Skip to content

symfonic.services.effects.admission

admission

EFX-L / EFX-F-3 — atomic effect admission at every classified port.

Three checkpoints, all conditional, all against the same lock the revert takes:

  • admit — before an effect attempt exists at all;
  • dispatch — the last check before the effect leaves the process;
  • commit — the check at commit time, which is the one EFX-F-3 actually requires, because a gap between "no fence observed" and "effect executed" has to be closed by the atomic operation rather than by ordering.

The fence is consulted before the lease at every checkpoint. Both are true after a revert (a revert revokes the leases it fences), and callers need to be able to tell "a fence stopped this" from "this lease simply aged out" — FencedEffectError subclasses LeaseInvalidError so a caller that only cares about denial still catches both.

EffectAdmissionController

EffectAdmissionController(*, store: InMemoryLeaseStore, fences: FenceLedger, ports: FencedPortRegistry, tracker: InFlightTracker, quarantine: ResultQuarantine, clock: Callable[[], float] = time.time)

One admission path for every classified effect port.

Source code in src/symfonic/services/effects/admission.py
def __init__(
    self,
    *,
    store: InMemoryLeaseStore,
    fences: FenceLedger,
    ports: FencedPortRegistry,
    tracker: InFlightTracker,
    quarantine: ResultQuarantine,
    clock: Callable[[], float] = time.time,
) -> None:
    self._store = store
    self._fences = fences
    self._ports = ports
    self._tracker = tracker
    self._quarantine = quarantine
    self._clock = clock

admit async

admit(*, lease_id: str, port_id: str, operation: str = '') -> EffectTicket

EFX-L-1 — no effect attempt exists without a live, unfenced lease.

Source code in src/symfonic/services/effects/admission.py
async def admit(
    self, *, lease_id: str, port_id: str, operation: str = ""
) -> EffectTicket:
    """EFX-L-1 — no effect attempt exists without a live, unfenced lease."""
    port = self._ports.port(port_id)
    lease = await self._verify(lease_id)
    ticket = EffectTicket(
        ticket_id=f"efx-{uuid.uuid4().hex[:12]}",
        lease_id=lease.lease_id,
        invocation_id=lease.invocation_id,
        port_id=port_id,
        operation=operation,
        admitted_at=self._clock(),
        irreversible=port.irreversible,
    )
    return self._tracker.open(ticket)

commit async

commit(ticket: EffectTicket, result: Any = None) -> CommitOutcome

EFX-F-3 — the conditional check that decides whether the write lands.

Source code in src/symfonic/services/effects/admission.py
async def commit(self, ticket: EffectTicket, result: Any = None) -> CommitOutcome:
    """EFX-F-3 — the conditional check that decides whether the write lands."""
    try:
        await self._verify(ticket.lease_id)
    except LeaseStoreUnavailableError as exc:
        return self._suppress(ticket, f"lease store unavailable: {exc}")
    except FencedEffectError as exc:
        return self._suppress(ticket, f"fence: {exc}")
    except LeaseInvalidError as exc:
        return self._suppress(ticket, f"lease invalid: {exc}")
    self._tracker.mark(ticket.ticket_id, TicketState.COMMITTED)
    return CommitOutcome(
        decision=CommitDecision.COMMITTED, ticket_id=ticket.ticket_id
    )

commit_or_raise async

commit_or_raise(ticket: EffectTicket, result: Any = None) -> CommitOutcome

The raising variant, for callers that treat suppression as fatal.

Source code in src/symfonic/services/effects/admission.py
async def commit_or_raise(
    self, ticket: EffectTicket, result: Any = None
) -> CommitOutcome:
    """The raising variant, for callers that treat suppression as fatal."""
    await self._verify(ticket.lease_id)
    return await self.commit(ticket, result)

deliver_result async

deliver_result(ticket: EffectTicket, result: Any) -> Any

Hand back a port's answer, or quarantine it if the fence has landed.

Source code in src/symfonic/services/effects/admission.py
async def deliver_result(self, ticket: EffectTicket, result: Any) -> Any:
    """Hand back a port's answer, or quarantine it if the fence has landed."""
    try:
        await self._verify(ticket.lease_id)
    except LeaseStoreUnavailableError as exc:
        self._quarantine.quarantine(ticket, f"lease store unavailable: {exc}")
        return None
    except FencedEffectError as exc:
        self._quarantine.quarantine(ticket, f"fence covered this result: {exc}")
        return None
    except LeaseInvalidError as exc:
        self._quarantine.quarantine(ticket, f"lease invalid: {exc}")
        return None
    return result

dispatch async

dispatch(ticket: EffectTicket) -> EffectTicket

The last check before the effect leaves. Raising here costs nothing.

Source code in src/symfonic/services/effects/admission.py
async def dispatch(self, ticket: EffectTicket) -> EffectTicket:
    """The last check before the effect leaves. Raising here costs nothing."""
    await self._verify(ticket.lease_id)
    self._tracker.mark_dispatched(ticket.ticket_id, self._clock())
    return ticket