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 ¶
- Variables
- type AppendDeliveryOutcome
- type AppendDeliveryReport
- type AppendDroppedRows
- type AppendErrorDetails
- type AppendFailurePolicy
- type AppendLastFailure
- type AppendRowError
- type AppendRowsResult
- type AppendState
- type AppendStream
- func (s *AppendStream) Flush(ctx context.Context) (AppendDeliveryReport, error)
- func (s *AppendStream) Send(ctx context.Context, row any) error
- func (s *AppendStream) Shutdown(ctx context.Context) (AppendDeliveryReport, error)
- func (s *AppendStream) Stats() AppendStreamStats
- func (s *AppendStream) TrySend(row any) error
- type AppendStreamOptions
- type AppendStreamState
- type AppendStreamStats
- type CatalogListOptions
- type CatalogPage
- type Client
- func (c *Client) Close()
- func (c *Client) FetchDatabase(ctx context.Context, database string) (DatabaseResource, error)
- func (c *Client) FetchSchema(ctx context.Context, database string, schema string) (SchemaResource, error)
- func (c *Client) FetchTable(ctx context.Context, database string, schema string, table string) (TableResource, error)
- func (c *Client) IngestStream(statement string, options IngestStreamOptions) (*IngestStream, error)
- func (c *Client) IterateDatabases(ctx context.Context, options CatalogListOptions) iter.Seq2[DatabaseResource, error]
- func (c *Client) IterateSchemas(ctx context.Context, database string, options CatalogListOptions) iter.Seq2[SchemaResource, error]
- func (c *Client) IterateTables(ctx context.Context, database string, schema string, ...) iter.Seq2[TableResourceSummary, error]
- func (c *Client) ListDatabases(ctx context.Context, options CatalogListOptions) (CatalogPage[DatabaseResource], error)
- func (c *Client) ListSchemas(ctx context.Context, database string, options CatalogListOptions) (CatalogPage[SchemaResource], error)
- func (c *Client) ListTables(ctx context.Context, database string, schema string, ...) (CatalogPage[TableResourceSummary], error)
- func (c *Client) Query(ctx context.Context, scopeql string) (*ResultSet, error)
- func (c *Client) Statement(stmt string) *Statement
- func (c *Client) StatementHandle(id uuid.UUID) *StatementHandle
- func (c *Client) Table(tableName string) *Table
- type Compression
- type Config
- type DataType
- type DatabaseResource
- type Error
- type ErrorKind
- type FieldSchema
- type IngestResult
- type IngestStream
- type IngestStreamOptions
- type ResultSet
- type Schema
- type SchemaResource
- type Statement
- type StatementCancelResult
- type StatementErrorCode
- type StatementErrorDetails
- type StatementHandle
- func (h *StatementHandle) Cancel(ctx context.Context) (StatementCancelResult, error)
- func (h *StatementHandle) ID() uuid.UUID
- func (h *StatementHandle) LastStatus() *StatementStatus
- func (h *StatementHandle) Progress() *StatementProgress
- func (h *StatementHandle) ResultSet() *ResultSet
- func (h *StatementHandle) Status(ctx context.Context) (StatementStatus, error)
- func (h *StatementHandle) Wait(ctx context.Context) (*ResultSet, error)
- type StatementProgress
- type StatementStatus
- type Table
- func (t *Table) AppendNDJSON(ctx context.Context, ndjson []byte) (AppendRowsResult, error)
- func (t *Table) AppendStream(options AppendStreamOptions) (*AppendStream, error)
- func (t *Table) Describe(ctx context.Context) (TableResource, error)
- func (t *Table) Drop(ctx context.Context) error
- func (t *Table) Identifier() string
- type TableColumnSpec
- type TableDistinctSpec
- type TableResource
- type TableResourceSummary
- type TableSpec
- type Value
Constants ¶
This section is empty.
Variables ¶
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 ¶
func (s *AppendStream) Flush(ctx context.Context) (AppendDeliveryReport, error)
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 ¶
func (s *AppendStream) Shutdown(ctx context.Context) (AppendDeliveryReport, error)
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 (*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 ¶
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) StatementHandle ¶
func (c *Client) StatementHandle(id uuid.UUID) *StatementHandle
StatementHandle creates a new StatementHandle with the given ID.
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 ¶
UnmarshalJSON accepts the canonical wire names and normalizes the historical unsigned-integer alias to UIntDataType.
type DatabaseResource ¶
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.
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 ¶
First returns the first row keyed by column name. The boolean is false when the result set is empty.
func (*ResultSet) RawRows ¶
RawRows returns the string-or-null cells from the JSON wire response without converting them to ScopeDB value types.
func (*ResultSet) ToObjects ¶
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 ¶
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 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.
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 ¶
func (h *StatementHandle) Cancel(ctx context.Context) (StatementCancelResult, error)
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 ¶
func (h *StatementHandle) Status(ctx context.Context) (StatementStatus, error)
Status fetches the latest status at most once. A cached terminal status is returned without making another request.
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 ¶
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 ¶
Drop drops the table from ScopeDB.
This method issues a DROP TABLE statement to ScopeDB and blocks until done.
func (*Table) Identifier ¶
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 ¶
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.
Source Files
¶
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. |