Module adcp.reporting.ledger.health
Deriving reporting health from immutable evidence and the clock.
Health is never stored. It is a pure function of the obligations, the revisions associated with them, and the snapshot's clock boundary – which is what makes it impossible for a stored health column to go stale, and what makes two readers of one snapshot agree.
The five states answer one question each:
waiting
Nothing is due yet.
ledger_as_of has not reached expected_at.
healthy
Everything due has been produced, and the scope is still open.
delayed
Something due is late and automated recovery is still running.
Bounded by
automated_recovery_deadline_at – a dead feed must never be parked here.
action_required
A human on the named responsible_party must act.
complete
The queried scope is closed and every final obligation is satisfied.
A committed zero-row revision satisfies an obligation exactly like any other. Absence is represented only by an empty association set, never by a zero.
Functions
def aggregate_reporting_health(obligation_health: Iterable[ReportingHealth],
*,
scope_closed: bool,
coverage_complete: bool) ‑> Literal['healthy', 'waiting', 'delayed', 'action_required', 'complete']-
Expand source code
def aggregate_reporting_health( obligation_health: Iterable[ReportingHealth], *, scope_closed: bool, coverage_complete: bool, ) -> ReportingHealth: """Roll obligation health up to a scope. Worst-case wins, with one addition: incomplete coverage is itself ``action_required`` regardless of the obligations underneath it. A scope whose denominator is not fully known cannot honestly be called healthy -- the obligations that *are* present may simply be the ones that happened to resolve. """ values = list(obligation_health) if not coverage_complete: return "action_required" if "action_required" in values: return "action_required" if "delayed" in values: return "delayed" if scope_closed and all(value == "complete" for value in values): return "complete" if any(value in {"healthy", "complete"} for value in values): return "healthy" return "waiting"Roll obligation health up to a scope.
Worst-case wins, with one addition: incomplete coverage is itself
action_requiredregardless of the obligations underneath it. A scope whose denominator is not fully known cannot honestly be called healthy – the obligations that are present may simply be the ones that happened to resolve. def current_required_revision(obligation: ReportingObligationRecord,
revisions: Sequence[ReportingRevisionRecord]) ‑> ReportingRevisionRecord | None-
Expand source code
def current_required_revision( obligation: ReportingObligationRecord, revisions: Sequence[ReportingRevisionRecord], ) -> ReportingRevisionRecord | None: """Source-compatible wrapper. Corrupt and not-ready histories both return None. New callers should consume ``select_reporting_revision``'s discriminated result so corruption can be parked for repair instead of retried as absence. """ result = select_reporting_revision( revisions, account_id=obligation.account_id, reporting_obligation_id=obligation.reporting_obligation_id, required_finality=obligation.required_finality, ) return result.revision if result.kind == "selected" else NoneSource-compatible wrapper. Corrupt and not-ready histories both return None.
New callers should consume
select_reporting_revision's discriminated result so corruption can be parked for repair instead of retried as absence. def issue_id_for(kind: str, *parts: object) ‑> str-
Expand source code
def issue_id_for(kind: str, *parts: object) -> str: """A stable, derived issue identity. Derived rather than stored because these conditions are monotone for one immutable obligation: once a qualifying revision is associated, ``REPORT_OVERDUE`` cannot recur. That makes a stateless identity correct without claiming a general issue-lifecycle ledger -- and it means a consumer polling twice deduplicates on the same id both times. """ digest = hashlib.sha256(canonical_json_utf8_v1([kind, *[str(part) for part in parts]])).digest() return "rpti_" + base64.urlsafe_b64encode(digest).decode("ascii").rstrip("=")[:32]A stable, derived issue identity.
Derived rather than stored because these conditions are monotone for one immutable obligation: once a qualifying revision is associated,
REPORT_OVERDUEcannot recur. That makes a stateless identity correct without claiming a general issue-lifecycle ledger – and it means a consumer polling twice deduplicates on the same id both times. def issue_id_for_occurrence(issue_key: str, generation: int) ‑> str-
Expand source code
def issue_id_for_occurrence(issue_key: str, generation: int) -> str: """The id for one *occurrence* of a stored, non-monotone condition. :func:`issue_id_for` is enough for conditions that cannot recur once satisfied. A consumer mismatch can: the buyer supersedes, the seller restates, the disagreement clears and comes back. AdCP 3.2.0-rc.3 requires a recurrence after retirement to get a *new* ``issue_id``, so identity has to include which occurrence this is -- the generation counter held by :class:`~adcp.reporting.ledger.models.ReportingIssueLifecycle`. Folding the generation into the digest rather than appending it keeps the id opaque, so nothing downstream can parse it back into "how many times has this buyer complained". """ digest = hashlib.sha256( canonical_json_utf8_v1(["core-issue-occurrence-v1", issue_key, generation]) ).digest() return "rpti_" + base64.urlsafe_b64encode(digest).decode("ascii").rstrip("=")[:32]The id for one occurrence of a stored, non-monotone condition.
:func:
issue_id_for()is enough for conditions that cannot recur once satisfied. A consumer mismatch can: the buyer supersedes, the seller restates, the disagreement clears and comes back. AdCP 3.2.0-rc.3 requires a recurrence after retirement to get a newissue_id, so identity has to include which occurrence this is – the generation counter held by :class:~adcp.reporting.ledger.models.ReportingIssueLifecycle.Folding the generation into the digest rather than appending it keeps the id opaque, so nothing downstream can parse it back into "how many times has this buyer complained".
def project_obligation_health(obligation: ReportingObligationRecord,
revisions: Sequence[ReportingRevisionRecord],
*,
ledger_as_of: datetime,
scope_closed: bool) ‑> ObligationProjection-
Expand source code
def project_obligation_health( obligation: ReportingObligationRecord, revisions: Sequence[ReportingRevisionRecord], *, ledger_as_of: datetime, scope_closed: bool, ) -> ObligationProjection: """Classify one obligation's immutable evidence at the snapshot's clock.""" boundary = _utc(ledger_as_of) selection = select_reporting_revision( revisions, account_id=obligation.account_id, reporting_obligation_id=obligation.reporting_obligation_id, required_finality=obligation.required_finality, ) current = selection.revision if selection.kind == "selected" else None if selection.kind == "corrupt": return ObligationProjection( health="action_required", production_status="published" if revisions else "pending", issues=( ReportingIssue( issue_id=issue_id_for( "core-revision-history-corrupt-v2", obligation.account_id, obligation.reporting_obligation_id, ), code="HISTORY_UNAVAILABLE", severity="action_required", responsible_party="seller", recommended_action="contact_seller", reporting_obligation_id=obligation.reporting_obligation_id, delivery_config_id=obligation.delivery_config_id, delivery_config_version=obligation.delivery_config_version, feed_purpose=obligation.feed_purpose, media_buy_ids=obligation.media_buy_ids, period_start=obligation.period.start, period_end=obligation.period.end, message=( "The retained revision history is inconsistent; seller repair is required." ), ), ), satisfied=False, current_revision=None, ) if obligation.currency is None: return ObligationProjection( health="action_required", production_status="published" if revisions else "pending", issues=( ReportingIssue( issue_id=issue_id_for( "core-currency-unresolved-v1", obligation.reporting_obligation_id ), code="HISTORY_UNAVAILABLE", severity="action_required", responsible_party="seller", recommended_action="contact_seller", reporting_obligation_id=obligation.reporting_obligation_id, delivery_config_id=obligation.delivery_config_id, delivery_config_version=obligation.delivery_config_version, feed_purpose=obligation.feed_purpose, media_buy_ids=obligation.media_buy_ids, period_start=obligation.period.start, period_end=obligation.period.end, message=( "Currency was not retained for this legacy obligation; " "verified historical evidence is required." ), ), ), satisfied=False, current_revision=current, ) if current is not None and current.readable: return ObligationProjection( health="complete" if scope_closed else "healthy", production_status="published", issues=(), satisfied=True, current_revision=current, ) if current is not None: # A qualifying revision exists but nothing is readable. Core's promise # is retained readability, so this is a seller repair, not a wait. return ObligationProjection( health="action_required", production_status="published", issues=(_unreadable_issue(obligation, current.reporting_revision_id),), satisfied=False, current_revision=current, ) if revisions: # Revisions exist but none meet the required finality: the seller has # published something, so production is not pending -- it is simply not # yet official. production_status: ReportingProductionStatus = "published" elif boundary < _utc(obligation.period.expected_at): production_status = "not_due" else: production_status = "pending" if boundary < _utc(obligation.period.expected_at): return ObligationProjection( health="waiting", production_status=production_status, issues=(), satisfied=False, current_revision=None, ) health: ReportingHealth = ( "delayed" if boundary < _utc(obligation.automated_recovery_deadline_at) else "action_required" ) return ObligationProjection( health=health, production_status=production_status, issues=(_overdue_issue(obligation, health),), satisfied=False, current_revision=None, )Classify one obligation's immutable evidence at the snapshot's clock.
Classes
class ObligationProjection (health: ReportingHealth,
production_status: ReportingProductionStatus,
issues: tuple[ReportingIssue, ...],
satisfied: bool,
current_revision: ReportingRevisionRecord | None)-
Expand source code
@dataclass(frozen=True) class ObligationProjection: """One obligation's derived status at a snapshot boundary.""" health: ReportingHealth production_status: ReportingProductionStatus issues: tuple[ReportingIssue, ...] satisfied: bool current_revision: ReportingRevisionRecord | NoneOne obligation's derived status at a snapshot boundary.
Instance variables
var current_revision : ReportingRevisionRecord | Nonevar health : Literal['healthy', 'waiting', 'delayed', 'action_required', 'complete']var issues : tuple[ReportingIssue, ...]var production_status : Literal['not_due', 'pending', 'published', 'failed']var satisfied : bool