bigquery

package module
v0.3.0 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Oct 7, 2026 License: MIT Imports: 36 Imported by: 0

README

dalgo2bigquery

An analytical BigQuery driver for bounded, explicitly approved reads. It compiles a restricted DALgo scalar query, verifies a reviewed native TABLE, performs a dry run, and executes one capped jobs.query submission. Results retain a schema and exact scalar values without inventing record keys. Same-job paging, partial-page resume, status, cancellation requests, and trusted same-subject reconnect use a persistent private ledger.

The governing accepted A0 contract and approved plan require both Go and JavaScript. DALgo v0.93.0, record v0.1.4 and Google API SDK v0.296.0 are pinned without replacements. The independently reviewed Go driver was released as v0.1.0 from 0a559674. v0.1.1, from 7df3d802, adds the independently reviewed Client.Snapshot API. Consumer integration, JavaScript acceptance and actual live acceptance remain required.

Operator integration

Job-free physical table source

NewReadOnlyDatabase(SourceConfig{...}) creates a DALgo handle for one explicitly selected BigQuery project and dataset. It exposes dbschema.SchemaReader, dbschema.SourceRowsReader, and dbschema.SourceViewReader. ListCollections lists physical tables; ListSourceViews exposes ordinary view names and columns as metadata (its CreateSQL is empty because BigQuery returns the SELECT body, not original CREATE VIEW DDL), while OpenSourceRows refuses view reads. DescribeCollection carries BigQuery's declared column types and primary/foreign-key metadata; BigQuery key declarations are not enforced. DescribeSourceFields also returns original field descriptions, which portable dbschema.FieldDef cannot represent. Search/vector indexes and reverse foreign-key lookups are reported unsupported, rather than reported absent.

The source uses only datasets.get, tables.list, tables.get, and tabledata.list GET requests. It never submits a query job, writes data, discovers credentials, or stores tokens. The caller supplies an existing Provider and transport with read access to the selected dataset. Rows stream by page. INT64 values stay exact in int64; NUMERIC/BIGNUMERIC stay decimal strings; BYTES become raw []byte, preserving empty bytes separately from SQL NULL; TIMESTAMP becomes UTC time.Time. Nested/repeated types, field configurations such as maxLength or defaultValueExpression, and non-table objects fail before row dispatch; this is not a complete BigQuery schema reader. SourceLimits are admission bounds: reaching one returns an error instead of a successful prefix.

The cursor rejects repeated page tokens, changing totalRows, short or excessive reads, and changed table metadata at the end of the scan. BigQuery tabledata.list does not provide an immutable cross-page snapshot. A write that changes and then restores a table between checks can still escape these guards; callers claiming historical parity should compare stable tables or use a separate snapshot workflow. The existing approved-plan Executor remains the bounded SQL read path, with its own approval and job budget.

The operator supplies a reviewed SourceProfile, a trusted Provider, an explicit job project/cap/session budget, and a Prepare policy bridge. Provider.Authorize attests a verified stable subject, current generation, expiry and granted capabilities; it authenticates above the supplied dispatch guard. Request JSON, token presence, email and ADC labels cannot attest identity. Providers must respect the supplied context, wrap the provided transport, and redact their own errors/logging. Preparation, identity, authentication middleware, transport and body work receive the original finite operation bound. Ignored cancellation returns control to the caller, without granting delayed work continuation authority. A process-wide ceiling of 64 dependency workers bounds abandoned work; a never-returning trusted dependency can retain a worker slot and its physical per-run HTTP gate until process exit. Worker admission refusal releases an unused physical gate and never records a dispatch. After admission, the same worker retains its slot through physical response Close; EOF alone does not release the gate, and cleanup needs no additional worker admission. Subsequent operations cannot bypass that gate. This ceiling is not a session-wide active-job limit. The driver never acquires credentials or requests broader scopes.

