The Catalog Interface enables VGI workers to expose database-like structures to clients, supporting DuckDB's ATTACH command for external catalogs.
While VGI functions provide computational capabilities (scalar transforms, table generation), the Catalog Interface provides metadata management - exposing catalogs, schemas, tables, views, and functions as first-class database objects.
-- Attach a VGI-backed catalog in DuckDB
ATTACH 'mydb' (TYPE 'vgi', LOCATION './my-worker');
-- Query tables from the attached catalog
SELECT * FROM mydb.main.users;
-- List available schemas
SELECT * FROM information_schema.schemata WHERE catalog_name = 'mydb';Key Characteristics:
| Aspect | Functions | Catalog Interface |
|---|---|---|
| Purpose | Compute data | Manage metadata |
| Protocol | Bind → Init → Stream | vgi_rpc typed dispatch |
| Stateful | Per-invocation | Per-attachment |
| Discovery | Worker.functions list |
CatalogInterface.catalogs() |
Catalog methods are dispatched via vgi_rpc typed protocol methods. Each method has its own typed request/response defined in vgi.protocol, with automatic Arrow serialization handled by the RPC layer. This is simpler than the function protocol — no bind/init phases.
VGI protocol 2.0 represents every schema location as a SchemaPath = list[str], ordered from the outermost schema to the innermost schema. The components are raw identifiers: ["a.b"] is one schema whose name contains a dot, while ["a", "b"] is a nested schema. Implementations must preserve those boundaries and must not flatten a path to a dotted string. On the Arrow wire these values are list<utf8>.
Schema paths are non-empty and contain only non-empty strings. Every nested path must have its parent path in the catalog. Schema enumeration preserves a worker's order among peers while ensuring that parents precede children.
CatalogAttachResult.default_schema remains a single string because it selects a root-level default/search-path schema. DuckDB 1.5 adapters—including the bundled Python transactor—consume protocol 2.0 by accepting only paths of length one and rejecting deeper paths explicitly.
from vgi.catalog import AttachOpaqueData, TransactionOpaqueData, SerializedSchema, SqlExpression
from vgi import SchemaPath
AttachOpaqueData = NewType("AttachOpaqueData", bytes) # Unique attachment identifier
TransactionOpaqueData = NewType("TransactionOpaqueData", bytes) # Transaction identifier
SerializedSchema = NewType("SerializedSchema", bytes) # Arrow schema bytes
SqlExpression = NewType("SqlExpression", str) # SQL expression string
# SchemaPath is list[str]: outermost-to-innermost raw schema identifiersReturned by catalog_attach() with attachment metadata:
| Field | Type | Description |
|---|---|---|
attach_opaque_data |
AttachOpaqueData |
Opaque per-attachment state the implementation owns (see note below) |
supports_transactions |
bool |
Whether transactions are supported |
supports_time_travel |
bool |
Whether time travel queries work |
catalog_version_frozen |
bool |
Whether metadata will change |
catalog_version |
int |
Current version (increments on changes) |
attach_opaque_data_required |
bool |
Whether attach_opaque_data must be persisted |
default_schema |
str |
Name of the default schema (usually "main") |
attach_opaque_data/transaction_opaque_dataare opaque, implementation-owned, and may carry secrets. They are not framework identifiers — they are arbitrarybytesyour implementation returns fromcatalog_attach()/catalog_transaction_begin(), and the client round-trips them back verbatim on every subsequent call. An implementation may pack connection handles, credentials, or any state into them. Never log either value raw — treat them like theoptionsdict. The worker already enforces this for its own catalog-lifecycle logs (it short-hashes both fields at a single chokepoint), and on HTTP transport it additionally seals each value in an AEAD envelope bound to the caller's identity, so a value minted for one principal cannot be replayed by another. Your implementation only ever sees the plaintext; the sealing and unsealing happen transparently in the worker.
Information about a schema in a catalog:
| Field | Type | Description |
|---|---|---|
attach_opaque_data |
AttachOpaqueData |
Parent attachment |
path |
SchemaPath |
Qualified schema path |
comment |
str | None |
Optional description |
tags |
dict[str, str] |
Key-value metadata |
Information about a table:
| Field | Type | Description |
|---|---|---|
name |
str |
Table name |
schema_path |
SchemaPath |
Parent schema path |
columns |
SerializedSchema |
Column definitions as Arrow schema bytes |
not_null_constraints |
list[int] |
Column indices with NOT NULL |
unique_constraints |
list[list[int]] |
Column index groups for UNIQUE |
check_constraints |
list[str] |
SQL check expressions |
comment |
str | None |
Optional description |
tags |
dict[str, str] |
Key-value metadata |
Information about a view:
| Field | Type | Description |
|---|---|---|
name |
str |
View name |
schema_path |
SchemaPath |
Parent schema path |
definition |
str |
SQL SELECT statement |
comment |
str | None |
Optional description |
tags |
dict[str, str] |
Key-value metadata |
Information about a function in a schema:
| Field | Type | Description |
|---|---|---|
name |
str |
Function name |
schema_path |
SchemaPath |
Parent schema path |
function_type |
FunctionType |
SCALAR or TABLE |
arguments |
SerializedSchema |
Argument schema as Arrow bytes |
output_schema |
SerializedSchema |
Output schema as Arrow bytes |
parameter_default_values |
pa.RecordBatch | None |
One-row typed defaults batch containing only defaulted parameters in signature order; a present null is an explicit NULL default |
argument_monotonicity |
list[ArgumentMonotonicity] | None |
Scalar-only claims aligned exactly with arguments; a vararg declaration occupies one slot |
comment |
str | None |
Optional description |
tags |
dict[str, str] |
Key-value metadata |
Enum for filtering objects in schema_contents():
| Value | Description |
|---|---|
TABLE |
Filter to return only tables |
VIEW |
Filter to return only views |
SCALAR_FUNCTION |
Filter to return only scalar functions |
TABLE_FUNCTION |
Filter to return only table functions |
Result from table_scan_function_get() that tells the VGI DuckDB extension which DuckDB function to call to obtain table data. This enables catalogs to delegate scanning to any DuckDB function (e.g., read_parquet, iceberg_scan, or a custom VGI table function) with appropriate arguments.
| Field | Type | Description |
|---|---|---|
function_name |
str |
The DuckDB function to call (e.g., "read_parquet", "iceberg_scan") |
positional_arguments |
list[pa.Scalar] |
Positional arguments to pass to the function |
named_arguments |
dict[str, pa.Scalar] |
Named arguments to pass to the function |
required_extensions |
list[str] |
DuckDB extensions that must be loaded before calling the function |
Example usage:
def table_scan_function_get(
self,
*,
attach_opaque_data: AttachOpaqueData,
transaction_opaque_data: TransactionOpaqueData | None,
schema_path: SchemaPath,
name: str,
at_unit: str | None,
at_value: str | None,
) -> ScanFunctionResult:
# Return a parquet scan for this table
return ScanFunctionResult(
function_name="read_parquet",
positional_arguments=[pa.scalar(f"s3://bucket/{'/'.join(schema_path)}/{name}/*.parquet")],
named_arguments={"hive_partitioning": pa.scalar(True)},
required_extensions=["parquet", "httpfs"],
)The CatalogInterface abstract base class defines all catalog operations. Subclass it and implement the abstract methods.
from abc import ABC, abstractmethod
from vgi.catalog import CatalogInterface, CatalogAttachResult, SchemaInfo, TableInfo, ViewInfo
class MyCatalog(CatalogInterface):
@abstractmethod
def catalogs(self) -> Iterable[str]:
"""List available catalog names."""
@abstractmethod
def catalog_attach(self, *, name: str, options: dict[str, Any]) -> CatalogAttachResult:
"""Attach to a catalog, returning attachment metadata."""
@abstractmethod
def schema_get(self, *, attach_opaque_data: AttachOpaqueData, transaction_opaque_data: TransactionOpaqueData | None, path: SchemaPath) -> SchemaInfo | None:
"""Get schema info, or None if not found."""
@abstractmethod
def table_get(self, *, attach_opaque_data: AttachOpaqueData, transaction_opaque_data: TransactionOpaqueData | None, schema_path: SchemaPath, name: str) -> TableInfo | None:
"""Get table info, or None if not found."""
@abstractmethod
def view_get(self, *, attach_opaque_data: AttachOpaqueData, transaction_opaque_data: TransactionOpaqueData | None, schema_path: SchemaPath, name: str) -> ViewInfo | None:
"""Get view info, or None if not found."""| Category | Method | Default Behavior |
|---|---|---|
| Catalog | catalog_create() |
NotImplementedError |
catalog_drop() |
NotImplementedError |
|
catalog_detach() |
No-op | |
catalog_version() |
Returns 0 |
|
| Transaction | catalog_transaction_begin() |
NotImplementedError |
catalog_transaction_commit() |
NotImplementedError |
|
catalog_transaction_rollback() |
NotImplementedError |
|
| Schema | schemas() |
Returns ["main"] |
schema_create() |
NotImplementedError |
|
schema_drop() |
NotImplementedError |
|
schema_contents() |
NotImplementedError |
|
| Table | table_create() |
NotImplementedError |
table_drop() |
NotImplementedError |
|
table_rename() |
NotImplementedError |
|
table_comment_set() |
NotImplementedError |
|
table_column_add() |
NotImplementedError |
|
table_column_drop() |
NotImplementedError |
|
table_column_rename() |
NotImplementedError |
|
table_column_type_change() |
NotImplementedError |
|
table_column_default_set() |
NotImplementedError |
|
table_column_default_drop() |
NotImplementedError |
|
table_not_null_set() |
NotImplementedError |
|
table_not_null_drop() |
NotImplementedError |
|
table_scan_function_get() |
NotImplementedError |
|
| View | view_create() |
NotImplementedError |
view_drop() |
NotImplementedError |
|
view_rename() |
NotImplementedError |
|
view_comment_set() |
NotImplementedError |
|
| Observability | loggable_attach_options() |
Returns {} (no options logged — see below) |
| Bulk load | catalog_contents() |
Composes schemas() + schema_contents() per kind; no etag (see below) |
A client normally learns a catalog lazily: catalog_schemas, then one
catalog_schema_contents_* call per schema and object kind. When the attach
result sets supports_catalog_contents (ReadOnlyCatalogInterface does), the
client may instead load everything with one catalog_contents RPC (protocol
2.1.0):
catalog_contents(attach_opaque_data: binary, if_none_match: utf8 nullable)
-> CatalogContentsResponse {
catalog_version: int64
etag: utf8 nullable -- null: the worker does not revalidate
not_modified: bool -- true: schemas is empty, keep what you have
schemas: list<struct<path: list<utf8>, schema: binary,
tables, views, scalar_functions, aggregate_functions,
table_functions, scalar_macros, table_macros,
indexes: list<binary>>> -- parents before children
}
Each struct row is one schema: path equals the SchemaInfo.path encoded in
schema, and every item is byte-identical to what the matching per-schema RPC
returns. There is no transaction parameter: it is the committed catalog at
catalog_version.
CatalogInterface.catalog_contents(*, attach_opaque_data, if_none_match=None)
returns a CatalogContentsResult(schemas, etag=None, not_modified=False). The
default composes schemas() and schema_contents() for every kind (skipping
kinds whose estimated_object_count is exactly 0) and returns no etag.
Revalidation. Return an etag to let a client revalidate a non-frozen
catalog with if_none_match instead of re-downloading it. Check it before
building, so an unchanged catalog costs almost nothing:
from vgi.catalog import CatalogContentsResult, CatalogInterface
class GenerationCatalog(CatalogInterface): # other methods omitted
generation = 1 # bumped by every DDL
def catalog_contents(self, *, attach_opaque_data, if_none_match=None):
etag = f"gen-{self.generation}"
if if_none_match == etag:
return CatalogContentsResult(etag=etag, not_modified=True) # nothing built
full = super().catalog_contents(attach_opaque_data=attach_opaque_data)
return CatalogContentsResult(schemas=full.schemas, etag=etag)The worker enforces the rules: not_modified requires an etag equal to
if_none_match and no schemas; a full answer whose etag equals
if_none_match is turned into not_modified; a catalog with no etag always
answers in full (it ignores if_none_match).
Content-hash etag. A catalog without a cheap validator can set
catalog_contents_etag = "content-hash": when it returns no etag, the worker
uses the hex SHA-256 of the serialized snapshot. The worker still builds the
snapshot on every call, but an unchanged catalog skips the transfer and the
client's decode. It is off by default, because for a non-frozen catalog the
client would revalidate with a full build on every transaction instead of a
cheap catalog_version poll.
Worker cache. When the catalog's version is frozen
(catalog_version_frozen) and its contents do not depend on the attach or the
caller (catalog_contents_attach_independent = True, the
ReadOnlyCatalogInterface default), the worker builds the response once per
catalog version and serves the same pre-serialized bytes to every call. A
ReadOnlyCatalogInterface subclass whose contents vary per caller must set
catalog_contents_attach_independent = False.
A worker declares the options its catalogs accept at ATTACH time in an AttachOptions inner class. Each option is advertised to clients at discovery (CatalogInfo.attach_option_specs), so DuckDB can cast and validate it, and frontends can render an options form.
from typing import Annotated
from vgi import Worker
from vgi.catalog.attach_option import AttachOption
class MyWorker(Worker):
class AttachOptions:
region: Annotated[str, AttachOption(desc="AWS region")] = "us-east-1"
# A credential: required, secret, and with no default.
api_key: Annotated[str, AttachOption(desc="API key", required=True, secret=True)]| Flag | Meaning |
|---|---|
required=True |
The caller must supply the option. Declare it with no class-level default; required plus a default is rejected. |
secret=True |
The option carries a credential. |
Credential options (API keys, tokens, passwords) MUST be declared secret=True. The caller passes them inline as attach options. Clients mask secret options, and the DuckDB extension redacts them from duckdb_databases(), keeps only a salted hash of them in its result-cache key, and never logs them. To keep a credential out of the SQL text (shell history, scripts, shared links), write the option as an expression:
ATTACH 'mydb' (TYPE vgi, LOCATION 'https://worker.example.com', api_key getenv('MYDB_API_KEY'));secret combines with required. It is allowed together with a default, but a secret option should normally have none: the default is advertised in plain text to every client at discovery.
On the wire, required and secret are nullable boolean columns appended after the four shared spec columns (name, description, type, default_value), in that order. Readers look columns up by name, so a spec from a peer that predates either column reads it as false, and older peers ignore columns they don't know.
secret does not change what the worker itself logs: attach options are never logged unless the catalog opts in through loggable_attach_options() (below), which must never return a secret option.
The worker emits structured _logger.info records and Sentry breadcrumbs for catalog lifecycle events (catalog.attach, catalog.detach, catalog.create, catalog.transaction.begin, catalog.transaction.commit, catalog.transaction.rollback). attach_opaque_data and transaction_opaque_data are short-hashed (12-char SHA-256 prefixes) before they reach the log record, the breadcrumb data, or the Sentry scope tags — the raw values never appear in observability output, since they may carry secrets. An operator correlates the short hash back to the catalog via these breadcrumbs.
The options dict passed to catalog_attach() and catalog_create() routinely carries credentials — passwords, tokens, OAuth secrets, connection strings. To avoid leaking these to logs and Sentry, the worker does not log option fields by default. Implementers opt in by overriding loggable_attach_options():
class MyCatalog(CatalogInterface):
def loggable_attach_options(self, options: Mapping[str, Any]) -> Mapping[str, Any]:
# Allowlist the keys you know are safe. Never include password / token / secret.
safe_keys = {"host", "region", "bucket", "database"}
return {k: v for k, v in options.items() if k in safe_keys}When the override returns an empty mapping (the default behaviour for catalogs that haven't opted in), the options field is omitted from the lifecycle event entirely — fail-closed: nothing is preferred over partial-leak.
The catalog name, attach id, transaction id, and version specs are always logged regardless.
A convenience base class for read-only catalogs that don't support DDL operations. All modification methods raise CatalogReadOnlyError.
The simplest way to expose VGI functions as a catalog:
from vgi.catalog import ReadOnlyCatalogInterface
from vgi import ScalarFunction, TableFunctionGenerator
class MyFunctionCatalog(ReadOnlyCatalogInterface):
catalog_name = "my_funcs" # Name for ATTACH
functions = [MyScalarFunction, MyTableFunction]
# Functions appear in the "main" schema:
# SELECT * FROM my_funcs.main.my_scalar_function(args);Implement abstract methods for more control:
from vgi.catalog import ReadOnlyCatalogInterface, CatalogAttachResult, SchemaInfo
class MyReadOnlyCatalog(ReadOnlyCatalogInterface):
def catalogs(self) -> Iterable[str]:
return ["readonly_db"]
def catalog_attach(self, *, name: str, options: dict[str, Any]) -> CatalogAttachResult:
return CatalogAttachResult(
attach_opaque_data=AttachOpaqueData(b"fixed-id"),
supports_transactions=False,
supports_time_travel=False,
catalog_version_frozen=True,
catalog_version=1,
attach_opaque_data_required=False,
)
def schema_get(self, *, attach_opaque_data, transaction_opaque_data, path) -> SchemaInfo | None:
if path == ["main"]:
return SchemaInfo(attach_opaque_data=attach_opaque_data, path=["main"], comment=None, tags={})
return None
# table_get, view_get return None by defaultA function name is not a global key. The worker resolves a call by the pair
(schema, function name), so the same name may be declared in more than one
schema of a catalog, and a schema-qualified call reaches the implementation in
that schema:
Catalog(
name="example",
default_schema="main",
schemas=[
Schema(path=["main"], functions=[ProdLookup]), # Meta.name = "lookup"
Schema(path=["staging"], functions=[StagingLookup]), # Meta.name = "lookup"
],
)SELECT example.main.lookup(1); -- ProdLookup
SELECT example.staging.lookup(1); -- StagingLookupThe DuckDB extension carries the owning schema on every bind request
(BindRequest.schema_path), taken from the schema entry the function was
registered into. Two consequences worth knowing:
- Overloads still work. Several classes sharing a name within one schema are overloads, disambiguated by argument signature as before. Only the cross-schema case is resolved by schema.
- Callers without a schema must be unambiguous. The pure-Python
Clientand the CLI send no schema. If the name is unique across the worker they resolve normally; if two schemas declare it, the worker raises an ambiguity error naming the schemas involved.
Functions declared via the legacy Worker.functions list have no schema of
their own, so they are registered into the catalog's default_schema — the
same schema DuckDB registers them into.
Across catalogs, the key is the attachment rather than the name: two
catalogs served by one worker process (see vgi.meta_worker.MetaWorker) may
each declare main.lookup, and each attachment's attach_opaque_data routes
its calls to the right catalog.
For most use cases, the declarative catalog API provides a simpler way to define catalogs using Python dataclasses instead of implementing CatalogInterface directly.
from vgi import Worker, TableFunctionGenerator
from vgi.catalog import Catalog, Schema, Table, ViewThe recommended pattern is to back tables with TableFunctionGenerator functions. The table schema is automatically derived from the function's output_schema, eliminating duplication:
import pyarrow as pa
from dataclasses import dataclass
from typing import ClassVar
from vgi import TableFunctionGenerator
from vgi.table_function import ProcessParams, OutputCollector
class UsersFunction(TableFunctionGenerator):
"""Generate user data."""
FIXED_SCHEMA: ClassVar[pa.Schema] = pa.schema([
("id", pa.int64()),
("name", pa.string()),
("active", pa.bool_()),
])
@classmethod
def process(cls, params, state, out: OutputCollector) -> None:
out.emit(pa.RecordBatch.from_pydict({
"id": [1, 2, 3],
"name": ["Alice", "Bob", "Carol"],
"active": [True, True, False],
}, schema=params.output_schema))
out.finish()
# Table with auto-derived schema
users_table = Table(
name="users",
function=UsersFunction, # Schema derived from FIXED_SCHEMA
not_null=["id"], # Constraint column names validated
unique=[["id"]],
comment="User accounts",
)from vgi import Worker
from vgi.catalog import Catalog, Schema, Table, View
class MyWorker(Worker):
catalog = Catalog(
name="myapp",
default_schema="main",
schemas=[
Schema(
path=["main"],
comment="Main application data",
tables=[users_table],
views=[
View(
name="active_users",
definition="SELECT * FROM users WHERE active = true",
comment="Active user accounts only",
),
],
functions=[UsersFunction],
),
Schema(
path=["analytics"],
comment="Analytics data",
tables=[events_table],
functions=[AggregateFunction],
),
],
)
if __name__ == "__main__":
MyWorker().run()| Feature | Description |
|---|---|
| No schema duplication | Function-backed tables derive schema automatically |
| Constraint validation | not_null, unique column names validated at definition time |
| Automatic scan handling | Function-backed tables don't need table_scan_function_get() |
| Type safety | Frozen dataclasses with runtime validation |
For tables not backed by functions, provide the schema explicitly:
# Explicit columns - requires table_scan_function_get() implementation
config_table = Table(
name="config",
columns=pa.schema([
("key", pa.string()),
("value", pa.string()),
]),
not_null=["key"],
unique=[["key"]],
)Note: Tables with explicit columns require the worker to implement table_scan_function_get() to tell DuckDB how to scan the data.
Declarative catalogs include comprehensive validation:
# Error: missing columns or function
Table(name="bad") # ValueError: must specify either 'columns' or 'function'
# Error: invalid constraint column
Table(
name="users",
columns=pa.schema([("id", pa.int64())]),
not_null=["nonexistent"], # ValueError: column 'nonexistent' not found
)
# Error: default_schema not in schemas
Catalog(
name="myapp",
default_schema="missing",
schemas=[Schema(path=["main"])], # ValueError: default_schema 'missing' not found
)Workers with functions automatically get a ReadOnlyCatalogInterface:
from vgi import Worker, ScalarFunction
class MyWorker(Worker):
functions = [MyFunction, OtherFunction]
catalog_name = "my_catalog" # Default: "functions"
# Automatically creates ReadOnlyCatalogInterface exposing functionsSet catalog_interface for full control:
from vgi import Worker
from vgi.catalog import CatalogInterface
class MyFullCatalog(CatalogInterface):
# ... implement abstract methods
class MyWorker(Worker):
catalog_interface = MyFullCatalog
functions = [] # Optional: functions can still be registeredTo disable the catalog interface entirely:
class MyWorker(Worker):
catalog_interface = None
catalog_name = None # Required to fully disable
functions = [...]A catalog-less worker is reachable only from the pure-Python [Client][], which
binds by (schema, function name) and needs no attachment — its functions
register into the default main schema. It is not reachable from DuckDB:
ATTACH ... (TYPE vgi) requires a catalog, and there is no standalone
call form. If the worker should be usable from SQL, give it a catalog.
The CatalogClientMixin adds catalog methods to the VGI Client:
from vgi.client import Client
from vgi.client.catalog_mixin import CatalogClientMixin
class CatalogClient(CatalogClientMixin, Client):
pass
# Connect and interact with catalog
client = CatalogClient("./my-worker")
# List catalogs
catalogs = client.catalogs() # ["my_catalog"]
# Attach to a catalog
result = client.catalog_attach(name="my_catalog", options={})
attach_opaque_data = result.attach_opaque_data
# List schemas
for schema in client.schemas(attach_opaque_data=attach_opaque_data):
print(f"Schema: {schema.path}")
# Get schema contents (tables, views, functions)
for obj in client.schema_contents(attach_opaque_data=attach_opaque_data, path=["main"]):
if isinstance(obj, TableInfo):
print(f"Table: {obj.name}")
elif isinstance(obj, ViewInfo):
print(f"View: {obj.name}")
elif isinstance(obj, FunctionInfo):
print(f"Function: {obj.name}")
# Or load the whole catalog at once. Uses one catalog_contents RPC when the
# attach result advertises supports_catalog_contents, else the per-schema RPCs.
snapshot = client.load_catalog(attach=result)
for entry in snapshot.schemas:
print(entry.schema.path, [t.name for t in entry.tables])
# Later: revalidate. A not_modified answer returns the same content cheaply.
snapshot = client.load_catalog(attach=result, previous=snapshot)
# Get only scalar functions using type filter
from vgi.catalog import SchemaObjectType
for obj in client.schema_contents(
attach_opaque_data=attach_opaque_data, path=["main"], type=SchemaObjectType.SCALAR_FUNCTION
):
print(f"Scalar Function: {obj.name}")
# Detach when done
client.catalog_detach(attach_opaque_data=attach_opaque_data)| Method | Description |
|---|---|
catalogs() |
List catalog names |
catalog_attach() |
Attach to a catalog |
catalog_detach() |
Detach from a catalog |
catalog_create() |
Create a new catalog |
catalog_drop() |
Drop a catalog |
catalog_version() |
Get catalog version |
schemas() |
List schemas |
schema_get() |
Get schema info |
schema_create() |
Create a schema |
schema_drop() |
Drop a schema |
schema_contents() |
List schema contents (optional type filter) |
load_catalog() |
Load the whole catalog: one catalog_contents call when the attach advertises it, else (or on failure, or inside a transaction) catalog_schemas + the per-schema RPCs; pass previous= to revalidate with if_none_match |
contents() |
The raw catalog_contents RPC (only valid when advertised; if_none_match revalidates) |
table_get() |
Get table info |
table_create() |
Create a table |
table_drop() |
Drop a table |
table_rename() |
Rename a table |
table_comment_set() |
Set table comment |
table_column_add() |
Add a column |
table_column_drop() |
Drop a column |
table_column_rename() |
Rename a column |
table_scan_function_get() |
Get scan function |
view_get() |
Get view info |
view_create() |
Create a view |
view_drop() |
Drop a view |
view_rename() |
Rename a view |
view_comment_set() |
Set view comment |
catalog_transaction_begin() |
Begin transaction |
catalog_transaction_commit() |
Commit transaction |
catalog_transaction_rollback() |
Rollback transaction |
For stateful catalogs, VGI provides CatalogStorage for persisting attachment and transaction state across worker processes.
SQLite-backed storage with WAL mode for concurrent access:
from vgi.catalog import CatalogStorageSqlite, AttachOpaqueData
# Use default location (~/.local/state/vgi/vgi_catalog.db)
storage = CatalogStorageSqlite()
# Or specify a custom path
storage = CatalogStorageSqlite("/path/to/catalog.db")
# Store attachment
attach_opaque_data = storage.generate_attach_opaque_data()
storage.attach_put(attach_opaque_data, catalog_name="mydb", options={"key": "value"})
# Retrieve attachment
result = storage.attach_get(attach_opaque_data) # ("mydb", {"key": "value"})
# List all attachments
all_ids = storage.attach_list()
# Delete attachment
storage.attach_delete(attach_opaque_data)from typing import Protocol
from vgi.catalog import AttachOpaqueData, TransactionOpaqueData
class CatalogStorage(Protocol):
# Attachment state
def attach_put(self, attach_opaque_data: AttachOpaqueData, catalog_name: str, options: dict) -> None: ...
def attach_get(self, attach_opaque_data: AttachOpaqueData) -> tuple[str, dict] | None: ...
def attach_delete(self, attach_opaque_data: AttachOpaqueData) -> None: ...
def attach_list(self) -> list[AttachOpaqueData]: ...
# Transaction state
def transaction_put(self, transaction_opaque_data: TransactionOpaqueData, attach_opaque_data: AttachOpaqueData, state: bytes) -> None: ...
def transaction_get(self, transaction_opaque_data: TransactionOpaqueData) -> tuple[AttachOpaqueData, bytes] | None: ...
def transaction_delete(self, transaction_opaque_data: TransactionOpaqueData) -> None: ...Catalogs can optionally support transactions:
class TransactionalCatalog(CatalogInterface):
def catalog_attach(self, *, name, options) -> CatalogAttachResult:
return CatalogAttachResult(
attach_opaque_data=...,
supports_transactions=True, # Enable transactions
...
)
def catalog_transaction_begin(self, *, attach_opaque_data) -> TransactionOpaqueData:
txn_id = self._create_transaction(attach_opaque_data)
return txn_id
def catalog_transaction_commit(self, *, attach_opaque_data, transaction_opaque_data) -> None:
self._commit_transaction(transaction_opaque_data)
def catalog_transaction_rollback(self, *, attach_opaque_data, transaction_opaque_data) -> None:
self._rollback_transaction(transaction_opaque_data)Transaction Guarantees:
- Transactions MAY span multiple worker processes
- Workers MUST treat
transaction_opaque_dataas opaque bytes - Workers MUST ensure idempotency of commit/rollback
- If
supports_transactions=False, transaction methods won't be called
Errors are returned as exceptions that propagate through the VGI protocol:
| Error | When Raised |
|---|---|
ValueError |
Invalid arguments, object not found |
NotImplementedError |
Method not supported |
CatalogReadOnlyError |
DDL on read-only catalog |
Example error handling:
from vgi.exceptions import CatalogReadOnlyError
class MyReadOnlyCatalog(ReadOnlyCatalogInterface):
def table_create(self, **kwargs) -> None:
# Automatically raises CatalogReadOnlyError
raise CatalogReadOnlyError("Cannot create table: catalog is read-only")from collections.abc import Iterable
from typing import Any
import uuid
from vgi import Worker
from vgi.catalog import (
AttachOpaqueData,
CatalogAttachResult,
CatalogInterface,
SchemaInfo,
TableInfo,
TransactionOpaqueData,
ViewInfo,
)
class SimpleCatalog(CatalogInterface):
"""A minimal catalog with a single schema."""
def __init__(self):
self._attachments: dict[AttachOpaqueData, str] = {}
def catalogs(self) -> Iterable[str]:
return ["simple_db"]
def catalog_attach(self, *, name: str, options: dict[str, Any]) -> CatalogAttachResult:
if name != "simple_db":
raise ValueError(f"Unknown catalog: {name}")
attach_opaque_data = AttachOpaqueData(uuid.uuid4().bytes)
self._attachments[attach_opaque_data] = name
return CatalogAttachResult(
attach_opaque_data=attach_opaque_data,
supports_transactions=False,
supports_time_travel=False,
catalog_version_frozen=True,
catalog_version=1,
attach_opaque_data_required=False,
)
def catalog_detach(self, *, attach_opaque_data: AttachOpaqueData) -> None:
self._attachments.pop(attach_opaque_data, None)
def schema_get(
self, *, attach_opaque_data: AttachOpaqueData, transaction_opaque_data: TransactionOpaqueData | None, path: SchemaPath
) -> SchemaInfo | None:
if path == ["main"]:
return SchemaInfo(
attach_opaque_data=attach_opaque_data,
path=["main"],
comment="Default schema",
tags={},
)
return None
def table_get(
self, *, attach_opaque_data: AttachOpaqueData, transaction_opaque_data: TransactionOpaqueData | None,
schema_path: SchemaPath, name: str
) -> TableInfo | None:
return None # No tables
def view_get(
self, *, attach_opaque_data: AttachOpaqueData, transaction_opaque_data: TransactionOpaqueData | None,
schema_path: SchemaPath, name: str
) -> ViewInfo | None:
return None # No views
class SimpleCatalogWorker(Worker):
catalog_interface = SimpleCatalog
functions = []
if __name__ == "__main__":
SimpleCatalogWorker().run()The current CatalogInterface has the following limitations:
- Functions: Cannot be created or dropped via catalog methods (use
Worker.functions) - Tags: Cannot be updated after object creation
- Schema metadata: Comments and tags cannot be updated on schemas
- Constraints: Only NOT NULL can be added/dropped (no ALTER for UNIQUE/CHECK)
INSERT, UPDATE, and DELETE are supported when the catalog supplies a write
function for the operation. For declarative catalogs, configure the
[Table][vgi.catalog.descriptors.Table] descriptor:
| Operation | Descriptor field | Requirement |
|---|---|---|
INSERT |
insert_function |
A TableInOutGenerator that accepts rows to insert |
UPDATE |
update_function |
A write function and a scan function that provides row IDs |
DELETE |
delete_function |
A write function and a scan function that provides row IDs |
Each field defaults to None, which leaves that operation unsupported. Custom
catalogs implement table_insert_function_get(), table_update_function_get(),
or table_delete_function_get() to return the corresponding ScanFunctionResult.
The worker's write functions implement the changes to the backing data source;
declaring a table alone does not make it writable.
ReadOnlyCatalogInterface rejects catalog DDL, but can expose these write functions
on declarative tables. Table writes do not require the bundled
transactor; that is an optional database-access service with
additional engine requirements.
Index metadata can be declared with [Index][vgi.catalog.descriptors.Index]. Custom
catalogs can implement index_get(), index_create(), and index_drop(); support
depends on the catalog implementation and backing data source.
- Function Lifecycle - Function execution phases
- Function Metadata - Function introspection