scopedb

package module
v0.6.2 Latest Latest
Warning

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

Go to latest
Published: Aug 23, 2026 License: Apache-2.0 Imports: 20 Imported by: 0

README

ScopeDB SDK for Go

Apache License, Version 2.0 Go Reference

The ScopeDB Go SDK supports ScopeQL statements, REST catalog discovery, and bounded asynchronous streaming writes. For application writes, start with Table.AppendStream.

Runtime and installation

The SDK requires Go 1.24 or later.

go get github.com/scopedb/goscopedb@latest

Create a client

Pass the ScopeDB endpoint and API key through application configuration. Keep API keys out of source control.

client, err := scopedb.NewClient(scopedb.Config{
	Endpoint: os.Getenv("SCOPEDB_ENDPOINT"),
	APIKey:   os.Getenv("SCOPEDB_API_KEY"),
})
if err != nil {
	return err
}
defer client.Close()

NewClient validates the endpoint and returns an error for invalid configuration. Set Config.HTTPClient when the application needs to own HTTP timeouts, proxies, TLS, or connection pooling. Client.Close closes idle connections only for the HTTP client created by the SDK; it never closes a caller-provided client.

Statement, transform-ingest, and AppendStream request bodies use zstd compression by default. Set Config.Compression to CompressionGzip when gzip is required. Direct caller-encoded AppendNDJSON requests remain identity-encoded.

ScopeQL documentation

The SDK executes ScopeQL but does not define the language. Start with the language documentation:

Query and results

Query submits a statement and waits for its result:

result, err := client.Query(ctx, "SELECT 1 AS ready")
if err != nil {
	return err
}

rows, err := result.ToObjects()
if err != nil {
	return err
}
fmt.Println(rows)

Use RawRows for unconverted wire values, ToValues for positional values, ToObjects for values keyed by column name, or First for an optional first row.

For a detached or long-running statement, keep its handle and choose between a local status snapshot, one remote status request, or waiting for the result:

handle, err := client.Statement("SELECT 1 AS ready").Submit(ctx)
if err != nil {
	return err
}

fmt.Println("statement ID:", handle.ID())
if cached := handle.LastStatus(); cached != nil {
	fmt.Println("cached status:", *cached) // No network request.
}

latest, err := handle.Status(ctx) // Fetches one remote snapshot while active.
if err != nil {
	return err
}
fmt.Println("latest status:", latest)

result, err = handle.Wait(ctx) // Polls until the statement terminates.
if err != nil {
	return err
}

Once a handle has a terminal status, Status returns the cached status without another request. Store handle.ID() and use client.StatementHandle(id) to resume the lifecycle in another process. Cancel returns the statement ID, creation time, status, and server message. If cancellation finds that a statement already finished or failed, Wait fetches the complete statement response needed for its result or structured failure details. A cancelled outcome uses the cancellation message directly.

Statement.ID and Statement.ExecTimeout are the only optional statement settings. Provide an ID when the application needs to choose the statement ID; otherwise ScopeDB generates one. StatementHandle.ID() always returns the ID confirmed by ScopeDB:

statement := client.Statement("FROM events")
statementID := uuid.New()
statement.ID = &statementID // Optional.
statement.ExecTimeout = "30s"
handle, err := statement.Submit(ctx)
if err != nil {
	return err
}
fmt.Println("statement ID:", handle.ID())

When Wait or Execute returns a *scopedb.Error with kind ErrorKindStatementFailed, StatementDetails preserves the server's structured error code, message, and code-specific JSON details. The outer error message remains the server's top-level statement message:

var scopeErr *scopedb.Error
if errors.As(err, &scopeErr) && scopeErr.StatementDetails != nil {
	fmt.Println("statement error code:", scopeErr.StatementDetails.Code)
	fmt.Println("statement error:", scopeErr.StatementDetails.Message)
	fmt.Println("details:", string(scopeErr.StatementDetails.Details))
}

Browse the REST catalog

List methods expose one explicit page. Iterators lazily request later pages and are the simpler choice for discovery:

for database, err := range client.IterateDatabases(ctx, scopedb.CatalogListOptions{
	PageSize: 100,
}) {
	if err != nil {
		return err
	}
	fmt.Println(database.Name)
}

for table, err := range client.IterateTables(
	ctx,
	"scopedb",
	"public",
	scopedb.CatalogListOptions{PageSize: 100},
) {
	if err != nil {
		return err
	}
	fmt.Println(table.Name)
}

Use ListDatabases, ListSchemas, or ListTables when the application owns page boundaries. Use FetchDatabase, FetchSchema, or FetchTable for a full resource.

Describe a table

The table helper defaults to database scopedb and schema public. Set both explicitly when the destination is application-configured:

table := client.Table("events")
table.Database = "scopedb"
table.Schema = "public"

description, err := table.Describe(ctx)
if err != nil {
	return err
}
fmt.Println(description.Columns)

Streaming writes

Table appends write rows to an existing destination table. For most applications, AppendStream is the recommended path: the SDK accepts typed rows and owns their encoding, bounded batching, backpressure, and request concurrency. Its current wire encoding is NDJSON, but callers do not construct the wire payload. Evaluate the examples against an explicitly selected disposable table before using a production destination.

Use AppendStream for normal application writes, including continuous and large producers. Send accepts typed rows and uses encoding/json to encode each value as one top-level JSON object. Standard JSON tags and custom MarshalJSON methods apply. The stream batches those objects by size or time, bounds pending bytes, and sends a bounded number of append requests concurrently. The zero-value options use bounded defaults; override them only when the workload needs a different delivery policy or resource bound.

Each AppendStream request contains at most 8 MiB of uncompressed NDJSON and 200,000 rows. The stream splits automatically at either limit.

type Event struct {
	ID   int    `json:"id"`
	Name string `json:"name"`
}

