yera.runtime.stream
Event stream implementation for inter-process app communication.
Symbols
await_input
await_input(
request_id: str,
expected_type: type[_InputDataT],
value_type: type[_InputValueT],
timeout: float | None = None,
) → tuple[_InputDataT, _InputValueT]Wait for a correlated typed response.
Parameters
Block ID of the request awaiting a response.
Typed input model expected for the request.
Python type represented by the serialized payload.
Maximum seconds to wait, or None to wait indefinitely.
Returns
Typed input data accepted for the request.
Raises
If the response targets another request.
If the response carries another input type.
If the submitted payload cannot be deserialized.
EventStream
Inter-process event stream backed by multiprocessing queues.
Methods
EventStream.push_output
push_output(
event: OutputEvent,
) → NonePut an output event onto the queue.
EventStream.push_input
push_input(
in_event: InputEvent,
) → NonePut an input event onto the queue.
EventStream.pop_output
pop_output(
timeout: float | None = None,
) → OutputEventRemove and return the next output event, blocking up to timeout.
EventStream.close
close() → NoneWake blocked output consumers and prevent duplicate close signals.
EventStream.pop_input
pop_input(
timeout: float | None = None,
) → InputEventRemove and return the next input event, blocking up to timeout.
EventStream.iter_output_blocking
iter_output_blocking(
timeout: float = 0.5,
should_continue: Callable[[], bool] | None = None,
) → Iterator[OutputEvent]Yield output events using blocking get(timeout). Reliable across processes.
Do not use iter_output() when the producer is in another process: Queue.empty() is unreliable across processes and the consumer may never see the event. This method blocks on get(timeout=...) so the exit event is received reliably.
EventStream.set_current
set_current() → NoneBind this stream to the current context variable.
EventStream.get_current
get_current() → EventStreamReturn the stream bound to the current context, raising if none is set.
EventStream.build
build() → EventStreamBuild an event stream object or return the one currently in-context.
EventStream.new
new() → EventStreamCreate a fresh stream, leaving the ambient context untouched.
Returns
A tuple of the new EventStream and the manager owning its queues. The caller must keep the manager alive and shut it down when done.
EventStream.drain_output
drain_output(
timeout: float = 0.1,
) → list[OutputEvent]Remove and return all immediately available output events.
Does not guarantee the queue is empty afterwards — a producer in another process may still be flushing events.
Returns
The events readable without blocking, in queue order.
failure_exit_event_already_pushed
failure_exit_event_already_pushed() → boolReturn True if a nested app already pushed a failure exit event.
is_restricted_stream
is_restricted_stream() → boolReturn whether the current event stream is restricted (only certain block types allowed).
push_input
push_input(
request_id: str,
data: InputData,
) → NonePush an input event to the current stream.
push_output
push_output(
event: OutputEvent,
) → NonePush an output event to the current stream.