Module adcp.reporting.materializer.service
Service-owned heartbeat around destination I/O outside database transactions.
Classes
class ReportingMaterializerService (store: ReportingMaterializerStore,
io: ReportingDestinationIO,
writer: ReportingDestinationWriter,
lease_seconds: int = 30,
io_timeout_seconds: int = 300)-
Expand source code
@dataclass(frozen=True) class ReportingMaterializerService: """One autonomous turn; schedule repeatedly, including after process restart. The destination must implement its advertised conditional/idempotent write semantics durably. A timeout or interrupted outcome commit retains the same immutable attempt. The SDK opens fresh authorization sessions for write and paginated readback; the service checks current ledger authority before each. """ store: ReportingMaterializerStore = field(repr=False) io: ReportingDestinationIO = field(repr=False) writer: ReportingDestinationWriter = field(repr=False) lease_seconds: int = 30 io_timeout_seconds: int = 300 def __post_init__(self) -> None: validate_lease_seconds(self.lease_seconds) if type(self.io_timeout_seconds) is not int or not 1 <= self.io_timeout_seconds <= 3600: raise ValueError("materializer I/O deadline requires 1..3600 seconds") if type(self.io) is not ReportingDestinationIO: raise failure("UNSUPPORTED_VERIFICATION") @materializer_errors async def run_once(self) -> ReportingMaterializerTurn: keys = tuple( v.key for v in self.io.registry.verifiers if v.key.capability in self.writer.capabilities ) reserved = await self.store.claim_materialization( keys=keys, lease_seconds=self.lease_seconds ) if isinstance(reserved, ReportingMaterializerTurn): return reserved cancel = asyncio.Event() heartbeat = _LeaseHeartbeat(self.store, reserved, self.lease_seconds, cancel) task = asyncio.create_task(heartbeat.run()) context = ReportingIOContext( datetime.now(timezone.utc) + timedelta(seconds=self.io_timeout_seconds), cancel, heartbeat, ) canceled = False turn = ReportingMaterializerTurn( "pending", "effect_unknown", reserved.attempt.reporting_materialization_id ) try: try: await self.store.authorize_materialization(reserved) target = reserved.context if target.delivery is None: raise failure("BINDING_MISMATCH") prepared = await self.io.registry.prepare( key=reserved.request.verification_key, binding=target.binding, delivery=target.delivery, obligation=target.obligation, revisions=target.revisions, attempt=reserved.attempt, reader=self.store, context=context, ) await self.store.authorize_materialization(reserved) locator = await self.io.write(prepared, context=context) await self.store.authorize_materialization(reserved) verified = await self.io.verify(prepared, locator, context=context) turn = await self.store.finish_materialization( reserved, prepared=prepared, verified=verified ) except ReportingWriterError as exc: record = exc.failure # The writer owns whether a failure is known and terminal. All # uncertain I/O and lost leases keep the external identity. if record.code in {"DEADLINE_EXCEEDED", "LEASE_LOST"}: record = ReportingWriterFailure(record.code, "same_identity", "unknown") turn = await self.store.finish_materialization(reserved, error=record) except asyncio.CancelledError: canceled = not heartbeat.lost except Exception: # Includes commit/connection loss and failed atomic outbox enqueue. # The lease expires; a restart reuses the original identity. turn = ReportingMaterializerTurn( "pending", "effect_unknown", reserved.attempt.reporting_materialization_id ) finally: heartbeat.stopped.set() task.cancel() canceled = await _join_tasks(task) or canceled if canceled: raise asyncio.CancelledError return turnOne autonomous turn; schedule repeatedly, including after process restart.
The destination must implement its advertised conditional/idempotent write semantics durably. A timeout or interrupted outcome commit retains the same immutable attempt. The SDK opens fresh authorization sessions for write and paginated readback; the service checks current ledger authority before each.
Instance variables
var io : ReportingDestinationIOvar io_timeout_seconds : intvar lease_seconds : intvar store : ReportingMaterializerStorevar writer : ReportingDestinationWriter
Methods
async def run_once(self) ‑> ReportingMaterializerTurn-
Expand source code
@materializer_errors async def run_once(self) -> ReportingMaterializerTurn: keys = tuple( v.key for v in self.io.registry.verifiers if v.key.capability in self.writer.capabilities ) reserved = await self.store.claim_materialization( keys=keys, lease_seconds=self.lease_seconds ) if isinstance(reserved, ReportingMaterializerTurn): return reserved cancel = asyncio.Event() heartbeat = _LeaseHeartbeat(self.store, reserved, self.lease_seconds, cancel) task = asyncio.create_task(heartbeat.run()) context = ReportingIOContext( datetime.now(timezone.utc) + timedelta(seconds=self.io_timeout_seconds), cancel, heartbeat, ) canceled = False turn = ReportingMaterializerTurn( "pending", "effect_unknown", reserved.attempt.reporting_materialization_id ) try: try: await self.store.authorize_materialization(reserved) target = reserved.context if target.delivery is None: raise failure("BINDING_MISMATCH") prepared = await self.io.registry.prepare( key=reserved.request.verification_key, binding=target.binding, delivery=target.delivery, obligation=target.obligation, revisions=target.revisions, attempt=reserved.attempt, reader=self.store, context=context, ) await self.store.authorize_materialization(reserved) locator = await self.io.write(prepared, context=context) await self.store.authorize_materialization(reserved) verified = await self.io.verify(prepared, locator, context=context) turn = await self.store.finish_materialization( reserved, prepared=prepared, verified=verified ) except ReportingWriterError as exc: record = exc.failure # The writer owns whether a failure is known and terminal. All # uncertain I/O and lost leases keep the external identity. if record.code in {"DEADLINE_EXCEEDED", "LEASE_LOST"}: record = ReportingWriterFailure(record.code, "same_identity", "unknown") turn = await self.store.finish_materialization(reserved, error=record) except asyncio.CancelledError: canceled = not heartbeat.lost except Exception: # Includes commit/connection loss and failed atomic outbox enqueue. # The lease expires; a restart reuses the original identity. turn = ReportingMaterializerTurn( "pending", "effect_unknown", reserved.attempt.reporting_materialization_id ) finally: heartbeat.stopped.set() task.cancel() canceled = await _join_tasks(task) or canceled if canceled: raise asyncio.CancelledError return turn