Documentation
¶
Overview ¶
Package bigquery implements bounded, lossless and explicitly approved analytical BigQuery reads through DALgo recordsets.
Index ¶
- Constants
- Variables
- func CanonicalJSON(raw []byte) ([]byte, error)
- func DecodeRows(raw []byte, fields []Field, limit int) ([][]Cell, error)
- func HashPayload(name string, raw []byte) (string, []byte, error)
- func OperationDeadline(now, executionDeadline, callerDeadline time.Time, httpLimit time.Duration, ...) (time.Time, error)
- func ParseJSON(raw []byte, limit int) (any, error)
- func ValidateQuery(query dal.Query) error
- type Approval
- type Bounds
- type CancelResult
- type Cell
- type Client
- func (c *Client) Approve(preview Preview, digest string) (Approval, error)
- func (c *Client) CancelJob(ctx context.Context, job JobRef) (CancelResult, error)
- func (c *Client) Execute(ctx context.Context, a Approval) (*Run, error)
- func (c *Client) Observe(ctx context.Context, p SourceProfile) (Observation, error)
- func (c *Client) ObserveConfigured(ctx context.Context, sourceDigest string, bounds Bounds) (Observation, error)
- func (c *Client) Preview(ctx context.Context, plan ReadPlan, execution Execution, bounds Bounds) (Preview, error)
- func (c *Client) RebindJob(ctx context.Context, receipt Receipt, encoded string) (RebindResult, error)
- func (c *Client) Resume(ctx context.Context, receipt Receipt, encoded string) (*Run, error)
- func (c *Client) Snapshot(ctx context.Context, receipt Receipt) (SnapshotResult, error)
- func (c *Client) Status(ctx context.Context, job JobRef) (JobStatus, error)
- type Clock
- type Config
- type Counters
- type DiscoveryAuthorize
- type DiscoveryBinding
- type DiscoveryClient
- type DiscoveryConfig
- type DiscoveryLocator
- type DiscoveryProvenance
- type Error
- type Execution
- type Executor
- type Field
- type FileLedger
- type Identity
- type JobRef
- type JobStatus
- type Ledger
- type LoadConfig
- type LoadJobRef
- type LoadOutcomeUnknownError
- type LoadReceipt
- type LoadWriter
- func (w *LoadWriter) CheckTablesAbsent(ctx context.Context, tableIDs []string) error
- func (w *LoadWriter) CreateTable(ctx context.Context, tableID string, schema []Field, ...) error
- func (w *LoadWriter) DatasetID() string
- func (w *LoadWriter) EnsureDataset(ctx context.Context) error
- func (w *LoadWriter) LoadTable(ctx context.Context, tableID string, schema []Field, ndjson io.Reader) (LoadReceipt, error)
- func (w *LoadWriter) Location() string
- func (w *LoadWriter) ProjectID() string
- func (w *LoadWriter) RecoverLoad(ctx context.Context, ref LoadJobRef) (LoadReceipt, error)
- func (w *LoadWriter) SetConstraints(ctx context.Context, tableID string, constraints *bq.TableConstraints) error
- type MemoryLedger
- type Observation
- type Order
- type Page
- type Parameter
- type Predicate
- type Prepare
- type Preview
- type Principal
- type Provider
- type PublicDiscovery
- type PublicSchemaField
- type ReadOnlyDatabase
- func (d *ReadOnlyDatabase) DescribeCollection(ctx context.Context, ref *dal.CollectionRef) (*dbschema.CollectionDef, error)
- func (d *ReadOnlyDatabase) DescribeSourceFields(ctx context.Context, ref *dal.CollectionRef) ([]Field, error)
- func (d *ReadOnlyDatabase) ListCollections(ctx context.Context, _ *record.Key) ([]dal.CollectionRef, error)
- func (d *ReadOnlyDatabase) ListConstraints(ctx context.Context, ref *dal.CollectionRef) ([]dbschema.ConstraintDef, error)
- func (d *ReadOnlyDatabase) ListIndexes(context.Context, *dal.CollectionRef) ([]dbschema.IndexDef, error)
- func (d *ReadOnlyDatabase) ListReferrers(context.Context, *dal.CollectionRef) ([]dbschema.Referrer, error)
- func (d *ReadOnlyDatabase) ListSourceViews(ctx context.Context) ([]dbschema.SourceViewDef, error)
- func (d *ReadOnlyDatabase) OpenSourceRows(ctx context.Context, ref *dal.CollectionRef) (dbschema.SourceRowCursor, error)
- type ReadPlan
- type RebindResult
- type Receipt
- type Run
- type SnapshotResult
- type SourceConfig
- type SourceLimits
- type SourceProfile
- type TableRef
Constants ¶
const BigqueryScope = "https://www.googleapis.com/auth/bigquery"
BigqueryScope is the OAuth scope required for dataset and load-job writes. The caller must obtain consent or configured local credentials for it.
const MaxDepth = 32
const MaxResponseBytes = 10 * 1024 * 1024
Variables ¶
var ( ErrDatasetLocationMismatch = errors.New("BigQuery dataset location does not match the requested location") ErrTableExists = errors.New("BigQuery destination table already exists") ErrTableNotOwned = errors.New("BigQuery table was not created by this load writer") ErrTableBusy = errors.New("BigQuery table already has a load or metadata update in progress") ErrTableOwnershipChanged = errors.New("BigQuery table ownership or schema changed") ErrLoadOutcomeUnknown = errors.New("BigQuery load outcome is unknown") )
Functions ¶
func CanonicalJSON ¶
CanonicalJSON implements RFC8785 for adapter-owned payloads restricted to safe integer JSON tokens. Warehouse integer/decimal/float values are string tokens.
func DecodeRows ¶
DecodeRows validates exact warehouse f/v shapes against an already validated projection schema. Schema metadata eligibility is a separate mandatory gate.
func HashPayload ¶
HashPayload uses explicit named projections, never recursive digest deletion. Non-bound observation times and approval nonce/times may appear at top level.
func OperationDeadline ¶
func OperationDeadline(now, executionDeadline, callerDeadline time.Time, httpLimit time.Duration, control bool, bytesRemaining int64) (time.Time, error)
OperationDeadline computes a bound from the original trusted-ledger deadline. It never starts or renews a run. Control means explicit status/cancel only; callers must still debit the unchanged cumulative byte ledger.
func ParseJSON ¶
ParseJSON validates bounded raw bytes without collapsing null or losing integers. Callers must bound decompressed reads before allocating this buffer.
func ValidateQuery ¶
ValidateQuery visits the whole supported AST before any policy wrapper or generic DALgo executor can choose a paid leaf scan fallback.
Types ¶
type Approval ¶
type Approval struct {
// contains filtered or unexported fields
}
Approval cannot be manufactured from JSON or a caller boolean. Execute reloads its nonce and immutable binding from the original private ledger.
type Bounds ¶
type Bounds struct {
PageSize int `json:"pageSize"`
MaxRows int `json:"maxRows"`
MaxPages int `json:"maxPages"`
ResponseBytes int `json:"responseBytes"`
TotalResponseBytes int64 `json:"totalResponseBytes"`
WallMs int `json:"wallMs"`
HTTPMs int `json:"httpMs"`
Concurrency int `json:"concurrency"`
}
func DefaultBounds ¶
func DefaultBounds() Bounds
type CancelResult ¶
type Client ¶
type Client struct {
// contains filtered or unexported fields
}
func (*Client) Observe ¶
func (c *Client) Observe(ctx context.Context, p SourceProfile) (Observation, error)
func (*Client) ObserveConfigured ¶ added in v0.1.2
func (c *Client) ObserveConfigured(ctx context.Context, sourceDigest string, bounds Bounds) (Observation, error)
ObserveConfigured validates metadata for a profile frozen by NewClient, selected by its SourceDigest (also exposed by Compile). Bounds apply to the entire inspection and both metadata responses. It uses authenticated datasets.get and tables.get only: no query preparation, dry run, job, rows, or ledger access. A successful observation does not approve execution or admit a public source.
func (*Client) RebindJob ¶
func (c *Client) RebindJob(ctx context.Context, receipt Receipt, encoded string) (RebindResult, error)
RebindJob verifies a fresh trusted generation for the SAME original subject. It does not renew approval, submit/fetch a job, or change receipt provenance, offsets, reservations, counters or the original execution deadline. After that deadline only existing bounded Status/CancelJob recovery remains usable.
func (*Client) Snapshot ¶ added in v0.1.1
Snapshot reads authoritative local recovery state without provider/policy preparation, HTTP, row delivery or ledger mutation. Immutable receipt authority must match the trusted run. Mutable state/counters come only from the ledger. The existing cursor is returned unchanged, never minted or renewed. Snapshot remains available after execution expiry; it grants no execution authority. Use the same caller control context for a control operation and its snapshot. The entire local operation has a ceiling of 15 seconds; a busy run lease fails immediately. Operator-owned filesystem primitives must still return.
type Config ¶
type Config struct {
Profiles []SourceProfile
Provider Provider
Transport http.RoundTripper
Ledger Ledger
Prepare Prepare
Clock Clock
}
type DiscoveryAuthorize ¶ added in v0.1.2
type DiscoveryAuthorize func(context.Context, DiscoveryLocator) (DiscoveryBinding, error)
DiscoveryAuthorize reads protected current state, not request claims. It must reject absent consent, sign-out or disallowed targets and honor cancellation.
type DiscoveryBinding ¶ added in v0.1.2
DiscoveryBinding is private operator-attested state. PolicyRevision must bind the current owner, consent, selected source and session; it is never published.
type DiscoveryClient ¶ added in v0.1.2
type DiscoveryClient struct {
// contains filtered or unexported fields
}
DiscoveryClient has no SourceProfile, query preparer, ledger or job methods.
func NewDiscoveryClient ¶ added in v0.1.2
func NewDiscoveryClient(cfg DiscoveryConfig) (*DiscoveryClient, error)
func (*DiscoveryClient) Discover ¶ added in v0.1.2
func (c *DiscoveryClient) Discover(ctx context.Context, sourceID string, bounds Bounds) (PublicDiscovery, error)
Discover gets one dataset's METADATA and one exact table's STORAGE_STATS. The existing guarded SDK transport supplies decoded per-response/cumulative byte, retry and deadline bounds. It does not cap encoded gzip wire bytes. No API credential is discovered.
type DiscoveryConfig ¶ added in v0.1.2
type DiscoveryConfig struct {
Allowlist []DiscoveryLocator
Authorize DiscoveryAuthorize
Provider Provider
Transport http.RoundTripper
Clock Clock
// EvidenceKind describes the trusted transport's origin, not an admission.
// Only provider-metadata and synthetic-fixture are supported.
EvidenceKind string
}
type DiscoveryLocator ¶ added in v0.1.2
type DiscoveryLocator struct {
SourceID string `json:"source_id"`
SourceProject string `json:"source_project"`
DatasetID string `json:"dataset_id"`
TableID string `json:"table_id"`
}
DiscoveryLocator is a protected, exact metadata target, never an execution profile. NewDiscoveryClient copies the allowlist; Discover accepts only its ID.
type DiscoveryProvenance ¶ added in v0.1.2
type Executor ¶
type Executor struct {
Client *Client
Profile SourceProfile
Approval Approval
Plan ReadPlan
}
Executor is a one-approved-plan DALgo bridge, not a DB or generic fallback.
func (Executor) ExecuteQueryToRecordsReader ¶
type FileLedger ¶
type FileLedger struct {
// contains filtered or unexported fields
}
func NewFileLedger ¶
func NewFileLedger(dir string) (*FileLedger, error)
type Identity ¶
Identity is attested by the operator's provider, never by request JSON. Read and Cancel describe already granted capabilities; the SDK requests no scopes.
type Ledger ¶
type Ledger interface {
// contains filtered or unexported methods
}
Ledger serializes durable updates and separately holds per-run leases during network and delivery operations. File leases are released by the OS on exit, so deliberate Resume retains state while recovering after process termination. Applications should reuse one private FileLedger path for the budget session.
type LoadConfig ¶ added in v0.2.0
type LoadConfig struct {
ProjectID string
DatasetID string
Location string
HTTPClient *http.Client
}
LoadConfig selects a destination and an authenticated HTTP client. The caller supplies credentials; this package does not discover credentials or request broader scopes on its own.
type LoadJobRef ¶ added in v0.2.0
type LoadJobRef struct {
JobID string `json:"jobId"`
ProjectID string `json:"projectId"`
DatasetID string `json:"datasetId"`
TableID string `json:"tableId"`
Location string `json:"location"`
}
LoadJobRef identifies a submitted load job without containing credentials or source data. Keep it when LoadTable returns an uncertain outcome.
type LoadOutcomeUnknownError ¶ added in v0.2.0
type LoadOutcomeUnknownError struct{ Job LoadJobRef }
LoadOutcomeUnknownError carries the durable job reference needed to recover the same load job without submitting it again.
func (*LoadOutcomeUnknownError) Error ¶ added in v0.2.0
func (e *LoadOutcomeUnknownError) Error() string
func (*LoadOutcomeUnknownError) Unwrap ¶ added in v0.2.0
func (e *LoadOutcomeUnknownError) Unwrap() error
type LoadReceipt ¶ added in v0.2.0
type LoadWriter ¶ added in v0.2.0
type LoadWriter struct {
// contains filtered or unexported fields
}
LoadWriter creates schema-bearing BigQuery tables and loads NDJSON data. It is separate from Client, whose approval and bounded-job rules are for reading public BigQuery sources.
func NewLoadWriter ¶ added in v0.2.0
func NewLoadWriter(ctx context.Context, cfg LoadConfig) (*LoadWriter, error)
NewLoadWriter builds a writer pinned to the BigQuery API origin. Redirects are refused so credentials cannot be forwarded to another host.
func (*LoadWriter) CheckTablesAbsent ¶ added in v0.2.0
func (w *LoadWriter) CheckTablesAbsent(ctx context.Context, tableIDs []string) error
CheckTablesAbsent checks every target identifier before the caller creates any of them. It protects existing tables and prevents a later collision from leaving a partially prepared destination schema.
func (*LoadWriter) CreateTable ¶ added in v0.2.0
func (w *LoadWriter) CreateTable(ctx context.Context, tableID string, schema []Field, constraints *bq.TableConstraints) error
CreateTable refuses every existing table, including empty tables. Schema and optional primary/foreign-key declarations are applied once, without a truncate or replace path. BigQuery records constraints as NOT ENFORCED.
func (*LoadWriter) DatasetID ¶ added in v0.2.0
func (w *LoadWriter) DatasetID() string
func (*LoadWriter) EnsureDataset ¶ added in v0.2.0
func (w *LoadWriter) EnsureDataset(ctx context.Context) error
EnsureDataset creates a missing dataset at the explicit location. An existing dataset is never modified and must already have that location.
func (*LoadWriter) LoadTable ¶ added in v0.2.0
func (w *LoadWriter) LoadTable(ctx context.Context, tableID string, schema []Field, ndjson io.Reader) (LoadReceipt, error)
LoadTable submits one atomic NEWLINE_DELIMITED_JSON load job. The table must have been created by CreateTable. It never appends to or truncates a table. The writer checks ownership and ETag immediately before submission, but the BigQuery load-job API has no destination-table ETag precondition; another principal can still replace the table between that check and job dispatch. A receipt confirms the identified load job's result, not exclusion of that external concurrent-writer race.
func (*LoadWriter) Location ¶ added in v0.2.0
func (w *LoadWriter) Location() string
func (*LoadWriter) ProjectID ¶ added in v0.2.0
func (w *LoadWriter) ProjectID() string
func (*LoadWriter) RecoverLoad ¶ added in v0.2.0
func (w *LoadWriter) RecoverLoad(ctx context.Context, ref LoadJobRef) (LoadReceipt, error)
RecoverLoad checks and, while the job remains pending, polls the same load job identified by ref. It never resubmits rows. Callers should use a fresh context after an earlier upload context was canceled.
func (*LoadWriter) SetConstraints ¶ added in v0.2.0
func (w *LoadWriter) SetConstraints(ctx context.Context, tableID string, constraints *bq.TableConstraints) error
SetConstraints updates only the primary/foreign-key metadata of a table that was created by this transfer. It is separate from CreateTable so cross-table references can be declared after every target primary key exists.
type MemoryLedger ¶
type MemoryLedger struct {
// contains filtered or unexported fields
}
func NewMemoryLedger ¶
func NewMemoryLedger() *MemoryLedger
type Observation ¶
type Observation struct {
Table TableRef `json:"table"`
Location string `json:"location"`
Type string `json:"type"`
Config map[string]any `json:"config"`
Schema []Field `json:"schema"`
Digest string `json:"digest"`
ObservedAt time.Time `json:"observedAt"`
Etag string `json:"etag,omitempty"`
LastModified string `json:"lastModified,omitempty"`
}
type Prepare ¶
Prepare reauthorizes the ORIGINAL protected query/context on every operation. The consumer supplies this trusted policy bridge; a cursor grants no access.
type Preview ¶
type Preview struct {
Nonce string `json:"nonce"`
Plan ReadPlan `json:"plan"`
Observation Observation `json:"observation"`
PolicyDigest string `json:"policyDigest"`
Execution Execution `json:"execution"`
Bounds Bounds `json:"bounds"`
EstimatedBytes string `json:"estimatedBytes"`
CreatedAt time.Time `json:"createdAt"`
ExpiresAt time.Time `json:"expiresAt"`
ApprovalDigest string `json:"approvalDigest"`
}
type Provider ¶
type Provider interface {
Authorize(context.Context, http.RoundTripper) (Identity, http.RoundTripper, error)
}
Provider composes authentication ABOVE the supplied guarded dispatch transport. Retry-capable authentication middleware cannot bypass its POST attempt gate.
type PublicDiscovery ¶ added in v0.1.2
type PublicDiscovery struct {
Format string `json:"format"`
SourceID string `json:"source_id"`
SourceProject string `json:"source_project"`
DatasetID string `json:"dataset_id"`
TableID string `json:"table_id"`
Location string `json:"location"`
ObjectType string `json:"object_type"`
ObservedAt string `json:"observed_at"`
Schema []PublicSchemaField `json:"schema"`
Projection string `json:"projection"`
Provenance DiscoveryProvenance `json:"provenance"`
SHA256 string `json:"sha256,omitempty"`
}
PublicDiscovery is an observation-only, partial registry projection. It contains no execution, rights, cost, retention or admission fields. SHA256 hashes canonical JSON of this envelope with only the sha256 property omitted.
type PublicSchemaField ¶ added in v0.1.2
type PublicSchemaField struct {
Name string `json:"name"`
Type string `json:"type"`
Mode string `json:"mode"`
Fields []PublicSchemaField `json:"fields,omitempty"`
}
PublicSchemaField is intentionally separate from executable Field. Optional descriptors, descriptions and security properties are never copied.
type ReadOnlyDatabase ¶ added in v0.3.0
ReadOnlyDatabase is a DALgo DB exposing schema and physical row read capabilities. Generic query execution and every write operation are refused.
func NewReadOnlyDatabase ¶ added in v0.3.0
func NewReadOnlyDatabase(cfg SourceConfig) (*ReadOnlyDatabase, error)
func (*ReadOnlyDatabase) DescribeCollection ¶ added in v0.3.0
func (d *ReadOnlyDatabase) DescribeCollection(ctx context.Context, ref *dal.CollectionRef) (*dbschema.CollectionDef, error)
func (*ReadOnlyDatabase) DescribeSourceFields ¶ added in v0.3.0
func (d *ReadOnlyDatabase) DescribeSourceFields(ctx context.Context, ref *dal.CollectionRef) ([]Field, error)
DescribeSourceFields retains BigQuery column descriptions, which the portable dbschema.FieldDef has no property for. The returned slice is independent data.
func (*ReadOnlyDatabase) ListCollections ¶ added in v0.3.0
func (d *ReadOnlyDatabase) ListCollections(ctx context.Context, _ *record.Key) ([]dal.CollectionRef, error)
func (*ReadOnlyDatabase) ListConstraints ¶ added in v0.3.0
func (d *ReadOnlyDatabase) ListConstraints(ctx context.Context, ref *dal.CollectionRef) ([]dbschema.ConstraintDef, error)
func (*ReadOnlyDatabase) ListIndexes ¶ added in v0.3.0
func (d *ReadOnlyDatabase) ListIndexes(context.Context, *dal.CollectionRef) ([]dbschema.IndexDef, error)
func (*ReadOnlyDatabase) ListReferrers ¶ added in v0.3.0
func (d *ReadOnlyDatabase) ListReferrers(context.Context, *dal.CollectionRef) ([]dbschema.Referrer, error)
func (*ReadOnlyDatabase) ListSourceViews ¶ added in v0.3.0
func (d *ReadOnlyDatabase) ListSourceViews(ctx context.Context) ([]dbschema.SourceViewDef, error)
func (*ReadOnlyDatabase) OpenSourceRows ¶ added in v0.3.0
func (d *ReadOnlyDatabase) OpenSourceRows(ctx context.Context, ref *dal.CollectionRef) (dbschema.SourceRowCursor, error)
type ReadPlan ¶
type ReadPlan struct {
Version int `json:"version"`
SourceDigest string `json:"sourceDigest"`
Projection []string `json:"projection"`
Where *Predicate `json:"where"`
Order []Order `json:"order"`
Limit int `json:"limit"`
Parameters []Parameter `json:"parameters"`
SQL string `json:"sql"`
Digest string `json:"digest"`
}
type RebindResult ¶
type Receipt ¶
type Receipt struct {
Version int `json:"version"`
RunID string `json:"runId"`
ApprovalDigest string `json:"approvalDigest"`
SourceDigest string `json:"sourceDigest"`
ObservationDigest string `json:"observationDigest"`
SchemaDigest string `json:"schemaDigest"`
Principal Principal `json:"principal"`
Job *JobRef `json:"job,omitempty"`
State string `json:"state"`
RunStartedAt time.Time `json:"runStartedAt"`
ExecutionDeadline time.Time `json:"executionDeadline"`
Bounds Bounds `json:"bounds"`
Counters Counters `json:"counters"`
ProcessedBytes *string `json:"processedBytes,omitempty"`
BilledBytes *string `json:"billedBytes,omitempty"`
CacheHit *bool `json:"cacheHit,omitempty"`
Warnings []string `json:"warnings"`
Reason string `json:"reason,omitempty"`
LocalStopped bool `json:"localStopped"`
ResidualSourceReplacementRace bool `json:"residualSourceReplacementRace"`
}
type Run ¶
type Run struct {
// contains filtered or unexported fields
}
type SnapshotResult ¶ added in v0.1.1
SnapshotResult contains current public receipt state and the existing opaque delivery cursor, if one was issued. It contains no rows or private plan values.
type SourceConfig ¶ added in v0.3.0
type SourceConfig struct {
ProjectID string
DatasetID string
Provider Provider
Transport http.RoundTripper
Clock Clock
Limits SourceLimits
}
SourceConfig selects exactly one BigQuery dataset and supplies an existing credential provider. The adapter does not discover or persist credentials.
type SourceLimits ¶ added in v0.3.0
type SourceLimits struct {
PageSize int
MaxPages int
MaxRows int64
MaxBytes int64
ResponseBytes int
WallTime time.Duration
}
SourceLimits bound a physical table read. Exceeding a bound returns an error; it never returns a successful, truncated source. Zero values use defaults.
type SourceProfile ¶
type SourceProfile struct {
Version int `json:"version"`
SourceID string `json:"sourceId"`
DescriptorDigest string `json:"descriptorDigest"`
LogicalCollection string `json:"logicalCollection"`
SourceProject string `json:"sourceProject"`
DatasetID string `json:"datasetId"`
TableID string `json:"tableId"`
Location string `json:"location"`
Schema []Field `json:"schema"`
PublisherReviewRef string `json:"publisherReviewRef"`
RightsReviewRef string `json:"rightsReviewRef"`
Use string `json:"use"`
}
Source Files
¶
Directories
¶
| Path | Synopsis |
|---|---|
|
examples
|
|
|
metadata-discovery
command
Command metadata-discovery emits one bounded public projection.
|
Command metadata-discovery emits one bounded public projection. |
|
operator-server
command
Command operator-server composes the approved loopback harness.
|
Command operator-server composes the approved loopback harness. |
|
server
Package server is an operator-owned loopback service harness.
|
Package server is an operator-owned loopback service harness. |