Module adcp.reporting.ledger.status_projection

Pure status projection shared by polling, source transactions and clock sweeps.

No connection, clock or checkpoint is consulted here. A provisional result contains closed lifecycle intents; callers apply them and reproject before publishing. Checkpoints remember notifications, never determine health.

Functions

def apply_intents_to_snapshot(snapshot: ReportingStatusSnapshot,
intents: Sequence[StatusLifecycleIntent]) ‑> ReportingStatusSnapshot
Expand source code
def apply_intents_to_snapshot(
    snapshot: ReportingStatusSnapshot, intents: Sequence[StatusLifecycleIntent]
) -> ReportingStatusSnapshot:
    """Pure replay helper. Durable callers persist the same intents first."""
    lifecycles = {(i.issue_key, i.generation): i for i in snapshot.lifecycles}
    scopes = dict(snapshot.issue_scopes)
    for intent in intents:
        issue = intent.lifecycle
        lifecycles[(issue.issue_key, issue.generation)] = issue
        scopes[issue.issue_id] = intent.scope
    return replace(
        snapshot, lifecycles=tuple(lifecycles.values()), issue_scopes=tuple(scopes.items())
    )

Pure replay helper. Durable callers persist the same intents first.

def bind_mismatch_waiver(snapshot: ReportingStatusSnapshot,
issue: ReportingIssueLifecycle) ‑> ReportingIssueLifecycle
Expand source code
def bind_mismatch_waiver(
    snapshot: ReportingStatusSnapshot, issue: ReportingIssueLifecycle
) -> ReportingIssueLifecycle:
    """Capture the exact current statement/diagnosis before the operator move.

    Bilateral consent and its private audit are the adopter's prerequisite for
    calling set_issue_state(waived); external_ref is never interpreted as proof.
    """
    if issue.issue_state == "waived":
        return issue
    latest = {i.issue_key: i for i in sorted(snapshot.lifecycles, key=lambda i: i.generation)}
    for status in snapshot.statuses:
        if status.superseded or status.consumer_id != issue.consumer_id:
            continue
        owner = next(
            (o for o in snapshot.obligations if status_matches_obligation(status, o)), None
        )
        revisions = tuple(
            r
            for r in snapshot.revisions
            if owner is not None and r.reporting_obligation_id == owner.reporting_obligation_id
        )
        required = current_required_revision(owner, revisions) if owner else None
        key, _ = mismatch_occurrence(status, owner, required, latest)
        if key == issue.issue_key and consumer_statement_conflicts(
            current=status, current_revision=required, revisions=revisions
        ):
            return replace(
                issue,
                waived_reporting_status_id=status.reporting_status_id,
                waived_conflict_sha256=mismatch_conflict_fingerprint(status, owner, required),
            )
    return issue

Capture the exact current statement/diagnosis before the operator move.

Bilateral consent and its private audit are the adopter's prerequisite for calling set_issue_state(waived); external_ref is never interpreted as proof.

def configuration_selected(configuration: ReportingConfiguration,
*,
delivery_config_ids: Sequence[str] = (),
media_buy_ids: Sequence[str] = (),
feed_purposes: Sequence[str] = ()) ‑> bool
Expand source code
def configuration_selected(
    configuration: ReportingConfiguration,
    *,
    delivery_config_ids: Sequence[str] = (),
    media_buy_ids: Sequence[str] = (),
    feed_purposes: Sequence[str] = (),
) -> bool:
    return (
        (not delivery_config_ids or configuration.delivery_config_id in delivery_config_ids)
        and (not feed_purposes or configuration.feed_purpose in feed_purposes)
        and (
            not media_buy_ids or bool(set(configuration.media_buy_ids).intersection(media_buy_ids))
        )
    )
