SAService Architect Back to designer
SAService ArchitectDocumentation
Framework reference Map | Filter | FlatMap | KeyBy | Process | Join | MultiJoin | Delay | Case

Operator functions

Operator contracts invoked by generated Rust streams. Implement only the contract selected by the stream type in the designer.

When to use

Use this page while implementing a generated business function or endpoint handler.

Behavior

  • Every invocation receives the current stream and input value; runtimes with explicit message contexts also pass the context separately.
  • Emission is explicit. Returning from Map, FlatMap, KeyBy, or Process without using an output intentionally produces no downstream value.
  • Process has separate normal and error outputs. Throwing or returning a runtime error is not a substitute for intentionally emitting the graph error type.
  • Join inputs are already grouped by key and retained by the configured join storage before the user callback runs.

Streams without a user callback

Input, Sink, Merge, Split, FlatMap Iterable, Cycle Link, Error, and individual When branches are primarily graph/runtime nodes. Their behavior is generated from topology and configuration; they do not all require a custom operator implementation.

Output and error semantics

Normal output follows persisted graph links. A Process error output follows the dedicated Error link. A callback can emit more than once unless the selected stream contract or downstream protocol imposes a stricter rule.

Operator contracts

User-facing methods selected by generated stream wiring.

Type or methodRoleDescription
MapFunction<T, R>::map(context, stream, value, out)Async trait

Transforms one input value and emits the resulting value through the normal output.

FilterFunction<T>::filter(context, stream, value) -> boolAsync trait

Returns whether the current value should continue to downstream streams.

FlatMapFunction<T, R>::flat_map(context, stream, value, out)Async trait

Emits zero, one, or many output values for one input value.

KeyByFunction<T, K, V>::key_by(context, stream, value, out)Async trait

Creates a key/value pair used by keyed operators and join storage.

ProcessFunction<T, R, E>::process(context, stream, value, out, error)Async trait

Runs general business logic with independent normal and error outputs.

JoinFunction<K, T1, T2, R>::join(context, stream, key, left, right, out) -> boolAsync trait

Receives the current key and the buffered values from the left and right inputs, emits results, and returns the framework completion decision for that key.

MultiJoinFunction<K, T, R>::multi_join(context, stream, key, values, out) -> boolAsync trait

Receives the current key and one buffered value collection per input, emits results, and returns the framework completion decision for that key.

DelayFunction<T>::duration / delay_errorAsync trait

Calculates how long the value should wait; the paired error method decides what to emit when scheduling the delay fails.

BuildSwitchFunction<T>::build_switch(stream, when_items)Case routing

Builds the runtime selector used by a Case stream to choose one of its typed When branches.

Generated extension pattern

#[async_trait]
impl ProcessFunction<Order, Order, String> for ValidateOrder {
  async fn process(&self, context: MessageContext, _stream: &dyn RuntimeStream, value: &Order, out: &Collector<Order>, _error: &Collector<String>) {
    out.collect(context, value.clone()).await;
  }
}