Prepare re-evaluates the original protected query and source context on every operation, returning its current effective ReadPlan and policy digest. ValidateQuery must run before any generic DALgo executor can choose a paid leaf fallback. The consumer remains responsible for its real policy wrapper and diagnostics. Compile supports one unaliased collection, explicit scalar projection, comparisons with typed constant parameters, null predicates, bounded IN, conjunction/disjunction, scalar ordering and a positive bounded limit. Joins, computed/star/aliased projections, groups, HAVING, offsets and unknown nodes fail before I/O. Reviewed profiles require explicit schema modes; omitted REST mode means NULLABLE. Nested/repeated types remain discoverable and fail scalar execution.

Use one NewFileLedger(absolutePrivateDirectory) for a session, including after process restart. Linux/macOS use private modes, symlink rejection, OS file locks and fsynced atomic replacement. Windows uses an owner-only DACL, reparse-point rejection, OS file locks and write-through atomic replacement. MemoryLedger is an injected test harness, not durable session accounting. Bounds.Concurrency=1 serializes actual HTTP operations for each run through physical response-body completion. Separately approved runs share an atomic session budget; no session-wide active-job-count limit is claimed. Budget reservations include unknown submissions and terminal jobs lacking authoritative billing. Terminal billed bytes replace the reservation and remain accounted as session expenditure.

Recordsets expose only already authorized and accounted rows, detached from private buffered cells. Accessor or row mutation cannot reveal or alter future delivery. Private filesystem lock acquisition observes the operation context; atomic snapshots permit read-only receipt recovery without waiting for a writer lock. Filesystem read/write/fsync/rename primitives rely on the operator-owned local filesystem returning; this driver does not sandbox a hung OS filesystem. After a timeout, a bounded synchronous settlement attempt may add up to 10 ms; if a contended ledger prevents settlement, its persisted reserved/submitting record retains the full cap. No abandoned dependency can write the ledger or deliver rows.

The API sequence is:

plan, err := bigquery.Compile(profile, protectedQuery) // first ValidateQuery before fallback
// Handle err, establish trusted provider/Prepare, and open a private FileLedger.
client, err := bigquery.NewClient(bigquery.Config{
    Profiles: []bigquery.SourceProfile{profile}, Provider: trustedProvider,
    Transport: operatorTransport, Ledger: persistentLedger, Prepare: prepare,
})
preview, err := client.Preview(ctx, plan, execution, bigquery.DefaultBounds())
// Present the exact preview and its source/policy/identity/cost binding to the operator.
approval, err := client.Approve(preview, operatorApprovedDigest)
run, err := client.Execute(ctx, approval) // retain any returned run receipt even on error
page, err := run.NextPage()             // or Run.Next(), a DALgo RecordsetReader
resumed, err := client.Resume(ctx, page.Receipt, page.Cursor)
status, err := client.Status(ctx, *page.Receipt.Job)
cancel, err := client.CancelJob(ctx, *page.Receipt.Job)

Explicit destination loads

NewLoadWriter is a separate, narrow API for applications that own a BigQuery destination. The caller chooses the project, dataset, location, and authenticated HTTP client; the package does not discover credentials. The exported BigqueryScope names the write scope the caller must arrange locally. The writer creates a missing dataset only at the requested location, refuses existing destination tables, and records a private ownership label, schema and ETag for each table it creates. SetConstraints and LoadTable refuse tables that this writer did not create or whose ownership, schema or ETag changed. Constraint patches send the tracked ETag as If-Match, so a stale metadata update is rejected by BigQuery. The load-job API has no destination-table ETag precondition: an external principal can replace a table between the writer's final ownership check and load dispatch. A load receipt verifies the identified job and its destination, but cannot fence that external concurrent-writer race. The authenticated transport checks every dispatched API, resumable-upload and status request against the fixed HTTPS Google API origins before credentials are added, including SDK-generated chunk requests. It submits explicit-schema NDJSON load jobs with CREATE_NEVER, WRITE_EMPTY, zero bad records, and no ignored unknown fields. An ambiguous submission returns *LoadOutcomeUnknownError with a durable, credential-free LoadJobRef. Persist that reference and call RecoverLoad with a fresh context to poll the same job; it never resubmits source rows. Completed receipts require the same job's project, ID, location, destination and load statistics, so a missing row count is an unknown outcome rather than zero. Loads are atomic per table; a caller copying multiple tables must account for partial progress across separate jobs. Any primary and foreign keys are BigQuery metadata only: BigQuery does not enforce them.