def lifecycle_intents(snapshot: ReportingStatusSnapshot) ‑> tuple[StatusLifecycleIntent, ...]
Expand source code
def lifecycle_intents(snapshot: ReportingStatusSnapshot) -> tuple[StatusLifecycleIntent, ...]:
    """Plan lifecycle from current chains, including pre-obligation statements.

    Repair old status-without-issue rows using their earliest provable recorded
    observation (or the first superseding revision for stale receipts), never
    a sweeper's later wall clock.
    """
    latest: dict[str, ReportingIssueLifecycle] = {}
    for issue in snapshot.lifecycles:
        if issue.generation >= getattr(latest.get(issue.issue_key), "generation", 0):
            latest[issue.issue_key] = issue
    scopes = dict(snapshot.issue_scopes)
    configurations = {c.generation_key: c for c in snapshot.configurations}
    result: list[StatusLifecycleIntent] = []
    for status in snapshot.statuses:
        if status.superseded:
            continue
        configuration = configurations.get(status.generation_key)
        if configuration is None:
            continue
        owner = next(
            (o for o in snapshot.obligations if status_matches_obligation(status, o)), None
        )
        revisions = tuple(
            r
            for r in snapshot.revisions
            if owner is not None and r.reporting_obligation_id == owner.reporting_obligation_id
        )
        required = None
        if owner is not None:
            selection = select_reporting_revision(
                revisions,
                account_id=owner.account_id,
                reporting_obligation_id=owner.reporting_obligation_id,
                required_finality=owner.required_finality,
            )
            if selection.kind == "corrupt":
                # Keep prior consumer occurrences intact until seller history
                # can establish agreement. Health exposes HISTORY_UNAVAILABLE.
                continue
            if selection.kind == "selected":
                required = selection.revision
        conflicts = consumer_statement_conflicts(
            current=status, current_revision=required, revisions=revisions
        )
        scope = ReportingStatusScope(
            snapshot.account_id,
            status.generation_key,
            owner.reporting_obligation_id if owner is not None else None,
            status.consumer_id,
            cast(Literal["pacing", "analytics", "billing"], configuration.feed_purpose),
        )
        key, previous = mismatch_occurrence(status, owner, required, latest)
        live = previous if previous is not None and previous.live else None
        if live is not None:
            existing_scope = scopes.get(live.issue_id)
            validate_scope_refinement(existing_scope, scope)
            if existing_scope != scope:
                result.append(StatusLifecycleIntent("refine_scope", live, scope))
        if not conflicts:
            if live is not None and live.issue_state != "waived":
                result.append(
                    StatusLifecycleIntent(
                        "retire_mismatch",
                        replace(live, issue_state="resolved", retired_at=snapshot.as_of),
                        scope,
                    )
                )
            continue
        if live is not None:
            continue
        generation = previous.generation + 1 if previous is not None else 1
        observed_at = status.recorded_at
        superseders = [
            r.created_at
            for r in revisions
            if r.supersedes_reporting_revision_id == status.reporting_revision_id
        ]
        if status.consumer_status == "received" and superseders:
            observed_at = max(observed_at, min(superseders))
        if previous is not None and previous.retired_at is not None:
            observed_at = max(observed_at, previous.retired_at)
        waived_parent = next(
            (
                i
                for i in latest.values()
                if i.issue_state == "waived" and condition_after_waiver(i) == key
            ),
            None,
        )
        if waived_parent is not None:
            # A distinct post-waiver diagnosis is observed at this source
            # boundary, not when the earlier immutable statement was filed.
            # Replayed boundaries retain their original as_of timestamp.
            observed_at = max(observed_at, snapshot.as_of)
        result.append(
            StatusLifecycleIntent(
                "ensure_mismatch",
                ReportingIssueLifecycle(
                    issue_key=key,
                    issue_id=issue_id_for_occurrence(key, generation),
                    account_id=snapshot.account_id,
                    consumer_id=status.consumer_id,
                    opened_at=observed_at,
                    generation=generation,
                ),
                scope,
            )
        )
    return tuple(result)

Plan lifecycle from current chains, including pre-obligation statements.

Repair old status-without-issue rows using their earliest provable recorded observation (or the first superseding revision for stale receipts), never a sweeper's later wall clock.

def mismatch_key(status: ConsumerStatusRecord) ‑> str
Expand source code
def mismatch_key(status: ConsumerStatusRecord) -> str:
    return consumer_mismatch_issue_key(
        account_id=status.account_id,
        consumer_id=status.consumer_id,
        delivery_config_id=status.delivery_config_id,
        delivery_config_version=status.delivery_config_version,
        report_definition_id=status.report_definition_id,
        period_start=status.period_start,
        period_end=status.period_end,
    )
