Module adcp.reporting.ledger.consumer_status

sync_reporting_status ingest and the consumer-mismatch projection.

This closes the operational loop between what a seller says it published and what an authenticated buyer could actually consume. Without it, a seller's ledger is a monologue: every obligation can look complete while the buyer has never successfully read a single revision.

Opt-in, and off by default. sync_reporting_status is an additive extension in AdCP 3.2.0-rc.2; it becomes required Core only in the next eligible minor after 2026-10-24. Turn the ingest on with::

import adcp.reporting.ledger.consumer_status as consumer_status
consumer_status.CONSUMER_STATUS_ENABLED = True

or per-instance via ConsumerStatusIngest(..., enabled=True). A seller that turns it on must advertise consumer_status_task; a seller that does not should leave it off, because a half-implemented status loop is worse than none – a buyer that can file statements nobody reads believes it has told you.

Wire types come from :mod:adcp.types; the conditional rules that datamodel-code-generator flattens away (received requires a revision id and digest, obligation_missing forbids them, and so on) are enforced directly against the bundled core/reporting-consumer-status.json through :func:get_named_validator(), so they stay in the schema rather than being restated in Python and drifting.

What a statement is, and is not

It is operational status: could the required reporting for this expected period actually be consumed? It is not measurement data, delivery totals, a billing receipt, or authority over the seller's ledger. received proves the buyer consumed the exact Core revision binding – nothing about materialization evidence or billing control totals.

Four statuses, four different problems:

============================ ==================================================== received The exact revision content was consumed. obligation_missing The period is required by the accepted configuration generation but the seller ledger omitted it. revision_missing The obligation exists but no required revision was available after expected_at. unreadable A revision was advertised but its content could not be consumed; failure_code says why. ============================ ====================================================

obligation_missing is deliberately keyed without a seller-issued obligation id. Requiring one would make the first missing report invisible again, which is the exact failure this loop exists to surface.

Attribution

Consumer status is separately attributed and never satisfies seller production health. A conflicting current statement makes only the submitting caller's view action_required with a stable CONSUMER_STATUS_MISMATCH issue. One buyer's assertion never changes another caller's view, and never touches seller-advertised reliability statistics without corroboration.

Global variables

var CONSUMER_STATUS_ENABLED

Module-level opt-in for the whole consumer-status surface. Off by default: the task is an additive extension a seller must deliberately advertise.

Functions

def consumer_mismatch_issue_key(*,
account_id: str,
consumer_id: str,
delivery_config_id: str,
delivery_config_version: int,
report_definition_id: str,
period_start: datetime,
period_end: datetime) ‑> str
Expand source code
def consumer_mismatch_issue_key(
    *,
    account_id: str,
    consumer_id: str,
    delivery_config_id: str,
    delivery_config_version: int,
    report_definition_id: str,
    period_start: datetime,
    period_end: datetime,
) -> str:
    """Identity of the *condition*, not of the statement that evidences it.

    Keyed on the logical status chain, deliberately not on
    ``reporting_status_id``. A buyer that supersedes one conflicting statement
    with another conflicting statement has not resolved anything, and rc.3
    anchors the escalation clock to the issue's ``opened_at``. Keying on the
    statement id would mint a fresh issue -- and a fresh clock -- on every
    re-file, letting an unresolved disagreement dodge escalation forever.

    Also deliberately not keyed on the obligation id: a chain that began as
    ``obligation_missing`` attaches to the repaired obligation later, and the
    spec requires that chain never be lost, forked, or reset.
    """
    generation = ReportingConfigurationGenerationKey(
        account_id=account_id,
        consumer_id=consumer_id,
        delivery_config_id=delivery_config_id,
        delivery_config_version=delivery_config_version,
    )
    payload = canonical_json_utf8_v1(
        [
            "core-consumer-status-mismatch-v1",
            generation.account_id,
            consumer_id,
            generation.delivery_config_id,
            generation.delivery_config_version,
            report_definition_id,
            _utc(period_start).isoformat(),
            _utc(period_end).isoformat(),
        ]
    )
    return "rpik_" + hashlib.sha256(payload).hexdigest()[:40]

Identity of the condition, not of the statement that evidences it.

Keyed on the logical status chain, deliberately not on reporting_status_id. A buyer that supersedes one conflicting statement with another conflicting statement has not resolved anything, and rc.3 anchors the escalation clock to the issue's opened_at. Keying on the statement id would mint a fresh issue – and a fresh clock – on every re-file, letting an unresolved disagreement dodge escalation forever.

Also deliberately not keyed on the obligation id: a chain that began as obligation_missing attaches to the repaired obligation later, and the spec requires that chain never be lost, forked, or reset.

