yera.runtime.stream

Event stream implementation for inter-process app communication.

Symbols

def await_input — Wait for a correlated typed response.
class EventStream — Inter-process event stream backed by multiprocessing queues.
def failure_exit_event_already_pushed — Return True if a nested app already pushed a failure exit event.
def is_restricted_stream — Return whether the current event stream is restricted (only certain block types allowed).
def push_input — Push an input event to the current stream.
def push_output — Push an output event to the current stream.

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

request_id
type: str

Block ID of the request awaiting a response.

expected_type
type: type[_InputDataT]

Typed input model expected for the request.

value_type
type: type[_InputValueT]

Python type represented by the serialized payload.

timeout
type: float | None = None

Maximum seconds to wait, or None to wait indefinitely.

Returns

type: tuple[_InputDataT, _InputValueT]

Typed input data accepted for the request.

Raises

InputRequestMismatchError

If the response targets another request.

InputTypeMismatchError

If the response carries another input type.

InputValueError

If the submitted payload cannot be deserialized.

EventStream

Inter-process event stream backed by multiprocessing queues.

Methods

push_output — Put an output event onto the queue.
push_input — Put an input event onto the queue.
pop_output — Remove and return the next output event, blocking up to timeout.
close — Wake blocked output consumers and prevent duplicate close signals.
pop_input — Remove and return the next input event, blocking up to timeout.
iter_output_blocking — Yield output events using blocking get(timeout). Reliable across processes.
set_current — Bind this stream to the current context variable.
get_current — Return the stream bound to the current context, raising if none is set.
build — Build an event stream object or return the one currently in-context.
new — Create a fresh stream, leaving the ambient context untouched.
drain_output — Remove and return all immediately available output events.

EventStream.push_output

push_output(
    event: OutputEvent,
) → None

Put an output event onto the queue.

EventStream.push_input

push_input(
    in_event: InputEvent,
) → None

Put an input event onto the queue.

EventStream.pop_output

pop_output(
    timeout: float | None = None,
) → OutputEvent

Remove and return the next output event, blocking up to timeout.

EventStream.close

close() → None

Wake blocked output consumers and prevent duplicate close signals.

EventStream.pop_input

pop_input(
    timeout: float | None = None,
) → InputEvent

Remove 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() → None

Bind this stream to the current context variable.

EventStream.get_current

get_current() → EventStream

Return the stream bound to the current context, raising if none is set.

EventStream.build

build() → EventStream

Build an event stream object or return the one currently in-context.

EventStream.new

new() → EventStream

Create a fresh stream, leaving the ambient context untouched.

Returns

type: EventStream

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

type: list[OutputEvent]

The events readable without blocking, in queue order.

failure_exit_event_already_pushed

failure_exit_event_already_pushed() → bool

Return True if a nested app already pushed a failure exit event.

is_restricted_stream

is_restricted_stream() → bool

Return whether the current event stream is restricted (only certain block types allowed).

push_input

push_input(
    request_id: str,
    data: InputData,
) → None

Push an input event to the current stream.

push_output

push_output(
    event: OutputEvent,
) → None

Push an output event to the current stream.