def mismatch_occurrence(status: ConsumerStatusRecord,
obligation: ReportingObligationRecord | None,
required: ReportingRevisionRecord | None,
latest: dict[str, ReportingIssueLifecycle]) ‑> tuple[str, ReportingIssueLifecycle | None]
Expand source code
def mismatch_occurrence(
    status: ConsumerStatusRecord,
    obligation: ReportingObligationRecord | None,
    required: ReportingRevisionRecord | None,
    latest: dict[str, ReportingIssueLifecycle],
) -> tuple[str, ReportingIssueLifecycle | None]:
    """Follow terminal exact waivers; keep an unwaived occurrence stable."""
    key = mismatch_key(status)
    seen: set[str] = set()
    while (previous := latest.get(key)) is not None and previous.issue_state == "waived":
        if key in seen:
            raise ReportingNotificationError("invalid_waiver_chain")
        seen.add(key)
        if waiver_covers_mismatch(previous, status, obligation, required):
            return key, previous
        key = condition_after_waiver(previous)
    return key, latest.get(key)

Follow terminal exact waivers; keep an unwaived occurrence stable.

def period_selected(start: datetime,
end: datetime,
requested_start: datetime | None,
requested_end: datetime | None) ‑> bool
Expand source code
def period_selected(
    start: datetime,
    end: datetime,
    requested_start: datetime | None,
    requested_end: datetime | None,
) -> bool:
    """One half-open intersection rule for records and aggregate projection."""
    return (requested_start is None or end > requested_start) and (
        requested_end is None or start < requested_end
    )

One half-open intersection rule for records and aggregate projection.

def project_status_scope(value: StatusProjectionInput) ‑> StatusProjectionResult
Expand source code
def project_status_scope(value: StatusProjectionInput) -> StatusProjectionResult:
    """Project exact semantics and only deadlines that can change those semantics."""
    result, candidates = _project(value)
    if result.intents:
        return result
    for at in sorted(candidates):
        if at <= value.snapshot.as_of:
            continue
        future, _ = _project(replace(value, snapshot=replace(value.snapshot, as_of=at)))
        if future.intents or future.fingerprint != result.fingerprint:
            return replace(result, next_due_at=at)
    return result

Project exact semantics and only deadlines that can change those semantics.

def projection_scopes(snapshot: ReportingStatusSnapshot) ‑> tuple[ReportingStatusScope, ...]
Expand source code
def projection_scopes(snapshot: ReportingStatusSnapshot) -> tuple[ReportingStatusScope, ...]:
    """Seller mutations affect every extant private view of that account.

    Conservative expansion deliberately includes configurations without a
    consumer statement yet. No opaque issue key is parsed to discover scope.
    """
    scopes = [
        ReportingStatusScope(
            snapshot.account_id,
            configuration.generation_key,
            consumer_id=configuration.consumer_id,
            feed_purpose=cast(
                Literal["pacing", "analytics", "billing"], configuration.feed_purpose
            ),
        )
        for configuration in snapshot.configurations
    ]
    scopes.extend(ReportingStatusScope.for_obligation(o) for o in snapshot.obligations)
    return tuple(sorted(scopes, key=lambda s: s.checkpoint_key))

Seller mutations affect every extant private view of that account.

Conservative expansion deliberately includes configurations without a consumer statement yet. No opaque issue key is parsed to discover scope.

def status_matches_obligation(status: ConsumerStatusRecord, obligation: ReportingObligationRecord) ‑> bool
Expand source code
def status_matches_obligation(
    status: ConsumerStatusRecord, obligation: ReportingObligationRecord
) -> bool:
    return status.account_id == obligation.account_id and (
        status.reporting_obligation_id == obligation.reporting_obligation_id
        or (
            status.generation_key == obligation.generation_key
            and status.report_definition_id == obligation.report_definition_id
            and status.period_start == obligation.period.start
            and status.period_end == obligation.period.end
        )
    )
def status_retained_from(configurations: Sequence[ReportingConfiguration], as_of: datetime) ‑> datetime.datetime
Expand source code
def status_retained_from(
    configurations: Sequence[ReportingConfiguration], as_of: datetime
) -> datetime:
    """The common provable horizon is the latest per-generation boundary."""
    return max(
        (
            (
                max(as_of - timedelta(days=c.status_retention_days), min(c.activated_at, as_of))
                if c.activated_at is not None
                else as_of - timedelta(days=c.status_retention_days)
            )
            for c in configurations
        ),
        default=as_of,
    )