def consumer_statement_conflicts(*,
current: ConsumerStatusRecord,
current_revision: ReportingRevisionRecord | None,
revisions: Sequence[ReportingRevisionRecord] = ()) ‑> bool
Expand source code
def consumer_statement_conflicts(
    *,
    current: ConsumerStatusRecord,
    current_revision: ReportingRevisionRecord | None,
    revisions: Sequence[ReportingRevisionRecord] = (),
) -> bool:
    """Whether this statement contradicts the seller's *record* of the period.

    Deliberately independent of the seller's health *label*. The two questions
    are different and conflating them is how an escalation clock gets reset:
    health answers "should this caller's view be degraded right now", while
    this answers "does the disagreement still stand" -- which is what decides
    whether the issue may be retired.

    The spec allows resolution two ways: the consumer supersedes with a
    statement that agrees, or the seller's own projection changes so the two no
    longer conflict. Both are record comparisons:

    * ``obligation_missing`` -- the obligation is in front of us, so the buyer's
      claim that the seller omitted the period is contradicted. Always conflicts.
    * ``revision_missing`` -- conflicts only while the seller actually has a
      current required revision. If the seller has none either, they agree.
    * ``unreadable`` -- conflicts while the named revision is still readable in
      the seller's record. A seller that has since marked it unreadable agrees.
    * ``content_mismatch`` -- conflicts while the named revision is still the
      one the seller requires, because that is the seller standing behind the
      content the buyer disputes.
    * ``received`` -- conflicts only once something superseded the revision the
      buyer named.
    """
    status = current.consumer_status
    if status == "obligation_missing":
        return True
    if status == "revision_missing":
        return current_revision is not None
    if status == "unreadable":
        named = _named_revision(current, revisions)
        return named is None or named.readable
    if status == "content_mismatch":
        return (
            current_revision is not None
            and current.reporting_revision_id == current_revision.reporting_revision_id
        )
    if status == "received":
        if current.reporting_revision_id is None:
            return False
        if (
            current_revision is not None
            and current.reporting_revision_id == current_revision.reporting_revision_id
        ):
            return False
        # Stale only if a restatement actually superseded it. With no
        # supersession and no current required revision there is nothing to
        # re-read, so nothing to forgive and nothing to page about.
        return _superseder(current.reporting_revision_id, revisions) is not None or (
            current_revision is not None
        )
    # No fallback return: every ``ConsumerStatusValue`` is handled above, so a
    # sixth status added to the schema fails mypy with "missing return
    # statement" rather than silently classifying itself as "no conflict".

Whether this statement contradicts the seller's record of the period.

Deliberately independent of the seller's health label. The two questions are different and conflating them is how an escalation clock gets reset: health answers "should this caller's view be degraded right now", while this answers "does the disagreement still stand" – which is what decides whether the issue may be retired.

The spec allows resolution two ways: the consumer supersedes with a statement that agrees, or the seller's own projection changes so the two no longer conflict. Both are record comparisons:

  • obligation_missing – the obligation is in front of us, so the buyer's claim that the seller omitted the period is contradicted. Always conflicts.
  • revision_missing – conflicts only while the seller actually has a current required revision. If the seller has none either, they agree.
  • unreadable – conflicts while the named revision is still readable in the seller's record. A seller that has since marked it unreadable agrees.
  • content_mismatch – conflicts while the named revision is still the one the seller requires, because that is the seller standing behind the content the buyer disputes.
  • received – conflicts only once something superseded the revision the buyer named.
def current_consumer_statement(statuses: Sequence[ConsumerStatusRecord]) ‑> ConsumerStatusRecord | None
Expand source code
def current_consumer_statement(
    statuses: Sequence[ConsumerStatusRecord],
) -> ConsumerStatusRecord | None:
    """This caller's one unsuperseded leaf, or ``None`` for an empty chain."""
    return next((item for item in statuses if not item.superseded), None)

This caller's one unsuperseded leaf, or None for an empty chain.