This shows API shape; applications must handle each error before the next call. Executor bridges one approved effective plan to DALgo recordsets and rejects keyed records or a changed query. HTTP errors contain only sanitized codes/reasons. The generated SDK runs with a discard logger; validated raw response bodies drive cell/state decoding. Fixed Google origin, redirect refusal, decompressed body limits, duplicate-key/Unicode validation and a lower dispatch gate apply beneath authentication middleware. POST replay is never permitted; GET retries share a three-attempt limit and debit actual response/page counters.

Run.Close stops local delivery and preserves the known running job and reservation. A cursor is an opaque, ledger-checked continuation of the original query, job, schema, page, offset, principal generation and cumulative counters. Resume retains the original absolute execution deadline. Expired result work stops before new dispatch or body I/O; explicit status/cancel retain a bounded control opportunity and the same response budget. A cancellation request alone is not confirmed cancellation; an ambiguous stopped reason remains a failure. Known jobs stay recoverable after later timeout or malformed rows; lost or late unvalidated initial submissions retain submission_unknown and full cap.

RebindJob(ctx, receipt, cursor) verifies a newly trusted generation for the original stable kind/subject, current read grant and unchanged policy/source binding. It atomically changes only the private active access principal and rotates cursor authority. Original receipt/approval principal provenance, offsets, page hashes, schema, counters, reservations and deadline remain intact. A stale cursor/generation or concurrent lease loses authority. It issues no BigQuery request, dry run or approval renewal. After execution expiry it only prepares existing bounded status/cancel recovery. This implements the accepted same-subject reconnection clause; actual browser GIS/CLI identity attestation remains consumer work.

Client.Snapshot(ctx, priorReceipt) returns SnapshotResult{Receipt, Cursor} from the trusted local ledger under the existing run lease. It validates immutable receipt authority and returns current billing, state and cumulative counters with defensive copies. The opaque cursor remains the exact previously issued delivery binding, or an empty string if none exists; control byte debits are reflected in the receipt without minting a new cursor. Snapshot performs no provider/policy preparation, HTTP, row reads or state changes. It remains available after execution expiry, without granting permission to resume expired results. Preserve each control operation's separate status/cancel result, including partial-error outcomes, alongside this receipt snapshot. Consumers should share the original bounded control context across control and Snapshot. A busy run lease refuses immediately; the local operation observes earlier caller cancellation or a 15-second ceiling and retains the documented filesystem limitations. There is no new HTTP route.

Metadata inspection

Client.ObserveConfigured(ctx, plan.SourceDigest, bounds) selects only a reviewed profile frozen in NewClient and validates its dataset location, native table configuration and schema using authenticated datasets.get and tables.get. Google documents tables.get as table structure metadata without table data. Explicit Bounds limit per-response and cumulative bytes, per-HTTP time and total wall time. Unknown profiles and invalid bounds fail before dispatch. The method does not invoke Prepare, perform a dry run, create a job, read rows or access the query ledger. It uses the same trusted Provider and transport guards as query operations; read permission is required, cancellation permission is not. The existing Observe(ctx, profile) API retains its caller-supplied profile and default bounds behavior.

NewDiscoveryClient(DiscoveryConfig) provides separate initial metadata discovery without an execution profile. It freezes at most 16 protected exact source/project/dataset/table locators; Discover(ctx, sourceID, bounds) accepts only a configured source ID. Its required DiscoveryAuthorize bridge reads the current owner/consent/source/session binding and private policy revision before each dispatch and delivery, including another policy check after provider identity callbacks return. The trusted Provider attests the matching principal, expiry and read capability. Neither a bearer token nor request JSON attests this state. A discovery-only physical transport wrapper repeats identity and policy attestation inside the admitted transport worker immediately before invoking the trusted transport, including consent changes while that worker was queued. The client holds no query preparer or ledger and exposes no job/row API.

