Module adcp.reporting.ledger.status
get_reporting_status: one authoritative read over the obligation ledger.
This is the buyer's answer to "do I have definitive reporting for this period – and if not, whose problem is it?". Polling it is the authoritative recovery path; a polling-only seller is fully conformant, and push is a doorbell, never the source of truth.
The handler is framework-agnostic on purpose.
It takes a plain request mapping
plus an authenticated caller and returns a plain response mapping, so the same
projection serves an :mod:adcp.server handler, a bare ASGI route, a CLI, or a
test.
Nothing here imports a web framework.
Three views:
summary
Health, coverage, data_through, obligation counts, and issues for a
scope.
No records, no cursor.
periods
The paginated flat union of obligations, revisions, adjustments, and (when
the preview ingest is on) this caller's consumer statuses.
Every page from
one cursor repeats the same ledger_snapshot_id and ledger_as_of.
revision
One exact revision plus its adjustments.
Caller isolation is not a filter applied at the end; it is a parameter to every lookup. Unknown and unauthorized identifiers deliberately produce the same shape, because a distinguishable "not found" is an enumeration oracle.
Classes
class ReportingStatusCaller (account_id: str, consumer_id: str)-
Expand source code
@dataclass(frozen=True) class ReportingStatusCaller: """The authenticated caller, resolved from transport, never from the body. ``consumer_id`` is what scopes consumer-status visibility. A request body that asserts a buyer principal is ignored: identity comes from the authenticated transport or it does not exist. """ account_id: str consumer_id: str def __post_init__(self) -> None: from adcp.reporting.evidence import consumer_reference, principal_reference principal_reference(self.account_id) consumer_reference(self.consumer_id)The authenticated caller, resolved from transport, never from the body.
consumer_idis what scopes consumer-status visibility. A request body that asserts a buyer principal is ignored: identity comes from the authenticated transport or it does not exist.Instance variables
var account_id : strvar consumer_id : str
class ReportingStatusHandler (store: ReportingLedgerStore,
*,
page_size: int = 100,
consumer_status_enabled: bool = False,
escalation: ReportingDeliveryEscalation | None = None)-
Expand source code
class ReportingStatusHandler: """Render one captured status snapshot using the same pure projection as push.""" def __init__( self, store: ReportingLedgerStore, *, page_size: int = _DEFAULT_PAGE_SIZE, consumer_status_enabled: bool = False, escalation: ReportingDeliveryEscalation | None = None, ) -> None: self._store = store self._page_size = page_size self._consumer_status_enabled = consumer_status_enabled self._escalation = escalation async def handle( self, request: dict[str, Any], *, caller: ReportingStatusCaller ) -> dict[str, Any]: snapshot = await self._load(caller) return self.render_snapshot(request, caller=caller, snapshot=snapshot) async def _load(self, caller: ReportingStatusCaller) -> ReportingStatusSnapshot: from adcp.reporting.ledger.status_snapshot import ReportingStatusParticipant if isinstance(self._store, ReportingStatusParticipant): return await self._store.read_status_snapshot(caller=caller) # Optional upgrade: custom stores keep their original structural API. # This fallback cannot establish durable status-notification readiness. boundary = await self._store.open_snapshot( caller=caller, filters_fingerprint="status-projection-v1" ) configurations = await self._store.list_configurations(caller=caller) obligations: list[ReportingObligationRecord] = [] revisions: list[ReportingRevisionRecord] = [] adjustments: list[ReportingAdjustmentRecord] = [] statuses: list[ConsumerStatusRecord] = [] offset = 0 changes: list[tuple[int, str, str, str]] = [] while True: page = await self._store.read_page( snapshot=boundary, consumer_id=caller.consumer_id, delivery_config_ids=None, media_buy_ids=None, offset=offset, limit=self._page_size, changes_after_sequence=None, ) obligations.extend(page.obligations) revisions.extend(page.revisions) adjustments.extend(page.adjustments) statuses.extend(page.consumer_statuses) if not page.has_more: break offset += self._page_size lifecycles = [] for key in sorted({mismatch_key(s) for s in statuses}): seen: set[str] = set() while ( issue := await self._store.get_issue(account_id=caller.account_id, issue_key=key) ) is not None: if issue.issue_id in seen or issue.issue_key != key: raise LedgerConflictError( "STATUS_PROJECTION_UNAVAILABLE", "invalid waiver chain" ) seen.add(issue.issue_id) lifecycles.append(issue) if issue.issue_state != "waived": break key = condition_after_waiver(issue) for kind, records, attribute in ( ("obligation", obligations, "reporting_obligation_id"), ("revision", revisions, "reporting_revision_id"), ("adjustment", adjustments, "reporting_adjustment_id"), ("consumer_status", statuses, "reporting_status_id"), ): # The legacy page API has a durable boundary but no per-record # ordinals. Replay its retained records when that boundary advances; # incremental_repair permits identity-deduplicated over-inclusion. # Positional ordinals would instead omit later records after rebuilds. changes.extend( (boundary.max_sequence, kind, getattr(record, attribute), caller.consumer_id) for record in records ) snapshot = ReportingStatusSnapshot( caller.account_id, boundary.ledger_as_of, tuple(configurations), tuple(obligations), tuple(revisions), tuple(statuses), tuple(lifecycles), adjustments=tuple(adjustments), changes=tuple(changes), ) for _ in range(3): intents = lifecycle_intents(snapshot) if not intents: return snapshot for intent in intents: issue = intent.lifecycle if intent.action == "ensure_mismatch": await self._store.ensure_issue_opened( issue_key=issue.issue_key, account_id=issue.account_id, consumer_id=issue.consumer_id, observed_at=issue.opened_at, ) elif intent.action == "retire_mismatch": await self._store.retire_issue( issue_key=issue.issue_key, account_id=issue.account_id, at=snapshot.as_of ) snapshot = apply_intents_to_snapshot(snapshot, intents) raise LedgerConflictError( "STATUS_PROJECTION_UNAVAILABLE", "status lifecycle did not converge" ) def render_snapshot( self, request: dict[str, Any], *, caller: ReportingStatusCaller, snapshot: ReportingStatusSnapshot, reconciliation: tuple[ReportingDeliveryRecord, ...] | None = None, revision_ownership: bool = False, ) -> dict[str, Any]: """Render captured database evidence without any store calls or clock reads.""" view = request.get("view", "summary") if view not in {"summary", "periods", "revision"}: raise LedgerConflictError("INVALID_VIEW", "unsupported reporting status view") if snapshot.account_id != caller.account_id: raise LedgerConflictError("LOOKUP_UNAVAILABLE", "status is unavailable to this caller") filters = _filters(request) version = normalize_to_release_precision( request.get("adcp_version") or resolve_adcp_version(None) ) complete_start_forecast = is_adcp_version_at_least(version, "3.2-rc.6") if complete_start_forecast: # Bind the new wire contract without relabelling explicit rc.3 snapshots. filters["adcp_version"] = version scope = ReportingStatusScope( caller.account_id, consumer_id=caller.consumer_id, ) from adcp.reporting.ledger.delivery_models import ReportingDeliveryPrincipal from adcp.reporting.materializer.capture import private_snapshot snapshot = private_snapshot( snapshot, ReportingDeliveryPrincipal(caller.account_id, caller.consumer_id) ) snapshot_id = ( "rpls_" + _fingerprint( [ caller.account_id, caller.consumer_id, filters, snapshot.max_sequence, ] )[:32] ) if reconciliation is not None: from adcp.reporting.ledger._delivery_state import fingerprint, principal reconciliation = tuple( r for r in reconciliation if principal(r).account_id == caller.account_id and principal(r).consumer_id == caller.consumer_id ) snapshot_id = ( "rpls_" + _fingerprint( [snapshot_id, [fingerprint(r) for r in reconciliation], _iso(snapshot.as_of)] )[:32] ) offset = 0 lower = _checkpoint_sequence(request.get("changes_after"), caller=caller) or 0 cursor = (request.get("pagination") or {}).get("cursor") if cursor: decoded = decode_cursor(cursor) _require_owned_cursor(decoded, caller=caller) if decoded.get("snapshot") != snapshot_id: raise LedgerConflictError( "CURSOR_SNAPSHOT_MISMATCH", "this cursor belongs to a different snapshot, caller or filter set;" " restart the walk", ) cursor_lower = decoded.get("lower") if not isinstance(cursor_lower, int): raise LedgerConflictError( "CURSOR_SNAPSHOT_MISMATCH", "this cursor does not retain its incremental lower bound; restart the walk", ) if "changes_after" in request and lower != cursor_lower: raise LedgerConflictError( "CURSOR_SNAPSHOT_MISMATCH", "this cursor belongs to a different incremental lower bound; restart the walk", ) lower = cursor_lower offset = int(decoded.get("offset", 0)) if "as_of" in decoded: snapshot = replace(snapshot, as_of=datetime.fromisoformat(decoded["as_of"])) value = StatusProjectionInput( snapshot, scope, self._escalation, tuple(filters["delivery_config_ids"] or ()), tuple(filters["media_buy_ids"] or ()), tuple(filters["feed_purposes"] or ()), _parse(filters["period_start"]), _parse(filters["period_end"]), reconciliation=reconciliation, consumer_status_enabled=self._consumer_status_enabled, ) result = project_status_scope(value) if result.intents: raise LedgerConflictError( "STATUS_PROJECTION_UNAVAILABLE", "status lifecycle is pending" ) common: dict[str, Any] = { "status": "completed", "view": view, "ledger_snapshot_id": snapshot_id, "ledger_as_of": _iso(snapshot.as_of), "account_id": caller.account_id, } if view == "revision": revision_id = request.get("reporting_revision_id") if not revision_id: raise LedgerConflictError( "MISSING_REVISION_ID", "a revision view requires reporting_revision_id" ) revision = next( (r for r in snapshot.revisions if r.reporting_revision_id == revision_id), None ) if revision is None: raise LedgerConflictError( "LOOKUP_UNAVAILABLE", "no such revision is available to this caller" ) owner = next( ( o for o in snapshot.obligations if o.reporting_obligation_id == revision.reporting_obligation_id ), None, ) response = { **common, "revision": _revision_to_wire(revision, owner), "adjustments": [ _adjustment_to_wire(a) for a in snapshot.adjustments if a.adjusts_reporting_revision_id == revision_id ], "materializations": [], "receipts": [], "pagination": {"total_count": 1, "has_more": False}, } if reconciliation is not None: from adcp.reporting.projection.wire import exact_revision_evidence response.update( exact_revision_evidence(snapshot, revision, owner, reconciliation, caller) ) if revision_ownership: from adcp.reporting.ownership import with_revision_ownership response = with_revision_ownership( response, {revision.reporting_revision_id: revision.reporting_obligation_id} ) return response common["scope"] = _scope_to_wire( result.configurations, ledger_as_of=snapshot.as_of, request=request, obligations=tuple(p.obligation for p in result.obligations), ) obligations = tuple(p.obligation for p in result.obligations) if view == "summary": states: list[str] = [p.projection.health for p in result.obligations] counts = { "total": len(obligations), **{ state: states.count(state) for state in ("waiting", "healthy", "delayed", "action_required", "complete") }, } if self._consumer_status_enabled: counts["consumer_status_pending"] = result.pending_count watermark = _scope_data_through(p.projection for p in result.obligations) configurations = tuple( c for c in result.configurations if (not request.get("finality") or c.required_finality in request["finality"]) and ( not request.get("health") or project_status_scope( replace( value, scope=ReportingStatusScope( c.account_id, c.generation_key, consumer_id=scope.consumer_id ), ) ).health in request["health"] ) ) selected_obligations = tuple( p.obligation for p in result.obligations if ( not request.get("finality") or p.obligation.required_finality in request["finality"] ) and (not request.get("health") or p.projection.health in request["health"]) ) next_expected = ( next_reporting_period_start(configurations, as_of=snapshot.as_of) if complete_start_forecast and result.health == "complete" else next_reporting_expectation( configurations, selected_obligations, as_of=snapshot.as_of, period_start=value.period_start, period_end=value.period_end, ) ) return { **common, "health": result.health, "coverage": _coverage_roll_up(obligations, as_of=snapshot.as_of), "data_through": _iso(watermark) if watermark else None, **({"next_expected_at": _iso(next_expected)} if next_expected is not None else {}), "obligation_counts": counts, "issues": [i.to_wire() for i in result.issues], } owners = {o.reporting_obligation_id: o for o in obligations} records: dict[tuple[str, str], Any] = { ("obligation", o.reporting_obligation_id): o for o in obligations } for r in snapshot.revisions: if r.reporting_obligation_id in owners: records[("revision", r.reporting_revision_id)] = r for a in snapshot.adjustments: if ("revision", a.adjusts_reporting_revision_id) in records: records[("adjustment", a.reporting_adjustment_id)] = a generation_keys = {c.generation_key for c in result.configurations} if self._consumer_status_enabled: for s in snapshot.statuses: if s.consumer_id == caller.consumer_id and s.generation_key in generation_keys: if (value.period_start is not None and s.period_end <= value.period_start) or ( value.period_end is not None and s.period_start >= value.period_end ): continue records[("consumer_status", s.reporting_status_id)] = s selected = [ (kind, records[(kind, record_id)]) for seq, kind, record_id, _ in sorted(snapshot.changes) if seq > lower and (kind, record_id) in records ] window = selected[offset : offset + self._page_size] has_more = offset + self._page_size < len(selected) periods = [] for kind, record in window: if kind != "obligation": continue scoped = project_status_scope( replace(value, scope=ReportingStatusScope.for_obligation(record, scope.consumer_id)) ) item = scoped.obligations[0] periods.append( _obligation_to_wire( record, revisions=item.revisions, health=scoped.health, production_status=item.projection.production_status, issues=scoped.issues, statuses=item.statuses, ) ) payload = { **common, "health": result.health, "issues": [issue.to_wire() for issue in result.issues], "changes_checkpoint": _encode_checkpoint(snapshot.max_sequence, caller=caller), "periods": periods, "revisions": [ _revision_to_wire(r, owners.get(r.reporting_obligation_id)) for kind, r in window if kind == "revision" ], "adjustments": [_adjustment_to_wire(a) for kind, a in window if kind == "adjustment"], "materializations": [], "receipts": [], "pagination": { "total_count": len(selected), "has_more": has_more, **( { "cursor": encode_cursor( { "ownership": 2, "account": caller.account_id, "consumer": caller.consumer_id, "snapshot": snapshot_id, "offset": offset + self._page_size, "lower": lower, "as_of": snapshot.as_of.isoformat(), } ) } if has_more else {} ), }, } if self._consumer_status_enabled: payload["consumer_statuses"] = [ _consumer_status_to_wire(s) for kind, s in window if kind == "consumer_status" ] return payloadRender one captured status snapshot using the same pure projection as push.
Methods
async def handle(self,
request: dict[str, Any],
*,
caller: ReportingStatusCaller) ‑> dict[str, typing.Any]-
Expand source code
async def handle( self, request: dict[str, Any], *, caller: ReportingStatusCaller ) -> dict[str, Any]: snapshot = await self._load(caller) return self.render_snapshot(request, caller=caller, snapshot=snapshot) def render_snapshot(self,
request: dict[str, Any],
*,
caller: ReportingStatusCaller,
snapshot: ReportingStatusSnapshot,
reconciliation: tuple[ReportingDeliveryRecord, ...] | None = None,
revision_ownership: bool = False) ‑> dict[str, typing.Any]-
Expand source code
def render_snapshot( self, request: dict[str, Any], *, caller: ReportingStatusCaller, snapshot: ReportingStatusSnapshot, reconciliation: tuple[ReportingDeliveryRecord, ...] | None = None, revision_ownership: bool = False, ) -> dict[str, Any]: """Render captured database evidence without any store calls or clock reads.""" view = request.get("view", "summary") if view not in {"summary", "periods", "revision"}: raise LedgerConflictError("INVALID_VIEW", "unsupported reporting status view") if snapshot.account_id != caller.account_id: raise LedgerConflictError("LOOKUP_UNAVAILABLE", "status is unavailable to this caller") filters = _filters(request) version = normalize_to_release_precision( request.get("adcp_version") or resolve_adcp_version(None) ) complete_start_forecast = is_adcp_version_at_least(version, "3.2-rc.6") if complete_start_forecast: # Bind the new wire contract without relabelling explicit rc.3 snapshots. filters["adcp_version"] = version scope = ReportingStatusScope( caller.account_id, consumer_id=caller.consumer_id, ) from adcp.reporting.ledger.delivery_models import ReportingDeliveryPrincipal from adcp.reporting.materializer.capture import private_snapshot snapshot = private_snapshot( snapshot, ReportingDeliveryPrincipal(caller.account_id, caller.consumer_id) ) snapshot_id = ( "rpls_" + _fingerprint( [ caller.account_id, caller.consumer_id, filters, snapshot.max_sequence, ] )[:32] ) if reconciliation is not None: from adcp.reporting.ledger._delivery_state import fingerprint, principal reconciliation = tuple( r for r in reconciliation if principal(r).account_id == caller.account_id and principal(r).consumer_id == caller.consumer_id ) snapshot_id = ( "rpls_" + _fingerprint( [snapshot_id, [fingerprint(r) for r in reconciliation], _iso(snapshot.as_of)] )[:32] ) offset = 0 lower = _checkpoint_sequence(request.get("changes_after"), caller=caller) or 0 cursor = (request.get("pagination") or {}).get("cursor") if cursor: decoded = decode_cursor(cursor) _require_owned_cursor(decoded, caller=caller) if decoded.get("snapshot") != snapshot_id: raise LedgerConflictError( "CURSOR_SNAPSHOT_MISMATCH", "this cursor belongs to a different snapshot, caller or filter set;" " restart the walk", ) cursor_lower = decoded.get("lower") if not isinstance(cursor_lower, int): raise LedgerConflictError( "CURSOR_SNAPSHOT_MISMATCH", "this cursor does not retain its incremental lower bound; restart the walk", ) if "changes_after" in request and lower != cursor_lower: raise LedgerConflictError( "CURSOR_SNAPSHOT_MISMATCH", "this cursor belongs to a different incremental lower bound; restart the walk", ) lower = cursor_lower offset = int(decoded.get("offset", 0)) if "as_of" in decoded: snapshot = replace(snapshot, as_of=datetime.fromisoformat(decoded["as_of"])) value = StatusProjectionInput( snapshot, scope, self._escalation, tuple(filters["delivery_config_ids"] or ()), tuple(filters["media_buy_ids"] or ()), tuple(filters["feed_purposes"] or ()), _parse(filters["period_start"]), _parse(filters["period_end"]), reconciliation=reconciliation, consumer_status_enabled=self._consumer_status_enabled, ) result = project_status_scope(value) if result.intents: raise LedgerConflictError( "STATUS_PROJECTION_UNAVAILABLE", "status lifecycle is pending" ) common: dict[str, Any] = { "status": "completed", "view": view, "ledger_snapshot_id": snapshot_id, "ledger_as_of": _iso(snapshot.as_of), "account_id": caller.account_id, } if view == "revision": revision_id = request.get("reporting_revision_id") if not revision_id: raise LedgerConflictError( "MISSING_REVISION_ID", "a revision view requires reporting_revision_id" ) revision = next( (r for r in snapshot.revisions if r.reporting_revision_id == revision_id), None ) if revision is None: raise LedgerConflictError( "LOOKUP_UNAVAILABLE", "no such revision is available to this caller" ) owner = next( ( o for o in snapshot.obligations if o.reporting_obligation_id == revision.reporting_obligation_id ), None, ) response = { **common, "revision": _revision_to_wire(revision, owner), "adjustments": [ _adjustment_to_wire(a) for a in snapshot.adjustments if a.adjusts_reporting_revision_id == revision_id ], "materializations": [], "receipts": [], "pagination": {"total_count": 1, "has_more": False}, } if reconciliation is not None: from adcp.reporting.projection.wire import exact_revision_evidence response.update( exact_revision_evidence(snapshot, revision, owner, reconciliation, caller) ) if revision_ownership: from adcp.reporting.ownership import with_revision_ownership response = with_revision_ownership( response, {revision.reporting_revision_id: revision.reporting_obligation_id} ) return response common["scope"] = _scope_to_wire( result.configurations, ledger_as_of=snapshot.as_of, request=request, obligations=tuple(p.obligation for p in result.obligations), ) obligations = tuple(p.obligation for p in result.obligations) if view == "summary": states: list[str] = [p.projection.health for p in result.obligations] counts = { "total": len(obligations), **{ state: states.count(state) for state in ("waiting", "healthy", "delayed", "action_required", "complete") }, } if self._consumer_status_enabled: counts["consumer_status_pending"] = result.pending_count watermark = _scope_data_through(p.projection for p in result.obligations) configurations = tuple( c for c in result.configurations if (not request.get("finality") or c.required_finality in request["finality"]) and ( not request.get("health") or project_status_scope( replace( value, scope=ReportingStatusScope( c.account_id, c.generation_key, consumer_id=scope.consumer_id ), ) ).health in request["health"] ) ) selected_obligations = tuple( p.obligation for p in result.obligations if ( not request.get("finality") or p.obligation.required_finality in request["finality"] ) and (not request.get("health") or p.projection.health in request["health"]) ) next_expected = ( next_reporting_period_start(configurations, as_of=snapshot.as_of) if complete_start_forecast and result.health == "complete" else next_reporting_expectation( configurations, selected_obligations, as_of=snapshot.as_of, period_start=value.period_start, period_end=value.period_end, ) ) return { **common, "health": result.health, "coverage": _coverage_roll_up(obligations, as_of=snapshot.as_of), "data_through": _iso(watermark) if watermark else None, **({"next_expected_at": _iso(next_expected)} if next_expected is not None else {}), "obligation_counts": counts, "issues": [i.to_wire() for i in result.issues], } owners = {o.reporting_obligation_id: o for o in obligations} records: dict[tuple[str, str], Any] = { ("obligation", o.reporting_obligation_id): o for o in obligations } for r in snapshot.revisions: if r.reporting_obligation_id in owners: records[("revision", r.reporting_revision_id)] = r for a in snapshot.adjustments: if ("revision", a.adjusts_reporting_revision_id) in records: records[("adjustment", a.reporting_adjustment_id)] = a generation_keys = {c.generation_key for c in result.configurations} if self._consumer_status_enabled: for s in snapshot.statuses: if s.consumer_id == caller.consumer_id and s.generation_key in generation_keys: if (value.period_start is not None and s.period_end <= value.period_start) or ( value.period_end is not None and s.period_start >= value.period_end ): continue records[("consumer_status", s.reporting_status_id)] = s selected = [ (kind, records[(kind, record_id)]) for seq, kind, record_id, _ in sorted(snapshot.changes) if seq > lower and (kind, record_id) in records ] window = selected[offset : offset + self._page_size] has_more = offset + self._page_size < len(selected) periods = [] for kind, record in window: if kind != "obligation": continue scoped = project_status_scope( replace(value, scope=ReportingStatusScope.for_obligation(record, scope.consumer_id)) ) item = scoped.obligations[0] periods.append( _obligation_to_wire( record, revisions=item.revisions, health=scoped.health, production_status=item.projection.production_status, issues=scoped.issues, statuses=item.statuses, ) ) payload = { **common, "health": result.health, "issues": [issue.to_wire() for issue in result.issues], "changes_checkpoint": _encode_checkpoint(snapshot.max_sequence, caller=caller), "periods": periods, "revisions": [ _revision_to_wire(r, owners.get(r.reporting_obligation_id)) for kind, r in window if kind == "revision" ], "adjustments": [_adjustment_to_wire(a) for kind, a in window if kind == "adjustment"], "materializations": [], "receipts": [], "pagination": { "total_count": len(selected), "has_more": has_more, **( { "cursor": encode_cursor( { "ownership": 2, "account": caller.account_id, "consumer": caller.consumer_id, "snapshot": snapshot_id, "offset": offset + self._page_size, "lower": lower, "as_of": snapshot.as_of.isoformat(), } ) } if has_more else {} ), }, } if self._consumer_status_enabled: payload["consumer_statuses"] = [ _consumer_status_to_wire(s) for kind, s in window if kind == "consumer_status" ] return payloadRender captured database evidence without any store calls or clock reads.