def project_consumer_mismatch(*,
obligation: ReportingObligationRecord,
current_revision: ReportingRevisionRecord | None,
statuses: Sequence[ConsumerStatusRecord],
seller_health: ReportingHealth,
revisions: Sequence[ReportingRevisionRecord] = (),
ledger_as_of: datetime | None = None,
delivery_sla: timedelta = datetime.timedelta(0),
automated_recovery_window: timedelta = datetime.timedelta(seconds=21600),
escalation: ReportingDeliveryEscalation | None = None,
lifecycle: ReportingIssueLifecycle | None = None) ‑> ConsumerMismatch | None
Expand source code
def project_consumer_mismatch(
    *,
    obligation: ReportingObligationRecord,
    current_revision: ReportingRevisionRecord | None,
    statuses: Sequence[ConsumerStatusRecord],
    seller_health: ReportingHealth,
    revisions: Sequence[ReportingRevisionRecord] = (),
    ledger_as_of: datetime | None = None,
    delivery_sla: timedelta = timedelta(0),
    automated_recovery_window: timedelta = timedelta(hours=6),
    escalation: ReportingDeliveryEscalation | None = None,
    lifecycle: ReportingIssueLifecycle | None = None,
) -> ConsumerMismatch | None:
    """Compare the seller's projection with this caller's current statement.

    Returns a mismatch only when they genuinely conflict. Five conflict kinds,
    four of them immediate:

    * ``obligation_missing``, ``revision_missing``, ``unreadable``, and
      ``content_mismatch`` against an otherwise ``healthy``/``complete`` seller
      projection are ``action_required`` at once; and
    * ``received`` naming a revision the seller has since superseded is
      ``delayed`` until its grace deadline, then ``action_required``.

    The carve-out exists because that buyer consumed exactly what the seller
    required at the time and has not yet had a bounded chance to re-read.
    Paging a human the instant a seller restates a provisional revision would
    train buyers to ignore the signal.

    An advertised ``consumer_mismatch_escalation`` boundary takes precedence
    when the two windows overlap: past ``opened_at`` plus that window the issue
    is ``action_required`` with a ``contact_*`` action regardless of the grace
    deadline. ``wait_for_retry`` must not survive that boundary -- an
    unattended mismatch is an escalation, not a retry.

    Missing consumer status stays *unknown*. It never excuses seller reporting
    and never by itself degrades seller health; silence is surfaced only
    through ``obligation_counts.consumer_status_pending``.
    """
    current = current_consumer_statement(statuses)
    if current is None:
        return None
    if lifecycle is not None and lifecycle.issue_state == "waived":
        if waiver_covers_mismatch(lifecycle, current, obligation, current_revision):
            return None
        # A caller must settle the new occurrence before publishing. Never
        # treat a legacy/unrelated waiver as permission to clear this health.
        lifecycle = None
    if not consumer_statement_conflicts(
        current=current, current_revision=current_revision, revisions=revisions
    ):
        return None
    if seller_health not in {"healthy", "complete"}:
        # The seller already knows something is wrong and has said so. Adding a
        # second issue for the same condition would double-count it.
        #
        # Note this is *only* a suppression of the emission. The condition is
        # still open, so a caller must not read ``None`` here as "resolved" --
        # see :func:`consumer_statement_conflicts`.
        return None

    boundary = _utc(ledger_as_of) if ledger_as_of is not None else None
    severity: Literal["delayed", "action_required"] = "action_required"
    grace_deadline: datetime | None = None

    if current.consumer_status == "received":
        assert current.reporting_revision_id is not None  # the conflict test proved it
        grace_deadline = stale_received_grace_deadline(
            received_revision_id=current.reporting_revision_id,
            revisions=revisions,
            delivery_sla=delivery_sla,
            automated_recovery_window=automated_recovery_window,
        )
        if grace_deadline is not None and boundary is not None and boundary < grace_deadline:
            severity = "delayed"
        responsible: ResponsibleParty = "buyer"
        message = (
            "the authenticated consumer received revision "
            f"{current.reporting_revision_id} but the current required revision is "
            f"{current_revision.reporting_revision_id if current_revision else 'unknown'}; "
            "the consumer is working from a superseded restatement"
        )
    else:
        responsible = _responsible_for(current)
        message = _negative_message(current)

    opened_at = _utc(lifecycle.opened_at) if lifecycle is not None else boundary
    escalated = False
    window = escalation.consumer_mismatch_escalation if escalation is not None else None
    if window is not None and opened_at is not None and boundary is not None:
        # Precedence: the escalation boundary wins when it overlaps the grace
        # window. Measured from the issue's opened_at, never from this poll, so
        # re-emission cannot reset it.
        if boundary >= opened_at + window:
            severity = "action_required"
            escalated = True

    issue = _mismatch_issue(
        obligation,
        current,
        responsible_party=responsible,
        message=message,
        severity=severity,
        escalated=escalated,
        lifecycle=lifecycle,
        opened_at=opened_at,
    )
    published = lifecycle is None or lifecycle.published
    return ConsumerMismatch(issue=issue, severity=severity, published=published)

Compare the seller's projection with this caller's current statement.

Returns a mismatch only when they genuinely conflict. Five conflict kinds, four of them immediate:

  • obligation_missing, revision_missing, unreadable, and content_mismatch against an otherwise healthy/complete seller projection are action_required at once; and
  • received naming a revision the seller has since superseded is delayed until its grace deadline, then action_required.

