Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion .github/workflows/build.yml
Original file line number Diff line number Diff line change
Expand Up @@ -2183,7 +2183,7 @@ jobs:
PYTHONPATH: '${{ env.PYTHONPATH }}:${{ github.workspace }}/test/python'
IS_SKIP_MCP_TEST: 'true'
if: success() && env.CI_IS_EXPRESS != 'true' && matrix.platform == 'linux/amd64' && env.BUILD_IMAGE_REQUIRED == 'true' && matrix.db_backend == 'sqlite'
timeout-minutes: ${{ vars.DEFAULT_STEP_TIMEOUT_MIN == '' && 20 || vars.DEFAULT_STEP_TIMEOUT_MIN }}
timeout-minutes: ${{ vars.DEFAULT_LONG_STEP_TIMEOUT_MIN == '' && 40 || vars.DEFAULT_LONG_STEP_TIMEOUT_MIN }}
run: |
sudo rm -rf test/tmp || true
mkdir -p test/tmp
Expand Down
2 changes: 2 additions & 0 deletions .gitignore
Original file line number Diff line number Diff line change
Expand Up @@ -46,3 +46,5 @@ mcp-audit-*.log

## Generic sink default name (used by pkg/sink tests / non-MCP callers)
sink_*.log

/*.db
66 changes: 65 additions & 1 deletion docs/preview.md
Original file line number Diff line number Diff line change
Expand Up @@ -36,7 +36,7 @@ _googleProject="stackql-demo" && \
### What streams, and what does not

- Use `--output jsonl` or `--output otel`. Each row is written and flushed as it arrives. The default table output holds every row until the query completes.
- `ORDER BY`, `GROUP BY`, `HAVING`, `DISTINCT` and aggregates are refused on these relations, since each needs every row before it can emit one.
- `ORDER BY`, `GROUP BY`, `HAVING`, `DISTINCT` and aggregates are refused on these relations, since each needs every row before it can emit one, unless [staging](#staging) is enabled.
- `LIMIT` without `ORDER BY` is pushed down and stops the requests early.
- Supported joins are `INNER JOIN` and `LEFT JOIN` with `ON`. A condition in `ON` that feeds a required input of the joined relation becomes a request per row of the left side.

Expand Down Expand Up @@ -155,3 +155,67 @@ from stackql_unstable_google.iam.service_accounts s
inner join stackql_unstable_google.iam.service_account_keys k on k.serviceAccountsId = s.email
where s.projectsId = '${_googleProject}' and k.projectsId = '${_googleProject}';"
```

## Staging

Add `"staging":true` to `--preview` to run `ORDER BY`, `GROUP BY`, `HAVING`, `DISTINCT`, aggregates and `OFFSET` over `stackql_unstable_*` relations. `omnisdk` still runs the joins and filters; its final result is staged in one query-owned table in the SQL backend, which evaluates the rest of the statement. Queries without those clauses stream as before and stage nothing. See [the staging design note](/docs/technical/omnisdk_staging.md).

Staged results are not streamed: the final result is read in full before it is returned. `SELECT *` and subqueries are refused when staging applies.

### Live test

Each query below has been run against the live GitHub API.

```bash
./build/stackql exec "registry pull github v26.08.00448;"

## Chuck these in ./cicd/vol/vendor-secrets/secrets.sh
## export STACKQL_GITHUB_USERNAME='<github username>'
## export STACKQL_GITHUB_PASSWORD='<github personal access token>'

source ./cicd/vol/vendor-secrets/secrets.sh

## Ordering, LIMIT and OFFSET prior and preview
./build/stackql exec --auth '{ "github": { "credentialsenvvar": "STACKQL_GITHUB_TOKEN", "type": "api_key", "valuePrefix": "Bearer " } }' --output csv \
"select name, stargazers_count
from github.repos.repos
where org = 'stackql'
order by stargazers_count desc limit 5 offset 1;"


./build/stackql exec --preview='{"unstable":true,"staging":true}' --auth '{ "github": { "credentialsenvvar": "STACKQL_GITHUB_TOKEN", "type": "api_key", "valuePrefix": "Bearer " } }' --output csv \
"select name, stargazers_count
from stackql_unstable_github.repos.repos
where org = 'stackql'
order by stargazers_count desc limit 5 offset 1;"

./build/stackql exec --preview='{"omni":"all","staging":true}' --auth '{ "github": { "credentialsenvvar": "STACKQL_GITHUB_TOKEN", "type": "api_key", "valuePrefix": "Bearer " } }' --output csv \
"select name, stargazers_count
from github.repos.repos
where org = 'stackql'
order by stargazers_count desc limit 5 offset 1;"


## Grouping and aggregation.
./build/stackql exec --preview='{"unstable":true,"staging":true}' --output csv \
"select language, count(*) as repo_count
from stackql_unstable_github.repos.repos
where org = 'stackql'
group by language order by repo_count desc;"

## Aggregate with no column references.
./build/stackql exec --preview='{"unstable":true,"staging":true}' --output csv \
"select count(*) as member_count
from stackql_unstable_github.orgs.members
where org = 'stackql';"
```

Without `"staging":true` the same queries are refused, for example with `ORDER BY cannot be applied to stackql_unstable_* relations`.

Staged tables are dropped as soon as each result has been read. To confirm, add `--sqlBackend='{"dsn":"file:/tmp/stackql-staging.db"}'` to the commands above, then:

```bash
sqlite3 /tmp/stackql-staging.db "select name from sqlite_master where name like '__iql__.queries.%';"
```

An empty result means every staged table was released.
91 changes: 91 additions & 0 deletions docs/technical/omnisdk_staging.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,91 @@

# `omnisdk` staging

`omnisdk` results are cursor streams. When the query plan can preserve SQL semantics incrementally, batches can flow directly to the end user. Otherwise, batches can be staged in an RDBMS for relational operations such as ordering, aggregation, and set operations.

Unlike the eager, per-row RDBMS ingestion used by the `any-sdk` path, Omni staging is conditional and batched. The staging tablespace sits alongside StackQL-owned relations in the selected RDBMS, whether embedded SQLite or TCP-routed PostgreSQL.

## High level details of staging

Let $Q$ be a StackQL query with query ID $q$. Omni executes the joins and exchanges in its own query plan and produces a single final relation $R_Q$. An Omni plan may contain multiple exchanges, but exchange boundaries do not create staging tables or query IDs.

The cursor partitions $R_Q$ into batches without changing its contents:

$$
R_Q = B_{1} \mathbin{\|} B_{2} \mathbin{\|} \cdots \mathbin{\|} B_{n}
$$

Here $\|$ denotes concatenation in cursor order. Batch size affects transport and insertion cost, not the relational result. The SQL operators that remain after Omni planning (for example global `ORDER BY`, aggregation, `DISTINCT`, set operations, or a `LIMIT` that depends on ordering) are evaluated over $R_Q$. If each of them can be evaluated incrementally within the chosen streaming-state bound, batches flow directly to the result consumer and no table is created. Otherwise $R_Q$ is staged, as a multiset, in exactly one query-owned table:

$$
T_{q} = \biguplus_{k=1}^{n} B_{k}
$$

The RDBMS then evaluates the remaining SQL over $T_q$, and its result cursor becomes the output stream. The planner, not the mere presence of an Omni exchange, determines whether materialization is needed.

A blocking operation whose inputs Omni cannot combine into one relation (for example a set operation over two independent Omni plans) is not staged as several tables under one query ID. Pre-analysis instead splits $Q$ into child queries $q_1, \ldots, q_m$; each child owns at most one table $T_{q_j}$, and the parent evaluates the operation over those tables. Child IDs therefore correspond only to real pre-analysis splits, and the invariant is one table per query ID.

Query IDs are allocated by the RDBMS itself (a PostgreSQL sequence, or an SQLite `AUTOINCREMENT` table), so they are collision free across concurrent queries and across processes sharing a backend. Rows from two queries never share a table, and cleanup is scoped to a query and its descendants.

Reading a batch and inserting it form a back-pressured pipeline: the next batch is requested only after the previous one is written, rather than eagerly loading every result into application memory. Cursor exhaustion marks completion; query ownership tracks staged tables for cleanup on success, error, or cancellation.

## Concrete descisions

### Decision 1: Savage cut in transaction control counters

The prior `any-sdk` implementation is supported by counters for:

- `Generation ID`.
- `Sessions ID`.
- `Transaction ID`.
- `Insert ID`.

The new, `omnisdk` implementation will require only one counter, `Query ID`. This is because staging happens only per query. We do reserve the right to split queries, so shall maintain a concurrency safe hierarchy store for `Query ID` parent-child associations.

### Decision 2: New tablespace for omnisdk staging



```bash
./build/stackql exec \
--sqlBackend='{"dsn":"file:./stackql.db"}' \
"SELECT
v.vpc_id,
s.subnet_id
FROM aws.ec2.vpcs AS v
INNER JOIN aws.ec2.subnets AS s
ON v.vpc_id = s.vpc_id
WHERE v.region = 'ap-southeast-2'
AND s.region = 'ap-southeast-2';"
```

**Figure MQ-1** Model query 1. A simple working query.

---

**Table T-1**: Tablespace comparison for model query MQ-1.

| sdk | RDBMS | tables |
|---|---|---|
| any-sdk | sqlite | `"aws.ec2.vpcs.generation_<v>"`, <br/> `"aws.ec2.subnets.generation_<s>"` |
| any-sdk | postgres | `"<table_schema>"."aws.ec2.vpcs.generation_<v>"`, <br/> `"<table_schema>"."aws.ec2.subnets.generation_<s>"` |
| omnisdk | sqlite | `"__iql__.queries.<Query ID>"` |
| omnisdk | postgres | `"<query_schema>"."<Query ID>"` |

`<query_schema>` is configured with `schemata.querySchema` in `--sqlBackend` and defaults to `stackql_queries`.

## Decision 3: GC simplification for omnisdk

For omnisdk, any materializations needed will be created eagerly, **in the same place for prior** and all query tables can be marked for deletion immediately where no cache is in operation, or at whatever future time if cacheing is in effect. We will need a robust mechanism in place from day 1, default to no cache.

This does imply a keyval store to look up tombstone times and find query ID by query plaintext. Intuitively, I favour a new GC mechanism with the old one phased out when `any-sdk` is decommissioned.


## Terse comparison prior any-sdk vs omnisdk

| Aspect | any-sdk | omnisdk | Commment |
|----|----|----|----|
| RDBMS Ingestion | Per API response "record" | Per query. AOT configurable and runtime responsive batching | omnisdk clearly better performance best case |
| SQL translation of relation names | Recursive: global and per API call | Per query and flat | omnisdk has appealing simplicity and debug property of thin layer queries |
| SQL control counters | Mutiple counters and recursively applied: global and per API call | Per query and single counter only | |
| SQL secure storage | - | - | We have decided not to address security or obfuscation at this time. This will happen in future versions |
38 changes: 28 additions & 10 deletions internal/stackql/intrinsic/doc.go
Original file line number Diff line number Diff line change
Expand Up @@ -28,17 +28,21 @@ const UnstablePrefix = aot.DefaultProviderPrefix
// into. They are documents read straight from disk, with none of the registry's
// curation behind them, so nothing exposes them until a caller asks.
func IsUnstableEnabled() bool {
return previewCfg.getUnstableEnabled()
return previewCfg.getUnstableEnabled() || previewCfg.getOmniAll()
}

// docProvider is the bundle behind an unstable provider name, or false.
// docProvider is the bundle behind an unstable provider name, or false. Once
// every provider is routed to omnisdk, an unprefixed name is one too.
func docProvider(name string) (string, bool) {
if !IsUnstableEnabled() {
return "", false
}
trimmed := strings.TrimSpace(name)
if !strings.HasPrefix(strings.ToLower(trimmed), UnstablePrefix) {
return "", false
if !previewCfg.getOmniAll() || trimmed == "" || strings.EqualFold(trimmed, ProviderName) {
return "", false
}
return trimmed, true
}
bundle := trimmed[len(UnstablePrefix):]
if bundle == "" {
Expand Down Expand Up @@ -159,6 +163,9 @@ func docSelectFunc(
node *sqlparser.Select,
currentProvider string,
) (func() internaldto.ExecutorOutput, bool) {
if previewCfg.getStagingEnabled() && needsStaging(node) {
return stagedSelectFunc(ctx, node, currentProvider), true
}
translated, err := translateSelect(node, currentProvider)
if err != nil {
return refuse(err), true
Expand Down Expand Up @@ -192,15 +199,16 @@ func docMutationFunc(

const mutationSuccessMessage = "The operation was despatched successfully"

// runDocQuery describes each relation, resolves the query and runs it.
func runDocQuery(ctx queryContext, translated docQuery) internaldto.ExecutorOutput {
// openDocQuery describes each relation, resolves the query and opens its
// cursor, returning it with the alias of the relation it reports.
func openDocQuery(ctx queryContext, translated docQuery) (omnisdk.Rows, string, error) {
registry := registryRoot(ctx)
q := translated.getQuery()
tables := make(map[string]omnisdk.Table, len(q.From())+1)
for _, join := range q.From() {
tbl, describeErr := omnisdk.DescribeTable(registry, join.Resource().Handle())
if describeErr != nil {
return internaldto.NewErroneousExecutorOutput(describeErr)
return nil, "", describeErr
}
tables[join.Resource().Alias()] = tbl
}
Expand All @@ -209,7 +217,7 @@ func runDocQuery(ctx queryContext, translated docQuery) internaldto.ExecutorOutp
tbl, describeErr := omnisdk.DescribeMutation(
registry, target.Resource().Handle(), target.Verb().String())
if describeErr != nil {
return internaldto.NewErroneousExecutorOutput(describeErr)
return nil, "", describeErr
}
tables[target.Resource().Alias()] = tbl
relation = target.Resource().Alias()
Expand All @@ -218,7 +226,7 @@ func runDocQuery(ctx queryContext, translated docQuery) internaldto.ExecutorOutp
}
res, resolveErr := omnisdk.Resolve(q, tables)
if resolveErr != nil {
return internaldto.NewErroneousExecutorOutput(resolveErr)
return nil, "", resolveErr
}
// omnisdk takes one credential per run: the first relation's cloud - a
// mutation's target - leaving the rest to the canonical environment
Expand All @@ -227,12 +235,22 @@ func runDocQuery(ctx queryContext, translated docQuery) internaldto.ExecutorOutp
args.Tuning.Limit = translated.getLimit()
plan, planErr := omnisdk.NewGraphSelectQuery(registry, res.Graph(), args)
if planErr != nil {
return internaldto.NewErroneousExecutorOutput(planErr)
return nil, "", planErr
}
rows, openErr := plan.Open(context.Background())
if openErr != nil {
return internaldto.NewErroneousExecutorOutput(openErr)
return nil, "", openErr
}
return rows, relation, nil
}

// runDocQuery runs the query and streams its rows back.
func runDocQuery(ctx queryContext, translated docQuery) internaldto.ExecutorOutput {
rows, relation, err := openDocQuery(ctx, translated)
if err != nil {
return internaldto.NewErroneousExecutorOutput(err)
}
q := translated.getQuery()
if q.Target() != nil && len(translated.getOutputs()) == 0 {
return drainMutation(rows)
}
Expand Down
3 changes: 3 additions & 0 deletions internal/stackql/intrinsic/intrinsic.go
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@ import (

"github.com/stackql/any-sdk/pkg/dto"
"github.com/stackql/any-sdk/public/formulation"
"github.com/stackql/any-sdk/public/sqlengine"
"github.com/stackql/stackql/internal/stackql/internal_data_transfer/internaldto"
"github.com/stackql/stackql/internal/stackql/typing"
"github.com/stackql/stackql/internal/stackql/util"
Expand Down Expand Up @@ -50,6 +51,8 @@ type queryContext interface {
GetTypingConfig() typing.Config
GetAuthContext(providerName string) (*dto.AuthCtx, error)
GetRuntimeContext() dto.RuntimeCtx
GetSQLEngine() sqlengine.SQLEngine
GetASTFormatter() sqlparser.NodeFormatter
}

func GeneratePrimitiveFunc(
Expand Down
23 changes: 23 additions & 0 deletions internal/stackql/intrinsic/omnisdk.go
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,8 @@ const defaultBatchSize = 100

const defaultFlushInterval = 50 * time.Millisecond

const omniAll = "all"

func relationName(path string) string {
return strings.ReplaceAll(path, ".", "_")
}
Expand Down Expand Up @@ -611,6 +613,12 @@ func omnisdkAuth(authCtx *dto.AuthCtx) *omnisdk.Auth {
UsernameEnvVar: authCtx.EnvVarUsername,
PasswordEnvVar: authCtx.EnvVarPassword,
}
if strings.EqualFold(authCtx.Type, "api_key") && auth.Name == "" {
auth.Name = "Authorization"
if auth.ValuePrefix == "" {
auth.ValuePrefix = "Bearer "
}
}
if credentials, credErr := authCtx.GetCredentialsBytes(); credErr == nil {
auth.SecretAccessKey = string(credentials)
auth.Credentials = string(credentials)
Expand Down Expand Up @@ -727,6 +735,8 @@ type backendInput interface {
getFlushInterval() time.Duration
getInsecureSkipTLSVerify() bool
getUnstableEnabled() bool
getStagingEnabled() bool
getOmniAll() bool
}

type standardBackendInput struct {
Expand All @@ -735,6 +745,8 @@ type standardBackendInput struct {
flushInterval time.Duration
insecureSkipTLSVerify bool
unstableEnabled bool
stagingEnabled bool
omniAll bool
}

// previewCfg is the parsed --preview argument. Cobra binds the raw string in
Expand All @@ -755,6 +767,11 @@ type previewCfgDTO struct {
Endpoint json.RawMessage `json:"endpoint"`
InsecureSkipTLSVerify bool `json:"insecureSkipTLSVerify"`
Unstable bool `json:"unstable"`
// Staging opts SELECTs over document-driven relations into RDBMS staging
// for the SQL omnisdk leaves unapplied, instead of refusing them.
Staging bool `json:"staging"`
// Omni "all" routes every provider through omnisdk, never any-sdk.
Omni string `json:"omni"`
}

func (c previewCfgDTO) endpoint() string {
Expand Down Expand Up @@ -786,6 +803,8 @@ func newBackendInput(cfg previewCfgDTO) backendInput {
flushInterval: defaultFlushInterval,
insecureSkipTLSVerify: cfg.InsecureSkipTLSVerify,
unstableEnabled: cfg.Unstable,
stagingEnabled: cfg.Staging,
omniAll: strings.EqualFold(cfg.Omni, omniAll),
}
if cfg.BatchSize > 0 {
rv.batchSize = cfg.BatchSize
Expand All @@ -806,6 +825,10 @@ func (b *standardBackendInput) getInsecureSkipTLSVerify() bool { return b.insecu

func (b *standardBackendInput) getUnstableEnabled() bool { return b.unstableEnabled }

func (b *standardBackendInput) getStagingEnabled() bool { return b.stagingEnabled }

func (b *standardBackendInput) getOmniAll() bool { return b.omniAll }

// sourceKey is the row key a column reads from: its own name, unless an alias
// renamed it.
func (c column) sourceKey() string {
Expand Down
20 changes: 20 additions & 0 deletions internal/stackql/intrinsic/omnisdk_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -303,3 +303,23 @@ func TestRowsReachOutputBeforeNextPage(t *testing.T) {
type writerFunc func([]byte) (int, error)

func (f writerFunc) Write(b []byte) (int, error) { return f(b) }

func TestOmnisdkAuthAPIKeyDefaultsMatchCanonicalProviders(t *testing.T) {
for _, tc := range []struct {
in dto.AuthCtx
wantName string
wantPrefix string
}{
{in: dto.AuthCtx{Type: "api_key", ValuePrefix: "token "}, wantName: "Authorization", wantPrefix: "token "},
{in: dto.AuthCtx{Type: "api_key"}, wantName: "Authorization", wantPrefix: "Bearer "},
{in: dto.AuthCtx{Type: "api_key", Name: "X-Api-Key"}, wantName: "X-Api-Key", wantPrefix: ""},
{in: dto.AuthCtx{Type: "bearer"}, wantName: "", wantPrefix: ""},
} {
authCtx := tc.in
got := omnisdkAuth(&authCtx)
if got.Name != tc.wantName || got.ValuePrefix != tc.wantPrefix {
t.Errorf("%+v: got name %q prefix %q, want %q %q",
tc.in, got.Name, got.ValuePrefix, tc.wantName, tc.wantPrefix)
}
}
}
Loading
Loading