Module adcp.reporting.ledger.delivery_changes
Consumer-scoped incremental reconciliation reads, independent of Core cursors.
Functions
def change_boundary(caller: ReportingDeliveryPrincipal,
maximum: int,
as_of: datetime,
*,
after: int = 0,
total_count: int | None = None,
filters: ReportingReconciliationFilter = ReportingReconciliationFilter(record_kinds=(), reporting_obligation_id=None)) ‑> ReportingReconciliationSnapshotToken-
Expand source code
def change_boundary( caller: ReportingDeliveryPrincipal, maximum: int, as_of: datetime, *, after: int = 0, total_count: int | None = None, filters: ReportingReconciliationFilter = ReportingReconciliationFilter(), ) -> ReportingReconciliationSnapshotToken: count = maximum - after if total_count is None else total_count if ( any(type(value) is not int for value in (after, maximum, count)) or not 0 <= after <= maximum or not 0 <= count <= maximum - after ): raise LedgerConflictError("INVALID_CHECKPOINT", "invalid reconciliation boundary") identity = hashlib.sha256( canonical_json_utf8_v1([_token_scope(caller, filters), after, maximum, count]) ).hexdigest() return ReportingReconciliationSnapshotToken( caller, "rprc_" + identity[:32], aware_utc(as_of), after, maximum, count, filters ) def change_page(caller: ReportingDeliveryPrincipal,
boundary: ReportingReconciliationSnapshotToken,
after: int,
changes: tuple[ReportingReconciliationChange, ...],
limit: int) ‑> ReportingReconciliationPage-
Expand source code
def change_page( caller: ReportingDeliveryPrincipal, boundary: ReportingReconciliationSnapshotToken, after: int, changes: tuple[ReportingReconciliationChange, ...], limit: int, ) -> ReportingReconciliationPage: validate_boundary(caller, boundary, boundary.max_sequence) if not boundary.min_sequence <= after <= boundary.max_sequence: raise LedgerConflictError("INVALID_CHECKPOINT", "reconciliation position is unavailable") previous = after for change in changes: if ( type(change.sequence) is not int or not previous < change.sequence <= boundary.max_sequence or principal(change.record) != caller or not boundary.filters.matches(change.record) ): fail("REPORTING_HISTORY_CORRUPT") previous = change.sequence window, more = changes[:limit], len(changes) > limit position = _token_scope(caller, boundary.filters) return ReportingReconciliationPage( caller, boundary, window, more, ( ReportingReconciliationCursor( encode_cursor( { **position, "type": "cursor", "seq": window[-1].sequence, "key": change_id(window[-1].record), "after": boundary.min_sequence, "through": boundary.max_sequence, "as_of": boundary.ledger_as_of.isoformat(), "count": boundary.total_count, "snapshot": boundary.snapshot_id, } ) ) if more else None ), ( None if more else ReportingReconciliationCheckpoint( encode_cursor({**position, "type": "checkpoint", "seq": boundary.max_sequence}) ) ), ) def read_position(caller: ReportingDeliveryPrincipal,
changes_after: ReportingReconciliationCheckpoint | None,
cursor: ReportingReconciliationCursor | None,
limit: int,
filters: ReportingReconciliationFilter) ‑> tuple[int, ReportingReconciliationSnapshotToken | None, str | None]-
Expand source code
def read_position( caller: ReportingDeliveryPrincipal, changes_after: ReportingReconciliationCheckpoint | None, cursor: ReportingReconciliationCursor | None, limit: int, filters: ReportingReconciliationFilter, ) -> tuple[int, ReportingReconciliationSnapshotToken | None, str | None]: if ( type(limit) is not int or not 1 <= limit <= 1000 or (changes_after is not None and cursor is not None) or type(caller) is not ReportingDeliveryPrincipal or type(filters) is not ReportingReconciliationFilter ): raise LedgerConflictError("INVALID_CHECKPOINT", "invalid reconciliation page request") if cursor is None and changes_after is None: return 0, None, None position = _decode_position( cursor if cursor is not None else changes_after or "", caller, filters ) sequence = position.get("seq") expected = {*_token_scope(caller, filters), "type", "seq"} if cursor is not None: expected.update(("after", "through", "as_of", "count", "snapshot", "key")) if ( set(position) != expected or type(sequence) is not int or sequence < 0 or position["type"] != ("cursor" if cursor is not None else "checkpoint") ): raise LedgerConflictError("INVALID_CHECKPOINT", "invalid reconciliation position") if cursor is None: return sequence, None, None through, after, as_of = position["through"], position["after"], position["as_of"] boundary = None if ( type(through) is int and type(after) is int and 0 <= after < sequence < through and type(as_of) is str and type(position["count"]) is int and type(position["key"]) is str and len(position["key"]) == 64 ): try: boundary = change_boundary( caller, through, datetime.fromisoformat(as_of), after=after, total_count=position["count"], filters=filters, ) except (ValueError, LedgerConflictError): # Reject below without exposing the parse or validation exception. pass if boundary is None or position["snapshot"] != boundary.snapshot_id: raise LedgerConflictError("INVALID_CHECKPOINT", "invalid reconciliation boundary") return sequence, boundary, position["key"] def validate_boundary(caller: ReportingDeliveryPrincipal,
boundary: ReportingReconciliationSnapshotToken,
maximum: int) ‑> None-
Expand source code
def validate_boundary( caller: ReportingDeliveryPrincipal, boundary: ReportingReconciliationSnapshotToken, maximum: int, ) -> None: valid = False try: if ( type(boundary) is ReportingReconciliationSnapshotToken and boundary.caller == caller and type(boundary.filters) is ReportingReconciliationFilter ): expected = change_boundary( caller, boundary.max_sequence, boundary.ledger_as_of, after=boundary.min_sequence, total_count=boundary.total_count, filters=boundary.filters, ) valid = boundary == expected and boundary.max_sequence <= maximum except (ValueError, TypeError, LedgerConflictError): # Keep the boundary invalid and classify outside the exception handler. pass if not valid: raise LedgerConflictError("INVALID_CHECKPOINT", "reconciliation boundary is unavailable")
Classes
class ReportingReconciliationChange (sequence: int, record: ReportingDeliveryRecord)-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingReconciliationChange: sequence: int record: ReportingDeliveryRecordReportingReconciliationChange(sequence: 'int', record: 'ReportingDeliveryRecord')
Instance variables
var record : ReportingDestinationBinding | ReportingObligationDeliveryRecord | ReportingMaterializationAttempt | ReportingMaterializationRecord | ReportingMaterializationCheck | ReportingRevisionReceiptRecord | ReportingAdjustmentReceiptRecord-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingReconciliationChange: sequence: int record: ReportingDeliveryRecord var sequence : int-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingReconciliationChange: sequence: int record: ReportingDeliveryRecord
class ReportingReconciliationFeedStore (*args, **kwargs)-
Expand source code
@runtime_checkable class ReportingReconciliationFeedStore(Protocol): """An optional seam; existing Core and reconciliation protocols are unchanged. Persist partial pages with their cursor. Only the final page issues a checkpoint, which opens a new walk via ``changes_after``. Tokens bind caller, version, filters, bounds and the last emitted key, and survive reader restarts. Neither token is an authorization grant or a Core status checkpoint. """ async def read_reconciliation_changes( self, *, caller: ReportingDeliveryPrincipal, changes_after: ReportingReconciliationCheckpoint | None = None, cursor: ReportingReconciliationCursor | None = None, limit: int = 100, filters: ReportingReconciliationFilter = ReportingReconciliationFilter(), ) -> ReportingReconciliationPage: ...An optional seam; existing Core and reconciliation protocols are unchanged.
Persist partial pages with their cursor. Only the final page issues a checkpoint, which opens a new walk via
changes_after. Tokens bind caller, version, filters, bounds and the last emitted key, and survive reader restarts. Neither token is an authorization grant or a Core status checkpoint.Ancestors
- typing.Protocol
- typing.Generic
Methods
async def read_reconciliation_changes(self,
*,
caller: ReportingDeliveryPrincipal,
changes_after: ReportingReconciliationCheckpoint | None = None,
cursor: ReportingReconciliationCursor | None = None,
limit: int = 100,
filters: ReportingReconciliationFilter = ReportingReconciliationFilter(record_kinds=(), reporting_obligation_id=None)) ‑> ReportingReconciliationPage-
Expand source code
async def read_reconciliation_changes( self, *, caller: ReportingDeliveryPrincipal, changes_after: ReportingReconciliationCheckpoint | None = None, cursor: ReportingReconciliationCursor | None = None, limit: int = 100, filters: ReportingReconciliationFilter = ReportingReconciliationFilter(), ) -> ReportingReconciliationPage: ...
class ReportingReconciliationChangeStore (*args, **kwargs)-
Expand source code
@runtime_checkable class ReportingReconciliationFeedStore(Protocol): """An optional seam; existing Core and reconciliation protocols are unchanged. Persist partial pages with their cursor. Only the final page issues a checkpoint, which opens a new walk via ``changes_after``. Tokens bind caller, version, filters, bounds and the last emitted key, and survive reader restarts. Neither token is an authorization grant or a Core status checkpoint. """ async def read_reconciliation_changes( self, *, caller: ReportingDeliveryPrincipal, changes_after: ReportingReconciliationCheckpoint | None = None, cursor: ReportingReconciliationCursor | None = None, limit: int = 100, filters: ReportingReconciliationFilter = ReportingReconciliationFilter(), ) -> ReportingReconciliationPage: ...An optional seam; existing Core and reconciliation protocols are unchanged.
Persist partial pages with their cursor. Only the final page issues a checkpoint, which opens a new walk via
changes_after. Tokens bind caller, version, filters, bounds and the last emitted key, and survive reader restarts. Neither token is an authorization grant or a Core status checkpoint.Ancestors
- typing.Protocol
- typing.Generic
Methods
async def read_reconciliation_changes(self,
*,
caller: ReportingDeliveryPrincipal,
changes_after: ReportingReconciliationCheckpoint | None = None,
cursor: ReportingReconciliationCursor | None = None,
limit: int = 100,
filters: ReportingReconciliationFilter = ReportingReconciliationFilter(record_kinds=(), reporting_obligation_id=None)) ‑> ReportingReconciliationPage-
Expand source code
async def read_reconciliation_changes( self, *, caller: ReportingDeliveryPrincipal, changes_after: ReportingReconciliationCheckpoint | None = None, cursor: ReportingReconciliationCursor | None = None, limit: int = 100, filters: ReportingReconciliationFilter = ReportingReconciliationFilter(), ) -> ReportingReconciliationPage: ...
class ReportingReconciliationFilter (record_kinds: tuple[ReportingReconciliationRecordKind, ...] = (),
reporting_obligation_id: str | None = None)-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingReconciliationFilter: record_kinds: tuple[ReportingReconciliationRecordKind, ...] = () reporting_obligation_id: str | None = None def __post_init__(self) -> None: if type(self.record_kinds) is not tuple or any( not _is_record_kind(kind) for kind in self.record_kinds ): raise ValueError("invalid reconciliation record filter") object.__setattr__(self, "record_kinds", tuple(sorted(set(self.record_kinds)))) if self.reporting_obligation_id is not None: reporting_identifier(self.reporting_obligation_id) @property def fingerprint(self) -> str: return hashlib.sha256(canonical_json_utf8_v1(asdict(self))).hexdigest() def matches(self, record: ReportingDeliveryRecord) -> bool: return (not self.record_kinds or record.kind in self.record_kinds) and ( self.reporting_obligation_id is None or ( not isinstance(record, ReportingDestinationBinding) and record.scope.reporting_obligation_id == self.reporting_obligation_id ) )ReportingReconciliationFilter(record_kinds: 'tuple[ReportingReconciliationRecordKind, …]' = (), reporting_obligation_id: 'str | None' = None)
Instance variables
prop fingerprint : str-
Expand source code
@property def fingerprint(self) -> str: return hashlib.sha256(canonical_json_utf8_v1(asdict(self))).hexdigest() var record_kinds : tuple[typing.Literal['destination_binding', 'obligation_delivery', 'materialization_attempt', 'materialization', 'materialization_check', 'revision_receipt', 'adjustment_receipt'], ...]-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingReconciliationFilter: record_kinds: tuple[ReportingReconciliationRecordKind, ...] = () reporting_obligation_id: str | None = None def __post_init__(self) -> None: if type(self.record_kinds) is not tuple or any( not _is_record_kind(kind) for kind in self.record_kinds ): raise ValueError("invalid reconciliation record filter") object.__setattr__(self, "record_kinds", tuple(sorted(set(self.record_kinds)))) if self.reporting_obligation_id is not None: reporting_identifier(self.reporting_obligation_id) @property def fingerprint(self) -> str: return hashlib.sha256(canonical_json_utf8_v1(asdict(self))).hexdigest() def matches(self, record: ReportingDeliveryRecord) -> bool: return (not self.record_kinds or record.kind in self.record_kinds) and ( self.reporting_obligation_id is None or ( not isinstance(record, ReportingDestinationBinding) and record.scope.reporting_obligation_id == self.reporting_obligation_id ) ) var reporting_obligation_id : str | None-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingReconciliationFilter: record_kinds: tuple[ReportingReconciliationRecordKind, ...] = () reporting_obligation_id: str | None = None def __post_init__(self) -> None: if type(self.record_kinds) is not tuple or any( not _is_record_kind(kind) for kind in self.record_kinds ): raise ValueError("invalid reconciliation record filter") object.__setattr__(self, "record_kinds", tuple(sorted(set(self.record_kinds)))) if self.reporting_obligation_id is not None: reporting_identifier(self.reporting_obligation_id) @property def fingerprint(self) -> str: return hashlib.sha256(canonical_json_utf8_v1(asdict(self))).hexdigest() def matches(self, record: ReportingDeliveryRecord) -> bool: return (not self.record_kinds or record.kind in self.record_kinds) and ( self.reporting_obligation_id is None or ( not isinstance(record, ReportingDestinationBinding) and record.scope.reporting_obligation_id == self.reporting_obligation_id ) )
Methods
def matches(self, record: ReportingDeliveryRecord) ‑> bool-
Expand source code
def matches(self, record: ReportingDeliveryRecord) -> bool: return (not self.record_kinds or record.kind in self.record_kinds) and ( self.reporting_obligation_id is None or ( not isinstance(record, ReportingDestinationBinding) and record.scope.reporting_obligation_id == self.reporting_obligation_id ) )
class ReportingReconciliationPage (caller: ReportingDeliveryPrincipal,
boundary: ReportingReconciliationSnapshotToken,
changes: tuple[ReportingReconciliationChange, ...],
has_more: bool,
cursor: ReportingReconciliationCursor | None,
changes_checkpoint: ReportingReconciliationCheckpoint | None)-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingReconciliationPage: caller: ReportingDeliveryPrincipal boundary: ReportingReconciliationSnapshotToken changes: tuple[ReportingReconciliationChange, ...] has_more: bool cursor: ReportingReconciliationCursor | None changes_checkpoint: ReportingReconciliationCheckpoint | None @property def total_count(self) -> int: return self.boundary.total_countReportingReconciliationPage(caller: 'ReportingDeliveryPrincipal', boundary: 'ReportingReconciliationSnapshotToken', changes: 'tuple[ReportingReconciliationChange, …]', has_more: 'bool', cursor: 'ReportingReconciliationCursor | None', changes_checkpoint: 'ReportingReconciliationCheckpoint | None')
Instance variables
var boundary : ReportingReconciliationSnapshotToken-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingReconciliationPage: caller: ReportingDeliveryPrincipal boundary: ReportingReconciliationSnapshotToken changes: tuple[ReportingReconciliationChange, ...] has_more: bool cursor: ReportingReconciliationCursor | None changes_checkpoint: ReportingReconciliationCheckpoint | None @property def total_count(self) -> int: return self.boundary.total_count var caller : ReportingDeliveryPrincipal-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingReconciliationPage: caller: ReportingDeliveryPrincipal boundary: ReportingReconciliationSnapshotToken changes: tuple[ReportingReconciliationChange, ...] has_more: bool cursor: ReportingReconciliationCursor | None changes_checkpoint: ReportingReconciliationCheckpoint | None @property def total_count(self) -> int: return self.boundary.total_count var changes : tuple[ReportingReconciliationChange, ...]-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingReconciliationPage: caller: ReportingDeliveryPrincipal boundary: ReportingReconciliationSnapshotToken changes: tuple[ReportingReconciliationChange, ...] has_more: bool cursor: ReportingReconciliationCursor | None changes_checkpoint: ReportingReconciliationCheckpoint | None @property def total_count(self) -> int: return self.boundary.total_count var changes_checkpoint : adcp.reporting.ledger.delivery_changes.ReportingReconciliationCheckpoint | None-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingReconciliationPage: caller: ReportingDeliveryPrincipal boundary: ReportingReconciliationSnapshotToken changes: tuple[ReportingReconciliationChange, ...] has_more: bool cursor: ReportingReconciliationCursor | None changes_checkpoint: ReportingReconciliationCheckpoint | None @property def total_count(self) -> int: return self.boundary.total_count var cursor : adcp.reporting.ledger.delivery_changes.ReportingReconciliationCursor | None-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingReconciliationPage: caller: ReportingDeliveryPrincipal boundary: ReportingReconciliationSnapshotToken changes: tuple[ReportingReconciliationChange, ...] has_more: bool cursor: ReportingReconciliationCursor | None changes_checkpoint: ReportingReconciliationCheckpoint | None @property def total_count(self) -> int: return self.boundary.total_count var has_more : bool-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingReconciliationPage: caller: ReportingDeliveryPrincipal boundary: ReportingReconciliationSnapshotToken changes: tuple[ReportingReconciliationChange, ...] has_more: bool cursor: ReportingReconciliationCursor | None changes_checkpoint: ReportingReconciliationCheckpoint | None @property def total_count(self) -> int: return self.boundary.total_count prop total_count : int-
Expand source code
@property def total_count(self) -> int: return self.boundary.total_count
class ReportingReconciliationSnapshotToken (caller: ReportingDeliveryPrincipal,
snapshot_id: str,
ledger_as_of: datetime,
min_sequence: int,
max_sequence: int,
total_count: int,
filters: ReportingReconciliationFilter)-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingReconciliationSnapshotToken: """A frozen boundary in one principal's feed, never a Core ledger sequence.""" caller: ReportingDeliveryPrincipal snapshot_id: str ledger_as_of: datetime min_sequence: int max_sequence: int total_count: int filters: ReportingReconciliationFilter @property def account_id(self) -> str: return self.caller.account_id @property def consumer_id(self) -> str: return self.caller.consumer_idA frozen boundary in one principal's feed, never a Core ledger sequence.
Instance variables
prop account_id : str-
Expand source code
@property def account_id(self) -> str: return self.caller.account_id var caller : ReportingDeliveryPrincipal-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingReconciliationSnapshotToken: """A frozen boundary in one principal's feed, never a Core ledger sequence.""" caller: ReportingDeliveryPrincipal snapshot_id: str ledger_as_of: datetime min_sequence: int max_sequence: int total_count: int filters: ReportingReconciliationFilter @property def account_id(self) -> str: return self.caller.account_id @property def consumer_id(self) -> str: return self.caller.consumer_id prop consumer_id : str-
Expand source code
@property def consumer_id(self) -> str: return self.caller.consumer_id var filters : ReportingReconciliationFilter-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingReconciliationSnapshotToken: """A frozen boundary in one principal's feed, never a Core ledger sequence.""" caller: ReportingDeliveryPrincipal snapshot_id: str ledger_as_of: datetime min_sequence: int max_sequence: int total_count: int filters: ReportingReconciliationFilter @property def account_id(self) -> str: return self.caller.account_id @property def consumer_id(self) -> str: return self.caller.consumer_id var ledger_as_of : datetime.datetime-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingReconciliationSnapshotToken: """A frozen boundary in one principal's feed, never a Core ledger sequence.""" caller: ReportingDeliveryPrincipal snapshot_id: str ledger_as_of: datetime min_sequence: int max_sequence: int total_count: int filters: ReportingReconciliationFilter @property def account_id(self) -> str: return self.caller.account_id @property def consumer_id(self) -> str: return self.caller.consumer_id var max_sequence : int-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingReconciliationSnapshotToken: """A frozen boundary in one principal's feed, never a Core ledger sequence.""" caller: ReportingDeliveryPrincipal snapshot_id: str ledger_as_of: datetime min_sequence: int max_sequence: int total_count: int filters: ReportingReconciliationFilter @property def account_id(self) -> str: return self.caller.account_id @property def consumer_id(self) -> str: return self.caller.consumer_id var min_sequence : int-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingReconciliationSnapshotToken: """A frozen boundary in one principal's feed, never a Core ledger sequence.""" caller: ReportingDeliveryPrincipal snapshot_id: str ledger_as_of: datetime min_sequence: int max_sequence: int total_count: int filters: ReportingReconciliationFilter @property def account_id(self) -> str: return self.caller.account_id @property def consumer_id(self) -> str: return self.caller.consumer_id var snapshot_id : str-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingReconciliationSnapshotToken: """A frozen boundary in one principal's feed, never a Core ledger sequence.""" caller: ReportingDeliveryPrincipal snapshot_id: str ledger_as_of: datetime min_sequence: int max_sequence: int total_count: int filters: ReportingReconciliationFilter @property def account_id(self) -> str: return self.caller.account_id @property def consumer_id(self) -> str: return self.caller.consumer_id var total_count : int-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingReconciliationSnapshotToken: """A frozen boundary in one principal's feed, never a Core ledger sequence.""" caller: ReportingDeliveryPrincipal snapshot_id: str ledger_as_of: datetime min_sequence: int max_sequence: int total_count: int filters: ReportingReconciliationFilter @property def account_id(self) -> str: return self.caller.account_id @property def consumer_id(self) -> str: return self.caller.consumer_id