The common provable horizon is the latest per-generation boundary.

def with_replay_lifecycles(snapshot: ReportingStatusSnapshot,
prior: ReportingStatusSnapshot | None) ‑> ReportingStatusSnapshot
Expand source code
def with_replay_lifecycles(
    snapshot: ReportingStatusSnapshot,
    prior: ReportingStatusSnapshot | None,
) -> ReportingStatusSnapshot:
    """Carry derived occurrences across old-writer boundaries lacking lifecycle.

    A later capture may contain an old writer's still-open occurrence after a
    matching statement already retired it in replay. Never regress that state,
    drop a derived recurrence, or broaden its validated persisted scope.
    """
    if prior is None:
        return snapshot
    issues = {(i.issue_key, i.generation): i for i in prior.lifecycles}
    rank = {"open": 0, "acknowledged": 1, "resolved": 2, "waived": 2}
    for issue in snapshot.lifecycles:
        key = (issue.issue_key, issue.generation)
        old = issues.get(key)
        if (
            old is None
            or (
                rank[issue.issue_state] >= rank[old.issue_state]
                and old.issue_state not in {"resolved", "waived"}
            )
            or (
                old.issue_state == issue.issue_state == "waived"
                and old.waived_reporting_status_id == issue.waived_reporting_status_id
                and old.waived_conflict_sha256 == issue.waived_conflict_sha256
                and old.retired_at == issue.retired_at
            )
        ):
            issues[key] = issue
    scopes = dict(prior.issue_scopes)
    for issue_id, scope in snapshot.issue_scopes:
        previous = scopes.get(issue_id)
        if previous is None:
            scopes[issue_id] = scope
            continue
        try:
            validate_scope_refinement(previous, scope)
            scopes[issue_id] = scope
        except ReportingNotificationError:
            validate_scope_refinement(scope, previous)
    return replace(snapshot, lifecycles=tuple(issues.values()), issue_scopes=tuple(scopes.items()))

Carry derived occurrences across old-writer boundaries lacking lifecycle.

A later capture may contain an old writer's still-open occurrence after a matching statement already retired it in replay. Never regress that state, drop a derived recurrence, or broaden its validated persisted scope.

Classes

class ReportingStatusSnapshot (account_id: str,
as_of: datetime,
configurations: tuple[ReportingConfiguration, ...] = (),
obligations: tuple[ReportingObligationRecord, ...] = (),
revisions: tuple[ReportingRevisionRecord, ...] = (),
statuses: tuple[ConsumerStatusRecord, ...] = (),
lifecycles: tuple[ReportingIssueLifecycle, ...] = (),
issue_scopes: tuple[tuple[str, ReportingStatusScope], ...] = (),
consumer_ids: tuple[str, ...] = (),
adjustments: tuple[ReportingAdjustmentRecord, ...] = (),
changes: tuple[tuple[int, str, str, str], ...] = ())
Expand source code
@dataclass(frozen=True)
class ReportingStatusSnapshot:
    account_id: str
    as_of: datetime
    configurations: tuple[ReportingConfiguration, ...] = ()
    obligations: tuple[ReportingObligationRecord, ...] = ()
    revisions: tuple[ReportingRevisionRecord, ...] = ()
    statuses: tuple[ConsumerStatusRecord, ...] = ()
    lifecycles: tuple[ReportingIssueLifecycle, ...] = ()
    issue_scopes: tuple[tuple[str, ReportingStatusScope], ...] = ()
    consumer_ids: tuple[str, ...] = ()
    adjustments: tuple[ReportingAdjustmentRecord, ...] = ()
    # Sequence, record kind, ID, authenticated generation/status owner.
    changes: tuple[tuple[int, str, str, str], ...] = ()

    @property
    def max_sequence(self) -> int:
        return max((row[0] for row in self.changes), default=0)