stream, err := table.AppendStream(scopedb.AppendStreamOptions{})
if err != nil {
	return err
}

for _, event := range []Event{
	{ID: 1, Name: "first"},
	{ID: 2, Name: "second"},
} {
	if err := stream.Send(ctx, event); err != nil {
		_, _ = stream.Shutdown(ctx)
		return err
	}
}

report, err := stream.Flush(ctx)
if err != nil {
	_, _ = stream.Shutdown(ctx)
	return err
}
fmt.Println("flush committed:", report.CommittedRows)

// Permanently closes admission and settles all remaining accepted rows.
_, err = stream.Shutdown(ctx)
return err

Send waits for bounded local admission capacity. A nil error means only that the row entered the local stream; it does not confirm a remote commit. Feed large sources one row at a time instead of starting one goroutine per row, which would move an unbounded backlog outside the stream.

JSON serialization validates only that each value encodes as an object. ScopeDB validates that object's fields and types against the destination table when it processes the batch. With the default stop policy, Flush or Shutdown returns the server error and structured row details. Continue mode reports failed rows through the barrier report and Stats().LastFailure. An earlier successful Send does not imply schema compatibility.

Send and TrySend are safe for concurrent producers. When source-side work benefits from parallelism, use a fixed worker pool; the append_stream example uses four producers. Remote batch commits are still unordered when MaxConcurrentBatches is greater than one.

Flush settles every row accepted before its barrier. Shutdown permanently closes admission and settles all accepted rows. Canceling either call's context stops that caller's wait after an enqueued barrier, but does not cancel remote settlement. Inspect Stats().LastReport if the wait is interrupted.

The default AppendFailureStop policy is strict: the first failed batch stops admission, and a successful barrier confirms that its accepted prefix committed. Concurrent batches have no defined commit order; set MaxConcurrentBatches: 1 when request submission must be serial.

Best-effort logs and telemetry

Logs and telemetry often cannot block a request path or stop forever after one remote failure. Opt into AppendFailureContinue, use TrySend, and inspect the settlement report and lifetime statistics:

telemetry, err := table.AppendStream(scopedb.AppendStreamOptions{
	FailurePolicy: scopedb.AppendFailureContinue,
	FlushInterval: time.Second,
})
if err != nil {
	return err
}

if err := telemetry.TrySend(map[string]any{
	"name":   "request.completed",
	"status": 200,
}); err != nil {
	// Send this diagnostic to a different sink.
	log.Printf("telemetry row dropped locally: %v", err)
}

report, err := telemetry.Shutdown(ctx)
if err != nil {
	return err
}
if report.Outcome != scopedb.AppendDeliveryOK {
	log.Printf("telemetry loss or ambiguity: %+v", report)
}
fmt.Printf("lifetime stats: %+v\n", telemetry.Stats())

TrySend does not wait for stream capacity. A nil error still means local admission only; an error can indicate invalid input, an oversized row, a full buffer, or a closed stream. Stats().DroppedByReason separates local loss causes.

Continue mode accounts for a failed batch and continues with later rows. A completed report separates committed, failed, unknown, and locally dropped rows. Stats().LastFailure preserves the latest HTTP status, request ID, retry metadata, and structured row errors for diagnostics. It is a settlement report, not a commit receipt for every row. The stream retries only an exact temporary batch that the server explicitly marks rejected. A timeout, transport failure, or malformed success response is unknown and is never automatically retried. Rows with an unknown outcome may already exist remotely, so never blindly replay them.

An in-memory stream is not a durable queue. Use an application-owned outbox and a reconciliation path when payloads must survive process failure or unknown outcomes.

Low-level: direct NDJSON append

Use AppendNDJSON only when the caller already owns one exact raw NDJSON body and its request boundary. The body contains one JSON object per non-empty line, not a JSON array:

ndjson := []byte("{\"id\":1,\"name\":\"first\"}\n{\"id\":2,\"name\":\"second\"}")
result, err := table.AppendNDJSON(ctx, ndjson)
if err != nil {
	return err
}
fmt.Println("committed rows:", result.NumRowsInserted)

One request is limited to 16 MiB and 200,000 rows.

Choose a delivery path
Workload Admission and delivery Example
Normal typed application writes SDK owns encoding and batches; strict barriers append_stream
Backfill or file import Sequential producer admission and bounded concurrent batches bulk_append
Long-running logs and events Non-blocking continue mode with observable loss telemetry
One exact raw NDJSON payload Caller encodes the body and owns the request boundary append_ndjson

Advanced: transform before writing

Use Client.IngestStream only when source JSON specifically needs a server-side ScopeQL transformation before it can match the destination table. For normal typed events, shape the row in the producer and use Table.AppendStream. See the guarded ingest_transform example for the advanced path.

This path is sequential and fail-fast. IngestStream.Send confirms local admission only, while Flush and Shutdown wait for the accepted prefix to settle when they succeed. If a remote ingest request returns an error, its commit outcome may be unknown. A nonzero result returned with that error counts only earlier confirmed batches; it is not a safe replay offset. Reconcile the failing batch before replaying records.

Structured errors

Server error messages pass through unchanged. scopedb.Error adds structured diagnostics without requiring applications to parse the message:

var scopeErr *scopedb.Error
if errors.As(err, &scopeErr) {
	log.Printf("kind=%s status=%d request_id=%s retryable=%t retry_after=%s",
		scopeErr.Kind,
		scopeErr.HTTPStatus,
		scopeErr.RequestID,
		scopeErr.Retryable,
		scopeErr.RetryAfter,
	)

	if details := scopeErr.AppendDetails; details != nil &&
		details.AppendState == scopedb.AppendStateUnknown {
		log.Print("append may have committed; reconcile before replaying")
	}
}

