bigquery

package module
v0.1.1 Latest Latest
Warning

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

Go to latest
Published: Oct 5, 2026 License: MIT Imports: 32 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.89.6, record v0.1.4 and Google API SDK v0.296.0 are pinned without replacements. This candidate is awaiting independent Task2 correction review; release, consumer integration and actual live acceptance remain required.

Operator integration

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)

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.

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 root interpretation remains subject to fresh independent review; 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.

Task1 corpus preparation and Task2 driver implementation advance here. 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 driver release, actual cloud query/spend, source admission, consumer cutover or dual-runtime live success is claimed. Root owns 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 MaxDepth = 32
View Source
const MaxResponseBytes = 10 * 1024 * 1024

Variables

This section is empty.

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