2. Your first table function
Step two: add a table function to the worker from
step 1. A scalar function transforms a column; a table function
produces rows, so it is called in a FROM clause. About 10 minutes.
Step 1 β Add the function
Section titled βStep 1 β Add the functionβReplace src/main.rs, keeping Double and adding Series:
src/main.rs
// Copyright 2025, 2026 Query Farm LLC - https://query.farm
//! The worker built across the vgi-rust tutorial: one scalar function and one
//! table function in a single catalog.
//!
//! The scalar `double` transforms a column in place. The table function `series`
//! *generates* rows from an argument, so it is called in a FROM clause rather
//! than an expression. One worker can serve any mix of shapes.
//!
//! ```text
//! cargo build --release --bin calc
//! # then, in a Haybarn shell:
//! ATTACH 'calc' (TYPE vgi, LOCATION './target/release/calc');
//! SELECT calc.double(21);
//! SELECT * FROM calc.series(3);
//! ```
use std::sync::Arc;
use arrow_array::cast::AsArray;
use arrow_array::types::Int64Type;
use arrow_array::{Array, ArrayRef, Int64Array, RecordBatch};
use arrow_schema::{DataType, Field, Schema, SchemaRef};
use vgi::catalog::CatalogModel;
use vgi::function::{
ArgSpec, BindParams, BindResponse, FunctionMetadata, ProcessParams, ScalarFunction,
};
use vgi::table_function::{TableFunction, TableProducer};
use vgi::vgi_rpc::OutputCollector;
use vgi::{Result, RpcError};
// ββ scalar: double(n) βββββββββββββββββββββββββββββββββββββββββββββββββββββββ
struct Double;
impl ScalarFunction for Double {
fn name(&self) -> &str {
"double"
}
fn metadata(&self) -> FunctionMetadata {
FunctionMetadata {
description: "Doubles a BIGINT".to_string(),
return_type: Some(DataType::Int64),
..Default::default()
}
}
fn argument_specs(&self) -> Vec<ArgSpec> {
vec![ArgSpec::column("n", 0, "int64", "Value to double")]
}
fn process(&self, params: &ProcessParams, batch: &RecordBatch) -> Result<RecordBatch> {
let n = batch.column(0).as_primitive::<Int64Type>();
let out: Int64Array = (0..n.len())
.map(|i| {
if n.is_valid(i) {
Some(n.value(i) * 2)
} else {
None
}
})
.collect();
RecordBatch::try_new(
params.output_schema.clone(),
vec![Arc::new(out) as ArrayRef],
)
.map_err(|e| RpcError::runtime_error(e.to_string()))
}
}
// ββ table: series(count) ββββββββββββββββββββββββββββββββββββββββββββββββββββ
const BATCH_SIZE: i64 = 1024;
/// The per-scan cursor. A table function is *pulled*: the engine calls
/// `next_batch` until it answers `None`, so whatever the function needs between
/// calls lives here rather than in the function itself (which is shared and
/// must stay `Sync`).
struct SeriesProducer {
schema: SchemaRef,
next: i64,
count: i64,
}
impl TableProducer for SeriesProducer {
fn next_batch(&mut self, _out: &mut OutputCollector) -> Result<Option<RecordBatch>> {
if self.next >= self.count {
// None is end-of-stream. Returning an empty batch instead would
// loop forever.
return Ok(None);
}
let end = (self.next + BATCH_SIZE).min(self.count);
let col: ArrayRef = Arc::new((self.next..end).collect::<Int64Array>());
self.next = end;
RecordBatch::try_new(self.schema.clone(), vec![col])
.map(Some)
.map_err(|e| RpcError::runtime_error(e.to_string()))
}
}
struct Series;
impl TableFunction for Series {
fn name(&self) -> &str {
"series"
}
fn metadata(&self) -> FunctionMetadata {
FunctionMetadata {
description: "Generates the integers 0..count-1".to_string(),
..Default::default()
}
}
fn argument_specs(&self) -> Vec<ArgSpec> {
// const_arg, not column: the value is fixed for the whole scan and read
// at bind. with_ge(0.0) is enforced by the framework before any row is
// produced, so series(-1) fails rather than returning nothing.
vec![ArgSpec::const_arg("count", 0, "int64", "How many numbers to generate").with_ge(0.0)]
}
/// Runs once per query, before any data moves. The output shape is fixed
/// here, so it never inspects the arguments; a function whose columns depend
/// on its arguments would build the schema from `params` instead.
fn on_bind(&self, _params: &BindParams) -> Result<BindResponse> {
Ok(BindResponse {
output_schema: Arc::new(Schema::new(vec![Field::new("n", DataType::Int64, true)])),
opaque_data: Vec::new(),
})
}
/// Runs once per scan, after bind. Arguments are fixed for the whole scan,
/// so this is where they are read β decoding them per batch would be waste.
fn producer(&self, params: &ProcessParams) -> Result<Box<dyn TableProducer>> {
Ok(Box::new(SeriesProducer {
schema: params.output_schema.clone(),
next: 0,
count: params.arguments.const_i64(0).unwrap_or(0),
}))
}
}
fn main() {
let mut worker = vgi::Worker::new();
worker.register_scalar(Double);
worker.register_table(Series);
worker.set_catalog(CatalogModel {
name: "calc".to_string(),
comment: Some("Tutorial worker: a scalar and a table function".to_string()),
..Default::default()
});
worker.run();
}
Five things to notice:
- A table function is two types.
TableFunctionis the registered, shared description βSend + Sync, one instance for the whole worker.TableProduceris the per-scan cursor, built fresh byproducer(). Anything that changes as rows are emitted lives on the producer, which is what lets the function itself stay immutable and shared. next_batchis the pull loop. DuckDB calls it until it answersOk(None). Returning an empty batch instead ofNonedoes not end the stream β it loops forever.- Arguments are read once, in
producer(). They are fixed for the whole scan, so decoding them per batch would be waste. const_arg, notcolumn. A table functionβs arguments are bind-time scalars, not per-row columns.params.arguments.const_i64(0)reads the first one.with_ge(0.0)is enforced. The framework validates declared constraints before your code runs, soseries(-1)fails at bind rather than looping or silently returning nothing.
A table function is pulled, not called: the engine asks the worker for output until it signals
completion. next_batch is the pull handler. series knows up front how many rows it owes, so its
producer is just a counter β but one reading a paginated API would keep its cursor in the same
struct and decide for itself when it is done.
Step 2 β Call it
Section titled βStep 2 β Call itβcargo build --release
Re-attach and query it in a FROM clause:
LOAD vgi;
ATTACH 'calc' (TYPE vgi, LOCATION './target/release/calc');
SELECT * FROM calc.series(3);
Output
| n |
|---|
| 0 |
| 1 |
| 2 |
Both functions live in one worker, so they compose in a single query:
SELECT calc.double(n) AS doubled FROM calc.series(3);
Output
| doubled |
|---|
| 0 |
| 2 |
| 4 |
The batching is real, not decorative β series(5000) streams through in chunks of 1024:
SELECT count(*), sum(n) FROM calc.series(5000);
Output
| count_star() | sum(n) |
|---|---|
| 5000 | 12497500 |
And the constraint does its job:
SELECT * FROM calc.series(-1);
Invalid Input Error: VGI Worker Exception: argument count: must be >= 0
What just happened: one calc worker now serves two functions. series ran in the FROM clause
β DuckDB pulled batches from next_batch until it answered None β and double transformed them,
all in the same process.
Binder Error: Failed to attach database: database with name "calc" already existsβcalcis still attached from step 1. RunDETACH calc;or open a fresh shell.Function "series" is a table function but it was used as a scalar functionβ table functions go inFROM, notSELECT.- The scan never ends β
next_batchmust eventually answerOk(None). An empty batch is not an end-of-stream signal. - Empty result β
series(0)is legitimately zero rows. Tryseries(5).
Next steps
Section titled βNext stepsβ- The other three shapes β Function patterns β table-in-out, aggregate and buffering, each with a runnable worker.
- When each callback fires β Function lifecycle.
- Exact contracts β Table functions.