Documentation
¶
Index ¶
- Variables
- func DecodeScanTasks(wire ScanTasks, metadata table.ScanPlanningMetadata, schema *iceberg.Schema, ...) ([]table.FileScanTask, error)
- type AuthManager
- type Catalog
- func (r *Catalog) CancelPlanning(ctx context.Context, ident table.Identifier, planID string) (err error)
- func (r *Catalog) CatalogType() catalog.Type
- func (r *Catalog) CheckFunctionExists(ctx context.Context, identifier table.Identifier) (bool, error)
- func (r *Catalog) CheckNamespaceExists(ctx context.Context, namespace table.Identifier) (bool, error)
- func (r *Catalog) CheckTableExists(ctx context.Context, identifier table.Identifier) (bool, error)
- func (r *Catalog) CheckViewExists(ctx context.Context, identifier table.Identifier) (bool, error)
- func (r *Catalog) Close() error
- func (r *Catalog) CommitTable(ctx context.Context, ident table.Identifier, requirements []table.Requirement, ...) (table.Metadata, string, error)
- func (r *Catalog) CommitTransaction(ctx context.Context, commits []table.TableCommit) error
- func (r *Catalog) CreateNamespace(ctx context.Context, namespace table.Identifier, props iceberg.Properties) error
- func (r *Catalog) CreateTable(ctx context.Context, identifier table.Identifier, schema *iceberg.Schema, ...) (*table.Table, error)
- func (r *Catalog) CreateView(ctx context.Context, identifier table.Identifier, version *view.Version, ...) (*view.View, error)
- func (r *Catalog) DropNamespace(ctx context.Context, namespace table.Identifier) error
- func (r *Catalog) DropTable(ctx context.Context, identifier table.Identifier) error
- func (r *Catalog) DropView(ctx context.Context, identifier table.Identifier) error
- func (r *Catalog) FetchPlanningResult(ctx context.Context, ident table.Identifier, planID string, ...) (result FetchPlanningResultResponse, err error)
- func (r *Catalog) FetchScanTasks(ctx context.Context, ident table.Identifier, req FetchScanTasksRequest) (result FetchScanTasksResponse, err error)
- func (r *Catalog) ListFunctions(ctx context.Context, namespace table.Identifier) iter.Seq2[table.Identifier, error]
- func (r *Catalog) ListNamespaces(ctx context.Context, parent table.Identifier) ([]table.Identifier, error)
- func (r *Catalog) ListTables(ctx context.Context, namespace table.Identifier) iter.Seq2[table.Identifier, error]
- func (r *Catalog) ListViews(ctx context.Context, namespace table.Identifier) iter.Seq2[table.Identifier, error]
- func (r *Catalog) LoadFunction(ctx context.Context, identifier table.Identifier) (*udf.UDF, error)
- func (r *Catalog) LoadNamespaceProperties(ctx context.Context, namespace table.Identifier) (iceberg.Properties, error)
- func (r *Catalog) LoadTable(ctx context.Context, identifier table.Identifier) (*table.Table, error)
- func (r *Catalog) LoadView(ctx context.Context, identifier table.Identifier) (*view.View, error)
- func (r *Catalog) Name() string
- func (r *Catalog) PlanFiles(ctx context.Context, req table.ScanPlanningRequest) (table.ScanPlanningResult, error)
- func (r *Catalog) PlanTableScan(ctx context.Context, ident table.Identifier, req PlanTableScanRequest) (result PlanTableScanResponse, err error)
- func (r *Catalog) PurgeTable(ctx context.Context, identifier table.Identifier) error
- func (r *Catalog) RegisterTable(ctx context.Context, identifier table.Identifier, metadataLoc string) (*table.Table, error)
- func (r *Catalog) RegisterView(ctx context.Context, identifier table.Identifier, metadataLoc string) (*view.View, error)
- func (r *Catalog) RenameTable(ctx context.Context, from, to table.Identifier) (*table.Table, error)
- func (r *Catalog) RenameView(ctx context.Context, from, to table.Identifier) (*view.View, error)
- func (r *Catalog) SetPageSize(ctx context.Context, sz int) context.Context
- func (r *Catalog) SupportsFullRemoteScanPlanning() bool
- func (r *Catalog) SupportsPlanTableScan() bool
- func (r *Catalog) SupportsRemoteScanPlanning() bool
- func (r *Catalog) UpdateNamespaceProperties(ctx context.Context, namespace table.Identifier, removals []string, ...) (catalog.PropertiesUpdateSummary, error)
- func (r *Catalog) UpdateTable(ctx context.Context, ident table.Identifier, requirements []table.Requirement, ...) (*table.Table, error)
- func (r *Catalog) UpdateView(ctx context.Context, ident table.Identifier, requirements []view.Requirement, ...) (*view.View, error)
- func (r *Catalog) WaitForPlan(ctx context.Context, ident table.Identifier, planID string, ...) (CompletedPlanningResult, error)
- type CompletedPlanningResult
- type ContextAuthManager
- type FetchPlanningResultOptions
- type FetchPlanningResultResponse
- type FetchScanTasksRequest
- type FetchScanTasksResponse
- type FunctionCatalog
- type Oauth2AuthManager
- type Option
- func WithAdditionalProps(props iceberg.Properties) Option
- func WithAudience(audience string) Option
- func WithAuthManager(authManager AuthManager) Option
- func WithAuthURI(uri *url.URL) Option
- func WithAwsConfig(cfg aws.Config) Option
- func WithCredential(cred string) Option
- func WithCustomTransport(transport http.RoundTripper) Option
- func WithHeaders(headers map[string]string) Option
- func WithMetadataLocation(loc string) Option
- func WithOAuthTLSConfig(config *tls.Config) Option
- func WithOAuthToken(token string) Option
- func WithPrefix(prefix string) Option
- func WithResource(resource string) Option
- func WithScope(scope string) Option
- func WithSigV4() Option
- func WithSigV4RegionSvc(region, service string) Option
- func WithTLSConfig(config *tls.Config) Option
- func WithTransportFactory(factory TransportFactory) Option
- func WithWarehouseLocation(loc string) Option
- type PlanFailedError
- type PlanStatus
- type PlanTableScanRequest
- type PlanTableScanResponse
- type PlanningError
- type RESTContentFile
- type RESTCountMap
- type RESTDataFile
- type RESTDeleteFile
- type RESTFileScanTask
- type RESTValueMap
- type ScanTasks
- type SnapshotMode
- type StorageCredential
- type TransportFactory
- type WaitForPlanOptions
Constants ¶
This section is empty.
Variables ¶
var ( ErrRESTError = errors.New("REST error") ErrBadRequest = fmt.Errorf("%w: bad request", ErrRESTError) ErrForbidden = fmt.Errorf("%w: forbidden", ErrRESTError) ErrAuthorizationExpired = fmt.Errorf("%w: authorization expired", ErrRESTError) ErrServerError = fmt.Errorf("%w: server error", ErrRESTError) // ErrCommitFailed wraps both ErrRESTError and table.ErrCommitFailed // so that callers can detect retryable commit conflicts via // errors.Is(err, table.ErrCommitFailed). ErrCommitFailed = fmt.Errorf("%w: %w", ErrRESTError, table.ErrCommitFailed) ErrCommitStateUnknown = fmt.Errorf("%w: commit failed due to unknown reason", ErrRESTError) ErrOAuthError = fmt.Errorf("%w: oauth error", ErrRESTError) )
var DefaultWaitForPlanOptions = WaitForPlanOptions{ MinDelay: 100 * time.Millisecond, MaxDelay: 5 * time.Second, CancelGracePeriod: 5 * time.Second, MaxRetries: 10, }
DefaultWaitForPlanOptions is the conservative polling backoff, cancel grace, and retry budget used when callers pass the zero-value WaitForPlanOptions. MaxRetries mirrors Java's maxRetries=10, but it applies only when the caller set no context deadline — with a deadline, the default cap is disabled so it cannot cut a long deadline short (see resolveWaitOptions). A caller expecting a long plan should therefore set a context deadline, not rely on this default.
var ErrEndpointNotSupported = errors.New("endpoint not supported by server")
ErrEndpointNotSupported means the server did not advertise an endpoint the operation needs.
var ErrNoSuchPlanTask = fmt.Errorf("%w: scan plan task not found", ErrRESTError)
ErrNoSuchPlanTask is returned when fetchScanTasks is called with a plan-task handle the server no longer knows about: a 404 whose error.type is NoSuchPlanTaskException. It is distinct from a table/namespace-gone 404 so a caller fanning out over plan-task handles can tell an expired handle from the table having vanished.
var ErrPlanCancelled = fmt.Errorf("%w: scan plan cancelled", ErrRESTError)
ErrPlanCancelled is returned by FetchPlanningResult when polling a plan that was cancelled (by this client or another). Like a failed plan it is terminal, so it surfaces as an error rather than a (resp, nil) the if-err idiom skips.
var ErrPlanExpired = fmt.Errorf("%w: scan plan expired", ErrRESTError)
ErrPlanExpired is returned when polling a plan that the server no longer knows about: a fetchPlanningResult 404 whose error.type is exactly NoSuchPlanIdException. It is distinct from a table/namespace-gone 404 (catalog.ErrNoSuchTable / catalog.ErrNoSuchNamespace) so the polling layer can tell retry-with-a-new-plan from abort. A bare or unrecognized 404 stays an ambiguous ErrRESTError rather than being guessed as an expiry.
var ErrPlanFailed = fmt.Errorf("%w: scan plan failed", ErrRESTError)
ErrPlanFailed is returned by PlanTableScan and FetchPlanningResult when the server reports a failed plan. The returned error is a *PlanFailedError, so the structured PlanningError detail is reachable via errors.As.
var ErrPlanPollExhausted = fmt.Errorf("%w: scan plan polling exhausted retries", ErrRESTError)
ErrPlanPollExhausted is returned by WaitForPlan when polling reaches WaitForPlanOptions.MaxRetries without the plan completing while the context is still live. It bounds a plan that stays submitted forever when the caller passes no context deadline, and is the Go analogue of Java's RemotePlanTimeoutException.
var ErrVendedCredentialsExpired = fmt.Errorf("%w: vended storage credentials expired", ErrRESTError)
ErrVendedCredentialsExpired is returned when a cached FileIO's vended creds expired with no endpoint to renew them (as a scan plan's own creds), so the caller sees this instead of undiagnosable storage 403s.
Functions ¶
func DecodeScanTasks ¶ added in v0.7.0
func DecodeScanTasks( wire ScanTasks, metadata table.ScanPlanningMetadata, schema *iceberg.Schema, fallbackResidual iceberg.BooleanExpression, ) ([]table.FileScanTask, error)
DecodeScanTasks converts one REST ScanTasks envelope into domain scan tasks. Delete-file references are scoped to one envelope, so callers must invoke this before combining inline and fetchScanTasks responses.
schema is the schema selected for the scan and is used to decode residual literals. When nil, metadata.CurrentSchema is used. fallbackResidual is used when a task omits residual-filter, as required by the REST specification; it is commonly the original scan filter.
Types ¶
type AuthManager ¶ added in v0.5.0
type AuthManager interface {
// AuthHeader returns the key and value for the authorization header.
AuthHeader() (string, string, error)
}
AuthManager is an interface for providing custom authorization headers.
type Catalog ¶
type Catalog struct {
// contains filtered or unexported fields
}
func NewCatalog ¶
func (*Catalog) CancelPlanning ¶ added in v0.7.0
func (r *Catalog) CancelPlanning(ctx context.Context, ident table.Identifier, planID string) (err error)
CancelPlanning cancels a server-side plan. Callers should cancel on context cancellation using a detached context with a short timeout. The spec supports idempotency and access-delegation headers on cancel; this low-level method deliberately defers those until a cancel options type is added, and suppresses the session-default access-delegation header (cancel vends no credentials). A 404 (already-expired or unknown plan) is not special-cased: cancel is best-effort, so the generic REST error is acceptable.
func (*Catalog) CatalogType ¶
func (*Catalog) CheckFunctionExists ¶ added in v0.7.0
func (r *Catalog) CheckFunctionExists(ctx context.Context, identifier table.Identifier) (bool, error)
CheckFunctionExists returns if the function exists. The REST spec defines no HEAD endpoint for functions, so existence is checked by loading the function; if loading is unsupported, ErrEndpointNotSupported surfaces rather than a bogus "not found". The identifier is validated once here, so an invalid one surfaces as an error and only a server-reported "not found" becomes (false, nil), matching the table and view existence checks.
func (*Catalog) CheckNamespaceExists ¶
func (*Catalog) CheckTableExists ¶
func (*Catalog) CheckViewExists ¶
func (*Catalog) Close ¶ added in v0.7.0
Close drains idle connections from transports created by the catalog and releases its metrics reporter. Closing the metrics dispatcher cancels any in-flight reports and drains its workers so outbound telemetry stops. Caller-provided transports remain caller-owned. Close is safe to call more than once. Callers holding a catalog.Catalog can reach this via a catalog.Closer type assertion.
func (*Catalog) CommitTable ¶
func (*Catalog) CommitTransaction ¶ added in v0.6.0
CommitTransaction atomically commits changes to multiple tables in a single request. It implements catalog.TransactionalCatalog.
The server applies all changes or none (all-or-nothing). On success (204 No Content) the method returns nil. Callers must LoadTable individually to obtain updated metadata.
func (*Catalog) CreateNamespace ¶
func (r *Catalog) CreateNamespace(ctx context.Context, namespace table.Identifier, props iceberg.Properties) error
func (*Catalog) CreateTable ¶
func (*Catalog) CreateView ¶ added in v0.3.0
func (r *Catalog) CreateView(ctx context.Context, identifier table.Identifier, version *view.Version, schema *iceberg.Schema, opts ...catalog.CreateViewOpt) (*view.View, error)
CreateView creates a new view in the catalog.
func (*Catalog) DropNamespace ¶
func (*Catalog) FetchPlanningResult ¶ added in v0.7.0
func (r *Catalog) FetchPlanningResult(ctx context.Context, ident table.Identifier, planID string, opts FetchPlanningResultOptions) (result FetchPlanningResultResponse, err error)
FetchPlanningResult polls a previously submitted plan. opts.AccessDelegation is sent as the X-Iceberg-Access-Delegation header so an async poll can still receive plan-scoped storage credentials: the spec defines data-access on this endpoint, and the completed-async result is where those credentials are vended.
completed and submitted return (resp, nil) and callers branch on resp.Status (done vs poll again). The terminal failure states surface as errors so an if-err poll loop cannot mistake them for an empty scan: failed returns a *PlanFailedError (errors.Is(err, ErrPlanFailed)) and cancelled returns ErrPlanCancelled. A 404 is split by the response error.type: a forgotten plan-id (NoSuchPlanIdException) is ErrPlanExpired, while a gone table or namespace is catalog.ErrNoSuchTable / catalog.ErrNoSuchNamespace, so the poller can tell retry-with-a-new-plan from abort. A bare or unrecognized 404 stays an ambiguous ErrRESTError rather than being guessed as an expiry.
func (*Catalog) FetchScanTasks ¶ added in v0.7.0
func (r *Catalog) FetchScanTasks(ctx context.Context, ident table.Identifier, req FetchScanTasksRequest) (result FetchScanTasksResponse, err error)
FetchScanTasks fetches the scan tasks for a plan-task handle returned by a completed plan. A 404 is split by the response error.type: an expired plan-task handle (NoSuchPlanTaskException) is ErrNoSuchPlanTask, while a gone table or namespace is catalog.ErrNoSuchTable / catalog.ErrNoSuchNamespace, so a caller fanning out over handles can tell a handle expiry from the table vanishing. A bare or unrecognized 404 stays an ambiguous ErrRESTError. An empty 200 body is rejected (requireBody) so a truncated response is not read as a successfully completed empty task set. Transient HTTP responses use the same retry policy as Catalog.PlanTableScan, retaining the idempotency key and plan-task handle across attempts.
func (*Catalog) ListFunctions ¶ added in v0.7.0
func (r *Catalog) ListFunctions(ctx context.Context, namespace table.Identifier) iter.Seq2[table.Identifier, error]
ListFunctions returns the function (SQL UDF) identifiers under a namespace, with the returned identifiers containing the information required to load the function via this catalog. The function endpoints are not part of the spec's assumed default endpoint set; if the server does not advertise the list-functions endpoint, it yields no functions rather than an error.
func (*Catalog) ListNamespaces ¶
func (r *Catalog) ListNamespaces(ctx context.Context, parent table.Identifier) ([]table.Identifier, error)
ListNamespaces lists namespaces under parent. If the server does not advertise the list-namespaces endpoint, it returns an empty slice rather than an error.
func (*Catalog) ListTables ¶
func (r *Catalog) ListTables(ctx context.Context, namespace table.Identifier) iter.Seq2[table.Identifier, error]
ListTables lists the tables in a namespace. If the server does not advertise the list-tables endpoint, it yields no tables rather than an error.
func (*Catalog) ListViews ¶
func (r *Catalog) ListViews(ctx context.Context, namespace table.Identifier) iter.Seq2[table.Identifier, error]
ListViews lists the views in a namespace. If the server does not advertise the list-views endpoint, it yields no views rather than an error.
func (*Catalog) LoadFunction ¶ added in v0.7.0
LoadFunction loads a function (SQL UDF) from the catalog. All overloaded definitions are included in the single metadata response. The function endpoints are not part of the spec's assumed default endpoint set; if the server does not advertise the load-function endpoint, LoadFunction fails with ErrEndpointNotSupported.
func (*Catalog) LoadNamespaceProperties ¶
func (r *Catalog) LoadNamespaceProperties(ctx context.Context, namespace table.Identifier) (iceberg.Properties, error)
func (*Catalog) LoadTable ¶
LoadTable loads a table from the catalog. It implements catalog.Catalog. When snapshot-loading-mode is set to "refs" in the catalog properties, only snapshots referenced by a named branch or tag are included in the response. Callers that need to time-travel to an unreferenced snapshot should use snapshot-loading-mode "all" or omit the property entirely.
func (*Catalog) PlanFiles ¶ added in v0.7.0
func (r *Catalog) PlanFiles(ctx context.Context, req table.ScanPlanningRequest) (table.ScanPlanningResult, error)
PlanFiles plans a scan server-side and returns tasks (and, optionally, a plan-scoped FileIO) for the table to read. It submits the plan, polls a submitted plan to completion, and expands any plan-task handles into their tasks before returning.
Each response is decoded as its own ScanTasks envelope before the results are combined. Delete-file references are scoped to the envelope that returned them, so flattening all responses first would attach deletes to the wrong data files when the server uses plan-task fanout.
func (*Catalog) PlanTableScan ¶ added in v0.7.0
func (r *Catalog) PlanTableScan(ctx context.Context, ident table.Identifier, req PlanTableScanRequest) (result PlanTableScanResponse, err error)
PlanTableScan submits a scan plan. A completed or submitted plan returns (resp, nil); callers branch on resp.Status for the plan-id (completed) vs poll (submitted) distinction. A failed plan returns a zero response and a non-nil *PlanFailedError (errors.Is(err, ErrPlanFailed)) so the if-err idiom does not mistake a failure for success; the server detail rides on the error. Any other status — including the empty status of a 200 with no body, which bypasses the response UnmarshalJSON validation — returns an ErrRESTError so a malformed response cannot masquerade as an empty completed plan.
Transient HTTP responses (408, 429, 500, 502, 503, 504) are retried up to three times with jittered backoff and the same idempotency key. Retry-After is honored, capped at five seconds when ctx has no deadline. Transport failures and malformed successful responses are not retried.
func (*Catalog) PurgeTable ¶
func (*Catalog) RegisterTable ¶
func (*Catalog) RegisterView ¶ added in v0.6.0
func (r *Catalog) RegisterView(ctx context.Context, identifier table.Identifier, metadataLoc string) (*view.View, error)
RegisterView registers an existing view in the catalog using its metadata file location. The metadata file must already be accessible to the catalog. This is the view equivalent of RegisterTable, using the REST endpoint POST /namespaces/{ns}/register-view defined in the Iceberg REST catalog specification.
func (*Catalog) RenameTable ¶
func (*Catalog) RenameView ¶ added in v0.7.0
func (*Catalog) SetPageSize ¶
func (*Catalog) SupportsFullRemoteScanPlanning ¶ added in v0.7.0
SupportsFullRemoteScanPlanning reports whether the server advertised the execution endpoints (plan, fetch-result, fetch-tasks), i.e. it can drive the async/fanout path, not just sync inline planning. Cancellation is optional cleanup and does not prevent a plan from producing tasks.
func (*Catalog) SupportsPlanTableScan ¶ added in v0.7.0
SupportsPlanTableScan reports whether the server advertised the plan submission endpoint.
func (*Catalog) SupportsRemoteScanPlanning ¶ added in v0.7.0
SupportsRemoteScanPlanning reports whether this catalog can submit a remote plan. Any continuation capability is validated against the response returned by the server; SupportsFullRemoteScanPlanning reports whether auto mode can safely use the complete continuation flow without a local fallback.
func (*Catalog) UpdateNamespaceProperties ¶
func (r *Catalog) UpdateNamespaceProperties(ctx context.Context, namespace table.Identifier, removals []string, updates iceberg.Properties, ) (catalog.PropertiesUpdateSummary, error)
func (*Catalog) UpdateTable ¶
func (*Catalog) UpdateView ¶ added in v0.5.0
func (r *Catalog) UpdateView(ctx context.Context, ident table.Identifier, requirements []view.Requirement, updates []view.Update) (*view.View, error)
UpdateView updates a view in the catalog.
func (*Catalog) WaitForPlan ¶ added in v0.7.0
func (r *Catalog) WaitForPlan(ctx context.Context, ident table.Identifier, planID string, opts WaitForPlanOptions) (CompletedPlanningResult, error)
WaitForPlan polls a submitted plan to completion using jittered backoff. It retries the idempotent-GET status set used by Java's REST transport (408/429/500/502/503/504); every other HTTP error is terminal. A positive Retry-After hint on a retried error response overrides the backoff, floored at MinDelay and — only when the caller set no context deadline — capped at MaxDelay; with a deadline the hint is honoured in full, bounded by that deadline (so a rate-limiting 429/503 is respected rather than clamped). The hint is read only from error responses, not from a 200 "submitted" body.
It returns the completed planning result verbatim: a completed plan may carry file-scan-tasks that are ready to read and/or plan-tasks (opaque handles that still need FetchScanTasks). WaitForPlan does not expand plan-tasks — driving that fanout is the caller's (PlanFiles') job — so "completed" does not imply every task is a file-scan-task.
Polling time is bounded by the caller's context deadline (the Go equivalent of Java's REST_SCAN_PLANNING_POLL_TIMEOUT_MS). opts.MaxRetries is a separate safety net for a caller that set no deadline: it caps the poll count and, when exhausted, cancels the plan server-side and returns ErrPlanPollExhausted. The default cap is disabled when the context has a deadline, so a longer deadline is never silently cut short by the retry count.
On context cancellation — or on retry exhaustion — WaitForPlan cancels the plan server-side and then returns (the wrapped context error, errors.Is context.Canceled / context.DeadlineExceeded; or ErrPlanPollExhausted). The cancel is synchronous so a short-lived caller cannot exit before the plan is released; it is best-effort and bounded by opts.CancelGracePeriod, so a stalled cancel endpoint delays the return by at most that grace (see abandonPlan).
It requires the fetchPlanningResult endpoint (returning ErrEndpointNotSupported otherwise) and a non-empty planID (ErrInvalidArgument otherwise); the cancel is best-effort, so cancelPlanning is not required. Other terminal errors from the poll — *PlanFailedError (errors.Is ErrPlanFailed), ErrPlanCancelled, ErrPlanExpired, or a gone table/namespace — are returned unchanged.
type CompletedPlanningResult ¶ added in v0.7.0
type CompletedPlanningResult struct {
Status PlanStatus `json:"status"`
ScanTasks
StorageCredentials []StorageCredential `json:"storage-credentials,omitempty"`
}
CompletedPlanningResult is the completed arm of the planning-result union. planTableScan carries the plan-id on PlanTableScanResponse; fetchPlanningResult omits it.
type ContextAuthManager ¶ added in v0.7.0
type ContextAuthManager interface {
AuthManager
// AuthHeaderWithContext returns the authorization header, honoring ctx.
AuthHeaderWithContext(ctx context.Context) (string, string, error)
}
ContextAuthManager is an optional interface an AuthManager may implement to honor a caller-supplied context (deadline and cancellation) while producing the authorization header. sessionTransport prefers AuthHeaderWithContext when the manager implements it, so a request's deadline also bounds the auth step rather than only the request/response cycle. A manager implementing only AuthManager still works, but its auth step cannot be interrupted mid-call — bounding that case is the implementer's responsibility.
type FetchPlanningResultOptions ¶ added in v0.7.0
type FetchPlanningResultOptions struct {
// AccessDelegation is sent as the X-Iceberg-Access-Delegation header; nil
// uses the catalog default. The spec defines data-access on this endpoint,
// so an async poll needs it to receive vended storage credentials.
AccessDelegation *string
}
FetchPlanningResultOptions carries per-call headers for fetchPlanningResult.
type FetchPlanningResultResponse ¶ added in v0.7.0
type FetchPlanningResultResponse struct {
Status PlanStatus `json:"status"`
Error *PlanningError `json:"error,omitempty"`
ScanTasks
StorageCredentials []StorageCredential `json:"storage-credentials,omitempty"`
}
FetchPlanningResultResponse is the GET .../plan/{plan-id} poll response. Same `status`-discriminated union (completed / submitted / cancelled / failed).
func (*FetchPlanningResultResponse) UnmarshalJSON ¶ added in v0.7.0
func (r *FetchPlanningResultResponse) UnmarshalJSON(data []byte) error
type FetchScanTasksRequest ¶ added in v0.7.0
type FetchScanTasksRequest struct {
// IdempotencyKey is sent as the Idempotency-Key header, not in the JSON body.
// If set, it must be a UUIDv7 string (RFC 9562). If nil, a fresh UUIDv7 is
// generated per call and reused for automatic retries. Supply an explicit
// key to preserve idempotency across separate calls.
IdempotencyKey *string `json:"-"`
PlanTask string `json:"plan-task"`
}
FetchScanTasksRequest is the POST .../tasks request body.
type FetchScanTasksResponse ¶ added in v0.7.0
type FetchScanTasksResponse struct {
ScanTasks
}
FetchScanTasksResponse is the POST .../tasks response. May itself return more plan-tasks for further fanout.
func (*FetchScanTasksResponse) UnmarshalJSON ¶ added in v0.7.0
func (r *FetchScanTasksResponse) UnmarshalJSON(data []byte) error
type FunctionCatalog ¶ added in v0.7.0
type FunctionCatalog interface {
// ListFunctions returns the function identifiers under a namespace,
// with the returned identifiers containing the information required
// to load the function via this catalog.
ListFunctions(ctx context.Context, namespace table.Identifier) iter.Seq2[table.Identifier, error]
// LoadFunction loads a function from the catalog. All overloaded
// definitions are included in the single metadata response.
LoadFunction(ctx context.Context, identifier table.Identifier) (*udf.UDF, error)
// CheckFunctionExists returns if the function exists. The REST spec
// defines no HEAD endpoint for functions, so existence is checked by
// loading the function.
CheckFunctionExists(ctx context.Context, identifier table.Identifier) (bool, error)
}
FunctionCatalog is an optional interface for catalogs that support the function (SQL UDF) read endpoints defined by the REST spec: list and load. The function endpoints are not part of the spec's assumed default endpoint set, so they are only used when the server advertises them: unsupported loads fail with ErrEndpointNotSupported and unsupported listings yield no results. Callers holding a catalog.Catalog can check for this capability via a type assertion:
if fc, ok := cat.(rest.FunctionCatalog); ok {
fn, err := fc.LoadFunction(ctx, ident)
}
type Oauth2AuthManager ¶ added in v0.5.0
type Oauth2AuthManager struct {
// contains filtered or unexported fields
}
Oauth2AuthManager is an implementation of the AuthManager interface which uses an oauth2.TokenSource to provide bearer tokens. The token source handles caching, thread-safe refresh, and expiry management.
func (*Oauth2AuthManager) AuthHeader ¶ added in v0.5.0
func (o *Oauth2AuthManager) AuthHeader() (string, string, error)
AuthHeader returns the authorization header with the bearer token.
func (*Oauth2AuthManager) AuthHeaderWithContext ¶ added in v0.7.0
AuthHeaderWithContext returns the authorization header, honoring ctx. A cached, unexpired token is returned without any network call, so a bounded caller (such as the metrics dispatcher) pays no auth latency in the common case. When a refresh is required it goes through the token-refresh HTTP client, whose Timeout bounds how long a stalled token endpoint can block; the ctx check here additionally short-circuits work for a caller whose deadline has already elapsed.
type Option ¶
type Option func(*options)
func WithAdditionalProps ¶
func WithAdditionalProps(props iceberg.Properties) Option
func WithAudience ¶ added in v0.6.0
WithAudience sets the audience parameter sent in OAuth token requests. See RFC 8693 for details on token exchange audiences.
func WithAuthManager ¶ added in v0.5.0
func WithAuthManager(authManager AuthManager) Option
func WithAuthURI ¶
WithAuthURI sets the OAuth2 token endpoint (oauth2-server-uri). When not set, the client falls back to {catalog}/v1/oauth/tokens. This is the programmatic equivalent of the oauth2-server-uri property (or its rest.authorization-url alias).
func WithAwsConfig ¶
func WithCredential ¶
func WithCustomTransport ¶ added in v0.4.0
func WithCustomTransport(transport http.RoundTripper) Option
WithCustomTransport replaces the internally configured http.Transport with the provided http.RoundTripper. Certain options such as WithTLSConfig which modify the default http.Transport will no longer work since the entire transport is replaced.
func WithHeaders ¶ added in v0.5.0
func WithMetadataLocation ¶
func WithOAuthTLSConfig ¶ added in v0.6.0
WithOAuthTLSConfig sets a separate TLS configuration for the HTTP client used to communicate with the OAuth2 server. This is useful when the OAuth2 server (oauth2-server-uri) is a different host than the catalog and requires different TLS settings (e.g. a different CA or client certificate).
If not set, the OAuth2 client reuses the catalog's HTTP client (and its TLS configuration).
func WithOAuthToken ¶
func WithPrefix ¶
func WithResource ¶ added in v0.6.0
WithResource sets the resource parameter sent in OAuth token requests. See RFC 8707 for details on resource indicators.
func WithSigV4RegionSvc ¶
func WithTLSConfig ¶
func WithTransportFactory ¶ added in v0.7.0
func WithTransportFactory(factory TransportFactory) Option
WithTransportFactory configures a catalog-owned transport factory. The returned cleanup function is called when the catalog is closed. Use WithCustomTransport when the caller owns the transport instead.
func WithWarehouseLocation ¶
type PlanFailedError ¶ added in v0.7.0
type PlanFailedError struct {
Detail *PlanningError
}
PlanFailedError carries the server's structured PlanningError detail for a failed plan. It satisfies errors.Is(err, ErrPlanFailed) so callers can branch on the failure with the if-err idiom while still reaching the detail via errors.As.
func (*PlanFailedError) Error ¶ added in v0.7.0
func (e *PlanFailedError) Error() string
func (*PlanFailedError) Unwrap ¶ added in v0.7.0
func (e *PlanFailedError) Unwrap() error
type PlanStatus ¶ added in v0.7.0
type PlanStatus string
PlanStatus is the status of a server-side plan.
const ( PlanStatusCompleted PlanStatus = "completed" PlanStatusSubmitted PlanStatus = "submitted" // PlanStatusCancelled is valid when polling a submitted plan, but invalid // as a planTableScan response. PlanTableScan and WaitForPlan should treat a // cancelled initial planning response as an error. PlanStatusCancelled PlanStatus = "cancelled" PlanStatusFailed PlanStatus = "failed" )
type PlanTableScanRequest ¶ added in v0.7.0
type PlanTableScanRequest struct {
// IdempotencyKey is sent as the Idempotency-Key header, not in the JSON body.
// If set, it must be a UUIDv7 string (RFC 9562). If nil, a fresh UUIDv7 is
// generated per call and reused for automatic retries. Supply an explicit
// key to preserve idempotency across separate calls.
IdempotencyKey *string `json:"-"`
// AccessDelegation is sent as the X-Iceberg-Access-Delegation header, not
// in the JSON body. Nil uses the catalog default.
AccessDelegation *string `json:"-"`
SnapshotID *int64 `json:"snapshot-id,omitempty"`
Select []string `json:"select,omitempty"`
Filter json.RawMessage `json:"filter,omitempty"`
MinRowsRequested *int64 `json:"min-rows-requested,omitempty"`
CaseSensitive *bool `json:"case-sensitive,omitempty"`
UseSnapshotSchema *bool `json:"use-snapshot-schema,omitempty"`
StatsFields []string `json:"stats-fields,omitempty"`
}
PlanTableScanRequest is the POST .../plan request body. Filter is the ExpressionParser-format JSON of the scan's row filter, produced by json.Marshal over the BooleanExpression (see planTableScanRequestFrom).
Point-in-time only for now: the spec's incremental start-snapshot-id / end-snapshot-id fields are deliberately omitted and land with the incremental phase, together with the matching fields on table.ScanPlanningRequest, so the wire type and the seam stay in agreement.
type PlanTableScanResponse ¶ added in v0.7.0
type PlanTableScanResponse struct {
Status PlanStatus `json:"status"`
PlanID *string `json:"plan-id,omitempty"`
Error *PlanningError `json:"error,omitempty"`
ScanTasks
StorageCredentials []StorageCredential `json:"storage-credentials,omitempty"`
}
PlanTableScanResponse is the POST .../plan response. The spec models this as a `status`-discriminated union; the flat struct carries every arm's fields with omitempty so none are discarded. Per the spec, plan-id is required for both completed (CompletedPlanningWithIDResult) and submitted (AsyncPlanningResult) responses here; the wire decoder must validate PlanID != nil at unmarshal rather than rely on the omitempty pointer. A cancelled status is invalid for this endpoint and must be treated as an error. A failed status decodes even when a non-compliant server omits or malforms Error; callers must branch on Status before dereferencing PlanID.
func (*PlanTableScanResponse) UnmarshalJSON ¶ added in v0.7.0
func (r *PlanTableScanResponse) UnmarshalJSON(data []byte) error
type PlanningError ¶ added in v0.7.0
type PlanningError struct {
Message string `json:"message"`
Type string `json:"type"`
Code int `json:"code"`
Stack []string `json:"stack,omitempty"`
}
PlanningError is the REST ErrorModel payload carried by the error arm of a failed planning result. It mirrors the package's internal error wire shape but is a dedicated exported struct so it renders cleanly in godoc and does not leak an unexported type or its unexported fields into the public API.
type RESTContentFile ¶ added in v0.7.0
type RESTContentFile struct {
SpecID int `json:"spec-id"`
Partition []json.RawMessage `json:"partition"`
Content string `json:"content"`
FilePath string `json:"file-path"`
FileFormat string `json:"file-format"`
FileSizeInBytes int64 `json:"file-size-in-bytes"`
RecordCount int64 `json:"record-count"`
KeyMetadata *string `json:"key-metadata,omitempty"`
SplitOffsets []int64 `json:"split-offsets,omitempty"`
SortOrderID *int `json:"sort-order-id,omitempty"`
}
RESTContentFile contains the fields common to REST data and delete files. Partition values remain raw until DecodeScanTasks resolves spec-id and the result type of each partition transform.
type RESTCountMap ¶ added in v0.7.0
RESTCountMap is the REST column-id-to-count representation. Keys and values are parallel arrays because JSON object keys cannot preserve integer IDs.
type RESTDataFile ¶ added in v0.7.0
type RESTDataFile struct {
RESTContentFile
FirstRowID *int64 `json:"first-row-id,omitempty"`
ColumnSizes *RESTCountMap `json:"column-sizes,omitempty"`
ValueCounts *RESTCountMap `json:"value-counts,omitempty"`
NullValueCounts *RESTCountMap `json:"null-value-counts,omitempty"`
NaNValueCounts *RESTCountMap `json:"nan-value-counts,omitempty"`
LowerBounds *RESTValueMap `json:"lower-bounds,omitempty"`
UpperBounds *RESTValueMap `json:"upper-bounds,omitempty"`
}
RESTDataFile is the data-file arm of the REST ContentFile union.
type RESTDeleteFile ¶ added in v0.7.0
type RESTDeleteFile struct {
RESTContentFile
EqualityIDs []int `json:"equality-ids,omitempty"`
// ReferencedDataFile is emitted by Java's ContentFileParser for
// file-scoped position deletes and deletion vectors. It is not currently
// declared by the REST OpenAPI DeleteFile schema, so clients must also
// support deriving a DV's target from its FileScanTask association.
ReferencedDataFile *string `json:"referenced-data-file,omitempty"`
ContentOffset *int64 `json:"content-offset,omitempty"`
ContentSizeInBytes *int64 `json:"content-size-in-bytes,omitempty"`
}
RESTDeleteFile is the position/equality arm of the REST ContentFile union. Content discriminates the variant. Puffin position deletes are decoded as deletion vectors and associated with the data file that references them.
type RESTFileScanTask ¶ added in v0.7.0
type RESTFileScanTask struct {
DataFile *RESTDataFile `json:"data-file"`
DeleteFileReferences []int `json:"delete-file-references,omitempty"`
ResidualFilter json.RawMessage `json:"residual-filter,omitempty"`
}
RESTFileScanTask is the REST FileScanTask wire payload. Delete references are zero-based indices into the DeleteFiles slice of the same ScanTasks envelope; they must be resolved before responses are combined.
type RESTValueMap ¶ added in v0.7.0
RESTValueMap is the REST column-id-to-binary representation used for lower and upper bounds. Java's ContentFileParser encodes each raw Iceberg binary value as a hexadecimal string; these are not typed JSON partition values.
type ScanTasks ¶ added in v0.7.0
type ScanTasks struct {
PlanTasks []string `json:"plan-tasks,omitempty"`
FileScanTasks []RESTFileScanTask `json:"file-scan-tasks,omitempty"`
DeleteFiles []RESTDeleteFile `json:"delete-files,omitempty"`
}
ScanTasks carries the task payload shared by completed planning responses and fetchScanTasks responses.
type SnapshotMode ¶ added in v0.7.0
type SnapshotMode string
SnapshotMode controls which snapshots are included in a loadTable response.
const ( // SnapshotModeAll requests all currently valid snapshots (server default). SnapshotModeAll SnapshotMode = "all" // SnapshotModeRefs requests only snapshots that are referenced by at least one named branch or tag. // The server omits unreferenced historical snapshots, which reduces response size for tables // with long snapshot histories. Note: time-travel to a snapshot not referenced by any branch // or tag will fail, as that snapshot will not be present in the returned metadata. SnapshotModeRefs SnapshotMode = "refs" )
type StorageCredential ¶ added in v0.7.0
type StorageCredential struct {
Prefix string `json:"prefix"`
Config iceberg.Properties `json:"config"`
}
StorageCredential carries REST-vended storage credentials scoped to matching object-location prefixes.
type TransportFactory ¶ added in v0.7.0
type TransportFactory func(*tls.Config) (http.RoundTripper, func())
TransportFactory creates a transport owned by the REST catalog and returns the cleanup function to run when the catalog is closed. The TLS config is the one configured with WithTLSConfig.
type WaitForPlanOptions ¶ added in v0.7.0
type WaitForPlanOptions struct {
MinDelay time.Duration
MaxDelay time.Duration
// CancelGracePeriod bounds the synchronous best-effort server-side cancel
// issued once the caller's context is done: it is the most WaitForPlan's
// return can be delayed past the deadline when the cancel endpoint stalls. A
// caller in a hurry can shrink it; a zero value uses DefaultWaitForPlanOptions.
CancelGracePeriod time.Duration
// MaxRetries caps the number of poll retries (attempts after the first) before
// WaitForPlan gives up with ErrPlanPollExhausted. It is the safety net for a
// caller that set no context deadline. Zero uses DefaultWaitForPlanOptions when
// the context has no deadline, and unlimited when it does (the deadline is the
// real bound). A negative value means unlimited; an explicit positive value is
// always honoured, even alongside a deadline.
MaxRetries int
// AccessDelegation is sent as the X-Iceberg-Access-Delegation header on each
// poll; nil uses the catalog default. Needed for async plans so the
// completed poll can return vended storage credentials.
AccessDelegation *string
}
WaitForPlanOptions tunes the polling backoff and bounds. The total polling wait is bounded by whichever comes first: the caller's context deadline (context.WithTimeout) or MaxRetries. There is deliberately no Timeout field — the context is the time bound, avoiding a duplicated deadline and its zero-value footgun. Zero MinDelay/MaxDelay/CancelGracePeriod values use DefaultWaitForPlanOptions. Zero MaxRetries uses the default only without a context deadline; with a deadline it is unlimited so the deadline remains the authoritative bound.