The main kinds are ErrorKindConfigInvalid, ErrorKindStatementFailed, ErrorKindAppendRowsFailed, and ErrorKindUnexpected. Transport and decoding causes support errors.Is and errors.As through Unwrap. A direct append context canceled before its request starts is returned directly.

Examples and development

The examples guide contains read-only discovery, guarded write examples, delivery contracts, and runnable commands.

go test ./...
go test -race ./...
go vet ./...

Release notes and the maintainer runbook are in CHANGELOG.md and RELEASE.md. See CONTRIBUTING.md to contribute.

License

This software is licensed under the Apache License, Version 2.0.

Documentation

Overview

Package scopedb provides a Go client for ScopeQL statements, REST catalog discovery, and streaming writes.

Create and close one Client for the application's connection pool:

client, err := scopedb.NewClient(scopedb.Config{
	Endpoint: os.Getenv("SCOPEDB_ENDPOINT"),
	APIKey:   os.Getenv("SCOPEDB_API_KEY"),
})
if err != nil {
	return err
}
defer client.Close()

Query waits for a statement result, which can be converted to keyed objects:

result, err := client.Query(ctx, "SELECT 1 AS ready")
if err != nil {
	return err
}
rows, err := result.ToObjects()
if err != nil {
	return err
}
fmt.Println(rows)

Use Statement.Submit when an application needs the statement ID, a local status snapshot, an explicit remote status request, or a separate wait. ID and ExecTimeout are the only optional statement settings; omit ID to let ScopeDB generate it, and read the confirmed ID from StatementHandle.ID. Failed statements expose structured server error details on Error.StatementDetails. Catalog iterators lazily traverse REST pages. Table.Describe returns table metadata. For application writes, Table.AppendStream is the recommended path: it accepts typed rows and adds bounded asynchronous batching that is safe for concurrent producers. Its request batches use the client's configured compression, which defaults to zstd. Table.AppendNDJSON is the low-level path for one caller-owned, identity-encoded NDJSON request.

AppendStream.Send uses encoding/json and accepts typed structs and other values that encode as a top-level JSON object. Send does not validate the destination table schema. ScopeDB validation failures surface when Flush or Shutdown settles: the stop policy returns the error, while continue mode reports failed rows in the delivery report and failure details in Stats().LastFailure.

Stream Send methods confirm local admission only. Flush and Shutdown wait for the accepted prefix. AppendStream reports rejected and unknown outcomes.

Client.IngestStream is an advanced path for source JSON that specifically needs a server-side ScopeQL transformation before it matches the destination table. An IngestStream error can follow a remote commit, so callers must reconcile before replaying the same records. A nonzero ingest result returned with an error covers only earlier confirmed batches, not a safe replay offset.

ScopeQL is documented separately in the quickstart, query guide, and language reference.

Index

Constants

This section is empty.

Variables

View Source
var (
	// ErrAppendStreamFull means TrySend could not reserve bounded local capacity.
	ErrAppendStreamFull = errors.New("append stream buffer is full")
	// ErrAppendStreamClosed means the stream no longer accepts rows.
	ErrAppendStreamClosed = errors.New("append stream is closed")
	// ErrAppendRowInvalid means a row did not encode to one JSON object.
	ErrAppendRowInvalid = errors.New("append row must encode to a JSON object")
	// ErrAppendRowTooLarge means one encoded row exceeds the append stream request limit.
	ErrAppendRowTooLarge = errors.New("append row exceeds the protocol limit")
)

Functions

This section is empty.

Types

type AppendDeliveryOutcome

type AppendDeliveryOutcome uint8

AppendDeliveryOutcome summarizes one completed delivery barrier.

const (
	// AppendDeliveryOK means the interval settled without failed, unknown, or dropped rows.
	AppendDeliveryOK AppendDeliveryOutcome = iota
	// AppendDeliveryPartial means some rows committed and some did not settle cleanly.
	AppendDeliveryPartial
	// AppendDeliveryFailed means no rows committed and at least one row failed or dropped.
	AppendDeliveryFailed
	// AppendDeliveryUnknown means no rows committed and at least one row may have committed.
	AppendDeliveryUnknown
)

func (AppendDeliveryOutcome) String

func (o AppendDeliveryOutcome) String() string

type AppendDeliveryReport

type AppendDeliveryReport struct {
	// Outcome summarizes delivery for this barrier interval.
	Outcome AppendDeliveryOutcome
	// AcceptedRows is the number of rows locally admitted in the interval.
	AcceptedRows uint64
	// CommittedRows is the number of rows confirmed committed.
	CommittedRows uint64
	// FailedRows is the number of rejected or unsent accepted rows.
	FailedRows uint64
	// UnknownRows is the number of rows whose commit outcome is unknown.
	UnknownRows uint64
	// DroppedRows is the number of TrySend calls rejected locally in the interval.
	DroppedRows      uint64
	CommittedBatches uint64
	FailedBatches    uint64
	UnknownBatches   uint64
	Retries          uint64
	Duration         time.Duration
}

AppendDeliveryReport contains settlement counters since the previous barrier.

type AppendDroppedRows

type AppendDroppedRows struct {
	// BufferFull counts rows rejected because immediate local admission was unavailable.
	BufferFull uint64
	// InvalidRow counts values that did not encode to a JSON object.
	InvalidRow uint64
	// RowTooLarge counts values over the single-request protocol limit.
	RowTooLarge uint64
	// Closed counts attempts after shutdown or a terminal failure.
	Closed uint64
}

AppendDroppedRows contains lifetime TrySend failures by reason.

type AppendErrorDetails

type AppendErrorDetails struct {
	AppendState        AppendState      `json:"append_state"`
	RowErrors          []AppendRowError `json:"row_errors"`
	RowErrorsTruncated bool             `json:"row_errors_truncated"`
}

AppendErrorDetails contains the structured outcome of a failed append.

type AppendFailurePolicy

type AppendFailurePolicy uint8

