Aggregate Functions
This page shows how to write a DuckDB aggregate function in Rust with quack-rs:
the callbacks DuckDB calls, their signatures, and how AggregateFunctionBuilder
registers them. An aggregate function reduces many rows to one value per group,
like SUM(), COUNT() or AVG(). Because DuckDB aggregates in parallel, it
also has a combine step that merges partial results from parallel workers.
Known DuckDB limitation
Out-of-bounds state reads. Two query shapes make DuckDB call every C-API aggregate's
updatewith a state array holding one state while passingcount > 1rows, so the callback readsstates[1..count]past the end of the array (undefined behaviour, in any C-API aggregate, whether built with quack-rs or by hand):
- Window aggregates whose frame is the whole partition, e.g.
agg(x) OVER ()—WindowConstantAggregator(src/function/window/window_constant_aggregator.cpp, ~lines 106 and 296–299 in DuckDB 1.5.5).- Ordered aggregates, e.g.
agg(x ORDER BY y)—src/function/aggregate/sorted_aggregate_function.cpp, ~lines 630–633.Both paths pass a
CONSTANT_VECTORof states because the function has nosimple_update(which the C API cannot set), andCAPIAggregateUpdate(src/main/capi/aggregate_function-c.cpp, ~lines 92–110) passes the vector's data pointer to the extension without flattening it. This is a defect in DuckDB's C API, not in quack-rs, and it cannot be detected from inside the callback: readingstates[1]to check is itself the out-of-bounds read. Until DuckDB fixes it, do not use C-API aggregates in those two query shapes.Reported upstream as duckdb/duckdb#26109.
The aggregate lifecycle
flowchart TD
REG["<b>Registration</b><br/>AggregateFunctionBuilder<br/>→ duckdb_register_aggregate_function"]
REG --> SIZE
SIZE --> INIT
INIT --> UPDATE
UPDATE --> COMBINE
COMBINE --> FINAL
FINAL --> DESTROY
SIZE["<b>state_size</b>()<br/>How many bytes to allocate per group?"]
INIT["<b>state_init</b>(state)<br/>Initialise a fresh state"]
UPDATE["<b>update</b>(chunk, states[])<br/>Process one input batch<br/>(NULL rows included — check is_valid)"]
COMBINE["<b>combine</b>(src[], tgt[], count)<br/>Merge partial results from parallel workers<br/>⚠️ Pitfall L1: target starts fresh — copy ALL config fields"]
FINAL["<b>finalize</b>(states[], out, count, offset)<br/>Write count results at out[offset..], once per result batch"]
DESTROY["<b>state_destroy</b>(states[], count)<br/>Free memory — after finalize and for combine<br/>sources after the merge (not every state: see Known Limitations)"]
style COMBINE fill:#fff3cd,stroke:#e6ac00,color:#333
DuckDB may call combine many times as it merges partial results. A combine
target holds whatever state_init set up, not a copy of the source, so combine
must carry every field across (Pitfall L1). state_size is called whenever an
operator sizes its state buffers, not once at registration, so it must always
return the same value. destroy runs after finalize, and on combine's source
states once they have been merged.
combine must leave its source states unchanged
(Pitfall L15).
A window's segment tree combines the same state into every frame that covers it,
from several threads at once. A combine that moves data out of its source
(mem::take, or zeroing a counter) is right for the first frame and wrong for the
rest: in quack-rs's regression test, a sliding-window sum written that way was
wrong on 4985 of 5000 rows. Read the source and copy or clone what the target
needs.
Registration
#![allow(unused)] fn main() { use libduckdb_sys::{duckdb_aggregate_state, duckdb_bind_info, duckdb_connection, duckdb_data_chunk, duckdb_function_info, duckdb_init_info, duckdb_vector, idx_t}; use quack_rs::prelude::*; #[derive(Default)] struct MyState { count: i64 } impl AggregateState for MyState {} unsafe extern "C" fn update(_: duckdb_function_info, _: duckdb_data_chunk, _: *mut duckdb_aggregate_state) {} unsafe extern "C" fn combine(_: duckdb_function_info, _: *mut duckdb_aggregate_state, _: *mut duckdb_aggregate_state, _: idx_t) {} unsafe extern "C" fn finalize(_: duckdb_function_info, _: *mut duckdb_aggregate_state, _: duckdb_vector, _: idx_t, _: idx_t) {} use quack_rs::aggregate::AggregateFunctionBuilder; use quack_rs::types::TypeId; unsafe fn register(con: duckdb_connection) -> Result<(), ExtensionError> { unsafe { AggregateFunctionBuilder::new("my_agg") .param(TypeId::BigInt) // input type(s) .returns(TypeId::BigInt) // output type .ffi_state::<MyState>() // state_size + init + destructor .update(update) .combine(combine) .finalize(finalize) .register(con)?; } Ok(()) } }
register returns an error if the return type or any of the five required
callbacks (state_size, init, update, combine, finalize) is missing. The
builder treats the destructor as optional (without one, register installs a
no-op destructor; see
Pitfall L13),
but FfiState<T> needs its destroy_callback to drop each T.
.ffi_state::<MyState>() sets state_size, init and destructor together from
FfiState<MyState>, so the three cannot describe different states; see
State Management. update,
combine and finalize read the state through FfiState::<MyState>::with_state /
with_state_mut with the same type.
Callback signatures
With FfiState<T> you do not write state_size, init or destroy yourself:
ffi_state::<T>() installs FfiState::<T>::size_callback, init_callback and
destroy_callback. The wrappers below show the signature DuckDB calls each one
with, and what it does. update, combine and finalize are yours to write.
state_size
#![allow(unused)] fn main() { use libduckdb_sys::{duckdb_aggregate_state, duckdb_bind_info, duckdb_connection, duckdb_data_chunk, duckdb_function_info, duckdb_init_info, duckdb_vector, idx_t}; use quack_rs::prelude::*; #[derive(Default)] struct MyState { config_field: i64, accumulator: i64 } impl MyState { fn accumulate(&mut self, v: i64) { self.accumulator += v; } fn result(&self) -> i64 { self.accumulator } } impl AggregateState for MyState {} unsafe extern "C" fn state_size(info: duckdb_function_info) -> idx_t { unsafe { FfiState::<MyState>::size_callback(info) } } }
Returns the number of bytes DuckDB allocates per group, FfiState::<MyState>::size():
a tag word followed by MyState itself, padded to whole words (or, for a state
larger than 256 bytes or aligned more strictly than usize, a Box<MyState>
pointer).
state_init
#![allow(unused)] fn main() { use libduckdb_sys::{duckdb_aggregate_state, duckdb_bind_info, duckdb_connection, duckdb_data_chunk, duckdb_function_info, duckdb_init_info, duckdb_vector, idx_t}; use quack_rs::prelude::*; #[derive(Default)] struct MyState { config_field: i64, accumulator: i64 } impl MyState { fn accumulate(&mut self, v: i64) { self.accumulator += v; } fn result(&self) -> i64 { self.accumulator } } impl AggregateState for MyState {} unsafe extern "C" fn state_init(info: duckdb_function_info, state: duckdb_aggregate_state) { unsafe { FfiState::<MyState>::init_callback(info, state) }; } }
Writes MyState::default() into the DuckDB-allocated state slot (or, for a boxed
state, a Box holding it), then writes the tag that marks the slot initialised. A
panic in default() is caught and reported to DuckDB as a query error.
update
#![allow(unused)] fn main() { use libduckdb_sys::{duckdb_aggregate_state, duckdb_bind_info, duckdb_connection, duckdb_data_chunk, duckdb_function_info, duckdb_init_info, duckdb_vector, idx_t}; use quack_rs::prelude::*; #[derive(Default)] struct MyState { config_field: i64, accumulator: i64 } impl MyState { fn accumulate(&mut self, v: i64) { self.accumulator += v; } fn result(&self) -> i64 { self.accumulator } } impl AggregateState for MyState {} unsafe extern "C" fn update( _info: duckdb_function_info, input: duckdb_data_chunk, states: *mut duckdb_aggregate_state, ) { let reader = unsafe { VectorReader::new(input, 0) }; let row_count = reader.row_count(); for row in 0..row_count { if unsafe { !reader.is_valid(row) } { continue; } let value = unsafe { reader.read_i64(row) }; let state_ptr = unsafe { *states.add(row) }; if let Some(st) = unsafe { FfiState::<MyState>::with_state_mut(state_ptr) } { st.accumulate(value); } } } }
states[row] is the state of that row's group; rows in the same group share one
state. update receives NULL rows too, so skip rows where is_valid is false (see
NULL Handling).
combine
#![allow(unused)] fn main() { use libduckdb_sys::{duckdb_aggregate_state, duckdb_bind_info, duckdb_connection, duckdb_data_chunk, duckdb_function_info, duckdb_init_info, duckdb_vector, idx_t}; use quack_rs::prelude::*; #[derive(Default)] struct MyState { config_field: i64, accumulator: i64 } impl MyState { fn accumulate(&mut self, v: i64) { self.accumulator += v; } fn result(&self) -> i64 { self.accumulator } } impl AggregateState for MyState {} unsafe extern "C" fn combine( _info: duckdb_function_info, source: *mut duckdb_aggregate_state, target: *mut duckdb_aggregate_state, count: idx_t, ) { for i in 0..count as usize { let src = unsafe { FfiState::<MyState>::with_state(*source.add(i)) }; let tgt = unsafe { FfiState::<MyState>::with_state_mut(*target.add(i)) }; if let (Some(s), Some(t)) = (src, tgt) { // ⚠️ MUST copy ALL fields — see Pitfall L1 t.config_field = s.config_field; // configuration t.accumulator += s.accumulator; // data } } } }
Pitfall L1 — critical: Target states are fresh states, set up by
state_init(withFfiState<T>::init_callback, aT::default()), not copies of the source. You must copy every field, including configuration fields set duringupdate. Forgetting even one config field produces silently wrong results. See Pitfall L1.
finalize
#![allow(unused)] fn main() { use libduckdb_sys::{duckdb_aggregate_state, duckdb_bind_info, duckdb_connection, duckdb_data_chunk, duckdb_function_info, duckdb_init_info, duckdb_vector, idx_t}; use quack_rs::prelude::*; #[derive(Default)] struct MyState { config_field: i64, accumulator: i64 } impl MyState { fn accumulate(&mut self, v: i64) { self.accumulator += v; } fn result(&self) -> i64 { self.accumulator } } impl AggregateState for MyState {} unsafe extern "C" fn finalize( _info: duckdb_function_info, source: *mut duckdb_aggregate_state, result: duckdb_vector, count: idx_t, offset: idx_t, ) { let mut writer = unsafe { VectorWriter::new(result) }; for i in 0..count as usize { let state_ptr = unsafe { *source.add(i) }; match unsafe { FfiState::<MyState>::with_state(state_ptr) } { Some(st) => unsafe { writer.write_i64(offset as usize + i, st.result()) }, None => unsafe { writer.set_null(offset as usize + i) }, } } } }
offset is non-zero when DuckDB writes the results into part of a larger vector.
Always add it to the output index.
state_destroy
#![allow(unused)] fn main() { use libduckdb_sys::{duckdb_aggregate_state, duckdb_bind_info, duckdb_connection, duckdb_data_chunk, duckdb_function_info, duckdb_init_info, duckdb_vector, idx_t}; use quack_rs::prelude::*; #[derive(Default)] struct MyState { config_field: i64, accumulator: i64 } impl AggregateState for MyState {} unsafe extern "C" fn state_destroy(states: *mut duckdb_aggregate_state, count: idx_t) { unsafe { FfiState::<MyState>::destroy_callback(states, count) }; } }
destroy_callback drops the T in each state whose tag matches (freeing its box,
if T is boxed), clearing the tag first, so a second call on the same state is a
no-op. See Pitfall L2.
Complex parameter and return types
For functions that accept or return parameterized types like LIST(BIGINT),
MAP(VARCHAR, INTEGER), or STRUCT(...), use param_logical and
returns_logical instead of param and returns:
#![allow(unused)] fn main() { use libduckdb_sys::{duckdb_aggregate_state, duckdb_bind_info, duckdb_connection, duckdb_data_chunk, duckdb_function_info, duckdb_init_info, duckdb_vector, idx_t}; use quack_rs::prelude::*; #[derive(Default)] struct MyState { count: i64 } impl AggregateState for MyState {} unsafe extern "C" fn update(_: duckdb_function_info, _: duckdb_data_chunk, _: *mut duckdb_aggregate_state) {} unsafe extern "C" fn combine(_: duckdb_function_info, _: *mut duckdb_aggregate_state, _: *mut duckdb_aggregate_state, _: idx_t) {} unsafe extern "C" fn finalize(_: duckdb_function_info, _: *mut duckdb_aggregate_state, _: duckdb_vector, _: idx_t, _: idx_t) {} use quack_rs::aggregate::AggregateFunctionBuilder; use quack_rs::types::{LogicalType, TypeId}; unsafe fn register(con: duckdb_connection) -> Result<(), ExtensionError> { unsafe { AggregateFunctionBuilder::new("retention") .param(TypeId::Boolean) .param(TypeId::Boolean) .returns_logical(LogicalType::list(TypeId::Boolean)) // LIST(BOOLEAN) .ffi_state::<MyState>() // state_size + init + destructor .update(update) .combine(combine) .finalize(finalize) .register(con)?; } Ok(()) } }
param_logical and param can be interleaved — the parameter position is
determined by the total number of calls made so far:
#![allow(unused)] fn main() { use libduckdb_sys::{duckdb_aggregate_state, duckdb_bind_info, duckdb_connection, duckdb_data_chunk, duckdb_function_info, duckdb_init_info, duckdb_vector, idx_t}; use quack_rs::prelude::*; fn demo() { let _ = AggregateFunctionBuilder::new("my_func") .param(TypeId::Varchar) // position 0: VARCHAR .param_logical(LogicalType::list(TypeId::BigInt)) // position 1: LIST(BIGINT) .param(TypeId::Integer) // position 2: INTEGER .returns(TypeId::BigInt) // ... ; } }
If both returns and returns_logical are called, the logical type takes precedence.
Extra info
extra_info attaches arbitrary data to an aggregate function, for example
configuration that parameterises its behaviour. DuckDB calls the destroy callback
to free the data when the function is dropped:
#![allow(unused)] fn main() { use libduckdb_sys::{duckdb_aggregate_state, duckdb_bind_info, duckdb_connection, duckdb_data_chunk, duckdb_function_info, duckdb_init_info, duckdb_vector, idx_t}; use quack_rs::prelude::*; #[derive(Default)] struct MyState { count: i64 } impl AggregateState for MyState {} unsafe extern "C" fn update(_: duckdb_function_info, _: duckdb_data_chunk, _: *mut duckdb_aggregate_state) {} unsafe extern "C" fn combine(_: duckdb_function_info, _: *mut duckdb_aggregate_state, _: *mut duckdb_aggregate_state, _: idx_t) {} unsafe extern "C" fn finalize(_: duckdb_function_info, _: *mut duckdb_aggregate_state, _: duckdb_vector, _: idx_t, _: idx_t) {} unsafe extern "C" fn my_destroy(p: *mut std::os::raw::c_void) { drop(unsafe { Box::from_raw(p.cast::<u64>()) }); } unsafe fn demo(con: duckdb_connection) -> Result<(), ExtensionError> { use std::os::raw::c_void; let config = Box::into_raw(Box::new(42u64)).cast::<c_void>(); unsafe { AggregateFunctionBuilder::new("my_agg") .param(TypeId::BigInt) .returns(TypeId::BigInt) .extra_info(config, Some(my_destroy)) .ffi_state::<MyState>() // state_size + init + destructor .update(update) .combine(combine) .finalize(finalize) .register(con)?; } Ok(()) } }
Inside callbacks, retrieve the extra info with AggregateFunctionInfo::get_extra_info().
AggregateFunctionInfo
AggregateFunctionInfo wraps the duckdb_function_info handle that DuckDB passes
to every aggregate callback except the destructor. It exposes:
get_extra_info() -> *mut c_void: the extra-info pointer set at registration.set_error(message): fails the current query withmessage. Called fromfinalize, it also leaves some states undestroyed; see Known Limitations.
#![allow(unused)] fn main() { use libduckdb_sys::{duckdb_aggregate_state, duckdb_bind_info, duckdb_connection, duckdb_data_chunk, duckdb_function_info, duckdb_init_info, duckdb_vector, idx_t}; use quack_rs::prelude::*; use quack_rs::aggregate::AggregateFunctionInfo; unsafe extern "C" fn update( info: duckdb_function_info, input: duckdb_data_chunk, states: *mut duckdb_aggregate_state, ) { let info = unsafe { AggregateFunctionInfo::new(info) }; let extra = unsafe { info.get_extra_info() }; // ... use extra info, or report errors via info.set_error("...") ... } }
Next steps
- State Management —
FfiState<T>,AggregateState, and lifecycle details - Overloading with Function Sets — register multiple signatures under one name