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 issueCapture 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 resultProject 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 : strvar adjustments : tuple[ReportingAdjustmentRecord, ...]var as_of : datetime.datetimevar 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: ReportingStatusScopeStatusLifecycleIntent(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 : ReportingIssueLifecyclevar 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 = NoneStatusObligationProjection(obligation: 'ReportingObligationRecord', projection: 'ObligationProjection', revisions: 'tuple[ReportingRevisionRecord, …]', statuses: 'tuple[ConsumerStatusRecord, …]', reconciliation: 'ReconciliationProjection | None' = None)
Instance variables
var obligation : ReportingObligationRecordvar projection : ObligationProjectionvar reconciliation : ReconciliationProjection | Nonevar 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 = TrueStatusProjectionInput(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 : boolvar delivery_config_ids : tuple[str, ...]var escalation : ReportingDeliveryEscalation | Nonevar feed_purposes : tuple[str, ...]var media_buy_ids : tuple[str, ...]var period_end : datetime.datetime | Nonevar period_start : datetime.datetime | Nonevar reconciliation : tuple[ReportingDestinationBinding | ReportingObligationDeliveryRecord | ReportingMaterializationAttempt | ReportingMaterializationRecord | ReportingMaterializationCheck | ReportingRevisionReceiptRecord | ReportingAdjustmentReceiptRecord, ...] | Nonevar scope : ReportingStatusScopevar 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 : strvar 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 | Nonevar obligations : tuple[StatusObligationProjection, ...]var pending_count : intprop 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], }