AppendFailurePolicy controls what a stream does after a remote batch fails.

const (
	// AppendFailureStop stops admission after the first failed batch.
	AppendFailureStop AppendFailurePolicy = iota
	// AppendFailureContinue accounts for the failed batch and processes later rows.
	AppendFailureContinue
)

func (AppendFailurePolicy) String

func (p AppendFailurePolicy) String() string

type AppendLastFailure

type AppendLastFailure struct {
	// At is when the client observed the failure.
	At time.Time
	// Message is the original error message.
	Message string
	// AppendState is rejected or unknown.
	AppendState AppendState
	// HTTPStatus is zero when no HTTP response was available.
	HTTPStatus int
	// RequestID is the request identifier reported by the service, when present.
	RequestID string
	// RetryAfter is the service-provided retry delay, when present.
	RetryAfter time.Duration
	// Retryable reports the service's retry classification. AppendStream retries
	// only exact rejected batches, regardless of this field alone.
	Retryable bool
	// RowErrors contains structured validation failures reported for the batch.
	RowErrors []AppendRowError
	// RowErrorsTruncated reports whether the service omitted additional row errors.
	RowErrorsTruncated bool
}

AppendLastFailure describes the last remotely failed batch.

type AppendRowError

type AppendRowError struct {
	RowIndex uint64 `json:"row_index"`
	Column   string `json:"column"`
	Message  string `json:"message"`
}

AppendRowError describes a validation error for one row in an append request.

type AppendRowsResult

type AppendRowsResult struct {
	AppendState     AppendState `json:"append_state"`
	NumRowsInserted int64       `json:"num_rows_inserted"`
}

AppendRowsResult is the result of a committed table append.

type AppendState

type AppendState string

AppendState describes the commit outcome of a table append request.

const (
	// AppendStateCommitted means all reported rows committed.
	AppendStateCommitted AppendState = "committed"
	// AppendStateRejected means no rows committed and the request is safe to correct and retry.
	AppendStateRejected AppendState = "rejected"
	// AppendStateUnknown means the client cannot determine whether rows committed.
	AppendStateUnknown AppendState = "unknown"
)

type AppendStream

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

AppendStream asynchronously batches rows into bounded NDJSON append requests. Send and TrySend are safe to call from concurrent producer goroutines.

func (*AppendStream) Flush

Flush dispatches all rows admitted before this barrier and waits for settlement. Canceling ctx stops only this wait; an already enqueued barrier keeps running. If ctx ends first, Flush returns a zero report and ctx.Err().

func (*AppendStream) Send

func (s *AppendStream) Send(ctx context.Context, row any) error

Send serializes row with encoding/json and admits it, waiting for bounded local capacity. Row may be a typed struct or any other value that encodes as one top-level JSON object. Success does not validate the destination table schema or confirm a remote commit.

func (*AppendStream) Shutdown

Shutdown permanently closes admission and settles every admitted row. It is idempotent; canceling ctx stops only the caller's wait. If ctx ends first, Shutdown returns a zero report and ctx.Err().

func (*AppendStream) Stats

func (s *AppendStream) Stats() AppendStreamStats

Stats returns a race-safe lifetime snapshot.

func (*AppendStream) TrySend

func (s *AppendStream) TrySend(row any) error

TrySend serializes row with encoding/json and attempts local admission without waiting for capacity. Row must encode as one top-level JSON object. Success does not validate the destination table schema or confirm a remote commit.

type AppendStreamOptions

type AppendStreamOptions struct {
	// FailurePolicy defaults to AppendFailureStop.
	FailurePolicy AppendFailurePolicy
	// TargetBatchBytes defaults to 8 MiB and cannot exceed the request limit.
	TargetBatchBytes int
	// MaxBatchRows defaults to 200,000 and cannot exceed the request limit.
	MaxBatchRows int
	// FlushInterval defaults to one second.
	FlushInterval time.Duration
	// MaxBufferedBytes bounds admitted rows not yet settled. It defaults to 64 MiB.
	MaxBufferedBytes int
	// MaxConcurrentBatches defaults to four, cannot exceed 1024, and may be set
	// to one for serial request submission.
	MaxConcurrentBatches int
	// AttemptTimeout bounds each HTTP attempt. It defaults to 30 seconds.
	// A timeout leaves that batch's commit outcome unknown.
	AttemptTimeout time.Duration
}

AppendStreamOptions configures bounded asynchronous table appends. Zero-valued fields use SDK defaults.

type AppendStreamState

type AppendStreamState uint8

AppendStreamState is the lifecycle state of an AppendStream.

const (
	// AppendStreamOpen means the stream accepts rows.
	AppendStreamOpen AppendStreamState = iota
	// AppendStreamClosing means shutdown has closed admission and is settling rows.
	AppendStreamClosing
	// AppendStreamClosed means shutdown settled all accepted rows without a fatal error.
	AppendStreamClosed
	// AppendStreamFailed means strict delivery stopped after a batch failure.
	AppendStreamFailed
)

func (AppendStreamState) String

func (s AppendStreamState) String() string

type AppendStreamStats

type AppendStreamStats struct {
	// State is the current stream lifecycle state.
	State           AppendStreamState
	AcceptedRows    uint64
	CommittedRows   uint64
	FailedRows      uint64
	UnknownRows     uint64
	DroppedRows     uint64
	DroppedByReason AppendDroppedRows
	Retries         uint64
	PendingRows     uint64
	PendingBytes    int
	InFlightBatches int
	// LastFailure is nil until a remote batch fails.
	LastFailure *AppendLastFailure
	// LastReport is nil until a delivery barrier completes.
	LastReport *AppendDeliveryReport
}

AppendStreamStats is a race-safe lifetime snapshot.

type CatalogListOptions

