Documentation
¶
Overview ¶
Package bigquery implements bounded, lossless and explicitly approved analytical BigQuery reads through DALgo recordsets.
Index ¶
- Constants
- 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) 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 Error
- type Execution
- type Executor
- type Field
- type FileLedger
- type Identity
- type JobRef
- type JobStatus
- type Ledger
- type MemoryLedger
- type Observation
- type Order
- type Page
- type Parameter
- type Predicate
- type Prepare
- type Preview
- type Principal
- type Provider
- type ReadPlan
- type RebindResult
- type Receipt
- type Run
- type SnapshotResult
- type SourceProfile
- type TableRef
Constants ¶
const MaxDepth = 32
const MaxResponseBytes = 10 * 1024 * 1024
Variables ¶
This section is empty.
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) 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 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 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 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 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
|
|
|
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. |