Skip to content

symfonic.kernel.backpressure

backpressure

Run-scoped bounded hand-off for event adapters (BP-1…BP-13).

Rendezvous adapters need no storage and are implemented by direct await or async iteration in :mod:symfonic.kernel.adapters. This module is the sole queue implementation for adapters whose compiled G9 row opts into decoupling.

AdapterMetrics dataclass

AdapterMetrics(high_watermark: int = 0, byte_high_watermark: int = 0, events_shed: dict[str, int] = dict(), blocked_seconds: float = 0.0, terminal_delivery_failed: bool = False, abandoned: bool = False)

Observable pressure for one adapter on one run (BP-12).

BoundedEventBuffer

BoundedEventBuffer(policy: EventAdapter, *, deadline_seconds: float | None = None)

A per-run queue bounded by both event count and approximate bytes.

One slot is held back from non-terminal events under reserve. Under preempt, a terminal event evicts only the oldest non-terminal entries. No module-global registry or storage is involved (BP-8/BP-13).

Source code in src/symfonic/kernel/backpressure.py
def __init__(
    self, policy: EventAdapter, *, deadline_seconds: float | None = None
) -> None:
    if policy.buffer != "bounded":
        raise ContractViolationError(
            f"{policy.name!r} is not declared as a bounded adapter (BP-1)."
        )
    self.policy = policy
    self.metrics = AdapterMetrics()
    self._items: deque[tuple[KernelEvent, int]] = deque()
    self._bytes = 0
    self._closed = False
    self._drops: dict[tuple[str, int], tuple[int, int, int, str]] = {}
    self._condition = asyncio.Condition()
    self._deadline = (
        time.monotonic() + deadline_seconds
        if deadline_seconds is not None
        else None
    )

close async

close(*, abandoned: bool = False) -> None

Detach this run's consumer and release all buffered references.

Source code in src/symfonic/kernel/backpressure.py
async def close(self, *, abandoned: bool = False) -> None:
    """Detach this run's consumer and release all buffered references."""
    async with self._condition:
        self.metrics.abandoned = abandoned
        self._closed = True
        self._items.clear()
        self._drops.clear()
        self._bytes = 0
        self._condition.notify_all()

get async

get() -> KernelEvent

Take the oldest surviving event, preserving emission order.

Source code in src/symfonic/kernel/backpressure.py
async def get(self) -> KernelEvent:
    """Take the oldest surviving event, preserving emission order."""
    async with self._condition:
        while not self._items and not self._drops:
            self._require_open()
            await self._condition.wait()
        if self._drops:
            first_drop = min(item[1] for item in self._drops.values())
            first_item = self._items[0][0].index if self._items else None
            if first_item is None or first_drop < first_item:
                return self._drop_notice()
        event, size = self._items.popleft()
        self._bytes -= size
        self._condition.notify_all()
        return event

put async

put(event: KernelEvent) -> bool

Put an event, returning False only for a declared shed.

Source code in src/symfonic/kernel/backpressure.py
async def put(self, event: KernelEvent) -> bool:
    """Put an event, returning ``False`` only for a declared shed."""
    size = _event_size(event)
    terminal = event.kind in _TERMINAL
    async with self._condition:
        self._require_open()
        if terminal and self.policy.terminal_policy == "preempt":
            self._preempt_until_fits(size)
        if self._fits(size, terminal):
            self._append(event, size)
            return True
        if not terminal and event.kind in self.policy.sheddable:
            self.metrics.record_shed(event.kind)
            self._record_drop(event)
            return False
        await self._wait_for_space(size, terminal)
        self._require_open()
        self._append(event, size)
        return True