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
1 change: 1 addition & 0 deletions .claude/.gitignore
Original file line number Diff line number Diff line change
@@ -1 +1,2 @@
*.pdf
local/
30 changes: 30 additions & 0 deletions docs/preview.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,30 @@

# Preview functionality

Bleeding edge `stackql` functionality can be exposed using the `--preview` CLI argument.


```bash
./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

./build/stackql exec --preview='{"unstable":true}' --output jsonl \
"select name, location, storageClass, timeCreated
from stackql_unstable_google.storage.buckets
where project = 'stackql-demo';"

```
2 changes: 1 addition & 1 deletion go.mod
Original file line number Diff line number Diff line change
Expand Up @@ -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-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
Expand Down
4 changes: 2 additions & 2 deletions go.sum
Original file line number Diff line number Diff line change
Expand Up @@ -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-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=
Expand Down
175 changes: 122 additions & 53 deletions internal/stackql/intrinsic/doc.go
Original file line number Diff line number Diff line change
Expand Up @@ -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 "<provider>.<service>.<resource>", 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 {
Expand Down Expand Up @@ -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
Expand Down
5 changes: 4 additions & 1 deletion internal/stackql/intrinsic/dynamic.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
9 changes: 7 additions & 2 deletions internal/stackql/intrinsic/intrinsic.go
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down
Loading
Loading