The carve-out exists because that buyer consumed exactly what the seller required at the time and has not yet had a bounded chance to re-read. Paging a human the instant a seller restates a provisional revision would train buyers to ignore the signal.

An advertised consumer_mismatch_escalation boundary takes precedence when the two windows overlap: past opened_at plus that window the issue is action_required with a contact_* action regardless of the grace deadline. wait_for_retry must not survive that boundary – an unattended mismatch is an escalation, not a retry.

Missing consumer status stays unknown. It never excuses seller reporting and never by itself degrades seller health; silence is surfaced only through obligation_counts.consumer_status_pending.

def stale_received_grace_deadline(*,
received_revision_id: str,
revisions: Sequence[ReportingRevisionRecord],
delivery_sla: timedelta,
automated_recovery_window: timedelta) ‑> datetime.datetime | None
Expand source code
def stale_received_grace_deadline(
    *,
    received_revision_id: str,
    revisions: Sequence[ReportingRevisionRecord],
    delivery_sla: timedelta,
    automated_recovery_window: timedelta,
) -> datetime | None:
    """When a stale ``received`` stops being excusable.

    The anchor is the ``created_at`` of the **first** revision that superseded
    the one the consumer named, plus the configuration generation's
    ``delivery_sla`` as exact elapsed UTC time. Anchoring to the first
    supersession rather than the newest one is what stops the window from
    becoming a hiding place: otherwise a seller could hold a genuinely
    unresolved mismatch below ``action_required`` by restating on a timer.

    When ``delivery_sla`` resolves to zero in any of its legal spellings, the
    seller uses ``automated_recovery_window`` instead, so a zero-SLA feed still
    yields a bounded re-read window rather than an instant escalation.

    Returns ``None`` when nothing superseded the named revision -- there is no
    staleness to forgive.
    """
    superseder = _superseder(received_revision_id, revisions)
    if superseder is None:
        return None
    window = delivery_sla if delivery_sla > timedelta(0) else automated_recovery_window
    return _utc(superseder.created_at) + window

When a stale received stops being excusable.

The anchor is the created_at of the first revision that superseded the one the consumer named, plus the configuration generation's delivery_sla as exact elapsed UTC time. Anchoring to the first supersession rather than the newest one is what stops the window from becoming a hiding place: otherwise a seller could hold a genuinely unresolved mismatch below action_required by restating on a timer.

When delivery_sla resolves to zero in any of its legal spellings, the seller uses automated_recovery_window instead, so a zero-SLA feed still yields a bounded re-read window rather than an instant escalation.

Returns None when nothing superseded the named revision – there is no staleness to forgive.

Classes

class ConsumerMismatch (issue: ReportingIssue,
severity: "Literal['delayed', 'action_required']",
published: bool = True)
Expand source code
@dataclass(frozen=True)
class ConsumerMismatch:
    """A consumer mismatch, with the severity it degrades the caller's view to.

    Severity is carried alongside the issue rather than read back off it so the
    caller cannot accidentally project ``delayed`` health from an
    ``action_required`` issue, or the reverse.
    """

    issue: ReportingIssue
    severity: Literal["delayed", "action_required"]
    #: Only published occurrences contribute health. A waived occurrence is
    #: omitted only while its exact statement and diagnosed conflict match
    #: the recorded bilateral waiver. Later disagreements are independent.
    published: bool = True

A consumer mismatch, with the severity it degrades the caller's view to.

Severity is carried alongside the issue rather than read back off it so the caller cannot accidentally project delayed health from an action_required issue, or the reverse.

Instance variables

var issue : ReportingIssue
var published : bool

Only published occurrences contribute health. A waived occurrence is omitted only while its exact statement and diagnosed conflict match the recorded bilateral waiver. Later disagreements are independent.

var severity : Literal['delayed', 'action_required']
class ConsumerStatusDisabledError (*args, **kwargs)
Expand source code
class ConsumerStatusDisabledError(RuntimeError):
    """The consumer-status ingest was called without being enabled."""

The consumer-status ingest was called without being enabled.

Ancestors

  • builtins.RuntimeError
  • builtins.Exception
  • builtins.BaseException
