From 9147717ef83ef0ad8cbfbbe430e5a922954c08df Mon Sep 17 00:00:00 2001 From: General Kroll Date: Fri, 25 Sep 2026 20:52:54 +1000 Subject: [PATCH 1/4] omnisdk-query-replacement-opt-in Summary: - Support for limited, lightweight streamed queries via `preview` opt in. - Pass each document-driven provider's configured stackql auth, basic included, to omnisdk, and report each method's SQL verb in SHOW METHODS. - Added robot test `Unstable Google Kms Key Rings Joined To Crypto Keys Jsonl Row Set Matches Expectation`. - Added robot test `Unstable Github Org Members Filtered By In List Jsonl Row Set Matches Expectation`. - Added robot test `Unstable Github Org Update Reports Despatch`. - Added robot test `Unstable Github Org Update Returning Jsonl Row Set Matches Expectation`. - Added robot test `Unstable Google Kms Key Ring Insert Returning Jsonl Row Set Matches Expectation`. - Pass each document-driven provider's configured stackql auth, basic included, to omnisdk, and report each method's SQL verb in SHOW METHODS. --- .claude/omnisdk-shortcomings.md | 44 ++ .claude/stackql-integration.md | 123 ++++ docs/preview.md | 23 + go.mod | 2 +- go.sum | 4 +- internal/stackql/intrinsic/doc.go | 175 +++-- internal/stackql/intrinsic/dynamic.go | 5 +- internal/stackql/intrinsic/intrinsic.go | 9 +- internal/stackql/intrinsic/omnisdk.go | 65 +- .../registry/fixture/v0.1.0/provider.yaml | 17 + .../fixture/v0.1.0/services/orgs.yaml | 54 ++ internal/stackql/intrinsic/translate.go | 684 ++++++++++++++++++ internal/stackql/intrinsic/translate_test.go | 299 ++++++++ internal/stackql/planbuilder/entrypoint.go | 35 +- .../stackql_test_tooling/flask/gcp/app.py | 10 + .../stackql_mocked_from_cmd_line.robot | 173 +++++ 16 files changed, 1633 insertions(+), 89 deletions(-) create mode 100644 .claude/omnisdk-shortcomings.md create mode 100644 .claude/stackql-integration.md create mode 100644 docs/preview.md create mode 100644 internal/stackql/intrinsic/testdata/registry/fixture/v0.1.0/provider.yaml create mode 100644 internal/stackql/intrinsic/testdata/registry/fixture/v0.1.0/services/orgs.yaml create mode 100644 internal/stackql/intrinsic/translate.go create mode 100644 internal/stackql/intrinsic/translate_test.go diff --git a/.claude/omnisdk-shortcomings.md b/.claude/omnisdk-shortcomings.md new file mode 100644 index 00000000..1c67f249 --- /dev/null +++ b/.claude/omnisdk-shortcomings.md @@ -0,0 +1,44 @@ + +# omnisdk shortcomings + +## High Level Expectations + +1. `omnisdk` will accept existing `any-sdk` auth structures and default configs in provider documents and behave as expected, can be observed in `any-sdk`. Behviours should be implemented protocol agnostic and best practice. Tech debt from `any-sdk` must not be inherited by copy or imitation. **`omnisdk` will never import from `any-sdk`, directly or indirectly.** +2. This is very important. Whereas `nay-sdk` has been tightly coupled to `http`, with some subprocess variants, by contrast `omnisdk` is **protocol agsnostc**. Full support for proticol buffers, `http`, subprocess calls, streaming protocols, whatever old and new transports, are intrinsic to `omnisdk`. +3. Any WAL, ledgers and the like should be configurable in location, substrate (eg: local vs, s3, different os...) and the like. Must be extensible and abstacted, no excuses. +4. `omnisdk` will support all functionality in `any-sdk`, but with clean implementation. Breaking changes will be specified ahead of time. The api is wildly different, but functional coverage will not be lesser. +5. `omnisdk` does not stage results in RDBMS or otherwise, it eagerly streams generated records. This is a significant and highly beneficial difference to `any-sdk`. Any ordering, aggregation or set operations (union and the like) is imposed post `omnisdk` query fulfilment. +6. `omnisdk` is required to support SQL extension funstions, including: scalar, redord and table valued ones. We do not expect the very first version to have total coverage and so some queries may need to be routed away from `omnisdk` at times, by logic within `stackql`. +7. `omnisdk` is to support all of the request and response trandsform grammars and shorthands already in `any-sdk`. +8. `omnisdk` will at some point support both SQL style saga rollbacks and also an IAC saga variant with fine grained locking and abstracted latches that are objectively superior to `terraform` style crude locks and failure modes. That said, early IAC forms are per stack locked. Patience is our watchword. + +## Migration plan at coarse grain + +- (a) Support joins, sql functions consuming only `omnisdk`, in `stackql_unstable_` namespace. It may take some time to acheive full covereage, in stages. +- (b) Cut versions and releases of both `omnisdk` and `stackql` along the way as useful milestones are reached. + + +## Current issues + +- (i) Auth has to function as expected. No excuses. The existing stackql patter with auth structures and docs simply **must** work. +- (ii) We need an orderly abstraction and catalogue of supported SQL extension functions in `omnisdk`. This will be in a discrete `pkg`. See below `Expected SQL extension functions` section. +- (iii) We need an orderly abstraction and catalogue of supported request and response processing grammars and shorthands in `omnisdk`. These are expected to mirror the `main` branch of `any-sdk` but be cleanly implemented in a discrete `pkg` with minimal dependencies and **zero** relation to `any-sdk`. See below `Expected Transformation Grammars and Shorthands` section. +- (iv) We want support for user/agent composed cross cloud rapid audit queries. +- (v) I want an SOC or whatever corporate audit query suite asap. +- (vi) `Args.Auth` is one struct shared by every node in the graph; each node reads the fields its scheme needs and falls back to env vars for any left empty. Two providers cannot carry distinct credentials in one query. Need per-provider credentials. Blocks (i) and (iv). +- (vii) `DescribeTable`/`DescribeMutation` drop parameters declared via `$ref` to `components/parameters` (e.g. github `org` on `orgs.members`, `username` on `users.users`), so joins and mutations cannot bind them. +- (viii) `DescribeMutation(dir, "stackql_unstable_google.storage.buckets", "insert")` fails with "read services: is a directory" while `DescribeTable` on the same address works; a `provider.yaml` service `$ref` to a missing file (`compute-v1.yaml`) gives the same message instead of naming the file. +- (ix) `Table` gives column names but no types, so every column reaches stackql as text. +- (x) Mutation outcomes ("was rejected" vs "may or may not have taken effect") are `fmt.Errorf` strings only; need sentinel errors for `errors.Is`. +- (xi) A mutation target carries one assignment set, so multi-row `INSERT ... VALUES` cannot be expressed. + +## Supporting information + +### Expected SQL extension functions + +The required SQL extension functions are precisely those that are robot tested anywhere in the `stackql` codebase. + +### Expected Transformation Grammars and Shorthands + +The required transformation grammars (eg for http request and response) are precisely those that are implemented in the `main` branch of the `any-sdk` codebase. + diff --git a/.claude/stackql-integration.md b/.claude/stackql-integration.md new file mode 100644 index 00000000..a8adea53 --- /dev/null +++ b/.claude/stackql-integration.md @@ -0,0 +1,123 @@ +# Routing stackql queries through omnisdk + +What stackql has to do to send a query to omnisdk and stream back its rows. omnisdk resolves the +query against the provider documents and runs it; stackql parses SQL, holds sessions, and does +whatever needs every row at once. + +## Division of work + +| stackql | omnisdk | +|---------|---------| +| Parse SQL | Choose each table's method from the documents | +| Split CTEs, subqueries and UNION into single queries, then combine results | Decide which conditions become request parameters, edges between tables, or row filters | +| ORDER BY, GROUP BY, aggregates, DISTINCT | Run requests, joins, `IN` fan-out, polls and mutations | +| Hold per-client sessions and turn them into per-query config | Return an eager, unordered, unaggregated row stream | + +## Per query + +### 1. Build a `query.Unresolved` (package `pkg/query`, standard library only) + +Handles are registry addresses: `stackql_unstable_aws.iam.users`. omnisdk does not map names. + +| SQL | Constructor | +|-----|-------------| +| `FROM t a` | `query.NewJoin(query.NewResource("a", "
"), query.Base)` | +| `[INNER] JOIN t b ON …` | `query.NewJoin(query.NewResource("b", …), query.Inner, on...)` | +| `LEFT JOIN t b ON …` | `query.NewJoin(…, query.Left, on...)` | +| `a.col` / unqualified `col` | `query.NewColumn("a", "col")` / `query.NewColumn("", "col")` | +| `'x'`, `1` | `query.NewLiteral(v)` | +| `f(x, y)` | `query.NewCall("f", x, y)` | +| `=`, `<>`, `<`, `<=`, `>`, `>=` | `query.NewCompare(query.Eq, l, r)` (and `Ne`, `Lt`, …) | +| `x IN (…)` | `query.NewIn(x, query.NewCollection(items...))` | +| `OR`, `NOT` | `query.NewOr(...)`, `query.NewNot(p)` | +| boolean function as condition | `query.NewTest(call)` | +| `SELECT expr AS name` | `query.NewOutput("name", expr)` | +| `SELECT *` / `a.*` | `query.NewOutput("", query.NewStar(""))` / `NewStar("a")` | +| `INSERT INTO t (c…) VALUES (…)` | `query.NewInsert(res, query.NewAssignment("c", v)...)` | +| `INSERT … SELECT` | same, with values reading the SELECT's tables, which go in `from` | +| `UPDATE t SET c = v` | `query.NewUpdate(res, assignments...)` | +| `DELETE FROM t` | `query.NewDelete(res)` | +| `RETURNING …` | the mutation's outputs | + +Split WHERE and each ON into conjuncts (top-level `AND`). Every output needs a name: give +unnamed expressions one. Then: + +```go +q, err := query.New(from, where, outputs) // a read +q, err := query.NewMutation(target, from, where, returning) // a mutation +``` + +Leave out ORDER BY, GROUP BY, aggregates, DISTINCT and HAVING. If ORDER BY or GROUP BY reads a +column that isn't selected, add it to the outputs so it comes back. + +Not supported: RIGHT, FULL and CROSS joins. + +### 2. Describe each table + +```go +tables := map[string]omnisdk.Table{} +for _, j := range q.From() { + t, err := omnisdk.DescribeTable(registry, j.Resource().Handle()) + tables[j.Resource().Alias()] = t +} +if tg := q.Target(); tg != nil { + t, err := omnisdk.DescribeMutation(registry, tg.Resource().Handle(), tg.Verb().String()) + tables[tg.Resource().Alias()] = t +} +``` + +A `Table` lists each method's parameters (name, location, required) and row columns, which is also +what `DESCRIBE` / `SHOW` can report. + +### 3. Resolve + +```go +res, err := omnisdk.Resolve(q, tables) +``` + +An error names the clause it could not place, for example a condition a mutation's method can't take +(it would only be checked after the effect), a left join whose preserved side needs the other side's +value, or an ambiguous `*`. Show it to the user as is. + +### 4. Run + +```go +args := omnisdk.Args{ + Params: merge(session.Params, res.Params()), // scope such as region, then the query's own + Auth: session.Auth, // nil falls back to the environment + Tuning: omnisdk.Tuning{Limit: limit}, // LIMIT without ORDER BY; otherwise apply it after sorting + Journal: session.Journal, // opt-in write-ahead intent for mutations + Redaction: session.Redaction, // nil drops credentials from rows +} +pl, err := omnisdk.NewGraphSelectQuery(registry, res.Graph(), args) +rows, err := pl.Open(ctx) +defer rows.Close() +for rows.Next() { + row := rows.Row() // map[string]any, keyed by output name +} +err = rows.Err() +``` + +Rows arrive as they are produced. A left join's unmatched row lacks the joined table's columns; +treat a missing column as NULL. + +## Per session + +- **Config:** keep it in the session and build a fresh `Args` for each query. omnisdk holds no + per-client state, so concurrent clients with different settings don't interfere. +- **Document changes for one client:** `EffectiveRegistry(registry, patches, omnisdk.DocCache{Dir: dir})` + returns a registry directory with that client's patches applied (RFC 7386 merge patches on service + documents). Use it as `registry` for that client's queries. `NewDocCache(parent)` gives a new cache + location; passing an existing `Dir` reuses its entries; `Fresh: true` rebuilds. +- **Mutations:** set `Args.Journal{State, RunID}` to record each effect before it is sent. A failure is + reported as "rejected" (the provider refused it) or "may or may not have taken effect" (no answer, or + a 5xx). Mutations are never retried. +- **Credentials:** dropped from result rows by default. `omnisdk.RedactNone()` keeps them, for a user + who needs them. + +## Once per process + +- `omnisdk.ConfigureDocumentCache(cache.Config{...})` bounds the parsed-document cache. The default + takes a quarter of the process's memory limit (`GOMEMLIMIT`, else the container's cgroup limit, + else 512 MiB). Parsed documents are keyed by the directory they are read from, so clients with + different patches never share entries. diff --git a/docs/preview.md b/docs/preview.md new file mode 100644 index 00000000..67ea400d --- /dev/null +++ b/docs/preview.md @@ -0,0 +1,23 @@ + +# Preview functionality + +Bleeding edge `stackql` functionality can be exposed using the `--preview` CLI argument. + + +```bash +./build/stackql shell --preview='{"unstable":true}' +``` + +## Relevant preview functionality + +Streaming high volume queries at low latency has releveance for audit and related use cases. + + +```bash + +./build/stackql exec --preview='{"unstable":true}' --output=jsonl \ + "select name, location, storageClass, timeCreated + from stackql_unstable_google.storage.buckets + where project = 'stackql-demo'" + +``` \ No newline at end of file diff --git a/go.mod b/go.mod index 9dfd7bd9..6717e0cf 100644 --- a/go.mod +++ b/go.mod @@ -19,7 +19,7 @@ require ( github.com/spf13/cobra v1.10.2 github.com/spf13/pflag v1.0.10 github.com/spf13/viper v1.10.1 - github.com/stackql-labs/omnisdk v0.1.2-beta02 + github.com/stackql-labs/omnisdk v0.1.3-alpha03 github.com/stackql/any-sdk v0.6.0-alpha01 github.com/stackql/go-suffix-map v0.0.1-alpha01 github.com/stackql/psql-wire v0.1.2-beta01 diff --git a/go.sum b/go.sum index b383fe0e..b3635cae 100644 --- a/go.sum +++ b/go.sum @@ -366,8 +366,8 @@ github.com/spf13/viper v1.10.1 h1:nuJZuYpG7gTj/XqiUwg8bA0cp1+M2mC3J4g5luUYBKk= github.com/spf13/viper v1.10.1/go.mod h1:IGlFPqhNAPKRxohIzWpI5QEy4kuI7tcl5WvR+8qy1rU= github.com/spiffe/go-spiffe/v2 v2.7.0 h1:uXe1MflJoHw58wAUvxVlcM7WpKtijWG7I1UidcGh6g4= github.com/spiffe/go-spiffe/v2 v2.7.0/go.mod h1:47Q0Q9/AqGha8QLHp+kxpH4Wca7X7EnOtlIJy3mxZ3U= -github.com/stackql-labs/omnisdk v0.1.2-beta02 h1:4UpcSdZtsqyPFhv5ZZ9Mz1siUvAXXeGah0R64tczMww= -github.com/stackql-labs/omnisdk v0.1.2-beta02/go.mod h1:WzvNj/bVv53yGFsVJpYWCJC1xAEdmQSFJl9eVpkRpCY= +github.com/stackql-labs/omnisdk v0.1.3-alpha03 h1:6bTEeF/FoBTsVjl/lk63a8plDS3VaPPwwAPByht6QE0= +github.com/stackql-labs/omnisdk v0.1.3-alpha03/go.mod h1:WzvNj/bVv53yGFsVJpYWCJC1xAEdmQSFJl9eVpkRpCY= github.com/stackql/any-sdk v0.6.0-alpha01 h1:mqy0bmZ1wghr7mUgLyaexoRU6rsUZCWXU9Mf60fC8dg= github.com/stackql/any-sdk v0.6.0-alpha01/go.mod h1:ahkRgwHHEn7RTmSfKrshPVk6i4X1BdHknSeBflSyRac= github.com/stackql/go-suffix-map v0.0.1-alpha01 h1:TDUDS8bySu41Oo9p0eniUeCm43mnRM6zFEd6j6VUaz8= diff --git a/internal/stackql/intrinsic/doc.go b/internal/stackql/intrinsic/doc.go index 2bbb591d..f18b2caf 100644 --- a/internal/stackql/intrinsic/doc.go +++ b/internal/stackql/intrinsic/doc.go @@ -132,68 +132,137 @@ func docMethods(ctx queryContext, bundle, service, resource string) ([]relationM } out := make([]relationMethod, 0, len(methods)) for _, method := range methods { - out = append(out, relationMethod{name: method.Name, description: method.OperationID}) + out = append(out, relationMethod{ + name: method.Name, + description: method.OperationID, + sqlVerb: docSQLVerb(method.SQLVerb), + }) } return out, nil } -// docSelectFunc routes a SELECT over a document-driven relation. The address is -// the bundle's own "..", and the plan it yields -// streams exactly as a hand-authored one does. +// docSQLVerb is the verb a document maps a method to, upper-cased as SHOW +// METHODS reports it. A method no verb maps is reachable only through EXEC. +func docSQLVerb(verb string) string { + if verb == "" { + return "EXEC" + } + return strings.ToUpper(verb) +} + +// docSelectFunc routes a SELECT over document-driven relations. omnisdk +// resolves the whole query against the provider documents - which method each +// relation runs, and whether a condition is a request parameter, an edge +// between relations or a row filter - and streams the rows back. func docSelectFunc( ctx queryContext, node *sqlparser.Select, - bundle, service, resource string, + currentProvider string, ) (func() internaldto.ExecutorOutput, bool) { - if unsupported := unsupportedClauses(node); len(unsupported) > 0 { - return refuse(fmt.Errorf( - "relation '%s%s.%s.%s' streams its rows, so %s cannot be applied; remove %s from the query", - UnstablePrefix, bundle, service, resource, - strings.Join(unsupported, ", "), pluralClause(len(unsupported)))), true - } - params, badPredicates := equalityPredicates(node.Where) - if len(badPredicates) > 0 { - return refuse(fmt.Errorf( - "relation '%s%s.%s.%s' streams its rows, so only equality predicates are applied; "+ - "%s cannot be honoured", - UnstablePrefix, bundle, service, resource, strings.Join(badPredicates, ", "))), true - } - address := fmt.Sprintf("%s%s.%s.%s", UnstablePrefix, bundle, service, resource) - return func() internaldto.ExecutorOutput { - input := previewCfg - dir, dirErr := docRoot(ctx, bundle) - if dirErr != nil { - return internaldto.NewErroneousExecutorOutput(dirErr) - } - plan, err := omnisdk.NewFromCatalog(dir, address, omnisdk.Args{ - Params: params, - Auth: omnisdkAuth(providerAuthContext(ctx, bundle)), - Endpoint: input.getEndpoint(), - InsecureSkipTLSVerify: input.getInsecureSkipTLSVerify(), - }) - if err != nil { - return internaldto.NewErroneousExecutorOutput(err) - } - rows, openErr := plan.Open(context.Background()) - if openErr != nil { - return internaldto.NewErroneousExecutorOutput(openErr) - } - // A document declares no egress schema, so the columns are those the - // first row carries; projection is applied over them. - stream := &rowStream{ - rows: rows, - batchSize: input.getBatchSize(), - flushInterval: input.getFlushInterval(), - table: sqldata.NewSQLTable(0, resource), - typCfg: ctx.GetTypingConfig(), - projection: node.SelectExprs, + translated, err := translateSelect(node, currentProvider) + if err != nil { + return refuse(err), true + } + return func() internaldto.ExecutorOutput { return runDocQuery(ctx, translated) }, true +} + +// docMutationFunc routes an INSERT, UPDATE or DELETE whose target is a +// document-driven relation. Without RETURNING the effects are driven to +// completion and reported as a message; with it, the returned rows stream back. +func docMutationFunc( + ctx queryContext, + stmt sqlparser.Statement, + currentProvider string, +) (func() internaldto.ExecutorOutput, bool) { + tables, isMutation := mutationTables(stmt) + if !isMutation { + return nil, false + } + if isDoc, err := fromDocProviders(tables, currentProvider); !isDoc { + return nil, false + } else if err != nil { + return refuse(err), true + } + translated, err := translateMutation(stmt, currentProvider) + if err != nil { + return refuse(err), true + } + return func() internaldto.ExecutorOutput { return runDocQuery(ctx, translated) }, true +} + +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 { + 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) } - primed, readErr := newPrimedStream(stream) - if readErr != nil { - return internaldto.NewErroneousExecutorOutput(readErr) + tables[join.Resource().Alias()] = tbl + } + var relation string + if target := q.Target(); target != nil { + tbl, describeErr := omnisdk.DescribeMutation( + registry, target.Resource().Handle(), target.Verb().String()) + if describeErr != nil { + return internaldto.NewErroneousExecutorOutput(describeErr) } - return internaldto.NewExecutorOutput(primed, nil, nil, nil, nil) - }, true + tables[target.Resource().Alias()] = tbl + relation = target.Resource().Alias() + } else { + relation = q.From()[0].Resource().Alias() + } + res, resolveErr := omnisdk.Resolve(q, tables) + if resolveErr != nil { + return internaldto.NewErroneousExecutorOutput(resolveErr) + } + // omnisdk takes one credential per run: the first relation's cloud - a + // mutation's target - leaving the rest to the canonical environment + // variables. + args := previewArgs(ctx, translated.getBundles()[0], res.Params()) + args.Tuning.Limit = translated.getLimit() + plan, planErr := omnisdk.NewGraphSelectQuery(registry, res.Graph(), args) + if planErr != nil { + return internaldto.NewErroneousExecutorOutput(planErr) + } + rows, openErr := plan.Open(context.Background()) + if openErr != nil { + return internaldto.NewErroneousExecutorOutput(openErr) + } + if q.Target() != nil && len(translated.getOutputs()) == 0 { + return drainMutation(rows) + } + input := previewCfg + stream := &rowStream{ + rows: rows, + batchSize: input.getBatchSize(), + flushInterval: input.getFlushInterval(), + table: sqldata.NewSQLTable(0, relation), + typCfg: ctx.GetTypingConfig(), + outputs: translated.getOutputs(), + } + primed, readErr := newPrimedStream(stream) + if readErr != nil { + return internaldto.NewErroneousExecutorOutput(readErr) + } + return internaldto.NewExecutorOutput(primed, nil, nil, nil, nil) +} + +// drainMutation pulls a mutation's cursor to the end, which is what sends its +// effects, and reports the outcome. +func drainMutation(rows omnisdk.Rows) internaldto.ExecutorOutput { + defer rows.Close() + for rows.Next() { //nolint:revive // draining sends the effects + } + if err := rows.Err(); err != nil { + return internaldto.NewErroneousExecutorOutput(err) + } + return internaldto.NewExecutorOutput(nil, nil, nil, + internaldto.NewBackendMessages([]string{mutationSuccessMessage}), nil) } func refuse(err error) func() internaldto.ExecutorOutput { @@ -255,7 +324,7 @@ func showDocMethods( row := map[string]interface{}{ "MethodName": method.name, "RequiredParams": strings.Join(method.requiredParams, ", "), - "SQLVerb": strings.ToUpper(selectMethodName), + "SQLVerb": method.sqlVerb, } if extended { row["description"] = method.description diff --git a/internal/stackql/intrinsic/dynamic.go b/internal/stackql/intrinsic/dynamic.go index fc2b4c59..47d6a847 100644 --- a/internal/stackql/intrinsic/dynamic.go +++ b/internal/stackql/intrinsic/dynamic.go @@ -115,7 +115,10 @@ func buildGraph(spec string) (omnisdk.Graph, error) { override.Address, override.ObjectKey, override.MediaType, override.ProgramType, override.ProgramBody)) } - return omnisdk.NewGraph(dtoSpec.Addresses, wirings, overrides...) + return omnisdk.NewGraph( + omnisdk.NodesOf(dtoSpec.Addresses...), + wirings, + overrides...) } // graphCloud is the cloud whose credential a graph runs under. omnisdk takes diff --git a/internal/stackql/intrinsic/intrinsic.go b/internal/stackql/intrinsic/intrinsic.go index 5a77add9..c23655ab 100644 --- a/internal/stackql/intrinsic/intrinsic.go +++ b/internal/stackql/intrinsic/intrinsic.go @@ -70,11 +70,16 @@ func GeneratePrimitiveFunc( return nil, false } +// GenerateStreamFunc plans a statement omnisdk runs end to end: a SELECT over +// its relations, or a mutation whose target is a document-driven relation. func GenerateStreamFunc( ctx queryContext, - node *sqlparser.Select, + stmt sqlparser.Statement, ) (func() internaldto.ExecutorOutput, bool) { - return selectFunc(ctx, node, ctx.GetCurrentProvider()) + if node, isSelect := stmt.(*sqlparser.Select); isSelect { + return selectFunc(ctx, node, ctx.GetCurrentProvider()) + } + return docMutationFunc(ctx, stmt, ctx.GetCurrentProvider()) } func IsProvider(name string) bool { diff --git a/internal/stackql/intrinsic/omnisdk.go b/internal/stackql/intrinsic/omnisdk.go index 5b47d7bd..7362c0a2 100644 --- a/internal/stackql/intrinsic/omnisdk.go +++ b/internal/stackql/intrinsic/omnisdk.go @@ -190,6 +190,7 @@ type rowStream struct { columns []column table sqldata.ISQLTable projection sqlparser.SelectExprs + outputs []string typCfg columnFactory done bool } @@ -261,7 +262,9 @@ func (rs *rowStream) startProducer() { } func (rs *rowStream) result(batch []omnisdk.Row) sqldata.ISQLResult { - if len(rs.columns) == 0 && len(batch) > 0 { + if len(rs.columns) == 0 && len(rs.outputs) > 0 { + rs.columns = outputColumns(rs.outputs, batch) + } else if len(rs.columns) == 0 && len(batch) > 0 { for _, name := range sortedKeys(batch[0]) { rs.columns = append(rs.columns, column{name: name}) } @@ -286,6 +289,37 @@ func (rs *rowStream) result(batch []omnisdk.Row) sqldata.ISQLResult { return sqldata.NewSQLResult(columns, uint64(len(rows)), 0, rows) } +// outputColumns lays the columns out in select-list order. A star is filled +// from the first row's keys that no named output claims, so it waits for a row; +// a list without one is known before any arrives. +func outputColumns(outputs []string, batch []omnisdk.Row) []column { + named := make(map[string]bool, len(outputs)) + hasStar := false + for _, name := range outputs { + if name == starOutput { + hasStar = true + continue + } + named[name] = true + } + if hasStar && len(batch) == 0 { + return nil + } + out := make([]column, 0, len(outputs)) + for _, name := range outputs { + if name != starOutput { + out = append(out, column{name: name}) + continue + } + for _, key := range sortedKeys(batch[0]) { + if !named[key] { + out = append(out, column{name: key}) + } + } + } + return out +} + func (rs *rowStream) Write(sqldata.ISQLResult) error { return fmt.Errorf("intrinsic: omnisdk result stream is read-only") } @@ -308,6 +342,12 @@ func selectFunc( node *sqlparser.Select, currentProvider string, ) (func() internaldto.ExecutorOutput, bool) { + if isDoc, err := fromDocProviders(node.From, currentProvider); isDoc { + if err != nil { + return refuse(err), true + } + return docSelectFunc(ctx, node, currentProvider) + } if len(node.From) != 1 { return nil, false } @@ -323,11 +363,6 @@ func selectFunc( tableName.Qualifier.GetRawVal(), currentProvider); isPreview { return previewSelectFunc(ctx, node, service, tableName.Name.GetRawVal()) } - if bundle, isDoc := docProvider( - resolveProvider(tableName.QualifierSecond.GetRawVal(), currentProvider)); isDoc { - return docSelectFunc(ctx, node, bundle, - tableName.Qualifier.GetRawVal(), tableName.Name.GetRawVal()) - } if !strings.EqualFold(tableName.Qualifier.GetRawVal(), auditService) || !IsProvider(resolveProvider(tableName.QualifierSecond.GetRawVal(), currentProvider)) { return nil, false @@ -476,6 +511,7 @@ type relationMethod struct { name string description string requiredParams []string + sqlVerb string } func (t table) methods() []relationMethod { @@ -533,7 +569,8 @@ func providerAuthContext(ctx queryContext, resourcePath string) *dto.AuthCtx { cloud, _, _ := strings.Cut(resourcePath, ".") providerName, ok := cloudProviders[cloud] if !ok { - return nil + // A document-driven provider is addressed by its own stackql name. + providerName = cloud } authCtx, err := ctx.GetAuthContext(providerName) if err != nil { @@ -547,11 +584,15 @@ func omnisdkAuth(authCtx *dto.AuthCtx) *omnisdk.Auth { return nil } auth := &omnisdk.Auth{ - Type: authCtx.Type, - ValuePrefix: authCtx.ValuePrefix, - Name: authCtx.Name, - Scopes: authCtx.Scopes, - TokenURL: authCtx.GetTokenURL(), + Type: authCtx.Type, + ValuePrefix: authCtx.ValuePrefix, + Name: authCtx.Name, + Scopes: authCtx.Scopes, + TokenURL: authCtx.GetTokenURL(), + Username: authCtx.Username, + Password: authCtx.Password, + UsernameEnvVar: authCtx.EnvVarUsername, + PasswordEnvVar: authCtx.EnvVarPassword, } if credentials, credErr := authCtx.GetCredentialsBytes(); credErr == nil { auth.SecretAccessKey = string(credentials) diff --git a/internal/stackql/intrinsic/testdata/registry/fixture/v0.1.0/provider.yaml b/internal/stackql/intrinsic/testdata/registry/fixture/v0.1.0/provider.yaml new file mode 100644 index 00000000..c54f8a0b --- /dev/null +++ b/internal/stackql/intrinsic/testdata/registry/fixture/v0.1.0/provider.yaml @@ -0,0 +1,17 @@ +id: fixture +name: fixture +version: v0.1.0 +providerServices: + orgs: + id: orgs:v0.1.0 + name: orgs + preferred: true + service: + $ref: fixture/v0.1.0/services/orgs.yaml + title: orgs API + version: v0.1.0 + description: organisation membership +config: + auth: + type: bearer + credentialsenvvar: FIXTURE_TOKEN diff --git a/internal/stackql/intrinsic/testdata/registry/fixture/v0.1.0/services/orgs.yaml b/internal/stackql/intrinsic/testdata/registry/fixture/v0.1.0/services/orgs.yaml new file mode 100644 index 00000000..616bb723 --- /dev/null +++ b/internal/stackql/intrinsic/testdata/registry/fixture/v0.1.0/services/orgs.yaml @@ -0,0 +1,54 @@ +openapi: 3.0.3 +info: + title: orgs API + version: v0.1.0 +servers: + - url: https://api.fixture.test +paths: + /orgs/{org}/members: + get: + operationId: orgs/list-members + parameters: + - name: org + in: path + required: true + schema: + type: string + responses: + '200': + description: members + content: + application/json: + schema: + type: array + items: + $ref: '#/components/schemas/member' +components: + schemas: + member: + type: object + properties: + login: + type: string + id: + type: integer + type: + type: string + x-stackQL-resources: + members: + id: fixture.orgs.members + name: members + title: Members + methods: + list_members: + operation: + $ref: '#/paths/~1orgs~1{org}~1members/get' + response: + mediaType: application/json + openAPIDocKey: '200' + sqlVerbs: + select: + - $ref: '#/components/x-stackQL-resources/members/methods/list_members' + insert: [] + update: [] + delete: [] diff --git a/internal/stackql/intrinsic/translate.go b/internal/stackql/intrinsic/translate.go new file mode 100644 index 00000000..908e1b80 --- /dev/null +++ b/internal/stackql/intrinsic/translate.go @@ -0,0 +1,684 @@ +package intrinsic + +// A SELECT over the document-driven providers is handed to omnisdk whole. This +// file turns the parsed statement into omnisdk's query.Unresolved: FROM becomes +// joins over registry addresses, WHERE and each ON become conjuncts, and the +// select list becomes named outputs. omnisdk decides which conditions become +// request parameters, edges between tables, or row filters. + +import ( + "fmt" + "strconv" + "strings" + + "github.com/stackql-labs/omnisdk/pkg/query" + + "github.com/stackql/stackql-parser/go/vt/sqlparser" +) + +// starOutput marks a select-list slot the first row's columns fill. +const starOutput = "*" + +// docQuery is a SELECT translated for omnisdk: the query itself, the output +// names in select-list order, and the LIMIT to push down. +type docQuery interface { + getQuery() query.Unresolved + getOutputs() []string + getLimit() int + getBundles() []string +} + +type standardDocQuery struct { + q query.Unresolved + outputs []string + limit int + bundles []string +} + +func newDocQuery(q query.Unresolved, outputs []string, limit int, bundles []string) docQuery { + return &standardDocQuery{q: q, outputs: outputs, limit: limit, bundles: bundles} +} + +func (d *standardDocQuery) getQuery() query.Unresolved { return d.q } + +func (d *standardDocQuery) getOutputs() []string { return d.outputs } + +func (d *standardDocQuery) getLimit() int { return d.limit } + +func (d *standardDocQuery) getBundles() []string { return d.bundles } + +// fromDocProviders reports whether a FROM clause reads the document-driven +// providers, and refuses one mixing them with anything else: omnisdk runs the +// whole query, so every relation in it has to be one it can resolve. +func fromDocProviders(from sqlparser.TableExprs, currentProvider string) (bool, error) { + var docs, others []string + var walk func(expr sqlparser.TableExpr) + walk = func(expr sqlparser.TableExpr) { + switch node := expr.(type) { + case *sqlparser.AliasedTableExpr: + tableName, isName := node.Expr.(sqlparser.TableName) + if !isName { + others = append(others, sqlparser.String(node)) + return + } + if _, isDoc := docProvider( + resolveProvider(tableName.QualifierSecond.GetRawVal(), currentProvider)); isDoc { + docs = append(docs, tableName.Name.GetRawVal()) + return + } + others = append(others, qualifiedName(tableName)) + case *sqlparser.JoinTableExpr: + walk(node.LeftExpr) + walk(node.RightExpr) + case *sqlparser.ParenTableExpr: + for _, inner := range node.Exprs { + walk(inner) + } + default: + others = append(others, sqlparser.String(node)) + } + } + for _, expr := range from { + walk(expr) + } + if len(docs) == 0 { + return false, nil + } + if len(others) > 0 { + return true, fmt.Errorf( + "%s relations cannot be combined with %s in one query", + UnstablePrefix+"*", strings.Join(others, ", ")) + } + return true, nil +} + +// translateSelect builds the omnisdk query for a SELECT over document-driven +// relations. What omnisdk leaves to the caller - ordering, grouping, +// aggregation, de-duplication - is refused until stackql applies it over the +// streamed rows. +func translateSelect(node *sqlparser.Select, currentProvider string) (docQuery, error) { + if unsupported := unsupportedDocClauses(node); len(unsupported) > 0 { + return nil, fmt.Errorf("%s cannot be applied to %s relations; remove %s from the query", + strings.Join(unsupported, ", "), UnstablePrefix+"*", pluralClause(len(unsupported))) + } + limit, err := pushedLimit(node.Limit) + if err != nil { + return nil, err + } + if len(node.From) != 1 { + return nil, fmt.Errorf("a comma-separated FROM cannot be applied to %s relations; use JOIN ... ON", + UnstablePrefix+"*") + } + t := &translator{currentProvider: currentProvider} + if err = t.from(node.From[0], query.Base, nil); err != nil { + return nil, err + } + var where []query.Predicate + if node.Where != nil { + if where, err = conjuncts(node.Where.Expr); err != nil { + return nil, err + } + } + outputs, names, err := selectOutputs(node.SelectExprs) + if err != nil { + return nil, err + } + q, err := query.New(t.joins, where, outputs) + if err != nil { + return nil, err + } + return newDocQuery(q, names, limit, t.bundles), nil +} + +// unsupportedDocClauses names what omnisdk returns unapplied: its row stream is +// unordered, ungrouped and not de-duplicated. LIMIT is pushed down instead. +func unsupportedDocClauses(node *sqlparser.Select) []string { + var out []string + if len(node.OrderBy) > 0 { + out = append(out, "ORDER BY") + } + if len(node.GroupBy) > 0 { + out = append(out, "GROUP BY") + } + if node.Having != nil { + out = append(out, "HAVING") + } + if node.Distinct { + out = append(out, "DISTINCT") + } + return out +} + +// pushedLimit is the row cap omnisdk applies. Without ORDER BY a LIMIT needs +// no sorted input, so it can stop the stream early. +func pushedLimit(limit *sqlparser.Limit) (int, error) { + if limit == nil { + return 0, nil + } + if limit.Offset != nil { + return 0, fmt.Errorf("OFFSET cannot be applied to %s relations", UnstablePrefix+"*") + } + val, isVal := limit.Rowcount.(*sqlparser.SQLVal) + if !isVal || val.Type != sqlparser.IntVal { + return 0, fmt.Errorf("LIMIT %s is not a row count", sqlparser.String(limit.Rowcount)) + } + n, err := strconv.Atoi(string(val.Val)) + if err != nil || n < 1 { + return 0, fmt.Errorf("LIMIT %s is not a positive row count", string(val.Val)) + } + return n, nil +} + +type translator struct { + currentProvider string + joins []query.Join + bundles []string +} + +// from appends a FROM item's joins in source order. A join's right side takes +// the join's form and ON; everything on its left keeps its own. +func (t *translator) from(expr sqlparser.TableExpr, form query.JoinForm, on []query.Predicate) error { + switch node := expr.(type) { + case *sqlparser.AliasedTableExpr: + tableName, isName := node.Expr.(sqlparser.TableName) + if !isName { + return fmt.Errorf("'%s' cannot be read from %s relations", sqlparser.String(node), UnstablePrefix+"*") + } + if len(t.joins) == 0 { + form = query.Base + } + t.joins = append(t.joins, query.NewJoin(t.resource(tableName, node.As), form, on...)) + return nil + case *sqlparser.ParenTableExpr: + if len(node.Exprs) != 1 { + return fmt.Errorf("'%s' cannot be read from %s relations", sqlparser.String(node), UnstablePrefix+"*") + } + return t.from(node.Exprs[0], form, on) + case *sqlparser.JoinTableExpr: + rightForm, err := joinForm(node.Join) + if err != nil { + return err + } + if len(node.Condition.Using) > 0 { + return fmt.Errorf("JOIN ... USING cannot be applied to %s relations; use ON", UnstablePrefix+"*") + } + if err = t.from(node.LeftExpr, form, on); err != nil { + return err + } + var rightOn []query.Predicate + if node.Condition.On != nil { + if rightOn, err = conjuncts(node.Condition.On); err != nil { + return err + } + } + return t.from(node.RightExpr, rightForm, rightOn) + default: + return fmt.Errorf("'%s' cannot be read from %s relations", sqlparser.String(node), UnstablePrefix+"*") + } +} + +// resource addresses a relation by its registry handle, aliased as written or +// else by its own name, and records the bundle whose credential it runs under. +func (t *translator) resource(tableName sqlparser.TableName, as sqlparser.TableIdent) query.Resource { + bundle, _ := docProvider(resolveProvider(tableName.QualifierSecond.GetRawVal(), t.currentProvider)) + t.bundles = append(t.bundles, bundle) + alias := as.GetRawVal() + if alias == "" { + alias = tableName.Name.GetRawVal() + } + address := fmt.Sprintf("%s%s.%s.%s", UnstablePrefix, bundle, + tableName.Qualifier.GetRawVal(), tableName.Name.GetRawVal()) + return query.NewResource(alias, address) +} + +func qualifiedName(tableName sqlparser.TableName) string { + var parts []string + for _, part := range []string{ + tableName.QualifierSecond.GetRawVal(), tableName.Qualifier.GetRawVal(), tableName.Name.GetRawVal(), + } { + if part != "" { + parts = append(parts, part) + } + } + return strings.Join(parts, ".") +} + +func joinForm(join string) (query.JoinForm, error) { + switch strings.ToLower(join) { + case sqlparser.JoinStr: + return query.Inner, nil + case sqlparser.LeftJoinStr, sqlparser.LeftOuterJoinStr: + return query.Left, nil + default: + return query.Base, fmt.Errorf("%s cannot be applied to %s relations", + strings.ToUpper(join), UnstablePrefix+"*") + } +} + +// conjuncts splits a condition at its top-level ANDs. +func conjuncts(expr sqlparser.Expr) ([]query.Predicate, error) { + if and, isAnd := expr.(*sqlparser.AndExpr); isAnd { + left, err := conjuncts(and.Left) + if err != nil { + return nil, err + } + right, err := conjuncts(and.Right) + if err != nil { + return nil, err + } + return append(left, right...), nil + } + p, err := predicate(expr) + if err != nil { + return nil, err + } + return []query.Predicate{p}, nil +} + +var compareOps = map[string]query.CompareOp{ //nolint:gochecknoglobals // fixed mapping + sqlparser.EqualStr: query.Eq, + sqlparser.NotEqualStr: query.Ne, + sqlparser.LessThanStr: query.Lt, + sqlparser.LessEqualStr: query.Le, + sqlparser.GreaterThanStr: query.Gt, + sqlparser.GreaterEqualStr: query.Ge, +} + +func predicate(expr sqlparser.Expr) (query.Predicate, error) { + switch node := expr.(type) { + case *sqlparser.AndExpr: + parts, err := conjuncts(node) + if err != nil { + return nil, err + } + // Under OR or NOT an AND is a single predicate: every part must hold. + negated := make([]query.Predicate, 0, len(parts)) + for _, part := range parts { + negated = append(negated, query.NewNot(part)) + } + return query.NewNot(query.NewOr(negated...)), nil + case *sqlparser.OrExpr: + left, err := predicate(node.Left) + if err != nil { + return nil, err + } + right, err := predicate(node.Right) + if err != nil { + return nil, err + } + return query.NewOr(left, right), nil + case *sqlparser.NotExpr: + inner, err := predicate(node.Expr) + if err != nil { + return nil, err + } + return query.NewNot(inner), nil + case *sqlparser.ComparisonExpr: + return comparison(node) + case *sqlparser.FuncExpr: + call, err := expression(node) + if err != nil { + return nil, err + } + return query.NewTest(call), nil + default: + return nil, fmt.Errorf("condition '%s' cannot be applied to %s relations", + sqlparser.String(expr), UnstablePrefix+"*") + } +} + +func comparison(node *sqlparser.ComparisonExpr) (query.Predicate, error) { + left, err := expression(node.Left) + if err != nil { + return nil, err + } + right, err := expression(node.Right) + if err != nil { + return nil, err + } + switch node.Operator { + case sqlparser.InStr: + return query.NewIn(left, right), nil + case sqlparser.NotInStr: + return query.NewNot(query.NewIn(left, right)), nil + } + op, known := compareOps[node.Operator] + if !known { + return nil, fmt.Errorf("condition '%s' cannot be applied to %s relations", + sqlparser.String(node), UnstablePrefix+"*") + } + return query.NewCompare(op, left, right), nil +} + +func expression(expr sqlparser.Expr) (query.Expr, error) { + switch node := expr.(type) { + case *sqlparser.ColName: + return query.NewColumn(node.Qualifier.Name.GetRawVal(), node.Name.GetRawVal()), nil + case *sqlparser.SQLVal: + return literal(node) + case sqlparser.BoolVal: + return query.NewLiteral(bool(node)), nil + case *sqlparser.NullVal: + return query.NewLiteral(nil), nil + case sqlparser.ValTuple: + items := make([]query.Expr, 0, len(node)) + for _, item := range node { + translated, err := expression(item) + if err != nil { + return nil, err + } + items = append(items, translated) + } + return query.NewCollection(items...), nil + case *sqlparser.FuncExpr: + if node.IsAggregate() || node.Distinct || node.Over != nil { + return nil, fmt.Errorf("'%s' cannot be applied to %s relations", + sqlparser.String(node), UnstablePrefix+"*") + } + args := make([]query.Expr, 0, len(node.Exprs)) + for _, arg := range node.Exprs { + aliased, isAliased := arg.(*sqlparser.AliasedExpr) + if !isAliased { + return nil, fmt.Errorf("argument '%s' cannot be applied to %s relations", + sqlparser.String(arg), UnstablePrefix+"*") + } + translated, err := expression(aliased.Expr) + if err != nil { + return nil, err + } + args = append(args, translated) + } + return query.NewCall(node.Name.Lowered(), args...), nil + default: + return nil, fmt.Errorf("'%s' cannot be applied to %s relations", + sqlparser.String(expr), UnstablePrefix+"*") + } +} + +func literal(val *sqlparser.SQLVal) (query.Expr, error) { + switch val.Type { //nolint:exhaustive // hex, bit and bind values have no omnisdk literal + case sqlparser.StrVal: + return query.NewLiteral(string(val.Val)), nil + case sqlparser.IntVal: + n, err := strconv.ParseInt(string(val.Val), 10, 64) + if err != nil { + return nil, err + } + return query.NewLiteral(n), nil + case sqlparser.FloatVal: + f, err := strconv.ParseFloat(string(val.Val), 64) + if err != nil { + return nil, err + } + return query.NewLiteral(f), nil + default: + return nil, fmt.Errorf("'%s' cannot be applied to %s relations", + sqlparser.String(val), UnstablePrefix+"*") + } +} + +// selectOutputs names every output: an alias, else a bare column's own name, +// else the expression as written. A star is left unnamed for omnisdk to expand. +func selectOutputs(exprs sqlparser.SelectExprs) ([]query.Output, []string, error) { + outputs := make([]query.Output, 0, len(exprs)) + names := make([]string, 0, len(exprs)) + for _, expr := range exprs { + switch node := expr.(type) { + case *sqlparser.StarExpr: + outputs = append(outputs, query.NewOutput("", query.NewStar(node.TableName.Name.GetRawVal()))) + names = append(names, starOutput) + case *sqlparser.AliasedExpr: + translated, err := expression(node.Expr) + if err != nil { + return nil, nil, err + } + name := node.As.GetRawVal() + if name == "" { + if col, isCol := node.Expr.(*sqlparser.ColName); isCol { + name = col.Name.GetRawVal() + } else { + name = sqlparser.String(node.Expr) + } + } + outputs = append(outputs, query.NewOutput(name, translated)) + names = append(names, name) + default: + return nil, nil, fmt.Errorf("'%s' cannot be applied to %s relations", + sqlparser.String(expr), UnstablePrefix+"*") + } + } + return outputs, names, nil +} + +// mutationTables is every relation a mutation reads or writes: the target +// first, then its sources. +func mutationTables(stmt sqlparser.Statement) (sqlparser.TableExprs, bool) { + switch node := stmt.(type) { + case *sqlparser.Insert: + tables := sqlparser.TableExprs{&sqlparser.AliasedTableExpr{Expr: node.Table}} + if sel, isSelect := node.Rows.(*sqlparser.Select); isSelect { + tables = append(tables, sel.From...) + } + return tables, true + case *sqlparser.Update: + return append(append(sqlparser.TableExprs{}, node.TableExprs...), updateSources(node)...), true + case *sqlparser.Delete: + return node.TableExprs, true + default: + return nil, false + } +} + +// translateMutation builds the omnisdk mutation for an INSERT, UPDATE or +// DELETE whose target is a document-driven relation. A RETURNING list becomes +// the mutation's outputs. +func translateMutation(stmt sqlparser.Statement, currentProvider string) (docQuery, error) { + t := &translator{currentProvider: currentProvider} + var ( + target query.Target + where []query.Predicate + returning sqlparser.SelectExprs + err error + ) + switch node := stmt.(type) { + case *sqlparser.Insert: + target, where, err = t.insert(node) + returning = node.SelectExprs + case *sqlparser.Update: + target, where, err = t.update(node) + returning = node.SelectExprs + case *sqlparser.Delete: + target, where, err = t.delete(node) + returning = node.SelectExprs + default: + return nil, fmt.Errorf("'%s' cannot be applied to %s relations", sqlparser.String(stmt), UnstablePrefix+"*") + } + if err != nil { + return nil, err + } + var outputs []query.Output + var names []string + if len(returning) > 0 { + if outputs, names, err = selectOutputs(returning); err != nil { + return nil, err + } + } + q, err := query.NewMutation(target, t.joins, where, outputs) + if err != nil { + return nil, err + } + return newDocQuery(q, names, 0, t.bundles), nil +} + +// insert takes its values from a single VALUES row, or from a SELECT whose +// relations become the mutation's sources and whose list is matched to the +// columns by position. +func (t *translator) insert(node *sqlparser.Insert) (query.Target, []query.Predicate, error) { + switch { + case !strings.EqualFold(node.Action, sqlparser.InsertStr): + return nil, nil, fmt.Errorf("%s cannot be applied to %s relations; use INSERT", + strings.ToUpper(node.Action), UnstablePrefix+"*") + case node.Ignore != "" || len(node.OnDup) > 0: + return nil, nil, fmt.Errorf("INSERT IGNORE and ON DUPLICATE KEY cannot be applied to %s relations", + UnstablePrefix+"*") + case len(node.Columns) == 0: + return nil, nil, fmt.Errorf("an INSERT into %s relations must name its columns", UnstablePrefix+"*") + } + resource := t.resource(node.Table, sqlparser.NewTableIdent("")) + var values []sqlparser.Expr + var where []query.Predicate + switch rows := node.Rows.(type) { + case sqlparser.Values: + if len(rows) != 1 { + return nil, nil, fmt.Errorf("an INSERT into %s relations takes exactly one VALUES row; got %d", + UnstablePrefix+"*", len(rows)) + } + values = rows[0] + case *sqlparser.Select: + var err error + if values, where, err = t.insertSelect(rows); err != nil { + return nil, nil, err + } + default: + return nil, nil, fmt.Errorf("'%s' cannot be inserted into %s relations", + sqlparser.String(node.Rows), UnstablePrefix+"*") + } + if len(values) != len(node.Columns) { + return nil, nil, fmt.Errorf("INSERT names %d columns but supplies %d values", len(node.Columns), len(values)) + } + assignments := make([]query.Assignment, 0, len(values)) + for i, value := range values { + translated, err := expression(value) + if err != nil { + return nil, nil, err + } + assignments = append(assignments, query.NewAssignment(node.Columns[i].GetRawVal(), translated)) + } + return query.NewInsert(resource, assignments...), where, nil +} + +// insertSelect reads an INSERT ... SELECT: its relations become the mutation's +// sources, and its list supplies the values, matched to the columns by position. +func (t *translator) insertSelect(rows *sqlparser.Select) ([]sqlparser.Expr, []query.Predicate, error) { + if unsupported := unsupportedDocClauses(rows); len(unsupported) > 0 || rows.Limit != nil { + return nil, nil, fmt.Errorf("an INSERT ... SELECT into %s relations takes no "+ + "ORDER BY, GROUP BY, HAVING, DISTINCT or LIMIT", UnstablePrefix+"*") + } + if len(rows.From) != 1 { + return nil, nil, fmt.Errorf("a comma-separated FROM cannot be applied to %s relations; use JOIN ... ON", + UnstablePrefix+"*") + } + if err := t.from(rows.From[0], query.Base, nil); err != nil { + return nil, nil, err + } + where, err := whereConjuncts(rows.Where) + if err != nil { + return nil, nil, err + } + values := make([]sqlparser.Expr, 0, len(rows.SelectExprs)) + for _, expr := range rows.SelectExprs { + aliased, isAliased := expr.(*sqlparser.AliasedExpr) + if !isAliased { + return nil, nil, fmt.Errorf("'%s' cannot be inserted into %s relations; name each value", + sqlparser.String(expr), UnstablePrefix+"*") + } + values = append(values, aliased.Expr) + } + return values, where, nil +} + +func (t *translator) update(node *sqlparser.Update) (query.Target, []query.Predicate, error) { + if !strings.EqualFold(node.Action, sqlparser.UpdateStr) { + return nil, nil, fmt.Errorf("%s cannot be applied to %s relations; use UPDATE", + strings.ToUpper(node.Action), UnstablePrefix+"*") + } + if len(node.OrderBy) > 0 || node.Limit != nil { + return nil, nil, fmt.Errorf("ORDER BY and LIMIT cannot be applied to an UPDATE of %s relations", + UnstablePrefix+"*") + } + resource, err := t.targetResource(node.TableExprs) + if err != nil { + return nil, nil, err + } + if err = t.sources(updateSources(node)); err != nil { + return nil, nil, err + } + assignments := make([]query.Assignment, 0, len(node.Exprs)) + for _, set := range node.Exprs { + translated, translateErr := expression(set.Expr) + if translateErr != nil { + return nil, nil, translateErr + } + assignments = append(assignments, query.NewAssignment(set.Name.Name.GetRawVal(), translated)) + } + where, err := whereConjuncts(node.Where) + if err != nil { + return nil, nil, err + } + return query.NewUpdate(resource, assignments...), where, nil +} + +func (t *translator) delete(node *sqlparser.Delete) (query.Target, []query.Predicate, error) { + switch { + case len(node.Targets) > 0: + return nil, nil, fmt.Errorf("a multi-table DELETE cannot be applied to %s relations", UnstablePrefix+"*") + case len(node.OrderBy) > 0 || node.Limit != nil: + return nil, nil, fmt.Errorf("ORDER BY and LIMIT cannot be applied to a DELETE of %s relations", + UnstablePrefix+"*") + } + resource, err := t.targetResource(node.TableExprs) + if err != nil { + return nil, nil, err + } + where, err := whereConjuncts(node.Where) + if err != nil { + return nil, nil, err + } + return query.NewDelete(resource), where, nil +} + +// targetResource is the single relation an UPDATE or DELETE writes. +func (t *translator) targetResource(exprs sqlparser.TableExprs) (query.Resource, error) { + if len(exprs) == 1 { + if aliased, isAliased := exprs[0].(*sqlparser.AliasedTableExpr); isAliased { + if tableName, isName := aliased.Expr.(sqlparser.TableName); isName { + return t.resource(tableName, aliased.As), nil + } + } + } + return nil, fmt.Errorf("'%s' cannot be written as one %s relation", sqlparser.String(exprs), UnstablePrefix+"*") +} + +// updateSources is an UPDATE's FROM clause. The parser stands in "dual" for an +// absent one, which reads nothing. +func updateSources(node *sqlparser.Update) sqlparser.TableExprs { + if len(node.From) == 1 { + if aliased, isAliased := node.From[0].(*sqlparser.AliasedTableExpr); isAliased { + if tableName, isName := aliased.Expr.(sqlparser.TableName); isName && + tableName.Qualifier.IsEmpty() && tableName.Name.GetRawVal() == "dual" { + return nil + } + } + } + return node.From +} + +// sources reads an UPDATE ... FROM clause as the mutation's source joins. +func (t *translator) sources(from sqlparser.TableExprs) error { + switch len(from) { + case 0: + return nil + case 1: + return t.from(from[0], query.Base, nil) + default: + return fmt.Errorf("a comma-separated FROM cannot be applied to %s relations; use JOIN ... ON", + UnstablePrefix+"*") + } +} + +func whereConjuncts(where *sqlparser.Where) ([]query.Predicate, error) { + if where == nil { + return nil, nil + } + return conjuncts(where.Expr) +} diff --git a/internal/stackql/intrinsic/translate_test.go b/internal/stackql/intrinsic/translate_test.go new file mode 100644 index 00000000..b6af24db --- /dev/null +++ b/internal/stackql/intrinsic/translate_test.go @@ -0,0 +1,299 @@ +package intrinsic //nolint:testpackage // tests unexported translation + +import ( + "fmt" + "path/filepath" + "reflect" + "testing" + + "github.com/stackql-labs/omnisdk/pkg/omnisdk" + "github.com/stackql-labs/omnisdk/pkg/query" + + "github.com/stackql/stackql-parser/go/vt/sqlparser" +) + +func parseSelect(t *testing.T, sql string) *sqlparser.Select { + t.Helper() + stmt, err := sqlparser.Parse(sql) + if err != nil { + t.Fatalf("parse %q: %v", sql, err) + } + sel, isSelect := stmt.(*sqlparser.Select) + if !isSelect { + t.Fatalf("%q is not a select", sql) + } + return sel +} + +func TestTranslateSelectJoins(t *testing.T) { + withUnstable(t, true) + sel := parseSelect(t, "select k.name as ring, c.name from stackql_unstable_google.cloudkms.key_rings k "+ + "inner join stackql_unstable_google.cloudkms.crypto_keys c on c.keyRingsId = split_part(k.name, '/', 6) "+ + "left join stackql_unstable_google.cloudkms.crypto_keys c2 on c2.name = c.name "+ + "where k.projectsId = 'p' and k.locationsId = 'global' limit 5") + dq, err := translateSelect(sel, "") + if err != nil { + t.Fatalf("translate: %v", err) + } + type joinShape struct { + alias, handle string + form query.JoinForm + on int + } + var got []joinShape + for _, j := range dq.getQuery().From() { + got = append(got, joinShape{j.Resource().Alias(), j.Resource().Handle(), j.Form(), len(j.On())}) + } + want := []joinShape{ + {"k", "stackql_unstable_google.cloudkms.key_rings", query.Base, 0}, + {"c", "stackql_unstable_google.cloudkms.crypto_keys", query.Inner, 1}, + {"c2", "stackql_unstable_google.cloudkms.crypto_keys", query.Left, 1}, + } + if !reflect.DeepEqual(got, want) { + t.Fatalf("joins: got %+v, want %+v", got, want) + } + if n := len(dq.getQuery().Where()); n != 2 { + t.Fatalf("where conjuncts: got %d, want 2", n) + } + if !reflect.DeepEqual(dq.getOutputs(), []string{"ring", "name"}) { + t.Fatalf("outputs: got %v", dq.getOutputs()) + } + if dq.getLimit() != 5 { + t.Fatalf("limit: got %d, want 5", dq.getLimit()) + } +} + +func TestTranslateSelectResolvesAgainstRegistry(t *testing.T) { + withUnstable(t, true) + sel := parseSelect(t, "select login from stackql_unstable_fixture.orgs.members "+ + "where org = 'dummyorg' and (type = 'User' or not id = 2) and login in ('a', 'b')") + dq, err := translateSelect(sel, "") + if err != nil { + t.Fatalf("translate: %v", err) + } + tables := map[string]omnisdk.Table{} + for _, j := range dq.getQuery().From() { + tbl, describeErr := omnisdk.DescribeTable(filepath.Join("testdata", "registry"), j.Resource().Handle()) + if describeErr != nil { + t.Fatalf("describe: %v", describeErr) + } + tables[j.Resource().Alias()] = tbl + } + res, err := omnisdk.Resolve(dq.getQuery(), tables) + if err != nil { + t.Fatalf("resolve: %v", err) + } + // org is the method's path parameter, so it binds to the node; the OR and + // the IN list are row filters. + nodes := res.Graph().Nodes() + if len(nodes) != 1 || !reflect.DeepEqual(nodes[0].Params(), map[string]string{"org": "dummyorg"}) { + t.Fatalf("nodes: got %d, first params %v", len(nodes), nodes[0].Params()) + } + if len(res.Params()) != 0 { + t.Fatalf("query-wide params: got %v, want none", res.Params()) + } + if n := len(res.Graph().Filters()); n != 2 { + t.Fatalf("filters: got %d, want 2", n) + } +} + +func TestTranslateSelectRefusals(t *testing.T) { + withUnstable(t, true) + for sql, want := range map[string]string{ + "select login from stackql_unstable_github.orgs.members order by login": "ORDER BY cannot be applied to " + + "stackql_unstable_* relations; remove it from the query", + "select distinct login from stackql_unstable_github.orgs.members group by login": "GROUP BY, DISTINCT cannot " + + "be applied to stackql_unstable_* relations; remove them from the query", + "select count(*) from stackql_unstable_github.orgs.members": "'count(*)' cannot be applied to " + + "stackql_unstable_* relations", + "select login from stackql_unstable_github.orgs.members limit 1, 2": "OFFSET cannot be applied to " + + "stackql_unstable_* relations", + "select login from stackql_unstable_github.orgs.members where login like 'a%'": "condition " + + "'`login` like 'a%'' cannot be applied to stackql_unstable_* relations", + "select a.login from stackql_unstable_github.orgs.members a, stackql_unstable_github.orgs.members b": "a " + + "comma-separated FROM cannot be applied to stackql_unstable_* relations; use JOIN ... ON", + "select a.login from stackql_unstable_github.orgs.members a right join " + + "stackql_unstable_github.orgs.members b on a.login = b.login": "RIGHT JOIN cannot be applied to " + + "stackql_unstable_* relations", + } { + _, err := translateSelect(parseSelect(t, sql), "") + if err == nil || err.Error() != want { + t.Errorf("%s:\n got %v\nwant %s", sql, err, want) + } + } +} + +func TestFromDocProvidersRefusesMixing(t *testing.T) { + withUnstable(t, true) + sel := parseSelect(t, "select a.login from stackql_unstable_github.orgs.members a "+ + "inner join github.orgs.members b on a.login = b.login") + isDoc, err := fromDocProviders(sel.From, "") + if !isDoc { + t.Fatalf("expected the query to be claimed") + } + want := "stackql_unstable_* relations cannot be combined with github.orgs.members in one query" + if err == nil || err.Error() != want { + t.Fatalf("got %v, want %s", err, want) + } + if isDoc, _ := fromDocProviders(parseSelect(t, "select 1 from github.orgs.members").From, ""); isDoc { + t.Fatalf("a registry-only query must not be claimed") + } +} + +func TestOutputColumns(t *testing.T) { + batch := []omnisdk.Row{{"b": 1, "a": 2, "x": 3}} + cols := outputColumns([]string{"x", starOutput}, batch) + var got []string + for _, col := range cols { + got = append(got, col.name) + } + if !reflect.DeepEqual(got, []string{"x", "a", "b"}) { + t.Fatalf("got %v", got) + } + if outputColumns([]string{starOutput}, nil) != nil { + t.Fatalf("a star waits for a row") + } + if n := len(outputColumns([]string{"x", "y"}, nil)); n != 2 { + t.Fatalf("named outputs are known before any row: got %d", n) + } +} + +func parseStatement(t *testing.T, sql string) sqlparser.Statement { + t.Helper() + stmt, err := sqlparser.Parse(sql) + if err != nil { + t.Fatalf("parse %q: %v", sql, err) + } + return stmt +} + +type assignmentShape struct { + column string + value string +} + +func assignmentShapes(target query.Target) []assignmentShape { + var out []assignmentShape + for _, a := range target.Set() { + value := "?" + switch v := a.Value().(type) { + case query.Literal: + value = fmt.Sprint(v.Value()) + case query.Column: + value = v.Qualifier() + "." + v.Name() + case query.Call: + value = v.Func() + "()" + } + out = append(out, assignmentShape{a.Column(), value}) + } + return out +} + +func TestTranslateMutations(t *testing.T) { + withUnstable(t, true) + for _, tc := range []struct { + sql string + verb query.Verb + alias string + assignments []assignmentShape + sources []string + where int + outputs []string + }{ + { + sql: "insert into stackql_unstable_google.cloudkms.key_rings (projectsId, locationsId, keyRingId) " + + "values ('p', 'global', 'r') returning name, createTime", + verb: query.Insert, + alias: "key_rings", + assignments: []assignmentShape{ + {"projectsId", "p"}, {"locationsId", "global"}, {"keyRingId", "r"}, + }, + outputs: []string{"name", "createTime"}, + }, + { + sql: "insert into stackql_unstable_google.cloudkms.crypto_keys (projectsId, locationsId, keyRingsId) " + + "select 'p', 'global', split_part(k.name, '/', 6) from stackql_unstable_google.cloudkms.key_rings k " + + "where k.projectsId = 'p' and k.locationsId = 'global'", + verb: query.Insert, + alias: "crypto_keys", + assignments: []assignmentShape{ + {"projectsId", "p"}, {"locationsId", "global"}, {"keyRingsId", "split_part()"}, + }, + sources: []string{"k"}, + where: 2, + }, + { + sql: "update stackql_unstable_github.orgs.orgs o set description = 'd' " + + "where o.org = 'dummyorg' returning login", + verb: query.Update, + alias: "o", + assignments: []assignmentShape{{"description", "d"}}, + where: 1, + outputs: []string{"login"}, + }, + { + sql: "delete from stackql_unstable_google.compute.firewalls where project = 'p' and firewall in ('a', 'b')", + verb: query.Delete, + alias: "firewalls", + where: 2, + }, + } { + dq, err := translateMutation(parseStatement(t, tc.sql), "") + if err != nil { + t.Errorf("%s: %v", tc.sql, err) + continue + } + target := dq.getQuery().Target() + if target.Verb() != tc.verb || target.Resource().Alias() != tc.alias { + t.Errorf("%s: target %s %s", tc.sql, target.Verb(), target.Resource().Alias()) + } + if got := assignmentShapes(target); !reflect.DeepEqual(got, tc.assignments) { + t.Errorf("%s: assignments %v, want %v", tc.sql, got, tc.assignments) + } + var sources []string + for _, j := range dq.getQuery().From() { + sources = append(sources, j.Resource().Alias()) + } + if !reflect.DeepEqual(sources, tc.sources) { + t.Errorf("%s: sources %v, want %v", tc.sql, sources, tc.sources) + } + if n := len(dq.getQuery().Where()); n != tc.where { + t.Errorf("%s: %d where conjuncts, want %d", tc.sql, n, tc.where) + } + if !reflect.DeepEqual(dq.getOutputs(), tc.outputs) { + t.Errorf("%s: outputs %v, want %v", tc.sql, dq.getOutputs(), tc.outputs) + } + } +} + +func TestTranslateMutationRefusals(t *testing.T) { + withUnstable(t, true) + for sql, want := range map[string]string{ + "insert into stackql_unstable_google.cloudkms.key_rings (projectsId) values ('a'), ('b')": "an INSERT " + + "into stackql_unstable_* relations takes exactly one VALUES row; got 2", + "insert into stackql_unstable_google.cloudkms.key_rings (projectsId, locationsId) values ('a')": "INSERT " + + "names 2 columns but supplies 1 values", + "update stackql_unstable_github.orgs.orgs set description = 'd' where org = 'x' limit 1": "ORDER BY and " + + "LIMIT cannot be applied to an UPDATE of stackql_unstable_* relations", + "delete from stackql_unstable_google.compute.firewalls where project = 'p' limit 1": "ORDER BY and " + + "LIMIT cannot be applied to a DELETE of stackql_unstable_* relations", + } { + _, err := translateMutation(parseStatement(t, sql), "") + if err == nil || err.Error() != want { + t.Errorf("%s:\n got %v\nwant %s", sql, err, want) + } + } +} + +func TestMutationTablesIgnoresImplicitDual(t *testing.T) { + withUnstable(t, true) + tables, isMutation := mutationTables(parseStatement(t, + "update stackql_unstable_github.orgs.orgs set description = 'd' where org = 'x'")) + if !isMutation || len(tables) != 1 { + t.Fatalf("got %d tables (mutation %v), want the target alone", len(tables), isMutation) + } + isDoc, err := fromDocProviders(tables, "") + if !isDoc || err != nil { + t.Fatalf("got isDoc %v, err %v", isDoc, err) + } +} diff --git a/internal/stackql/planbuilder/entrypoint.go b/internal/stackql/planbuilder/entrypoint.go index 08882ad6..5dbb9a01 100644 --- a/internal/stackql/planbuilder/entrypoint.go +++ b/internal/stackql/planbuilder/entrypoint.go @@ -107,25 +107,24 @@ func (pb *standardPlanBuilder) BuildPlanFromContext(handlerCtx handler.HandlerCo // An omnisdk data relation streams its rows straight to the output writer, // so it is not backed by a view and analysis would try, and fail, to resolve // it as a registry-backed relation. Plan it here, ahead of that analysis. - if sel, isSelect := statement.(*sqlparser.Select); isSelect { - if executor, isStream := intrinsic.GenerateStreamFunc(handlerCtx, sel); isStream { - qPlan.SetType(sqlparser.StmtSelect) - qPlan.SetReadOnly(true) - qPlan.SetCacheable(false) - qPlan.SetStatement(statement) - pGBuilder.getPlanGraphHolder().CreatePrimitiveNode( - primitive.NewLocalPrimitive( - func(_ primitive.IPrimitiveCtx) internaldto.ExecutorOutput { - return executor() - }, - ), - ) - qPlan.SetInstructions(pGBuilder.getPlanGraphHolder()) - if optimiseErr := qPlan.GetInstructions().GetPrimitiveGraph().Optimise(); optimiseErr != nil { - return createErroneousPlan(handlerCtx, qPlan, rowSort, optimiseErr) - } - return qPlan, nil + if executor, isStream := intrinsic.GenerateStreamFunc(handlerCtx, statement); isStream { + stmtType := sqlparser.ASTToStatementType(statement) + qPlan.SetType(stmtType) + qPlan.SetReadOnly(stmtType == sqlparser.StmtSelect) + qPlan.SetCacheable(false) + qPlan.SetStatement(statement) + pGBuilder.getPlanGraphHolder().CreatePrimitiveNode( + primitive.NewLocalPrimitive( + func(_ primitive.IPrimitiveCtx) internaldto.ExecutorOutput { + return executor() + }, + ), + ) + qPlan.SetInstructions(pGBuilder.getPlanGraphHolder()) + if optimiseErr := qPlan.GetInstructions().GetPrimitiveGraph().Optimise(); optimiseErr != nil { + return createErroneousPlan(handlerCtx, qPlan, rowSort, optimiseErr) } + return qPlan, nil } primitiveGenerator := primitivegenerator.NewRootPrimitiveGenerator( diff --git a/test/python/stackql_test_tooling/flask/gcp/app.py b/test/python/stackql_test_tooling/flask/gcp/app.py index 137d7dc9..f2bb0681 100644 --- a/test/python/stackql_test_tooling/flask/gcp/app.py +++ b/test/python/stackql_test_tooling/flask/gcp/app.py @@ -121,6 +121,16 @@ def projects_testing_project_global_network_detail(project_name: str, network_na network_name=network_name ), 200, {'Content-Type': 'application/json'} +@app.route('/v1/projects//locations//keyRings', methods=['POST']) +def v1_projects_locations_keyRings_create(project_name: str, location_name: str): + key_ring_id = request.args.get('keyRingId') + if not key_ring_id: + return '{"msg": "Invalid request: keyRingId not supplied"}', 400, {'Content-Type': 'application/json'} + return jsonify({ + 'name': f'projects/{project_name}/locations/{location_name}/keyRings/{key_ring_id}', + 'createTime': '2022-02-02T02:02:02.02000000Z', + }), 200 + @app.route('/v1/projects/testing-project-three/locations/global/keyRings/testing-three/cryptoKeys', methods=['GET']) def v1_projects_testing_project_three_locations_global_keyRings_testing_three_cryptoKeys(): return render_template('route_1_template.json'), 200, {'Content-Type': 'application/json'} diff --git a/test/robot/functional/stackql_mocked_from_cmd_line.robot b/test/robot/functional/stackql_mocked_from_cmd_line.robot index 603eb284..c23d8225 100644 --- a/test/robot/functional/stackql_mocked_from_cmd_line.robot +++ b/test/robot/functional/stackql_mocked_from_cmd_line.robot @@ -11180,6 +11180,179 @@ Unstable Google Kms Key Rings Jsonl Row Set Matches Expectation ... stdout=${CURDIR}${/}tmp${/}Unstable-Google-Kms-Key-Rings.tmp ... stderr=${CURDIR}${/}tmp${/}Unstable-Google-Kms-Key-Rings-stderr.tmp +Unstable Google Kms Key Rings Joined To Crypto Keys Jsonl Row Set Matches Expectation + [Documentation] A join over document-driven relations, resolved by + ... omnisdk: the ON condition feeds each key ring's id into + ... the crypto keys request, so the join runs as an edge + ... between the two exchanges rather than a row filter. + [Setup] Write Gcp Service Account ${OMNISDK_MOCK_GCP_SA_HOST} + [Teardown] Remove Preview Mock Environment + ${gcp_sa} = Set Variable If "${EXECUTION_PLATFORM}" == "docker" + ... /opt/test/tmp/omnisdk-gcp-sa.json ${OMNISDK_MOCK_GCP_SA_HOST} + Set Environment Variable GOOGLE_APPLICATION_CREDENTIALS ${gcp_sa} + ${preview} = Catenate SEPARATOR= + ... {"endpoint":"https://${LOCAL_HOST_ALIAS}:${MOCKSERVER_PORT_GOOGLE}", + ... "insecureSkipTLSVerify":true,"unstable":true} + ${auth} = Set Variable {"google":{"credentialsfilepath":"${gcp_sa}"}} + ${query} = Catenate SEPARATOR=${SPACE} + ... select k.name as key_ring, c.name as key_name + ... from stackql_unstable_google.cloudkms.key_rings k + ... inner join stackql_unstable_google.cloudkms.crypto_keys c + ... on c.keyRingsId = split_part(k.name, '/', 6) + ... where k.projectsId = 'testing-project' and k.locationsId = 'global' + ... and c.projectsId = 'testing-project' and c.locationsId = 'global'; + ${expected} = Catenate SEPARATOR= + ... {"key_ring":"projects/testing-project/locations/global/keyRings/testing", + ... "key_name":"projects/testing-project/locations/global/keyRings/testing/cryptoKeys/testing-demo-key"} + Should StackQL Exec Inline Jsonl Set Equal + ... ${STACKQL_EXE} + ... ${OKTA_SECRET_STR} + ... ${GITHUB_SECRET_STR} + ... ${K8S_SECRET_STR} + ... ${REGISTRY_NO_VERIFY_CFG_STR} + ... ${auth} + ... ${SQL_BACKEND_CFG_STR_CANONICAL} + ... ${query} + ... ${expected} + ... --preview\=${preview} + ... stdout=${CURDIR}${/}tmp${/}Unstable-Google-Kms-Key-Rings-Joined-To-Crypto-Keys.tmp + ... stderr=${CURDIR}${/}tmp${/}Unstable-Google-Kms-Key-Rings-Joined-To-Crypto-Keys-stderr.tmp + +Unstable Github Org Members Filtered By In List Jsonl Row Set Matches Expectation + [Documentation] A condition no request parameter can carry is applied by + ... omnisdk as a filter on the streamed rows. + [Teardown] Remove Preview Mock Environment + ${preview} = Catenate SEPARATOR= + ... {"endpoint":"https://${LOCAL_HOST_ALIAS}:${MOCKSERVER_PORT_GITHUB}", + ... "insecureSkipTLSVerify":true,"unstable":true} + ${expected} = Catenate SEPARATOR=\n + ... {"login":"some-jimbo-7","type":"User"} + ... {"login":"some-jimbo-3","type":"User"} + ${query} = Catenate SEPARATOR=${SPACE} + ... select login, type from stackql_unstable_github.orgs.members + ... where org = 'dummyorg' and login in ('some-jimbo-3', 'some-jimbo-7'); + Should StackQL Exec Inline Jsonl Set Equal + ... ${STACKQL_EXE} + ... ${OKTA_SECRET_STR} + ... ${GITHUB_SECRET_STR} + ... ${K8S_SECRET_STR} + ... ${REGISTRY_NO_VERIFY_CFG_STR} + ... ${AUTH_CFG_STR} + ... ${SQL_BACKEND_CFG_STR_CANONICAL} + ... ${query} + ... ${expected} + ... --preview\=${preview} + ... stdout=${CURDIR}${/}tmp${/}Unstable-Github-Org-Members-Filtered-By-In-List.tmp + ... stderr=${CURDIR}${/}tmp${/}Unstable-Github-Org-Members-Filtered-By-In-List-stderr.tmp + +Unstable Github Org Update Reports Despatch + [Documentation] A document-driven UPDATE without RETURNING: omnisdk sends + ... the effect and stackql reports it. The mock refuses any + ... other description, so success shows the SET value reached + ... the request body. + [Teardown] Remove Preview Mock Environment + ${preview} = Catenate SEPARATOR= + ... {"endpoint":"https://${LOCAL_HOST_ALIAS}:${MOCKSERVER_PORT_GITHUB}", + ... "insecureSkipTLSVerify":true,"unstable":true} + ${query} = Catenate SEPARATOR=${SPACE} + ... update stackql_unstable_github.orgs.orgs + ... set description = 'Some silly description.' + ... where org = 'dummyorg'; + Should StackQL Exec Inline Equal Stderr + ... ${STACKQL_EXE} + ... ${OKTA_SECRET_STR} + ... ${GITHUB_SECRET_STR} + ... ${K8S_SECRET_STR} + ... ${REGISTRY_NO_VERIFY_CFG_STR} + ... ${AUTH_CFG_STR} + ... ${SQL_BACKEND_CFG_STR_CANONICAL} + ... ${query} + ... The operation was despatched successfully + ... --preview\=${preview} + ... stdout=${CURDIR}${/}tmp${/}Unstable-Github-Org-Update-Reports-Despatch.tmp + ... stderr=${CURDIR}${/}tmp${/}Unstable-Github-Org-Update-Reports-Despatch-stderr.tmp + +Unstable Github Org Update Returning Jsonl Row Set Matches Expectation + [Documentation] A document-driven UPDATE with RETURNING streams back the + ... updated resource. + [Teardown] Remove Preview Mock Environment + ${preview} = Catenate SEPARATOR= + ... {"endpoint":"https://${LOCAL_HOST_ALIAS}:${MOCKSERVER_PORT_GITHUB}", + ... "insecureSkipTLSVerify":true,"unstable":true} + ${query} = Catenate SEPARATOR=${SPACE} + ... update stackql_unstable_github.orgs.orgs + ... set description = 'Some silly description.' + ... where org = 'dummyorg' + ... returning login, email; + ${expected} = Set Variable {"login":"dummyorg","email":"info@dummyorg.io"} + Should StackQL Exec Inline Jsonl Set Equal + ... ${STACKQL_EXE} + ... ${OKTA_SECRET_STR} + ... ${GITHUB_SECRET_STR} + ... ${K8S_SECRET_STR} + ... ${REGISTRY_NO_VERIFY_CFG_STR} + ... ${AUTH_CFG_STR} + ... ${SQL_BACKEND_CFG_STR_CANONICAL} + ... ${query} + ... ${expected} + ... --preview\=${preview} + ... stdout=${CURDIR}${/}tmp${/}Unstable-Github-Org-Update-Returning.tmp + ... stderr=${CURDIR}${/}tmp${/}Unstable-Github-Org-Update-Returning-stderr.tmp + +Unstable Google Kms Key Ring Insert Returning Jsonl Row Set Matches Expectation + [Documentation] A document-driven INSERT ... VALUES with RETURNING: path + ... and query parameters come from the column list, and the + ... created key ring streams back. + [Setup] Write Gcp Service Account ${OMNISDK_MOCK_GCP_SA_HOST} + [Teardown] Remove Preview Mock Environment + ${gcp_sa} = Set Variable If "${EXECUTION_PLATFORM}" == "docker" + ... /opt/test/tmp/omnisdk-gcp-sa.json ${OMNISDK_MOCK_GCP_SA_HOST} + Set Environment Variable GOOGLE_APPLICATION_CREDENTIALS ${gcp_sa} + ${preview} = Catenate SEPARATOR= + ... {"endpoint":"https://${LOCAL_HOST_ALIAS}:${MOCKSERVER_PORT_GOOGLE}", + ... "insecureSkipTLSVerify":true,"unstable":true} + ${auth} = Set Variable {"google":{"credentialsfilepath":"${gcp_sa}"}} + ${query} = Catenate SEPARATOR=${SPACE} + ... insert into stackql_unstable_google.cloudkms.key_rings + ... (projectsId, locationsId, keyRingId) + ... values ('testing-project', 'us-central1', 'fresh-ring') + ... returning name, createTime; + ${expected} = Catenate SEPARATOR= + ... {"name":"projects/testing-project/locations/us-central1/keyRings/fresh-ring", + ... "createTime":"2022-02-02T02:02:02.02000000Z"} + Should StackQL Exec Inline Jsonl Set Equal + ... ${STACKQL_EXE} + ... ${OKTA_SECRET_STR} + ... ${GITHUB_SECRET_STR} + ... ${K8S_SECRET_STR} + ... ${REGISTRY_NO_VERIFY_CFG_STR} + ... ${auth} + ... ${SQL_BACKEND_CFG_STR_CANONICAL} + ... ${query} + ... ${expected} + ... --preview\=${preview} + ... stdout=${CURDIR}${/}tmp${/}Unstable-Google-Kms-Key-Ring-Insert-Returning.tmp + ... stderr=${CURDIR}${/}tmp${/}Unstable-Google-Kms-Key-Ring-Insert-Returning-stderr.tmp + +Unstable Show Methods Reports Each Method Sql Verb + [Documentation] SHOW METHODS on a document-driven relation reports the + ... verb the document maps each method to, and EXEC for a + ... method no verb maps. + Should StackQL Exec Inline Equal + ... ${STACKQL_EXE} + ... ${OKTA_SECRET_STR} + ... ${GITHUB_SECRET_STR} + ... ${K8S_SECRET_STR} + ... ${REGISTRY_NO_VERIFY_CFG_STR} + ... ${AUTH_CFG_STR} + ... ${SQL_BACKEND_CFG_STR_CANONICAL} + ... show methods in stackql_unstable_github.orgs.orgs; + ... MethodName,RequiredParams,SQLVerb\nget,,SELECT\nlist,,SELECT\nlist_for_authenticated_user,,EXEC\nlist_for_user,,SELECT\nupdate,,UPDATE + ... \-o\=csv + ... --preview\={"unstable":true} + ... stdout=${CURDIR}${/}tmp${/}Unstable-Show-Methods-Reports-Each-Method-Sql-Verb.tmp + ... stderr=${CURDIR}${/}tmp${/}Unstable-Show-Methods-Reports-Each-Method-Sql-Verb-stderr.tmp + OTel Output Emits One Record Per Row Plus Completion [Documentation] Issue #738: --output otel writes one OTLP/JSON LogsData per ... row and a completion record with the row values as typed From 454bb7ba22eb2fc8aed26f0721ae4abee3f0153e Mon Sep 17 00:00:00 2001 From: General Kroll Date: Mon, 28 Sep 2026 18:52:21 +1000 Subject: [PATCH 2/4] remove-boilerplate --- .claude/.gitignore | 1 + .claude/omnisdk-shortcomings.md | 44 ------------ .claude/stackql-integration.md | 123 -------------------------------- 3 files changed, 1 insertion(+), 167 deletions(-) delete mode 100644 .claude/omnisdk-shortcomings.md delete mode 100644 .claude/stackql-integration.md diff --git a/.claude/.gitignore b/.claude/.gitignore index a1363379..04199acc 100644 --- a/.claude/.gitignore +++ b/.claude/.gitignore @@ -1 +1,2 @@ *.pdf +local/ \ No newline at end of file diff --git a/.claude/omnisdk-shortcomings.md b/.claude/omnisdk-shortcomings.md deleted file mode 100644 index 1c67f249..00000000 --- a/.claude/omnisdk-shortcomings.md +++ /dev/null @@ -1,44 +0,0 @@ - -# omnisdk shortcomings - -## High Level Expectations - -1. `omnisdk` will accept existing `any-sdk` auth structures and default configs in provider documents and behave as expected, can be observed in `any-sdk`. Behviours should be implemented protocol agnostic and best practice. Tech debt from `any-sdk` must not be inherited by copy or imitation. **`omnisdk` will never import from `any-sdk`, directly or indirectly.** -2. This is very important. Whereas `nay-sdk` has been tightly coupled to `http`, with some subprocess variants, by contrast `omnisdk` is **protocol agsnostc**. Full support for proticol buffers, `http`, subprocess calls, streaming protocols, whatever old and new transports, are intrinsic to `omnisdk`. -3. Any WAL, ledgers and the like should be configurable in location, substrate (eg: local vs, s3, different os...) and the like. Must be extensible and abstacted, no excuses. -4. `omnisdk` will support all functionality in `any-sdk`, but with clean implementation. Breaking changes will be specified ahead of time. The api is wildly different, but functional coverage will not be lesser. -5. `omnisdk` does not stage results in RDBMS or otherwise, it eagerly streams generated records. This is a significant and highly beneficial difference to `any-sdk`. Any ordering, aggregation or set operations (union and the like) is imposed post `omnisdk` query fulfilment. -6. `omnisdk` is required to support SQL extension funstions, including: scalar, redord and table valued ones. We do not expect the very first version to have total coverage and so some queries may need to be routed away from `omnisdk` at times, by logic within `stackql`. -7. `omnisdk` is to support all of the request and response trandsform grammars and shorthands already in `any-sdk`. -8. `omnisdk` will at some point support both SQL style saga rollbacks and also an IAC saga variant with fine grained locking and abstracted latches that are objectively superior to `terraform` style crude locks and failure modes. That said, early IAC forms are per stack locked. Patience is our watchword. - -## Migration plan at coarse grain - -- (a) Support joins, sql functions consuming only `omnisdk`, in `stackql_unstable_` namespace. It may take some time to acheive full covereage, in stages. -- (b) Cut versions and releases of both `omnisdk` and `stackql` along the way as useful milestones are reached. - - -## Current issues - -- (i) Auth has to function as expected. No excuses. The existing stackql patter with auth structures and docs simply **must** work. -- (ii) We need an orderly abstraction and catalogue of supported SQL extension functions in `omnisdk`. This will be in a discrete `pkg`. See below `Expected SQL extension functions` section. -- (iii) We need an orderly abstraction and catalogue of supported request and response processing grammars and shorthands in `omnisdk`. These are expected to mirror the `main` branch of `any-sdk` but be cleanly implemented in a discrete `pkg` with minimal dependencies and **zero** relation to `any-sdk`. See below `Expected Transformation Grammars and Shorthands` section. -- (iv) We want support for user/agent composed cross cloud rapid audit queries. -- (v) I want an SOC or whatever corporate audit query suite asap. -- (vi) `Args.Auth` is one struct shared by every node in the graph; each node reads the fields its scheme needs and falls back to env vars for any left empty. Two providers cannot carry distinct credentials in one query. Need per-provider credentials. Blocks (i) and (iv). -- (vii) `DescribeTable`/`DescribeMutation` drop parameters declared via `$ref` to `components/parameters` (e.g. github `org` on `orgs.members`, `username` on `users.users`), so joins and mutations cannot bind them. -- (viii) `DescribeMutation(dir, "stackql_unstable_google.storage.buckets", "insert")` fails with "read services: is a directory" while `DescribeTable` on the same address works; a `provider.yaml` service `$ref` to a missing file (`compute-v1.yaml`) gives the same message instead of naming the file. -- (ix) `Table` gives column names but no types, so every column reaches stackql as text. -- (x) Mutation outcomes ("was rejected" vs "may or may not have taken effect") are `fmt.Errorf` strings only; need sentinel errors for `errors.Is`. -- (xi) A mutation target carries one assignment set, so multi-row `INSERT ... VALUES` cannot be expressed. - -## Supporting information - -### Expected SQL extension functions - -The required SQL extension functions are precisely those that are robot tested anywhere in the `stackql` codebase. - -### Expected Transformation Grammars and Shorthands - -The required transformation grammars (eg for http request and response) are precisely those that are implemented in the `main` branch of the `any-sdk` codebase. - diff --git a/.claude/stackql-integration.md b/.claude/stackql-integration.md deleted file mode 100644 index a8adea53..00000000 --- a/.claude/stackql-integration.md +++ /dev/null @@ -1,123 +0,0 @@ -# Routing stackql queries through omnisdk - -What stackql has to do to send a query to omnisdk and stream back its rows. omnisdk resolves the -query against the provider documents and runs it; stackql parses SQL, holds sessions, and does -whatever needs every row at once. - -## Division of work - -| stackql | omnisdk | -|---------|---------| -| Parse SQL | Choose each table's method from the documents | -| Split CTEs, subqueries and UNION into single queries, then combine results | Decide which conditions become request parameters, edges between tables, or row filters | -| ORDER BY, GROUP BY, aggregates, DISTINCT | Run requests, joins, `IN` fan-out, polls and mutations | -| Hold per-client sessions and turn them into per-query config | Return an eager, unordered, unaggregated row stream | - -## Per query - -### 1. Build a `query.Unresolved` (package `pkg/query`, standard library only) - -Handles are registry addresses: `stackql_unstable_aws.iam.users`. omnisdk does not map names. - -| SQL | Constructor | -|-----|-------------| -| `FROM t a` | `query.NewJoin(query.NewResource("a", "
"), query.Base)` | -| `[INNER] JOIN t b ON …` | `query.NewJoin(query.NewResource("b", …), query.Inner, on...)` | -| `LEFT JOIN t b ON …` | `query.NewJoin(…, query.Left, on...)` | -| `a.col` / unqualified `col` | `query.NewColumn("a", "col")` / `query.NewColumn("", "col")` | -| `'x'`, `1` | `query.NewLiteral(v)` | -| `f(x, y)` | `query.NewCall("f", x, y)` | -| `=`, `<>`, `<`, `<=`, `>`, `>=` | `query.NewCompare(query.Eq, l, r)` (and `Ne`, `Lt`, …) | -| `x IN (…)` | `query.NewIn(x, query.NewCollection(items...))` | -| `OR`, `NOT` | `query.NewOr(...)`, `query.NewNot(p)` | -| boolean function as condition | `query.NewTest(call)` | -| `SELECT expr AS name` | `query.NewOutput("name", expr)` | -| `SELECT *` / `a.*` | `query.NewOutput("", query.NewStar(""))` / `NewStar("a")` | -| `INSERT INTO t (c…) VALUES (…)` | `query.NewInsert(res, query.NewAssignment("c", v)...)` | -| `INSERT … SELECT` | same, with values reading the SELECT's tables, which go in `from` | -| `UPDATE t SET c = v` | `query.NewUpdate(res, assignments...)` | -| `DELETE FROM t` | `query.NewDelete(res)` | -| `RETURNING …` | the mutation's outputs | - -Split WHERE and each ON into conjuncts (top-level `AND`). Every output needs a name: give -unnamed expressions one. Then: - -```go -q, err := query.New(from, where, outputs) // a read -q, err := query.NewMutation(target, from, where, returning) // a mutation -``` - -Leave out ORDER BY, GROUP BY, aggregates, DISTINCT and HAVING. If ORDER BY or GROUP BY reads a -column that isn't selected, add it to the outputs so it comes back. - -Not supported: RIGHT, FULL and CROSS joins. - -### 2. Describe each table - -```go -tables := map[string]omnisdk.Table{} -for _, j := range q.From() { - t, err := omnisdk.DescribeTable(registry, j.Resource().Handle()) - tables[j.Resource().Alias()] = t -} -if tg := q.Target(); tg != nil { - t, err := omnisdk.DescribeMutation(registry, tg.Resource().Handle(), tg.Verb().String()) - tables[tg.Resource().Alias()] = t -} -``` - -A `Table` lists each method's parameters (name, location, required) and row columns, which is also -what `DESCRIBE` / `SHOW` can report. - -### 3. Resolve - -```go -res, err := omnisdk.Resolve(q, tables) -``` - -An error names the clause it could not place, for example a condition a mutation's method can't take -(it would only be checked after the effect), a left join whose preserved side needs the other side's -value, or an ambiguous `*`. Show it to the user as is. - -### 4. Run - -```go -args := omnisdk.Args{ - Params: merge(session.Params, res.Params()), // scope such as region, then the query's own - Auth: session.Auth, // nil falls back to the environment - Tuning: omnisdk.Tuning{Limit: limit}, // LIMIT without ORDER BY; otherwise apply it after sorting - Journal: session.Journal, // opt-in write-ahead intent for mutations - Redaction: session.Redaction, // nil drops credentials from rows -} -pl, err := omnisdk.NewGraphSelectQuery(registry, res.Graph(), args) -rows, err := pl.Open(ctx) -defer rows.Close() -for rows.Next() { - row := rows.Row() // map[string]any, keyed by output name -} -err = rows.Err() -``` - -Rows arrive as they are produced. A left join's unmatched row lacks the joined table's columns; -treat a missing column as NULL. - -## Per session - -- **Config:** keep it in the session and build a fresh `Args` for each query. omnisdk holds no - per-client state, so concurrent clients with different settings don't interfere. -- **Document changes for one client:** `EffectiveRegistry(registry, patches, omnisdk.DocCache{Dir: dir})` - returns a registry directory with that client's patches applied (RFC 7386 merge patches on service - documents). Use it as `registry` for that client's queries. `NewDocCache(parent)` gives a new cache - location; passing an existing `Dir` reuses its entries; `Fresh: true` rebuilds. -- **Mutations:** set `Args.Journal{State, RunID}` to record each effect before it is sent. A failure is - reported as "rejected" (the provider refused it) or "may or may not have taken effect" (no answer, or - a 5xx). Mutations are never retried. -- **Credentials:** dropped from result rows by default. `omnisdk.RedactNone()` keeps them, for a user - who needs them. - -## Once per process - -- `omnisdk.ConfigureDocumentCache(cache.Config{...})` bounds the parsed-document cache. The default - takes a quarter of the process's memory limit (`GOMEMLIMIT`, else the container's cgroup limit, - else 512 MiB). Parsed documents are keyed by the directory they are read from, so clients with - different patches never share entries. From 7bc679bf68b6b530c029aac765d1b2d99327ca25 Mon Sep 17 00:00:00 2001 From: General Kroll Date: Mon, 28 Sep 2026 19:15:38 +1000 Subject: [PATCH 3/4] decent-preview-guide --- docs/preview.md | 15 +++++++++++---- go.mod | 2 +- go.sum | 4 ++-- 3 files changed, 14 insertions(+), 7 deletions(-) diff --git a/docs/preview.md b/docs/preview.md index 67ea400d..dee5ae08 100644 --- a/docs/preview.md +++ b/docs/preview.md @@ -8,16 +8,23 @@ Bleeding edge `stackql` functionality can be exposed using the `--preview` CLI a ./build/stackql shell --preview='{"unstable":true}' ``` +For any provider you wish to consume, you will need to pull, for example in support of what follows: + +```sql +registry pull google v26.08.00446; + +``` + ## Relevant preview functionality Streaming high volume queries at low latency has releveance for audit and related use cases. -```bash +```sql -./build/stackql exec --preview='{"unstable":true}' --output=jsonl \ - "select name, location, storageClass, timeCreated +./build/stackql exec --preview='{"unstable":true}' --output jsonl \ +"select name, location, storageClass, timeCreated from stackql_unstable_google.storage.buckets - where project = 'stackql-demo'" + where project = 'stackql-demo';" ``` \ No newline at end of file diff --git a/go.mod b/go.mod index 6717e0cf..4033ff4c 100644 --- a/go.mod +++ b/go.mod @@ -19,7 +19,7 @@ require ( github.com/spf13/cobra v1.10.2 github.com/spf13/pflag v1.0.10 github.com/spf13/viper v1.10.1 - github.com/stackql-labs/omnisdk v0.1.3-alpha03 + github.com/stackql-labs/omnisdk v0.1.3-alpha05 github.com/stackql/any-sdk v0.6.0-alpha01 github.com/stackql/go-suffix-map v0.0.1-alpha01 github.com/stackql/psql-wire v0.1.2-beta01 diff --git a/go.sum b/go.sum index b3635cae..eb9f280e 100644 --- a/go.sum +++ b/go.sum @@ -366,8 +366,8 @@ github.com/spf13/viper v1.10.1 h1:nuJZuYpG7gTj/XqiUwg8bA0cp1+M2mC3J4g5luUYBKk= github.com/spf13/viper v1.10.1/go.mod h1:IGlFPqhNAPKRxohIzWpI5QEy4kuI7tcl5WvR+8qy1rU= github.com/spiffe/go-spiffe/v2 v2.7.0 h1:uXe1MflJoHw58wAUvxVlcM7WpKtijWG7I1UidcGh6g4= github.com/spiffe/go-spiffe/v2 v2.7.0/go.mod h1:47Q0Q9/AqGha8QLHp+kxpH4Wca7X7EnOtlIJy3mxZ3U= -github.com/stackql-labs/omnisdk v0.1.3-alpha03 h1:6bTEeF/FoBTsVjl/lk63a8plDS3VaPPwwAPByht6QE0= -github.com/stackql-labs/omnisdk v0.1.3-alpha03/go.mod h1:WzvNj/bVv53yGFsVJpYWCJC1xAEdmQSFJl9eVpkRpCY= +github.com/stackql-labs/omnisdk v0.1.3-alpha05 h1:pAqzYLtmowbDuo2B9ssvrhbyISsu+kTajog6p6JnxX4= +github.com/stackql-labs/omnisdk v0.1.3-alpha05/go.mod h1:WzvNj/bVv53yGFsVJpYWCJC1xAEdmQSFJl9eVpkRpCY= github.com/stackql/any-sdk v0.6.0-alpha01 h1:mqy0bmZ1wghr7mUgLyaexoRU6rsUZCWXU9Mf60fC8dg= github.com/stackql/any-sdk v0.6.0-alpha01/go.mod h1:ahkRgwHHEn7RTmSfKrshPVk6i4X1BdHknSeBflSyRac= github.com/stackql/go-suffix-map v0.0.1-alpha01 h1:TDUDS8bySu41Oo9p0eniUeCm43mnRM6zFEd6j6VUaz8= From 1e8323a6997b0144119ff772744e1e48a36bdac2 Mon Sep 17 00:00:00 2001 From: General Kroll Date: Mon, 28 Sep 2026 19:15:59 +1000 Subject: [PATCH 4/4] format-correction --- docs/preview.md | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/docs/preview.md b/docs/preview.md index dee5ae08..5afae38d 100644 --- a/docs/preview.md +++ b/docs/preview.md @@ -20,7 +20,7 @@ registry pull google v26.08.00446; Streaming high volume queries at low latency has releveance for audit and related use cases. -```sql +```bash ./build/stackql exec --preview='{"unstable":true}' --output jsonl \ "select name, location, storageClass, timeCreated