ReportingStatusSnapshot(account_id: 'str', as_of: 'datetime', configurations: 'tuple[ReportingConfiguration, …]' = (), obligations: 'tuple[ReportingObligationRecord, …]' = (), revisions: 'tuple[ReportingRevisionRecord, …]' = (), statuses: 'tuple[ConsumerStatusRecord, …]' = (), lifecycles: 'tuple[ReportingIssueLifecycle, …]' = (), issue_scopes: 'tuple[tuple[str, ReportingStatusScope], …]' = (), consumer_ids: 'tuple[str, …]' = (), adjustments: 'tuple[ReportingAdjustmentRecord, …]' = (), changes: 'tuple[tuple[int, str, str, str], …]' = ())

Instance variables

var account_id : str
var adjustments : tuple[ReportingAdjustmentRecord, ...]
var as_of : datetime.datetime
var changes : tuple[tuple[int, str, str, str], ...]
var configurations : tuple[ReportingConfiguration, ...]
var consumer_ids : tuple[str, ...]
var issue_scopes : tuple[tuple[str, ReportingStatusScope], ...]
var lifecycles : tuple[ReportingIssueLifecycle, ...]
prop max_sequence : int
Expand source code
@property
def max_sequence(self) -> int:
    return max((row[0] for row in self.changes), default=0)
var obligations : tuple[ReportingObligationRecord, ...]
var revisions : tuple[ReportingRevisionRecord, ...]
var statuses : tuple[ConsumerStatusRecord, ...]
class StatusLifecycleIntent (action: "Literal['ensure_mismatch', 'retire_mismatch', 'refine_scope']",
lifecycle: ReportingIssueLifecycle,
scope: ReportingStatusScope)
Expand source code
@dataclass(frozen=True)
class StatusLifecycleIntent:
    action: Literal["ensure_mismatch", "retire_mismatch", "refine_scope"]
    lifecycle: ReportingIssueLifecycle
    scope: ReportingStatusScope

StatusLifecycleIntent(action: "Literal['ensure_mismatch', 'retire_mismatch', 'refine_scope']", lifecycle: 'ReportingIssueLifecycle', scope: 'ReportingStatusScope')

Instance variables

var action : Literal['ensure_mismatch', 'retire_mismatch', 'refine_scope']
var lifecycle : ReportingIssueLifecycle
var scope : ReportingStatusScope
class StatusObligationProjection (obligation: ReportingObligationRecord,
projection: ObligationProjection,
revisions: tuple[ReportingRevisionRecord, ...],
statuses: tuple[ConsumerStatusRecord, ...],
reconciliation: ReconciliationProjection | None = None)
Expand source code
@dataclass(frozen=True)
class StatusObligationProjection:
    obligation: ReportingObligationRecord
    projection: ObligationProjection
    revisions: tuple[ReportingRevisionRecord, ...]
    statuses: tuple[ConsumerStatusRecord, ...]
    reconciliation: ReconciliationProjection | None = None

StatusObligationProjection(obligation: 'ReportingObligationRecord', projection: 'ObligationProjection', revisions: 'tuple[ReportingRevisionRecord, …]', statuses: 'tuple[ConsumerStatusRecord, …]', reconciliation: 'ReconciliationProjection | None' = None)

Instance variables

var obligation : ReportingObligationRecord
var projection : ObligationProjection
var reconciliation : ReconciliationProjection | None
var revisions : tuple[ReportingRevisionRecord, ...]
var statuses : tuple[ConsumerStatusRecord, ...]
class StatusProjectionInput (snapshot: ReportingStatusSnapshot,
scope: ReportingStatusScope,
escalation: ReportingDeliveryEscalation | None = None,
delivery_config_ids: tuple[str, ...] = (),
media_buy_ids: tuple[str, ...] = (),
feed_purposes: tuple[str, ...] = (),
period_start: datetime | None = None,
period_end: datetime | None = None,
reconciliation: tuple[ReportingDeliveryRecord, ...] | None = None,
consumer_status_enabled: bool = True)
Expand source code
@dataclass(frozen=True)
class StatusProjectionInput:
    snapshot: ReportingStatusSnapshot
    scope: ReportingStatusScope
    escalation: ReportingDeliveryEscalation | None = None
    delivery_config_ids: tuple[str, ...] = ()
    media_buy_ids: tuple[str, ...] = ()
    feed_purposes: tuple[str, ...] = ()
    period_start: datetime | None = None
    period_end: datetime | None = None
    # None selects the immutable legacy C representation. Versioned callers
    # pass the complete captured account history, even when it is empty.
    reconciliation: tuple[ReportingDeliveryRecord, ...] | None = None
    consumer_status_enabled: bool = True