class ConsumerStatusIngest (store: ReportingLedgerStore,
enabled: bool | None = None,
clock: Callable[[], datetime] | None = None)
Expand source code
@dataclass(frozen=True)
class ConsumerStatusIngest:
    """Records authenticated consumer statements into the ledger.

    Framework-agnostic like :class:`~adcp.reporting.ledger.status.ReportingStatusHandler`:
    a request mapping plus an authenticated caller in, a response mapping out.
    """

    store: ReportingLedgerStore
    enabled: bool | None = None
    #: Supplies ``recorded_at`` -- when the seller durably recorded the
    #: statement -- and therefore the issue's ``opened_at``. Defaults to wall
    #: clock. A deployment whose ledger boundary comes from the database should
    #: pass the same source here, or an issue's escalation clock and the
    #: snapshot that reads it will disagree about what time it is.
    clock: Callable[[], datetime] | None = None

    def _require_enabled(self) -> None:
        active = CONSUMER_STATUS_ENABLED if self.enabled is None else self.enabled
        if not active:
            raise ConsumerStatusDisabledError(
                "sync_reporting_status is an additive opt-in extension. Set "
                "adcp.reporting.ledger.consumer_status.CONSUMER_STATUS_ENABLED = True, or "
                "pass enabled=True, and advertise consumer_status_task before accepting "
                "live statements."
            )

    async def handle(
        self, request: dict[str, Any], *, account_id: str, consumer_id: str
    ) -> dict[str, Any]:
        """Serve one ``sync_reporting_status`` batch.

        Returns one ``recorded`` / ``unchanged`` / ``failed`` result per
        submitted statement.  A failure in one statement does not fail the
        batch: a buyer reporting five periods should not lose four good
        statements to one stale supersession pointer.
        """
        self._require_enabled()

        statements = request.get("statuses") or []
        if not statements:
            raise LedgerConflictError(
                "EMPTY_BATCH", "a sync_reporting_status batch must carry at least one statement"
            )
        seen_chains: set[tuple[Any, ...]] = set()
        results: list[dict[str, Any]] = []
        for statement in statements:
            status_id = statement.get("reporting_status_id", "")
            problems = _validate_consumer_status_wire(statement)
            if problems:
                results.append(_failed(status_id, "INVALID_CONSUMER_STATUS", problems[0]))
                continue
            try:
                record = self._to_record(
                    statement,
                    account_id=account_id,
                    consumer_id=consumer_id,
                    recorded_at=self._now(),
                )
            except ValueError as error:
                results.append(_failed(status_id, "INVALID_CONSUMER_STATUS", str(error)))
                continue
            if record.chain_key in seen_chains:
                # "A batch contains at most one update for each logical status
                # chain" -- two updates in one batch have no defined order, so
                # the second would silently win or silently lose.
                results.append(
                    _failed(
                        status_id,
                        "DUPLICATE_CHAIN_IN_BATCH",
                        "a batch may carry at most one statement per logical status chain",
                    )
                )
                continue
            seen_chains.add(record.chain_key)
            try:
                # Answer "exact retry?" before validating anything, because
                # some of that validation is time-dependent. "content_mismatch
                # must name the revision the seller currently requires" has an
                # answer that changes when the seller restates, so a retry of a
                # statement this ledger already accepted would otherwise be
                # rejected for a reason that did not apply when it was first
                # recorded. The spec's result_mapping requires `unchanged`, and
                # a buyer retrying after a transport failure has no way to tell
                # a genuine rejection from that.
                replay = await self.store.resolve_consumer_status_replay(record)
                if replay is not None:
                    results.append({"result": "unchanged", "consumer_status": _to_wire(replay)})
                    continue
                await self._validate_against_configuration(record)
                from adcp.reporting.ledger.status_snapshot import ReportingStatusParticipant

                atomic = isinstance(self.store, ReportingStatusParticipant)
                if atomic and isinstance(self.store, ReportingStatusParticipant):
                    stored, recorded = await self.store.record_consumer_status_with_lifecycle(
                        record
                    )
                else:
                    stored, recorded = await self.store.record_consumer_status(record)
            except LedgerConflictError as error:
                results.append(_failed(status_id, error.code, str(error)))
                continue
            if recorded and not atomic:
                await self._open_issue_on_first_observation(stored)
            results.append(
                {
                    "result": "recorded" if recorded else "unchanged",
                    "consumer_status": _to_wire(stored),
                }
            )
        return {"status": "completed", "results": results}

    async def _open_issue_on_first_observation(self, stored: ConsumerStatusRecord) -> None:
        """Start the escalation clock when the statement arrives, not when polled.

        ``opened_at`` is "when the seller first observed this logical
        condition", and the seller observes it here -- a statement landing on
        this task *is* the observation. Deferring to the first
        ``get_reporting_status`` would tie the escalation clock to whether
        anyone happened to poll: a mismatch filed and never read would sit at
        ``opened_at = None`` indefinitely and could never cross
        ``consumer_mismatch_escalation_seconds``, which is precisely the
        unattended case the boundary exists for.

        Only opens; never retires and never changes an existing occurrence.
        ``ensure_issue_opened`` is idempotent, so a supersession that keeps the
        same condition open keeps the same ``opened_at``.
        """
        obligation: ReportingObligationRecord | None = None
        if stored.reporting_obligation_id is not None:
            obligation = await self.store.get_obligation(
                account_id=stored.account_id,
                reporting_obligation_id=stored.reporting_obligation_id,
            )
        if obligation is None:
            obligation = await self.store.find_obligation(
                account_id=stored.account_id,
                consumer_id=stored.consumer_id,
                delivery_config_id=stored.delivery_config_id,
                delivery_config_version=stored.delivery_config_version,
                period_start=stored.period_start,
                period_end=stored.period_end,
            )
        revisions = (
            await self.store.list_revisions(
                account_id=stored.account_id,
                reporting_obligation_id=obligation.reporting_obligation_id,
            )
            if obligation is not None
            else ()
        )
        current = None
        if obligation is not None:
            selection = select_reporting_revision(
                revisions,
                account_id=obligation.account_id,
                reporting_obligation_id=obligation.reporting_obligation_id,
                required_finality=obligation.required_finality,
            )
            if selection.kind == "corrupt":
                return  # Corruption cannot establish or resolve consumer disagreement.
            if selection.kind == "selected":
                current = selection.revision
        if not consumer_statement_conflicts(
            current=stored, current_revision=current, revisions=revisions
        ):
            return
        key = consumer_mismatch_issue_key(
            account_id=stored.account_id,
            consumer_id=stored.consumer_id,
            delivery_config_id=stored.delivery_config_id,
            delivery_config_version=stored.delivery_config_version,
            report_definition_id=stored.report_definition_id,
            period_start=stored.period_start,
            period_end=stored.period_end,
        )
        seen: set[str] = set()
        while (
            issue := await self.store.get_issue(account_id=stored.account_id, issue_key=key)
        ) is not None and issue.issue_state == "waived":
            if issue.issue_id in seen or issue.issue_key != key:
                raise LedgerConflictError("STATUS_PROJECTION_UNAVAILABLE", "invalid waiver chain")
            seen.add(issue.issue_id)
            if waiver_covers_mismatch(issue, stored, obligation, current):
                return
            key = condition_after_waiver(issue)
        await self.store.ensure_issue_opened(
            issue_key=key,
            account_id=stored.account_id,
            consumer_id=stored.consumer_id,
            observed_at=stored.recorded_at,
        )

    async def _validate_against_configuration(self, record: ConsumerStatusRecord) -> None:
        """Resolve every identifier the statement names, within caller + account.

        Deliberately does *not* require an obligation to exist *when none is
        named*: the whole point of ``obligation_missing`` is that it is filed
        when the seller's ledger omitted the period, and the configuration
        generation is what makes that period legitimate.

        But an obligation or revision id the caller *does* supply must resolve.
        The spec's ``batch_identity`` rule is explicit that every optional
        obligation, revision, and superseded status MUST resolve within the
        authenticated caller and account or fail with the same unavailable
        result. Skipping that lets a statement name a foreign, nonexistent, or
        long-superseded revision and still be recorded -- and then projected as
        ``action_required`` against a seller who never published the thing
        being disputed.
        """
        configurations = await self.store.list_configurations(
            caller=ReportingDeliveryPrincipal(record.account_id, record.consumer_id),
            delivery_config_ids=[record.delivery_config_id],
        )
        generation = next(
            (item for item in configurations if item.generation_key == record.generation_key),
            None,
        )
        if generation is None:
            raise LedgerConflictError(
                "UNKNOWN_CONFIGURATION_GENERATION",
                f"{record.delivery_config_id}@{record.delivery_config_version} is not an "
                "accepted configuration generation for this account",
            )
        if generation.report_definition_id != record.report_definition_id:
            raise LedgerConflictError(
                "REPORT_DEFINITION_MISMATCH",
                "the statement names a different report definition than the configuration "
                "generation accepted; unlike reporting promises must not share a status chain",
            )
        if (record.seller_ledger_snapshot_id is None) != (record.seller_ledger_as_of is None):
            raise LedgerConflictError(
                "SELLER_SNAPSHOT_EVIDENCE_INCOMPLETE",
                "seller_ledger_snapshot_id and seller_ledger_as_of are present together or "
                "not at all",
            )
        await self._resolve_named_records(record)
        validate_consumer_status_timing(record, generation, as_of=self._now())

    async def _resolve_named_records(self, record: ConsumerStatusRecord) -> None:
        """Resolve the optional obligation and revision ids the caller supplied.

        Unknown, unauthorized, cross-account, and cross-caller identifiers all
        produce the same result, so this surface cannot be used as an oracle to
        discover which ids exist on another tenant.
        """
        obligation: ReportingObligationRecord | None = None
        if record.reporting_obligation_id is not None:
            obligation = await self.store.get_obligation(
                account_id=record.account_id,
                reporting_obligation_id=record.reporting_obligation_id,
            )
            if obligation is None:
                raise LedgerConflictError(
                    "LOOKUP_UNAVAILABLE",
                    f"reporting_obligation_id {record.reporting_obligation_id} does not "
                    "resolve for this caller and account",
                )
            if (
                obligation.generation_key != record.generation_key
                or obligation.report_definition_id != record.report_definition_id
                or _utc(obligation.period.start) != _utc(record.period_start)
                or _utc(obligation.period.end) != _utc(record.period_end)
            ):
                # The id resolves but describes a different period or
                # generation. Recording it would attach this chain to the wrong
                # obligation, which the immutability rule forbids.
                raise LedgerConflictError(
                    "OBLIGATION_IDENTITY_MISMATCH",
                    f"reporting_obligation_id {record.reporting_obligation_id} names a "
                    "different configuration generation, report definition, or period than "
                    "this statement",
                )

        if obligation is None:
            obligation = await self.store.find_obligation(
                account_id=record.account_id,
                consumer_id=record.consumer_id,
                delivery_config_id=record.delivery_config_id,
                delivery_config_version=record.delivery_config_version,
                period_start=record.period_start,
                period_end=record.period_end,
            )
        required = None
        if obligation is not None:
            revisions = await self.store.list_revisions(
                account_id=record.account_id,
                reporting_obligation_id=obligation.reporting_obligation_id,
            )
            selection = select_reporting_revision(
                revisions,
                account_id=obligation.account_id,
                reporting_obligation_id=obligation.reporting_obligation_id,
                required_finality=obligation.required_finality,
            )
            if selection.kind == "corrupt":
                raise LedgerConflictError(
                    "HISTORY_UNAVAILABLE", "the revision history requires repair"
                )
            if selection.kind == "selected":
                required = selection.revision
        if record.reporting_revision_id is None:
            return

        revision = await self.store.get_revision(
            account_id=record.account_id, reporting_revision_id=record.reporting_revision_id
        )
        if revision is None or (
            obligation is not None
            and revision.reporting_obligation_id != obligation.reporting_obligation_id
        ):
            raise LedgerConflictError(
                "LOOKUP_UNAVAILABLE",
                f"reporting_revision_id {record.reporting_revision_id} does not resolve for "
                "this caller, account, and period",
            )

        if record.consumer_status != "content_mismatch" or obligation is None:
            return

        # ``content_mismatch`` is valid only against a revision the seller
        # *currently requires* for the period. Against a superseded revision the
        # dispute is already moot -- the seller has replaced the content -- and
        # recording it would degrade the caller's own view over bytes neither
        # party stands behind any more. The buyer's move there is to re-read and
        # either accept or dispute the current revision.
        if required is None or required.reporting_revision_id != record.reporting_revision_id:
            raise LedgerConflictError(
                "REVISION_NOT_CURRENTLY_REQUIRED",
                f"content_mismatch names revision {record.reporting_revision_id}, which is "
                "not the revision this seller currently requires for the period; re-read the "
                "current revision and file against that",
            )

    def _now(self) -> datetime:
        return _utc(self.clock() if self.clock is not None else datetime.now(timezone.utc))

    @staticmethod
    def _to_record(
        statement: dict[str, Any],
        *,
        account_id: str,
        consumer_id: str,
        recorded_at: datetime,
    ) -> ConsumerStatusRecord:
        from adcp.types import ReportingConsumerStatus

        parsed = ReportingConsumerStatus.model_validate(statement)
        period = parsed.period
        return ConsumerStatusRecord(
            reporting_status_id=parsed.reporting_status_id,
            # Identity comes from authenticated transport. A body that asserts
            # a buyer or consumer principal is ignored, not trusted.
            account_id=account_id,
            consumer_id=consumer_id,
            delivery_config_id=parsed.delivery_config_id,
            delivery_config_version=parsed.delivery_config_version,
            report_definition_id=parsed.report_definition_id,
            period_start=_utc(period.start),
            period_end=_utc(period.end),
            period_source_timezone=period.source_timezone,
            consumer_status=_narrow_recorded_status(parsed.consumer_status.value),
            status_as_of=_utc(parsed.status_as_of),
            recorded_at=recorded_at,
            supersedes_reporting_status_id=parsed.supersedes_reporting_status_id,
            reporting_obligation_id=parsed.reporting_obligation_id,
            reporting_revision_id=parsed.reporting_revision_id,
            observed_revision_content_sha256=parsed.observed_revision_content_sha256,
            failure_code=parsed.failure_code.value if parsed.failure_code else None,
            mismatch_code=parsed.mismatch_code.value if parsed.mismatch_code else None,
            consumer_commit_ref=parsed.consumer_commit_ref,
            seller_ledger_snapshot_id=parsed.seller_ledger_snapshot_id,
            seller_ledger_as_of=(
                _utc(parsed.seller_ledger_as_of) if parsed.seller_ledger_as_of else None
            ),
        )

