Skip to content

symfonic.core.streaming.event_processor

event_processor

Model-agnostic stream event processor.

process_stream_events async

process_stream_events(
    events: AsyncIterator[StreamEvent],
    *,
    on_event: Callable[[StreamEvent], Any] | None = None,
    filter_types: set[type] | None = None,
) -> list[StreamEvent]

Process a stream of events, optionally filtering and calling back.

  • Preserves order: events in -> events out in same order
  • Never drops events (ExtensionEvent included)
  • Never raises on unknown event types

Parameters:

Name Type Description Default
events AsyncIterator[StreamEvent]

async iterator of StreamEvent

required
on_event Callable[[StreamEvent], Any] | None

optional callback for each event

None
filter_types set[type] | None

if set, only include events of these types in output

None

Returns:

Type Description
list[StreamEvent]

List of all processed events (filtered if filter_types specified)

Source code in src/symfonic/core/streaming/event_processor.py
async def process_stream_events(
    events: AsyncIterator[StreamEvent],
    *,
    on_event: Callable[[StreamEvent], Any] | None = None,
    filter_types: set[type] | None = None,
) -> list[StreamEvent]:
    """Process a stream of events, optionally filtering and calling back.

    - Preserves order: events in -> events out in same order
    - Never drops events (ExtensionEvent included)
    - Never raises on unknown event types

    Args:
        events: async iterator of StreamEvent
        on_event: optional callback for each event
        filter_types: if set, only include events of these types in output

    Returns:
        List of all processed events (filtered if filter_types specified)
    """
    result: list[StreamEvent] = []
    async for event in events:
        if on_event is not None:
            on_event(event)
        if filter_types is None or type(event) in filter_types:
            result.append(event)
    return result