type CatalogListOptions struct {
	// PageSize is the maximum number of resources returned in one page.
	// Zero uses the server default; otherwise the value must be from 1 to 1000.
	PageSize int
	// PageToken is the opaque continuation token returned by the previous page.
	PageToken string
}

CatalogListOptions configures one catalog list request.

type CatalogPage

type CatalogPage[T any] struct {
	Items         []T    `json:"items"`
	NextPageToken string `json:"next_page_token,omitempty"`
}

CatalogPage is one page of catalog resources.

type Client

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

Client provides access to ScopeDB APIs.

func NewClient

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

NewClient creates a new ScopeDB client with the given configuration.

func (*Client) Close

func (c *Client) Close()

Close releases idle connections owned by the client. It does not close a caller-provided HTTP client.

func (*Client) FetchDatabase

func (c *Client) FetchDatabase(ctx context.Context, database string) (DatabaseResource, error)

FetchDatabase returns one database resource.

func (*Client) FetchSchema

func (c *Client) FetchSchema(
	ctx context.Context,
	database string,
	schema string,
) (SchemaResource, error)

FetchSchema returns one schema resource.

func (*Client) FetchTable

func (c *Client) FetchTable(
	ctx context.Context,
	database string,
	schema string,
	table string,
) (TableResource, error)

FetchTable returns one table resource with its complete specification.

func (*Client) IngestStream

func (c *Client) IngestStream(
	statement string,
	options IngestStreamOptions,
) (*IngestStream, error)

IngestStream creates a bounded transform-ingest stream and starts its worker.

func (*Client) IterateDatabases

func (c *Client) IterateDatabases(
	ctx context.Context,
	options CatalogListOptions,
) iter.Seq2[DatabaseResource, error]

IterateDatabases lazily iterates databases across all catalog pages.

func (*Client) IterateSchemas

func (c *Client) IterateSchemas(
	ctx context.Context,
	database string,
	options CatalogListOptions,
) iter.Seq2[SchemaResource, error]

IterateSchemas lazily iterates schemas across all catalog pages.

func (*Client) IterateTables

func (c *Client) IterateTables(
	ctx context.Context,
	database string,
	schema string,
	options CatalogListOptions,
) iter.Seq2[TableResourceSummary, error]

IterateTables lazily iterates tables across all catalog pages.

func (*Client) ListDatabases

func (c *Client) ListDatabases(
	ctx context.Context,
	options CatalogListOptions,
) (CatalogPage[DatabaseResource], error)

ListDatabases returns one page of databases.

func (*Client) ListSchemas

func (c *Client) ListSchemas(
	ctx context.Context,
	database string,
	options CatalogListOptions,
) (CatalogPage[SchemaResource], error)

ListSchemas returns one page of schemas in a database.

func (*Client) ListTables

func (c *Client) ListTables(
	ctx context.Context,
	database string,
	schema string,
	options CatalogListOptions,
) (CatalogPage[TableResourceSummary], error)

ListTables returns one page of tables in a schema.

func (*Client) Query

func (c *Client) Query(ctx context.Context, scopeql string) (*ResultSet, error)

Query executes a ScopeQL statement and waits for all result rows.

func (*Client) Statement

func (c *Client) Statement(stmt string) *Statement

Statement creates a new statement with the given ScopeQL statement.

func (*Client) StatementHandle

func (c *Client) StatementHandle(id uuid.UUID) *StatementHandle

StatementHandle creates a new StatementHandle with the given ID.

func (*Client) Table

func (c *Client) Table(tableName string) *Table

Table creates a new Table object with the given name.

type Compression

type Compression string

Compression defines the compression used for statement, transform-ingest, and AppendStream request bodies. Direct table appends send identity-encoded NDJSON.

const (
	// CompressionZstd uses Zstandard compression.
	CompressionZstd Compression = "zstd"
	// CompressionGzip uses gzip compression.
	CompressionGzip Compression = "gzip"
)

type Config

type Config struct {
	// Endpoint is the URL of the ScopeDB service.
	Endpoint string `json:"endpoint"`
	// APIKey is the API key used for authentication.
	//
	// When provided, the client sends it as the Authorization header using the
	// Bearer scheme.
	APIKey string `json:"api_key"`
	// Compression controls how statement, transform-ingest, and AppendStream
	// request bodies are compressed. It does not affect direct table appends.
	//
	// The default is CompressionZstd. Set this to CompressionGzip to talk to
	// older deployments that do not support zstd yet.
	Compression Compression `json:"compression"`
	// HTTPClient is the HTTP client used to send requests.
	//
	// When nil, the SDK creates and owns an independent HTTP client. A supplied
	// client remains owned by the caller and is never closed by the SDK.
	HTTPClient *http.Client `json:"-"`
}

Config defines the configuration for the client.

type DataType

type DataType string

DataType is the type of field.

const (
	// StringDataType indicates the data is of string data type.
	StringDataType DataType = "string"
	// BinaryDataType indicates the data is of binary data type.
	BinaryDataType DataType = "binary"
	// IntDataType indicates the data is of int data type.
	IntDataType DataType = "int"
	// UIntDataType indicates the data is of uint data type.
	UIntDataType DataType = "uint"
	// FloatDataType indicates the data is of float data type.
	FloatDataType DataType = "float"
	// BooleanDataType indicates the data is of bool data type.
	BooleanDataType DataType = "boolean"
	// TimestampDataType indicates the data is of timestamp data type.
	TimestampDataType DataType = "timestamp"
	// IntervalDataType indicates the data is of interval data type.
	IntervalDataType DataType = "interval"
	// ArrayDataType indicates the data is of array data type.
	ArrayDataType DataType = "array"
	// ObjectDataType indicates the data is of object data type.
	ObjectDataType DataType = "object"
	// AnyDataType indicates the data is of any data type.
	AnyDataType DataType = "any"
	// NullDataType indicates the data is of null data type.
	NullDataType DataType = "null"
)

