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_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.

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 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.

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_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.

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 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".

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 | None

One obligation's derived status at a snapshot boundary.

Instance variables

var current_revision : ReportingRevisionRecord | None
var health : Literal['healthy', 'waiting', 'delayed', 'action_required', 'complete']
var issues : tuple[ReportingIssue, ...]
var production_status : Literal['not_due', 'pending', 'published', 'failed']
var satisfied : bool