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 turn

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.

Instance variables

var io : ReportingDestinationIO
var io_timeout_seconds : int
var lease_seconds : int
var store : ReportingMaterializerStore
var 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