func (*DataType) UnmarshalJSON

func (d *DataType) UnmarshalJSON(data []byte) error

UnmarshalJSON accepts the canonical wire names and normalizes the historical unsigned-integer alias to UIntDataType.

type DatabaseResource

type DatabaseResource struct {
	Name    string  `json:"name"`
	Comment *string `json:"comment"`
}

DatabaseResource describes a database in the catalog.

type Error

type Error struct {
	Kind          ErrorKind
	Message       string
	HTTPStatus    int
	RequestID     string
	RetryAfter    time.Duration
	Retryable     bool
	AppendDetails *AppendErrorDetails
	// StatementDetails contains the structured server failure for a failed
	// statement, when available.
	StatementDetails *StatementErrorDetails
	// contains filtered or unexported fields
}

Error represents a ScopeDB client or API error.

func (*Error) Error

func (e *Error) Error() string

Error returns the server message unchanged when the error came from ScopeDB.

func (*Error) Unwrap

func (e *Error) Unwrap() error

Unwrap returns the transport, decoding, or configuration error that caused this error.

type ErrorKind

type ErrorKind string

ErrorKind identifies the operation-level category of a ScopeDB error.

const (
	// ErrorKindUnexpected is an error that has no more specific classification.
	ErrorKindUnexpected ErrorKind = "Unexpected"
	// ErrorKindConfigInvalid indicates invalid client configuration or arguments.
	ErrorKindConfigInvalid ErrorKind = "ConfigInvalid"
	// ErrorKindStatementFailed indicates a failed or cancelled statement.
	ErrorKindStatementFailed ErrorKind = "StatementFailed"
	// ErrorKindAppendRowsFailed indicates a rejected append or an unknown commit outcome.
	ErrorKindAppendRowsFailed ErrorKind = "AppendRowsFailed"
)

type FieldSchema

type FieldSchema struct {
	// Name is the field name.
	Name string
	// Type is the field data type.
	Type DataType
}

FieldSchema describes a single field.

type IngestResult

type IngestResult struct {
	NumRowsInserted int64
}

IngestResult reports rows confirmed inserted by requests covered by a delivery barrier. When returned with an error, it counts only earlier confirmed batches; the failing batch can have an unknown commit outcome.

type IngestStream

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

IngestStream serializes JSON objects into bounded, sequential transform-ingest requests. It is safe for concurrent producers.

func (*IngestStream) Flush

func (s *IngestStream) Flush(ctx context.Context) (IngestResult, error)

Flush dispatches every row admitted before this barrier and waits for their remote outcome. Cancelling ctx stops only this caller's wait after the barrier has been enqueued; the worker continues settling the accepted prefix. A nonzero result returned with an error excludes the failing batch and is not a safe offset for replay.

func (*IngestStream) Send

func (s *IngestStream) Send(ctx context.Context, row any) error

Send serializes and admits one JSON object, waiting for bounded local capacity. Success confirms local admission only; use Flush or Shutdown as a remote commit barrier.

func (*IngestStream) Shutdown

func (s *IngestStream) Shutdown(ctx context.Context) (IngestResult, error)

Shutdown closes admission, settles every accepted row, and returns the final barrier result. It is idempotent. Cancelling ctx stops only the caller's wait; shutdown continues in the background and a later call can await it. A nonzero result returned with an error excludes the failing batch and is not a safe offset for replay.

type IngestStreamOptions

type IngestStreamOptions struct {
	// TargetBatchBytes is the target NDJSON body size. A single row may exceed
	// this target up to the 16 MiB request limit.
	TargetBatchBytes int
	// MaxBatchRows is the maximum number of rows in one ingest request.
	MaxBatchRows int
	// FlushInterval is the maximum delay from the first buffered row until its
	// batch is dispatched.
	FlushInterval time.Duration
	// MaxBufferedBytes bounds serialized rows that have been admitted but have
	// not finished their remote request.
	MaxBufferedBytes int
	// AttemptTimeout bounds each transform-ingest request. It defaults to 30
	// seconds. A timeout can leave the remote commit outcome unknown.
	AttemptTimeout time.Duration
}

IngestStreamOptions configures transform-oriented JSON ingest batching. Zero values use the documented defaults.

type ResultSet

type ResultSet struct {
	// TotalRows is the total number of rows in the result set.
	TotalRows uint64
	// Schema is the schema of the result set.
	Schema Schema
	// contains filtered or unexported fields
}

ResultSet stores the result of a statement execution.

func (*ResultSet) First

func (rs *ResultSet) First() (map[string]Value, bool, error)

First returns the first row keyed by column name. The boolean is false when the result set is empty.

func (*ResultSet) RawRows

func (rs *ResultSet) RawRows() ([][]*string, error)

RawRows returns the string-or-null cells from the JSON wire response without converting them to ScopeDB value types.

func (*ResultSet) ToObjects

func (rs *ResultSet) ToObjects() ([]map[string]Value, error)

ToObjects returns rows keyed by result column name.

Duplicate output column names are rejected because representing them as a map would silently discard values.

func (*ResultSet) ToValues

func (rs *ResultSet) ToValues() ([][]Value, error)

ToValues reads the result set and returns the rows as a 2D array of values, i.e., rows of value lists.

Binary cells are returned as []byte, timestamps as time.Time, and intervals as time.Duration. Array, object, and any cells remain JSON strings, while null cells are returned as nil.

This method is only valid if the result set is of the JSON format.

type Schema

type Schema []*FieldSchema

Schema describes the fields in a table or query result.

type SchemaResource

type SchemaResource struct {
	Database string  `json:"database"`
	Name     string  `json:"name"`
	Comment  *string `json:"comment"`
}

SchemaResource describes a schema in the catalog.

type Statement