Discovery explicitly requests datasets.get?datasetView=METADATA, excluding the default ACL view, then tables.get?view=STORAGE_STATS, pinning the existing audited metadata view rather than the future-expanding FULL view. Storage statistics may arrive transiently but are discarded; no rows or table data are requested. Exact references and observed location must agree. The public projection is deliberately partial and can describe TABLE, VIEW, EXTERNAL, MATERIALIZED_VIEW or SNAPSHOT without making any of them executable. Unsupported flexible identifiers fail closed. Source IDs match the registry's segmented lowercase alphanumeric syntax and 80-character ceiling. There is no table listing or inferred table selection; WDI still needs an independently documented exact candidate table before this API can inspect it.

The output follows ovdb-bigquery-observation/draft-1: exact source refs, location/object type, UTC observation time at second precision, recursive schema (name, native type, explicit mode, nested fields), fixed partial-projection marker, credential-free provenance and canonical SHA-256. Depth is at most 8, schema fields total at most 500 and the complete envelope at most 64 KiB. Descriptions, etags, ACL/security/governance objects, default/generated expressions, optional type descriptors, counts, actor/consent IDs, credentials, job project and raw bodies are excluded recursively. Unknown properties are omitted; the output is not a complete native schema or an admission-ready profile. The SHA-256 covers canonical JSON of the entire public envelope with only sha256 omitted, including observation time and provenance; it differs from executable Observation digests. Consumers must retain inactive/blocked query, cost, rights and retention gates. Review and import the public envelope separately; never persist raw responses or the private operator binding.

The offline CLI harness demonstrates the production discovery path with hostile nested metadata fixtures. Its default live mode requires an injected operator factory and refuses when unconfigured. No WDI profile, source admission or real metadata access proof is included.

Metadata observation grants no query authority. Execution still requires an explicit job project, verified execution principal, current policy preparation, dry-run estimate, positive byte cap/session budget and exact approval. Source rights and provider-result retention authorization remain separate operator gates: BigQuery materializes query results in destination or temporary tables, and disabling cache retrieval does not prevent that storage. Go metadata tests do not establish browser OAuth or CORS readiness.

Local server harness

examples/server.NewHandler(Config) provides /preview, /execute, /page, /rebind, /status and /cancel with the injected operator client, fixed query/cost principal and explicit trusted browser origin. Closed bounded JSON requests cannot choose credentials, SQL or cost identity. The server rejects non-loopback clients and unknown origins. The runnable operator composition scaffold constructs the client, private ledger, handler and signal-owned server from a build-time injected trusted operator factory. Unconfigured invocation refuses before ledger creation or a listener; the factory must honor its bounded setup context. An operator application can also call server.Serve(ctx, "127.0.0.1:YOUR_PORT", handler) after establishing its trusted provider and private ledger; the context owns shutdown. No import or test starts a listener. The production handler tests use httptest.NewRecorder and an injected SDK transport for metadata → dry run → approved execution → two same-job pages → status.

Shared corpus and validation

testdata/contract/manifest.json binds corpus revision3 by SHA256 and records its immutable revision2 predecessor. All original 70 cases are byte-identical. Additions exercise named hash mutations and raw HTTP/job state, headers, bodies, dispatch counts, persistent-clock controls and reconnect. JavaScript must vendor exact accepted committed bytes, bind the source commit/manifest hash, and report every scenario through its production path. Private provisional comparisons do not establish final dual-runtime acceptance.

One rejected decompressed overflow scenario records expected_runtime_observations.go.bytes and .js.bytes: Go exposes a bounded prefix plus one overflow byte, while fetch exposes an entire chunk. Counters truthfully debit exposed bytes. The case explicitly names this streaming variance; error/state/cap/job identity/dispatch behavior and all other fields remain identical. A separate one-byte chunk schedule proves identical counters. This narrow interpretation was reviewed for the released Go component; JavaScript review and joint parity acceptance remain required. Counters are never omitted or clamped to manufacture parity.

Raw body, input and body_hex values must reach the production parser without lossy JSON conversion. HTTP advance_ms moves the injected clock before response headers arrive; gzip compresses the literal body for the harness. Cases distinguish SQL NULL, JSON text null, empty strings, omitted v, malformed UTF-8/surrogates and scalar numeric wire strings. FLOAT64 uses the exact shared entire-string decimal/exponent grammar; no Go literal syntax or JS coercion widens it. RFC8785 hashes accept only the frozen named projections and safe integer numeric JSON tokens, preserving UTF-16 key ordering and exact strings. Warehouse numerics remain strings.

