Skip to main content

clika_runtime.runtime functions

astream​

astream(pipe: 'Pipeline', request: 'Request') -> 'AsyncIterator[ResponseChunk]'

Stream request through pipe as an async iterator of chunks.

Each blocking read of the underlying :class:ResponseStream runs in the event loop's default executor: the read parks on a condition variable with the interpreter lock released, so the worker thread costs nothing while it waits and wakes the moment a chunk lands (a polled non-blocking read would trade wake latency for CPU). One worker thread is occupied per stream being awaited. Canceling the coroutine closes the stream; the session itself runs to completion. A failed session raises RuntimeError from the iteration, as the blocking form does.

Example::

async for chunk in astream(pipe, request):
print(chunk.outputs)

from_graph​

from_graph(graph: "'ModelGraph | TracedGraph'", *, kv_batch: 'int' = 1, kv_past_len: 'int' = 0) -> 'Node'

A pipeline node over a compiled graph.

graph is a ModelGraph, or the TracedGraph :func:clika_runtime.trace returns (its .graph is the ModelGraph). kv_batch and kv_past_len size a KV-cached graph's session state.

function_node​

function_node(*args, **kwargs)

function_node(schema: clika_runtime._core.runtime.Schema, *, run_once: dict | None = None, iterate: dict | None = None, make_state: object | None = None) -> clika_runtime._core.runtime.Node

A node served by Python callables. run_once maps a phase name to fn(ctx: PhaseContext) -> None (read ctx.inputs, write ctx.outputs); iterate maps a phase name to fn(ctx: StepContext) -> list[StepResult] (one per ctx.slots entry, in order). Phases run in the order given: every run_once phase, then every iterate phase. make_state(stream) returns the per-session state object (a stateful schema); an object exposing occupancy() and capacity() (and optionally step_tokens()) joins the continuous batcher. Callables run on the engine's thread while the request executes; an exception fails that session and surfaces where its result is read. The context objects are valid only during the call.

pipeline​

pipeline(steps: 'Sequence[Mapping[str, Any]]', *, external_inputs: 'Sequence[Any]' = (), external_outputs: 'Sequence[Any]' = ()) -> 'Pipeline'

Build a pipeline from step dicts (keys: name, node, and optionally exec, input_map, output_map) plus the external I/O schema (TensorSpec lists).

A step's node is a :class:Node, or a graph to wrap in one on the way: a ModelGraph or the TracedGraph :func:clika_runtime.trace returns (see :func:from_graph). input_map / output_map take a dict or a sequence of (name, name) pairs: an input_map pair is (pipeline input, the step input it feeds); an output_map pair is (step output, the pipeline output it publishes as). For a single step with no maps, omitted external specs default to that step's own schema. The pipeline holds its nodes alive.