type Statement struct {

	// ID is an optional caller-provided statement ID.
	//
	// When nil, ScopeDB generates the statement ID.
	ID *uuid.UUID
	// ExecTimeout is the maximum time allowed for statement execution.
	//
	// If the total execution time exceeds this value, the statement is failed
	// as timed out.
	//
	// Values use duration strings such as "1h".
	ExecTimeout string
	// contains filtered or unexported fields
}

Statement configures a ScopeQL statement before submission.

func (*Statement) Execute

func (s *Statement) Execute(ctx context.Context) (*ResultSet, error)

Execute submits the statement to ScopeDB for execution and waits for its completion.

func (*Statement) Submit

func (s *Statement) Submit(ctx context.Context) (*StatementHandle, error)

Submit submits the statement to ScopeDB for execution.

type StatementCancelResult

type StatementCancelResult struct {
	StatementID uuid.UUID       `json:"statement_id"`
	CreatedAt   time.Time       `json:"created_at"`
	Status      StatementStatus `json:"status"`
	Message     string          `json:"message"`
}

StatementCancelResult reports the server's complete cancellation outcome.

type StatementErrorCode

type StatementErrorCode string

StatementErrorCode identifies the reason a statement failed. Unknown values are preserved for forward compatibility.

const (
	// StatementErrorCodePrepareError indicates that statement preparation failed.
	StatementErrorCodePrepareError StatementErrorCode = "prepare_error"
	// StatementErrorCodeExecuteError indicates that statement execution failed.
	StatementErrorCodeExecuteError StatementErrorCode = "execute_error"
	// StatementErrorCodePendingTimeout indicates that the statement timed out before execution.
	StatementErrorCodePendingTimeout StatementErrorCode = "pending_timeout"
	// StatementErrorCodeExecutionTimeout indicates that the statement exceeded its execution timeout.
	StatementErrorCodeExecutionTimeout StatementErrorCode = "execution_timeout"
	// StatementErrorCodeHeartbeatLost indicates that the statement worker stopped reporting progress.
	StatementErrorCodeHeartbeatLost StatementErrorCode = "heartbeat_lost"
	// StatementErrorCodeRowLimitExceeded indicates that the statement exceeded a server-enforced row limit.
	StatementErrorCodeRowLimitExceeded StatementErrorCode = "row_limit_exceeded"
	// StatementErrorCodeScanLimitExceeded indicates that the statement exceeded a server-enforced scan limit.
	StatementErrorCodeScanLimitExceeded StatementErrorCode = "scan_limit_exceeded"
)

type StatementErrorDetails

type StatementErrorDetails struct {
	Code    StatementErrorCode `json:"code"`
	Message string             `json:"message"`
	Details json.RawMessage    `json:"details,omitempty"`
}

StatementErrorDetails contains the structured failure returned for a failed statement. Details is code-specific JSON and is nil when the server did not provide additional details.

type StatementHandle

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

StatementHandle is a handle to a statement that has been submitted to ScopeDB.

func (*StatementHandle) Cancel

Cancel cancels the statement if it is running or pending.

func (*StatementHandle) ID

func (h *StatementHandle) ID() uuid.UUID

ID returns the statement ID represented by this handle.

func (*StatementHandle) LastStatus

func (h *StatementHandle) LastStatus() *StatementStatus

LastStatus returns the latest locally cached status without making a request.

func (*StatementHandle) Progress

func (h *StatementHandle) Progress() *StatementProgress

Progress returns the last seen progress of the statement.

func (*StatementHandle) ResultSet

func (h *StatementHandle) ResultSet() *ResultSet

ResultSet returns the result set of the statement if available.

func (*StatementHandle) Status

Status fetches the latest status at most once. A cached terminal status is returned without making another request.

func (*StatementHandle) Wait

func (h *StatementHandle) Wait(ctx context.Context) (*ResultSet, error)

Wait polls until the statement is finished, failed, or cancelled.

When the statement is finished, the result set is returned. Otherwise, an error is returned.

type StatementProgress

type StatementProgress struct {
	// TotalPercentage denotes the total progress in percentage: [0.0, 100.0].
	TotalPercentage float64 `json:"total_percentage"`
	// NanosFromSubmitted denotes the duration in nanoseconds since the statement is submitted.
	NanosFromSubmitted int64 `json:"nanos_from_submitted"`
	// NanosFromStarted denotes the duration in nanoseconds since the statement is started.
	NanosFromStarted int64 `json:"nanos_from_started"`
	// TotalStages denotes the total number of stages to execute.
	TotalStages int64 `json:"total_stages"`
	// TotalPartitions denotes the estimated total number of partitions to scan.
	TotalPartitions int64 `json:"total_partitions"`
	// TotalRows denotes the estimated total number of rows to scan.
	TotalRows int64 `json:"total_rows"`
	// TotalCompressedBytes denotes the estimated total number of compressed bytes to scan.
	TotalCompressedBytes int64 `json:"total_compressed_bytes"`
	// TotalUncompressedBytes denotes the estimated total number of uncompressed bytes to scan.
	TotalUncompressedBytes int64 `json:"total_uncompressed_bytes"`
	// ScannedStages denotes the total number of stages executed.
	ScannedStages int64 `json:"scanned_stages"`
	// ScannedPartitions denotes the number of partitions scanned.
	ScannedPartitions int64 `json:"scanned_partitions"`
	// ScannedRows denotes the number of rows scanned.
	ScannedRows int64 `json:"scanned_rows"`
	// ScannedCompressedBytes denotes the number of compressed bytes scanned.
	ScannedCompressedBytes int64 `json:"scanned_compressed_bytes"`
	// ScannedUncompressedBytes denotes the number of uncompressed bytes scanned.
	ScannedUncompressedBytes int64 `json:"scanned_uncompressed_bytes"`
	// SkippedPartitions denotes the number of partitions skipped by pruning.
	SkippedPartitions int64 `json:"skipped_partitions"`
	// SkippedRows denotes the number of rows skipped by pruning.
	SkippedRows int64 `json:"skipped_rows"`
	// SkippedCompressedBytes denotes the number of compressed bytes skipped by pruning.
	SkippedCompressedBytes int64 `json:"skipped_compressed_bytes"`
	// SkippedUncompressedBytes denotes the number of uncompressed bytes skipped by pruning.
	SkippedUncompressedBytes int64 `json:"skipped_uncompressed_bytes"`
}