Run go test -race ./..., go vet ./... and gofmt validation. CI runs those checks on Linux, macOS and Windows; local cross-compiles also cover current CLI release targets. BIGQUERY_CONTRACT_REPORT=/absolute/private/report.json go test -run TestSharedCorpus -count=1 . emits production HTTP observations for private parity review, excluding only random run IDs/ledger references.

The Go component of Task2 is independently reviewed and released through driver v0.1.1. Shared-corpus acceptance across both runtimes, Task3 browser parity, Task4 actual CLI/app wiring and diagnostics, Task5 reviewed source-access profile, and Task6 both live journeys/rights/admission/cost gates remain open. No actual cloud query/spend, source admission, consumer cutover or dual-runtime live success is claimed. Root owns further landing and release after independent review and exact-head checks.

Documentation

Overview

Package bigquery implements bounded, lossless and explicitly approved analytical BigQuery reads through DALgo recordsets.

Index

Constants

View Source
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.

View Source
const MaxDepth = 32
View Source
const MaxResponseBytes = 10 * 1024 * 1024

Variables

View Source
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

func CanonicalJSON(raw []byte) ([]byte, error)

CanonicalJSON implements RFC8785 for adapter-owned payloads restricted to safe integer JSON tokens. Warehouse integer/decimal/float values are string tokens.

func DecodeRows

func DecodeRows(raw []byte, fields []Field, limit int) ([][]Cell, error)

DecodeRows validates exact warehouse f/v shapes against an already validated projection schema. Schema metadata eligibility is a separate mandatory gate.

func HashPayload

func HashPayload(name string, raw []byte) (string, []byte, error)

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

func ParseJSON(raw []byte, limit int) (any, error)

ParseJSON validates bounded raw bytes without collapsing null or losing integers. Callers must bound decompressed reads before allocating this buffer.

func ValidateQuery

func ValidateQuery(query dal.Query) error

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 CancelResult struct {
	Job   JobRef `json:"job"`
	State string `json:"state"`
}

type Cell

type Cell struct {
	Type  string `json:"type"`
	Value any    `json:"value"`
}

func NormalizeScalar

func NormalizeScalar(field Field, value any) (Cell, error)

NormalizeScalar preserves SQL NULL separately from JSON text "null".

type Client

type Client struct {
	// contains filtered or unexported fields
}

func NewClient

func NewClient(cfg Config) (*Client, error)

func (*Client) Approve

func (c *Client) Approve(preview Preview, digest string) (Approval, error)

func (*Client) CancelJob

func (c *Client) CancelJob(ctx context.Context, job JobRef) (CancelResult, error)

func (*Client) Execute

func (c *Client) Execute(ctx context.Context, a Approval) (*Run, error)

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) Preview

func (c *Client) Preview(ctx context.Context, plan ReadPlan, execution Execution, bounds Bounds) (Preview, error)

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) Resume

func (c *Client) Resume(ctx context.Context, receipt Receipt, encoded string) (*Run, error)

func (*Client) Snapshot added in v0.1.1

func (c *Client) Snapshot(ctx context.Context, receipt Receipt) (SnapshotResult, error)

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.

func (*Client) Status

func (c *Client) Status(ctx context.Context, job JobRef) (JobStatus, error)

type Clock

type Clock interface {
	Now() time.Time
	Sleep(context.Context, time.Duration) error
}

type Config

type Config struct {
	Profiles  []SourceProfile
	Provider  Provider
	Transport http.RoundTripper
	Ledger    Ledger
	Prepare   Prepare
	Clock     Clock
}

type Counters

type Counters struct {
	Rows  int   `json:"rows"`
	Pages int   `json:"pages"`
	Bytes int64 `json:"bytes"`
}

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

type DiscoveryBinding struct {
	Principal      Principal
	PolicyRevision string
}

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 DiscoveryProvenance struct {
	Kind     string `json:"kind"`
	Method   string `json:"method"`
	Verifier string `json:"verifier"`
}

type Error

type Error struct {
	Code   string
	Reason string
}

Error contains a stable, sanitized code, never query values or provider bodies.

func (*Error) Error

