streaming_reducers

codegen.streaming_reducers

Codegen registry: pipeline reducer callable -> backend streaming reducer.

A streaming reducer replaces the O(window * N^2)/step rolled-buffer recompute of a sliding-window source observation (e.g. compute_fc recomputed over the whole BOLD window every online-tuning iteration) with an incremental O(N^2)/step reducer over the (add, evict, emit, resync) protocol. The reducer is authored declaratively (:class:StreamingReducerSpec) and lowered to backend source by :mod:tvbo.codegen.reducers, so tvbo emits the realization into the generated experiment — no backend ships one.

The registry is per backend: a backend that has a registered streaming form for a given pipeline reducer uses the streaming path in generated code; a backend that registers nothing (tvb / pyrates / julia) keeps the recompute path byte-for-byte unchanged. The whole mechanism is transparent to the recipe / schema — a study keeps declaring pipeline: [compute_fc] and the algorithm codegen swaps in the streaming reducer only where one is registered.

Registering a new streaming reducer (e.g. for dFC / metastability) is a new database/reducers/*.yaml recipe — state + add/evict/resync/emit assignment strings — loaded at import; no code change.

Classes

Name Description
StreamingReducerSpec Declarative, backend-agnostic spec of a windowed streaming reducer (metadata).

StreamingReducerSpec

codegen.streaming_reducers.StreamingReducerSpec(
    state,
    add,
    evict,
    resync,
    emit,
    resync_masked=(),
    emit_kind='window',
)

Declarative, backend-agnostic spec of a windowed streaming reducer (metadata).

The reducer’s math is authored as symbolic assignment recurrences (strings over the state vars, the arriving/leaving sample v, and the window x); the general resolver in :mod:tvbo.codegen.reducers parses them and the sympy printers lower them per backend, so tvbo emits the reducer into the generated experiment and no backend ships a reducer realization (works on any backend/tvboptim version — nothing external to resolve). Sequential reassignment encodes the data flow — a later RHS sees the value assigned to a state var above it — so no new_* temporaries are needed. The recipe is data here; the template that emits it is fully general.

Attributes

Name Type Description
state tuple accumulator state variables, unpacked from acc and returned.
add tuple ordered (lhs, rhs) assignments folding one arriving sample v.
evict tuple ordered assignments dropping one leaving sample v (sliding window).
resync tuple ordered assignments rebuilding the state exactly from the window x.
resync_masked tuple optional assignments rebuilding the state from a max-sized ring buffer whose valid window is the last L rows, selected by the boolean mask m (row-vector over the full buffer) with count L. Present so a reducer can be resynced over a fixed-shape buffer with a traced window length — the tuning scan then compiles once across varying window sizes. Empty when the reducer has no masked form (the masked path is not offered).
emit str readout expression over the final state (the reduced value).
emit_kind str "window" (emit every step; wired today) or "stride" (emit a per-window observation for a downstream reduction — the dFC / FCD follow-on; :func:streaming_capable returns False for it until the template’s stride branch lands, recompute fallback meanwhile).

Functions

Name Description
is_windowed_reducer True when a pipeline reducer is a registered windowed reducer for backend.
lookup_streaming_reducer Return the :class:StreamingReducerSpec for a pipeline reducer, or None.
register_streaming_reducer Register spec as backend’s streaming form of the reducer callable.
streaming_capable True when a step-wise ("window") streaming reducer is registered.

is_windowed_reducer

codegen.streaming_reducers.is_windowed_reducer(
    reducer_module,
    reducer_name,
    backend='tvboptim',
)

True when a pipeline reducer is a registered windowed reducer for backend.

The registered reducer (compute_fc for tvboptim) is a windowed correlation reduction — undefined over a < 2-sample window, where jnp.corrcoef collapses to a scalar — so codegen routes it through a degenerate-window guard. Unlike :func:lookup_streaming_reducer, this is a pure registry-membership test with no factory-availability or TVBO_STREAMING_REDUCERS gating: the guard applies whether or not the streaming path is active. It shares the one registered reducer set so the guard’s routing and the streaming lookup cannot drift.

lookup_streaming_reducer

codegen.streaming_reducers.lookup_streaming_reducer(
    backend,
    reducer_module,
    reducer_name,
)

Return the :class:StreamingReducerSpec for a pipeline reducer, or None.

Returns None (-> recompute fallback) when the backend registers nothing for this reducer, or when streaming is disabled via the TVBO_STREAMING_REDUCERS=0 environment switch (used to codegen the recompute reference for validation and as an escape hatch). The reducer realization is tvbo-emitted, so there is no external factory to resolve.

register_streaming_reducer

codegen.streaming_reducers.register_streaming_reducer(backend, reducer, spec)

Register spec as backend’s streaming form of the reducer callable.

Parameters

Name Type Description Default
backend str Codegen backend key (e.g. "tvboptim"). required
reducer str Fully-qualified pipeline reducer callable the spec replaces (e.g. "tvboptim.observations.observation.compute_fc"). The bare callable name is also matched as a fallback. required
spec StreamingReducerSpec The streaming reducer descriptor. required

streaming_capable

codegen.streaming_reducers.streaming_capable(
    backend,
    reducer_module,
    reducer_name,
)

True when a step-wise ("window") streaming reducer is registered.

The algorithm template gates its streaming branch on this; "stride" (dFC / FCD) specs return False here until the template’s stride branch lands, so registering them is safe (recompute fallback stays in force).