StatusProjectionInput(snapshot: 'ReportingStatusSnapshot', scope: 'ReportingStatusScope', escalation: 'ReportingDeliveryEscalation | None' = None, delivery_config_ids: 'tuple[str, …]' = (), media_buy_ids: 'tuple[str, …]' = (), feed_purposes: 'tuple[str, …]' = (), period_start: 'datetime | None' = None, period_end: 'datetime | None' = None, reconciliation: 'tuple[ReportingDeliveryRecord, …] | None' = None, consumer_status_enabled: 'bool' = True)

Instance variables

var consumer_status_enabled : bool
var delivery_config_ids : tuple[str, ...]
var escalation : ReportingDeliveryEscalation | None
var feed_purposes : tuple[str, ...]
var media_buy_ids : tuple[str, ...]
var period_end : datetime.datetime | None
var period_start : datetime.datetime | None
var reconciliation : tuple[ReportingDestinationBinding | ReportingObligationDeliveryRecord | ReportingMaterializationAttempt | ReportingMaterializationRecord | ReportingMaterializationCheck | ReportingRevisionReceiptRecord | ReportingAdjustmentReceiptRecord, ...] | None
var scope : ReportingStatusScope
var snapshot : ReportingStatusSnapshot
class StatusProjectionResult (scope: ReportingStatusScope,
health: ReportingHealth,
issues: tuple[ReportingIssue, ...],
obligations: tuple[StatusObligationProjection, ...],
configurations: tuple[ReportingConfiguration, ...],
pending_count: int,
next_due_at: datetime | None,
intents: tuple[StatusLifecycleIntent, ...],
fingerprint: str)
Expand source code
@dataclass(frozen=True)
class StatusProjectionResult:
    scope: ReportingStatusScope
    health: ReportingHealth
    issues: tuple[ReportingIssue, ...]
    obligations: tuple[StatusObligationProjection, ...]
    configurations: tuple[ReportingConfiguration, ...]
    pending_count: int
    next_due_at: datetime | None
    intents: tuple[StatusLifecycleIntent, ...]
    fingerprint: str

    @property
    def publishable(self) -> bool:
        return not self.intents and (
            any(i.severity == self.health for i in self.issues)
            if self.health in {"delayed", "action_required"}
            else not self.issues
        )

    @property
    def issue_ids(self) -> tuple[str, ...]:
        return tuple(sorted({issue.issue_id for issue in self.issues}))[:16]

    def canonical(self) -> dict[str, Any]:
        return {
            "health": self.health,
            "issues": [issue.to_wire() for issue in self.issues],
        }

StatusProjectionResult(scope: 'ReportingStatusScope', health: 'ReportingHealth', issues: 'tuple[ReportingIssue, …]', obligations: 'tuple[StatusObligationProjection, …]', configurations: 'tuple[ReportingConfiguration, …]', pending_count: 'int', next_due_at: 'datetime | None', intents: 'tuple[StatusLifecycleIntent, …]', fingerprint: 'str')

Instance variables

var configurations : tuple[ReportingConfiguration, ...]
var fingerprint : str
var health : Literal['healthy', 'waiting', 'delayed', 'action_required', 'complete']
var intents : tuple[StatusLifecycleIntent, ...]
prop issue_ids : tuple[str, ...]
Expand source code
@property
def issue_ids(self) -> tuple[str, ...]:
    return tuple(sorted({issue.issue_id for issue in self.issues}))[:16]
var issues : tuple[ReportingIssue, ...]
var next_due_at : datetime.datetime | None
var obligations : tuple[StatusObligationProjection, ...]
var pending_count : int
prop publishable : bool
Expand source code
@property
def publishable(self) -> bool:
    return not self.intents and (
        any(i.severity == self.health for i in self.issues)
        if self.health in {"delayed", "action_required"}
        else not self.issues
    )
var scope : ReportingStatusScope

Methods

def canonical(self) ‑> dict[str, typing.Any]
Expand source code
def canonical(self) -> dict[str, Any]:
    return {
        "health": self.health,
        "issues": [issue.to_wire() for issue in self.issues],
    }