StatementProgress is a struct that represents the progress of a statement.

type StatementStatus

type StatementStatus string

StatementStatus is a string that represents the status of a statement.

const (
	// StatementStatusPending indicates the query is not started yet.
	StatementStatusPending StatementStatus = "pending"
	// StatementStatusRunning indicates the query is not finished yet.
	StatementStatusRunning StatementStatus = "running"
	// StatementStatusFinished indicates the query is finished.
	StatementStatusFinished StatementStatus = "finished"
	// StatementStatusFailed indicates the query is failed.
	StatementStatusFailed StatementStatus = "failed"
	// StatementStatusCancelled indicates the query is cancelled.
	StatementStatusCancelled StatementStatus = "cancelled"
)

func (StatementStatus) Finished

func (s StatementStatus) Finished() bool

Finished returns true if the statement is finished.

func (StatementStatus) Terminated

func (s StatementStatus) Terminated() bool

Terminated returns true if the statement is finished, failed, or cancelled.

type Table

type Table struct {

	// Database is the database name. Empty uses "scopedb" for REST APIs.
	Database string
	// Schema is the schema name. Empty uses "public" for REST APIs.
	Schema string
	// Name is the name of the table.
	Name string
	// contains filtered or unexported fields
}

Table references a ScopeDB table.

func (*Table) AppendNDJSON

func (t *Table) AppendNDJSON(ctx context.Context, ndjson []byte) (AppendRowsResult, error)

AppendNDJSON sends one caller-encoded NDJSON request to this table. The body contains one JSON object per non-empty line, not a JSON array.

func (*Table) AppendStream

func (t *Table) AppendStream(options AppendStreamOptions) (*AppendStream, error)

AppendStream creates a bounded asynchronous append stream for this table.

func (*Table) Describe

func (t *Table) Describe(ctx context.Context) (TableResource, error)

Describe returns the complete REST catalog resource for this table.

func (*Table) Drop

func (t *Table) Drop(ctx context.Context) error

Drop drops the table from ScopeDB.

This method issues a DROP TABLE statement to ScopeDB and blocks until done.

func (*Table) Identifier

func (t *Table) Identifier() string

Identifier returns the quoted table identifier.

type TableColumnSpec

type TableColumnSpec struct {
	Name     string   `json:"name"`
	DataType DataType `json:"data_type"`
	Comment  *string  `json:"comment"`
}

TableColumnSpec describes one table column.

type TableDistinctSpec

type TableDistinctSpec struct {
	On []string `json:"on"`
	By []string `json:"by"`
}

TableDistinctSpec describes the table's distinct-key configuration.

type TableResource

type TableResource struct {
	Database          string            `json:"database"`
	Schema            string            `json:"schema"`
	Name              string            `json:"name"`
	Columns           []TableColumnSpec `json:"columns"`
	PartitionBy       []string          `json:"partition_by"`
	ClusterBy         []string          `json:"cluster_by"`
	DistinctOn        TableDistinctSpec `json:"distinct_on"`
	DataRetentionDays *int32            `json:"data_retention_days"`
	Comment           *string           `json:"comment"`
}

TableResource describes a table and its complete catalog specification.

func (TableResource) Spec

func (resource TableResource) Spec() TableSpec

Spec returns the reusable table specification without its catalog identity.

type TableResourceSummary

type TableResourceSummary struct {
	Database string  `json:"database"`
	Schema   string  `json:"schema"`
	Name     string  `json:"name"`
	Comment  *string `json:"comment"`
}

TableResourceSummary describes a table without its full specification.

type TableSpec

type TableSpec struct {
	Columns           []TableColumnSpec `json:"columns"`
	PartitionBy       []string          `json:"partition_by"`
	ClusterBy         []string          `json:"cluster_by"`
	DistinctOn        TableDistinctSpec `json:"distinct_on"`
	DataRetentionDays *int32            `json:"data_retention_days"`
	Comment           *string           `json:"comment"`
}

TableSpec describes the reusable portion of a table definition.

type Value

type Value any

Value stores the contents of a single cell from a ScopeDB statement result.

Directories

Path Synopsis
examples
append_ndjson command
Package main demonstrates one caller-encoded raw NDJSON table append.
Package main demonstrates one caller-encoded raw NDJSON table append.
append_stream command
Package main demonstrates strict bounded asynchronous table appends.
Package main demonstrates strict bounded asynchronous table appends.
catalog command
Package main demonstrates lazy REST catalog discovery.
Package main demonstrates lazy REST catalog discovery.
ingest_transform command
Package main demonstrates the advanced IngestStream path for source JSON that needs a server-side ScopeQL transformation before it matches a table.
Package main demonstrates the advanced IngestStream path for source JSON that needs a server-side ScopeQL transformation before it matches a table.
internal/exampleutil
Package exampleutil contains configuration shared by the runnable examples.
Package exampleutil contains configuration shared by the runnable examples.
patterns/bulk_append command
Package main demonstrates bounded streaming writes for a bulk source.
Package main demonstrates bounded streaming writes for a bulk source.
patterns/telemetry command
Package main demonstrates observable best-effort telemetry delivery.
Package main demonstrates observable best-effort telemetry delivery.
statement command
Package main demonstrates synchronous and asynchronous statement queries.
Package main demonstrates synchronous and asynchronous statement queries.

Jump to

Keyboard shortcuts

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