func (e *Error) Error() string

type Execution

type Execution struct {
	JobProject         string    `json:"jobProject"`
	Principal          Principal `json:"principal"`
	MaximumBytesBilled string    `json:"maximumBytesBilled"`
	SessionBudgetBytes string    `json:"sessionBudgetBytes"`
}

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

func (e Executor) ExecuteQueryToRecordsReader(context.Context, dal.Query) (dal.RecordsReader, error)

func (Executor) ExecuteQueryToRecordsetReader

func (e Executor) ExecuteQueryToRecordsetReader(ctx context.Context, q dal.Query, options ...recordset.Option) (dal.RecordsetReader, error)

type Field

type Field struct {
	Name        string  `json:"name"`
	Type        string  `json:"type"`
	Mode        string  `json:"mode"`
	Description string  `json:"description,omitempty"`
	Fields      []Field `json:"fields,omitempty"`
	Precision   string  `json:"precision,omitempty"`
	Scale       string  `json:"scale,omitempty"`
}

type FileLedger

type FileLedger struct {
	// contains filtered or unexported fields
}

func NewFileLedger

func NewFileLedger(dir string) (*FileLedger, error)

type Identity

type Identity struct {
	Principal    Principal
	ExpiresAt    time.Time
	Read, Cancel bool
}

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 JobRef

type JobRef struct {
	ProjectID string `json:"projectId"`
	JobID     string `json:"jobId"`
	Location  string `json:"location"`
}

type JobStatus

type JobStatus struct {
	Job         JobRef   `json:"job"`
	State       string   `json:"state"`
	BilledBytes *string  `json:"billedBytes,omitempty"`
	Warnings    []string `json:"warnings"`
}

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 LoadReceipt struct {
	JobID     string
	ProjectID string
	DatasetID string
	TableID   string
	Location  string
	Rows      uint64
}

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 Order

type Order struct {
	Column    string `json:"column"`
	Direction string `json:"direction"`
}

type Page

type Page struct {
	Schema  []Field  `json:"schema"`
	Rows    [][]Cell `json:"rows"`
	Receipt Receipt  `json:"receipt"`
	Cursor  string   `json:"cursor"`
}

type Parameter

type Parameter struct {
	Name  string `json:"name"`
	Type  string `json:"type"`
	Value any    `json:"value"`
}

type Predicate

type Predicate struct {
	Op        string      `json:"op"`
	Column    string      `json:"column,omitempty"`
	Parameter string      `json:"parameter,omitempty"`
	Children  []Predicate `json:"children,omitempty"`
}

type Prepare

type Prepare func(context.Context) (ReadPlan, string, error)

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 Principal

type Principal struct {
	Kind       string `json:"kind"`
	Subject    string `json:"subject"`
	Generation string `json:"generation"`
}

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

type ReadOnlyDatabase struct {
	dal.DB
	// contains filtered or unexported fields
}

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 (*ReadOnlyDatabase) ListReferrers added in v0.3.0

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

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"`
}

func Compile

func Compile(profile SourceProfile, query dal.Query) (ReadPlan, error)

type RebindResult

type RebindResult struct {
	Receipt Receipt `json:"receipt"`
	Cursor  string  `json:"cursor"`
}

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
}

func (*Run) Close

func (r *Run) Close() error

func (*Run) Cursor

func (r *Run) Cursor() (string, error)

func (*Run) Next

func (r *Run) Next() (recordset.Row, recordset.Recordset, error)

func (*Run) NextPage

func (r *Run) NextPage() (Page, error)

func (*Run) Receipt

func (r *Run) Receipt() Receipt

func (*Run) Recordset

func (r *Run) Recordset() recordset.Recordset

Recordset exposes only previously authorized delivery, detached from private cells.

func (*Run) Schema

func (r *Run) Schema() []Field

type SnapshotResult added in v0.1.1

type SnapshotResult struct {
	Receipt Receipt `json:"receipt"`
	Cursor  string  `json:"cursor"`
}

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"`
}

type TableRef

type TableRef struct {
	ProjectID string `json:"projectId"`
	DatasetID string `json:"datasetId"`
	TableID   string `json:"tableId"`
}

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.

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL