Skip to content
Query-farmPublic

About

Go SDK for the VGI (Vector Gateway Interface) protocol — host DuckDB scalar, table, aggregate, and table-in-out functions in an external worker over Arrow IPC.

Topics

Resources

Stars

1 star

Watchers

0 watching

Forks

Repository files navigation

Vector Gateway Interface logo

VGI for Go

Add your own functions and tables to DuckDB with Go and Apache Arrow.
Built by 🚜 Query.Farm

CI Go Reference License

Go SDK for the VGI (Vector Gateway Interface) protocol. VGI lets DuckDB call user-defined scalar / table / aggregate / table-in-out functions hosted in an external worker process over Arrow IPC.

  • Sibling reference port (Python): vgi-python
  • DuckDB extension: vgi

Install

go get github.com/Query-farm/vgi-go/vgi

Requires Go 1.25+. The default stdio and HTTP transports run on upstream github.com/apache/arrow-go/v18.

Shared-memory transport (optional)

The zero-copy shared-memory side channel (VGI_RPC_SHM_SIZE_BYTES) requires the Query.Farm Arrow fork, which adds RecordBatch custom-metadata support used by the shm pointer-batch protocol. Because Go replace directives do not propagate to importers, a project that enables SHM must add the replace to its own go.mod:

replace github.com/apache/arrow-go/v18 => github.com/Query-farm/arrow-go/v18 v18.0.0-20260220022719-2d45cbd918a4

stdio and HTTP need no fork.

Quickstart

A minimum-viable worker — a scalar that adds two integers:

package main

import (
    "context"

    "github.com/Query-farm/vgi-go/vgi"
    "github.com/apache/arrow-go/v18/arrow"
    "github.com/apache/arrow-go/v18/arrow/array"
)

type AddInts struct{}

type addArgs struct {
    A int64 `vgi:"pos=0,const=false,doc=Left operand"`
    B int64 `vgi:"pos=1,const=false,doc=Right operand"`
}

func (*AddInts) Name() string               { return "add_ints" }
func (*AddInts) Metadata() vgi.FunctionMetadata {
    return vgi.FunctionMetadata{
        ArgumentMonotonicity: []vgi.ArgumentMonotonicity{
            vgi.ArgumentMonotonicityStrictlyIncreasing,
            vgi.ArgumentMonotonicityStrictlyIncreasing,
        },
    }
}

func (*AddInts) OnBindTyped(_ *addArgs, _ *vgi.BindParams) (*vgi.BindResponse, error) {
    return vgi.BindResult(arrow.PrimitiveTypes.Int64)
}

func (*AddInts) ProcessTyped(_ context.Context, _ *addArgs, params *vgi.ProcessParams, batch arrow.RecordBatch) (arrow.RecordBatch, error) {
    return vgi.MapAllColumns(params, batch, array.NewInt64Builder,
        func(cols []arrow.Array, i int) int64 {
            return vgi.GetInt64Value(cols[0], i) + vgi.GetInt64Value(cols[1], i)
        })
}

func main() {
    w := vgi.NewWorker(vgi.WithCatalogName("demo"))
    w.RegisterScalar(vgi.AsScalarFunction[addArgs](&AddInts{}))
    w.RunStdio()
}

ArgumentMonotonicity is optional and scalar-only. When present, its entries follow the ordered argument declarations: fixed, defaulted, and constant arguments each use one slot, and a vararg declaration uses one slot regardless of call-time expansion. Named SQL argument order does not reorder the list.

Build it (go build -o my-worker .), then install the VGI extension, attach the worker, and call the function from DuckDB:

INSTALL vgi FROM community;
LOAD vgi;

-- Attach the worker as a catalog. The first argument is the worker's catalog
-- name (WithCatalogName above); LOCATION is the command DuckDB runs to launch it.
ATTACH 'demo' AS demo (TYPE vgi, LOCATION './my-worker');

SELECT demo.add_ints(2, 3); -- => 5

LOCATION also accepts http://…/https://… for an HTTP worker, or a launch: command for the AF_UNIX transport.

Schemas and catalogs

Every registered function has exactly one home: a (catalog, schema) pair. A plain Register* call homes it in the worker's own catalog and default schema (main); registration can name either dimension instead:

w.RegisterScalar(&MainImpl{})                     // demo.main.collide
w.RegisterScalarInSchema("data", &DataImpl{})     // demo.data.collide
w.RegisterScalarForCatalog("other", &OtherImpl{}) // other.main.collide

A function name is not a unique key, so both the catalog listing and bind dispatch resolve the whole (catalog, schema, name) triple, exactly — there is no "visible everywhere" tier. demo.data.collide(x) reaches the data implementation and never the main one, and a call arriving through catalog other reaches only what other owns. Same-name registrations within one schema stay overloads, resolved by argument signature as before.

The one bind that legitimately names no schema is a COPY handler's: COPY formats are advertised at catalog level, not inside a schema. Such a bind resolves by name within the catalog, and errors naming the schemas involved if that name is declared in more than one.

Attach options

vgi.WithAttachOptions(...) declares the options a caller may pass at ATTACH time. Required: true marks one the catalog cannot be attached without. Credential options (API keys, tokens, passwords) must be declared Secret: true. Credentials are passed inline as attach options. The DuckDB extension redacts a secret option's value from duckdb_databases(), keeps only a salted hash of it in its cache key, and never logs it; clients mask it and keep it out of exported configuration. To keep the credential out of the SQL text, write it as an expression:

w := vgi.NewWorker(vgi.WithAttachOptions(
	vgi.AttachOptionSpec{Name: "api_key", Description: "API key",
		Type: arrow.BinaryTypes.String, Required: true, Secret: true},
))
ATTACH 'sales' (TYPE vgi, LOCATION 'https://worker.example.com',
    api_key getenv('SALES_API_KEY'));

Secret combines with Required. A secret option may declare a default, but normally has none, since the default is published to every client.

vgi_attach_ticket is reserved for attach tickets: a worker declaring an option of that name (any case) refuses to start.

Function shapes

Shape Interface Use case
Scalar ScalarFunction, TypedScalarFunc[A] 1:1 row mapping
Table generator TableFunction, TypedTableFunc[S] Row generator (no streamed input)
Table-in-out TableInOutFunction Stream input rows → output rows
Table-buffering TableBufferingFunction Sort / aggregate / join-style buffer
Aggregate AggregateFunction Cumulative state + final emit

Pick the typed variants (TypedScalarFunc, TypedTableFunc) when you want the framework to derive ArgumentSpecs from a Go struct with vgi:"..." tags. See examples/scalar/add_values.go and examples/table/sequence.go.

Logging

Workers emit structured logs through named slog loggers (vgi, vgi.worker, vgi.catalog, vgi.rpc, vgi.client, vgi.filter_pushdown). Configure them with the standard CLI flags by registering them in main():

fs := flag.CommandLine
logFlags := vgi.RegisterLoggingFlags(fs)
flag.Parse()
if err := logFlags.Apply(); err != nil {
    log.Fatal(err)
}

Then:

./my-worker --log-level=debug --log-format=json
./my-worker --log-logger=vgi.catalog --log-logger=vgi.rpc
VGI_WORKER_DEBUG=1 ./my-worker        # equivalent to --debug

Env-var fallbacks: VGI_LOG_LEVEL, VGI_LOG_FORMAT, VGI_LOG_LOGGER, VGI_WORKER_DEBUG.

Errors

The SDK defines a few error types that surface clean RPC-error messages to DuckDB rather than the generic RuntimeError:

  • ArgumentError — bad / missing argument at bind time
  • SchemaValidationError — schema mismatch with per-field detail
  • TypeBoundError — column type doesn't satisfy a declared type predicate
  • WorkerPanicError — captured panic from user code; worker stays alive

Panics inside user functions during bind, init, cardinality, and statistics dispatch are recovered automatically.

Hosting more protocols

A worker always hosts vgi.v2 and vgi_rpc.Reflection.v1. To host another application protocol beside them -- a reporting protocol, say -- return it from the hosted-protocols hook. Define it the way vgi-rpc-go defines any protocol:

w := vgi.NewWorker(
    vgi.WithHostedProtocols(func() ([]*vgirpc.Server, error) {
        reports := vgirpc.NewProtocol("acme.Reports.v1")
        vgirpc.Unary(reports, "status", reportStatus)
        return []*vgirpc.Server{reports}, nil
    }),
)

The hook is called once when the worker's server is built and may read configuration; what it returns is hosted after vgi.v2 on every transport (stdio, unix, TCP, HTTP) and fixed for the life of the process. It cannot change vgi.v2: every request routes on its protocol name. A returned protocol named under the reserved vgi_rpc. prefix, named vgi.v2, or named twice stops the worker from starting, with an error naming the hook.

Identity (vgi_rpc.Identity.v1)

A worker can resolve opaque bearer credentials to principals for a reverse proxy (introspect_token), and mint short-lived grants for the calling user (issue_grant), by opting in with vgi.WithIdentity:

w := vgi.NewWorker(vgi.WithIdentity(vgirpc.IdentityConfig{
    ResolveToken: func(token string) (vgirpc.TokenIdentity, bool, error) {
        row, err := apiKeys.Lookup(token)
        if err != nil {
            // "I could not find out" -- not "the credential is bad".
            return vgirpc.TokenIdentity{}, false, &vgirpc.AuthUnavailableError{Detail: err.Error(), RetryAfter: 5}
        }
        if row == nil {
            return vgirpc.TokenIdentity{}, false, nil // the store answered: unknown
        }
        return vgirpc.NewTokenIdentity(row.Principal), true, nil
    },
    IntrospectPrincipals: []string{"proxy@example.com"},
}))
  • The protocol is hosted on HTTP only, and only the methods whose hook is set: with neither hook it is absent, not hosted-and-refusing.
  • introspect_token needs an allowlist of principals permitted to ask -- IntrospectPrincipals, or VGI_INTROSPECT_PRINCIPALS (comma-separated). There is no permissive default; a worker with ResolveToken and no allowlist refuses to start, because "any authenticated caller" lets any user resolve any other user's credential to its owner.
  • For a transient failure (a store or sidecar is down), return *vgirpc.AuthUnavailableError with a RetryAfter -- the error an authenticator returns for an outage -- or *vgirpc.IdentityUnavailableError. The framework sends either as identity_unavailable / UNAVAILABLE with that retry hint as RetryInfo, which tells the caller not to negative-cache the answer. Return ok=false with no error only when the store answered and the credential is unknown. Never report an outage as a ValueError *RpcError: ChainAuthenticate reads that as "not my credential, try the next", so a thirty-second blip becomes a 401 for everyone.

Sealed grants

issue_grant mints a credential for unattended automation to present later as an ordinary bearer. Configure grant keys and the worker closes that loop itself: over HTTP it hosts issue_grant (minting sealed grants, unless WithIdentity supplies a MintGrant) and accepts its own grants back as Bearer credentials, after your SetAuthenticate authenticator.

VGI_RPC_GRANT_KEYS=$(openssl rand -base64 32) ./my-worker --http
./my-worker --http --grant-key "$NEW_KEY" --grant-key "$OLD_KEY"   # first mints, all verify

or vgi.WithGrantKeys(keys) in code. Optional VGI_RPC_GRANT_AUDIENCE and VGI_RPC_GRANT_MAX_TTL_SECONDS (default 7 days). A malformed key stops the worker at startup.

  • A grant authenticates as domain grant, its owner's principal, and claims grant_id, scopes, purpose -- with no auth_time, so a grant can never mint another grant.
  • Grants are not individually revocable: keep the max TTL short, and remove a key to revoke everything it minted.
  • A ResolveToken hook (WithIdentity) is consulted for bearers too: resolved means authenticated (domain token); unknown falls through to 401; an outage is a 503 with your RetryAfter.
  • Your authenticator must return a ValueError *vgirpc.RpcError for a bearer it does not recognise, or the grant verifier is never reached.

Attach tickets (vgi.attach_tickets.v1)

A grant says who attaches; an attach ticket says what. While a user is attached and logged in, a client calls seal_attach and the worker seals the options the user attached with -- secret ones included -- into a vgia1. ticket only this worker can open. Later a runner holding the user's grant reattaches with that single option, and never sees an option:

