Lossless in-process projections of the one internal event stream.
Every class attaches through its compiled G9 row. None of these creates
storage, so their coupling is rendezvous and upstream progress follows
consumer progress. The one adapter that does own storage and a worker โ the
callback fan-out โ lives in :mod:symfonic.kernel.fanout, because owning
work is a different responsibility from projecting a stream.
CallbackEventAdapter
CallbackEventAdapter(policy: EventAdapter, callbacks: Iterable[Callable[[KernelEvent], Any]], *, deadline_seconds: float | None, context: RequestContext | None = None)
Awaited, ordered, error-isolated callback event delivery.
Source code in src/symfonic/kernel/fanout.py
| def __init__(
self,
policy: EventAdapter,
callbacks: Iterable[Callable[[KernelEvent], Any]],
*,
deadline_seconds: float | None,
context: RequestContext | None = None,
) -> None:
self._policy = policy
self._callbacks = tuple(callbacks)
self._deadline_seconds = deadline_seconds
self._context = context
self._buffer = (
BoundedEventBuffer(policy, deadline_seconds=deadline_seconds)
if policy.buffer == "bounded"
else None
)
self._worker: asyncio.Task[Any] | None = None
self._terminal_owed = False
self.failures = 0
|
ResultCollector
ResultCollector(policy: EventAdapter, response: Any)
The single assembler for the public blocking result.
Source code in src/symfonic/kernel/adapters.py
| def __init__(self, policy: EventAdapter, response: Any) -> None:
self._policy = policy
self._response = response
|
StructuredOutputAdapter
StructuredOutputAdapter(policy: EventAdapter, response: Any)
Run terminal extraction through the response port under its G9 row.
Source code in src/symfonic/kernel/adapters.py
| def __init__(self, policy: EventAdapter, response: Any) -> None:
self._policy = policy
self._response = response
|
TextStreamAdapter
TextStreamAdapter(policy: EventAdapter)
Project answer deltas to strings while consuming the terminal event.
Source code in src/symfonic/kernel/adapters.py
| def __init__(self, policy: EventAdapter) -> None:
self._policy = policy
|
TypedStreamAdapter
TypedStreamAdapter(policy: EventAdapter, response: Any)
Project kernel events through ResponsePort.build_event in order.
Source code in src/symfonic/kernel/adapters.py
| def __init__(self, policy: EventAdapter, response: Any) -> None:
self._policy = policy
self._response = response
|