Records authenticated consumer statements into the ledger.

Framework-agnostic like :class:~adcp.reporting.ledger.status.ReportingStatusHandler: a request mapping plus an authenticated caller in, a response mapping out.

Instance variables

var clock : collections.abc.Callable[[], datetime.datetime] | None

Supplies recorded_at – when the seller durably recorded the statement – and therefore the issue's opened_at. Defaults to wall clock. A deployment whose ledger boundary comes from the database should pass the same source here, or an issue's escalation clock and the snapshot that reads it will disagree about what time it is.

var enabled : bool | None
var store : ReportingLedgerStore

Methods

async def handle(self, request: dict[str, Any], *, account_id: str, consumer_id: str) ‑> dict[str, typing.Any]
Expand source code
async def handle(
    self, request: dict[str, Any], *, account_id: str, consumer_id: str
) -> dict[str, Any]:
    """Serve one ``sync_reporting_status`` batch.

    Returns one ``recorded`` / ``unchanged`` / ``failed`` result per
    submitted statement.  A failure in one statement does not fail the
    batch: a buyer reporting five periods should not lose four good
    statements to one stale supersession pointer.
    """
    self._require_enabled()

    statements = request.get("statuses") or []
    if not statements:
        raise LedgerConflictError(
            "EMPTY_BATCH", "a sync_reporting_status batch must carry at least one statement"
        )
    seen_chains: set[tuple[Any, ...]] = set()
    results: list[dict[str, Any]] = []
    for statement in statements:
        status_id = statement.get("reporting_status_id", "")
        problems = _validate_consumer_status_wire(statement)
        if problems:
            results.append(_failed(status_id, "INVALID_CONSUMER_STATUS", problems[0]))
            continue
        try:
            record = self._to_record(
                statement,
                account_id=account_id,
                consumer_id=consumer_id,
                recorded_at=self._now(),
            )
        except ValueError as error:
            results.append(_failed(status_id, "INVALID_CONSUMER_STATUS", str(error)))
            continue
        if record.chain_key in seen_chains:
            # "A batch contains at most one update for each logical status
            # chain" -- two updates in one batch have no defined order, so
            # the second would silently win or silently lose.
            results.append(
                _failed(
                    status_id,
                    "DUPLICATE_CHAIN_IN_BATCH",
                    "a batch may carry at most one statement per logical status chain",
                )
            )
            continue
        seen_chains.add(record.chain_key)
        try:
            # Answer "exact retry?" before validating anything, because
            # some of that validation is time-dependent. "content_mismatch
            # must name the revision the seller currently requires" has an
            # answer that changes when the seller restates, so a retry of a
            # statement this ledger already accepted would otherwise be
            # rejected for a reason that did not apply when it was first
            # recorded. The spec's result_mapping requires `unchanged`, and
            # a buyer retrying after a transport failure has no way to tell
            # a genuine rejection from that.
            replay = await self.store.resolve_consumer_status_replay(record)
            if replay is not None:
                results.append({"result": "unchanged", "consumer_status": _to_wire(replay)})
                continue
            await self._validate_against_configuration(record)
            from adcp.reporting.ledger.status_snapshot import ReportingStatusParticipant

            atomic = isinstance(self.store, ReportingStatusParticipant)
            if atomic and isinstance(self.store, ReportingStatusParticipant):
                stored, recorded = await self.store.record_consumer_status_with_lifecycle(
                    record
                )
            else:
                stored, recorded = await self.store.record_consumer_status(record)
        except LedgerConflictError as error:
            results.append(_failed(status_id, error.code, str(error)))
            continue
        if recorded and not atomic:
            await self._open_issue_on_first_observation(stored)
        results.append(
            {
                "result": "recorded" if recorded else "unchanged",
                "consumer_status": _to_wire(stored),
            }
        )
    return {"status": "completed", "results": results}

Serve one sync_reporting_status batch.

Returns one recorded / unchanged / failed result per submitted statement. A failure in one statement does not fail the batch: a buyer reporting five periods should not lose four good statements to one stale supersession pointer.