ATTACH 'sales' AS s (TYPE vgi, LOCATION 'https://worker.example.com',
    vgi_attach_ticket '<ticket>');   -- authenticated by Bearer <grant>
  • Hosting. HTTP only, and only when the signing key is configured explicitly (VGI_SIGNING_KEY, or vgi.WithHttpSigningKey) and the worker can issue grants (grant keys, or a MintGrant hook). Otherwise the protocol is absent, which a client sees through reflection. A per-process generated key never enables it: every ticket would die on restart.
  • Format. "vgia1." + base64url(0x01 || nonce || XChaCha20-Poly1305), keyed by VGI_SIGNING_KEY, AAD "vgi.attach_ticket.v1\0" + principal. The principal only, not the login domain, so a ticket sealed under a JWT login opens under that user's grant. A ticket carries no authority: it opens only for the same principal.
  • Redemption. catalog_attach replaces a request whose options contain vgi_attach_ticket (any case) with the sealed catalog name, options and version specs before any catalog code runs, and before routing to a sub-catalog. Any other option beside it is invalid_request; a ticket that does not open is attach_ticket_invalid; one outside its lifetime is attach_ticket_expired.
  • Lifetime. The grant maximum (VGI_RPC_GRANT_MAX_TTL_SECONDS / the grant keys' max TTL); ttl_seconds 0 asks for all of it. No maximum means no expiry (expires_at is +Inf).
  • Reserved name. Declaring an attach option named vgi_attach_ticket (any case) stops the worker at startup.
  • Rotating VGI_SIGNING_KEY invalidates every ticket.

vgi.MintAttachTicket and vgi.OpenAttachTicket expose the format. The spec is vgi-python's docs/protocol/vgi-attach-tickets.md; vgi/testdata/attach_ticket_vectors.json pins it byte for byte. The example worker serves the cross-SDK ticket_probe fixture catalog (examples/ticket_probe), and hosts the protocol over HTTP when both VGI_SIGNING_KEY and VGI_RPC_GRANT_KEYS are set. Its test bearers vgi-test-alice / vgi-test-bob count as fresh logins, so they can call issue_grant.

Examples

The cmd/vgi-example-worker binary registers every example function via examples/all.RegisterAll(w). Browse:

  • examples/scalar/ — 20+ scalar examples (typed + classic)
  • examples/table/ — 40+ table generators (partitioned, paged, etc.)
  • examples/table_in_out/ — transform / buffering / aggregation patterns
  • examples/aggregate/ — cumulative aggregates
  • examples/schema_reconcile/, examples/versioned_tables/ — catalog + write paths

Build & test

make build              # build all example worker binaries
make fmt                # gofmt
make vet                # go vet
make lint               # golangci-lint (requires golangci-lint in PATH)
make test               # full integration suite over stdio (requires $VGI_DIR)
make test-http          # full suite over HTTP transport
make test-all           # both transports
make test-single TEST=test/sql/integration/scalar/add_values.test
go test ./...           # pure Go unit tests

Integration tests live in the DuckDB VGI extension repo, located by VGI_DIR (default: the sibling checkout ../vgi; e.g. make test VGI_DIR=/path/to/vgi), and use the DuckDB sqllogictest format.

See Iroh operations for bridge-ready raw and HTTP workers.

Repo layout

vgi/                            # SDK package
examples/scalar/                # scalar example functions
examples/table/                 # table example functions
examples/table_in_out/          # table-in-out + buffering examples
examples/aggregate/             # aggregate examples
examples/schema_reconcile/      # catalog-handlers fixture
examples/twin_catalogs/         # two catalogs, one worker, colliding names
examples/all/                   # RegisterAll(w) helper
cmd/vgi-example-worker/         # fixture worker (used by integration tests)
cmd/vgi-example-versioned-worker/
cmd/vgi-example-versioned-tables-worker/
cmd/vgi-example-attach-options-worker/

License

Copyright 2025, 2026 Query Farm LLC.

Licensed under the Query Farm Source-Available License, Version 1.0 — see LICENSE for the full terms. In brief, you may use, modify, and redistribute the software freely for non-production use, and for production use except where it would constitute a Competing Offering or a Commercial Marketplace as defined in the license. Each version converts to the Apache License, Version 2.0 on the tenth anniversary of its public release.

For uses not permitted under this license, contact hello@query.farm for a commercial license.

About

Go SDK for the VGI (Vector Gateway Interface) protocol — host DuckDB scalar, table, aggregate, and table-in-out functions in an external worker over Arrow IPC.

Topics

Resources

Stars

1 star

Watchers

0 watching

Forks

Releases

Packages

Contributors

Languages