Module adcp.reporting.testing
Deterministic adapter, failure, and service lifecycle test helpers.
Adapter authors can exercise the exact production wrapper without constructing manifests or staged objects themselves. The conformance runner executes the same frozen slice twice and verifies replay identity plus every staged byte. Named barriers work with synchronous provider SDKs as well as async adapters. Store wrappers and restart assertions use the public memory/PG store contracts.
Functions
async def assert_reporting_lifecycle_replay(expected: ReportingLifecycleSnapshot,
store: ReportingLedgerStore,
*,
page_size: int = 500,
max_pages: int = 1000) ‑> None-
Expand source code
async def assert_reporting_lifecycle_replay( expected: ReportingLifecycleSnapshot, store: ReportingLedgerStore, *, page_size: int = 500, max_pages: int = 1000, ) -> None: """Assert no lost, changed, or extra publication after restart and retry.""" actual = await capture_reporting_lifecycle( store, account_id=expected.obligation.account_id, reporting_obligation_id=expected.obligation.reporting_obligation_id, page_size=page_size, max_pages=max_pages, ) if actual != expected: raise AssertionError("reporting obligation, revisions, or row bytes changed after replay")Assert no lost, changed, or extra publication after restart and retry.
async def capture_reporting_lifecycle(store: ReportingLedgerStore,
*,
account_id: str,
reporting_obligation_id: str,
page_size: int = 500,
max_pages: int = 1000) ‑> adcp.reporting._testing_lifecycle.ReportingLifecycleSnapshot-
Expand source code
async def capture_reporting_lifecycle( store: ReportingLedgerStore, *, account_id: str, reporting_obligation_id: str, page_size: int = 500, max_pages: int = 1000, ) -> ReportingLifecycleSnapshot: """Read and verify one published obligation through public store methods. Works with in-memory, PostgreSQL, or adopter stores. All pages must repeat the exact revision identity/count and exhaust within ``max_pages``. The complete walk is hashed against the committed Core binding, including the typed control totals when the publisher supplied managed evidence. Keep workers quiescent during capture/replay comparison. For memory tests, retain the original stores across service instances; that models service recreation only. PostgreSQL tests should create fresh store instances over the retained database. No private state is copied or repaired by the kit. """ if type(page_size) is not int or not 1 <= page_size <= 500: raise ValueError("page_size must be an integer between 1 and 500") if type(max_pages) is not int or max_pages <= 0: raise ValueError("max_pages must be a positive integer") obligation = await store.get_obligation( account_id=account_id, reporting_obligation_id=reporting_obligation_id ) if obligation is None: raise AssertionError("reporting obligation did not survive lifecycle recovery") if (obligation.account_id, obligation.reporting_obligation_id) != ( account_id, reporting_obligation_id, ): raise AssertionError("obligation read returned a different account or identity") revisions = tuple( sorted( await store.list_revisions( account_id=account_id, reporting_obligation_id=reporting_obligation_id ), key=lambda item: item.reporting_revision_id, ) ) if not revisions: raise AssertionError("lifecycle conformance requires at least one published revision") if len({item.reporting_revision_id for item in revisions}) != len(revisions): raise AssertionError("revision listing repeated an identity") contents: list[bytes] = [] for revision in revisions: revision_id = revision.reporting_revision_id if (revision.account_id, revision.reporting_obligation_id) != ( account_id, reporting_obligation_id, ): raise AssertionError("revision read returned a different account or obligation") exact = await store.get_revision(account_id=account_id, reporting_revision_id=revision_id) if exact != revision: raise AssertionError("revision listing disagrees with the exact revision read") rows: list[dict[str, Any]] = [] cursor: str | None = None seen: set[str] = set() for _ in range(max_pages): page = await store.read_revision_rows( account_id=account_id, reporting_revision_id=revision_id, cursor=cursor, limit=page_size, ) if ( page.reporting_revision_id != revision_id or type(page.total_count) is not int or page.total_count != revision.row_count or type(page.has_more) is not bool or len(page.rows) > page_size ): raise AssertionError("revision page changed its identity, count, or page bound") rows.extend(page.rows) if len(rows) > revision.row_count: raise AssertionError("revision walk returned more rows than its binding") if not page.has_more: if page.cursor is not None: raise AssertionError("final revision page retained a continuation cursor") break if not page.rows or not isinstance(page.cursor, str) or not page.cursor: raise AssertionError("revision continuation made no progress") if page.cursor in seen: raise AssertionError("revision continuation repeated a cursor") seen.add(page.cursor) cursor = page.cursor else: raise AssertionError("revision walk exceeded max_pages") if len(rows) != revision.row_count: raise AssertionError("revision walk ended before all bound rows were read") digest = revision_content_sha256( reporting_revision_id=revision_id, row_count=revision.row_count, control_totals=revision.control_totals, control_total_evidence=revision.managed_control_totals, reporting_rows=rows, ) if digest != revision.revision_content_sha256: raise AssertionError("revision rows do not match the committed content binding") contents.append(canonical_json_utf8_v1(rows)) return ReportingLifecycleSnapshot(obligation, revisions, tuple(contents))Read and verify one published obligation through public store methods.
Works with in-memory, PostgreSQL, or adopter stores. All pages must repeat the exact revision identity/count and exhaust within
max_pages. The complete walk is hashed against the committed Core binding, including the typed control totals when the publisher supplied managed evidence.Keep workers quiescent during capture/replay comparison. For memory tests, retain the original stores across service instances; that models service recreation only. PostgreSQL tests should create fresh store instances over the retained database. No private state is copied or repaired by the kit.
async def run_reporting_adapter_conformance(adapter: ReportingAdapter,
request: ReportingSourceSliceRequestV1,
*,
clock: Callable[[], datetime] | None = None,
staging: ReportingStagingStore | None = None,
seals: ReportingSealStore | None = None) ‑> SourceBatchManifestV1-
Expand source code
async def run_reporting_adapter_conformance( adapter: ReportingAdapter, request: ReportingSourceSliceRequestV1, *, clock: Callable[[], datetime] | None = None, staging: ReportingStagingStore | None = None, seals: ReportingSealStore | None = None, ) -> SourceBatchManifestV1: """Wrap one minimal adapter and run the SDK's full replay conformance. The adapter needs only ``capabilities`` and ``fetch_slice``. This helper uses fresh memory stores by default. Supply retained memory stores or freshly constructed PostgreSQL stores to test recovery with the same request; create their schemas first. Stores and the adapter are borrowed. The same clock drives both execution timestamps and conformance deadlines. """ staging = staging if staging is not None else InMemoryStagingStore() seals = seals if seals is not None else InMemorySealStore() test_clock = clock or DeterministicReportingClock(request.period.source_read_cutoff_at) executor = InlineReportingSource( capabilities=adapter.capabilities, fetch=adapter.fetch_slice, staging=staging, seals=seals, clock=test_clock, ) return await run_reporting_source_replay_conformance( executor=executor, request=request, object_reader=staging, clock=test_clock, )Wrap one minimal adapter and run the SDK's full replay conformance.
The adapter needs only
capabilitiesandfetch_slice. This helper uses fresh memory stores by default. Supply retained memory stores or freshly constructed PostgreSQL stores to test recovery with the same request; create their schemas first. Stores and the adapter are borrowed. The same clock drives both execution timestamps and conformance deadlines.
Classes
class DeterministicReportingClock (now: datetime)-
Expand source code
class DeterministicReportingClock: """A small aware UTC clock that tests can advance without sleeping.""" def __init__(self, now: datetime) -> None: self.set(now) def __call__(self) -> datetime: return self._now def set(self, value: datetime) -> datetime: if value.tzinfo is None or value.utcoffset() is None: raise ValueError("DeterministicReportingClock requires an aware datetime") self._now = value.astimezone(timezone.utc) return self._now def advance(self, delta: timedelta) -> datetime: if delta < timedelta(0): raise ValueError("a deterministic reporting clock cannot move backward") return self.set(self._now + delta)A small aware UTC clock that tests can advance without sleeping.
Methods
def advance(self, delta: timedelta) ‑> datetime.datetime-
Expand source code
def advance(self, delta: timedelta) -> datetime: if delta < timedelta(0): raise ValueError("a deterministic reporting clock cannot move backward") return self.set(self._now + delta) def set(self, value: datetime) ‑> datetime.datetime-
Expand source code
def set(self, value: datetime) -> datetime: if value.tzinfo is None or value.utcoffset() is None: raise ValueError("DeterministicReportingClock requires an aware datetime") self._now = value.astimezone(timezone.utc) return self._now
class FaultInjectingReportingSealStore (store: ReportingSealStore,
failures: ReportingFailurePlan)-
Expand source code
class FaultInjectingReportingSealStore: """Borrow a seal store with ``seal.get`` and ``seal`` before/after points. In particular, ``seal.after`` models a lost acknowledgement: replay must discover the original seal instead of refetching or publishing new bytes. """ def __init__(self, store: ReportingSealStore, failures: ReportingFailurePlan) -> None: self._store = store self._failures = failures async def get(self, *, account_id: str, source_execution_key: str) -> SealedSlice | None: await self._failures.hit("seal.get.before") result = await self._store.get( account_id=account_id, source_execution_key=source_execution_key ) await self._failures.hit("seal.get.after") return result async def put( self, *, account_id: str, source_execution_key: str, sealed: SealedSlice ) -> SealedSlice: await self._failures.hit("seal.before") result = await self._store.put( account_id=account_id, source_execution_key=source_execution_key, sealed=sealed ) await self._failures.hit("seal.after") return resultBorrow a seal store with
seal.getandsealbefore/after points.In particular,
seal.aftermodels a lost acknowledgement: replay must discover the original seal instead of refetching or publishing new bytes.Methods
async def get(self, *, account_id: str, source_execution_key: str) ‑> SealedSlice | None-
Expand source code
async def get(self, *, account_id: str, source_execution_key: str) -> SealedSlice | None: await self._failures.hit("seal.get.before") result = await self._store.get( account_id=account_id, source_execution_key=source_execution_key ) await self._failures.hit("seal.get.after") return result async def put(self, *, account_id: str, source_execution_key: str, sealed: SealedSlice) ‑> SealedSlice-
Expand source code
async def put( self, *, account_id: str, source_execution_key: str, sealed: SealedSlice ) -> SealedSlice: await self._failures.hit("seal.before") result = await self._store.put( account_id=account_id, source_execution_key=source_execution_key, sealed=sealed ) await self._failures.hit("seal.after") return result
class FaultInjectingReportingStagingStore (store: ReportingStagingStore,
failures: ReportingFailurePlan)-
Expand source code
class FaultInjectingReportingStagingStore: """Borrow a memory or durable store with ``stage`` / ``read`` fault points. Each operation visits ``<operation>.before`` and ``<operation>.after``. A fault after staging retains the real store's write. Schema creation and resource ownership stay with the caller; this wrapper copies no state. """ def __init__(self, store: ReportingStagingStore, failures: ReportingFailurePlan) -> None: self._store = store self._failures = failures async def stage( self, *, account_id: str, source_execution_key: str, ordinal: int, payload: bytes ) -> tuple[str, str]: await self._failures.hit("stage.before") result = await self._store.stage( account_id=account_id, source_execution_key=source_execution_key, ordinal=ordinal, payload=payload, ) await self._failures.hit("stage.after") return result async def read( self, *, object_ref: str, object_generation: str, account_id: str, source_scope: Mapping[str, Any], cancel: asyncio.Event, ) -> bytes: await self._failures.hit("read.before") result = await self._store.read( object_ref=object_ref, object_generation=object_generation, account_id=account_id, source_scope=source_scope, cancel=cancel, ) await self._failures.hit("read.after") return resultBorrow a memory or durable store with
stage/readfault points.Each operation visits
<operation>.beforeand<operation>.after. A fault after staging retains the real store's write. Schema creation and resource ownership stay with the caller; this wrapper copies no state.Methods
async def read(self,
*,
object_ref: str,
object_generation: str,
account_id: str,
source_scope: Mapping[str, Any],
cancel: asyncio.Event) ‑> bytes-
Expand source code
async def read( self, *, object_ref: str, object_generation: str, account_id: str, source_scope: Mapping[str, Any], cancel: asyncio.Event, ) -> bytes: await self._failures.hit("read.before") result = await self._store.read( object_ref=object_ref, object_generation=object_generation, account_id=account_id, source_scope=source_scope, cancel=cancel, ) await self._failures.hit("read.after") return result async def stage(self, *, account_id: str, source_execution_key: str, ordinal: int, payload: bytes) ‑> tuple[str, str]-
Expand source code
async def stage( self, *, account_id: str, source_execution_key: str, ordinal: int, payload: bytes ) -> tuple[str, str]: await self._failures.hit("stage.before") result = await self._store.stage( account_id=account_id, source_execution_key=source_execution_key, ordinal=ordinal, payload=payload, ) await self._failures.hit("stage.after") return result
class ReportingFailurePlan-
Expand source code
class ReportingFailurePlan: """FIFO faults and barriers at named operation boundaries. ``at("seal.after", OSError("lost acknowledgement"))`` fails once *after* the seal store commits. ``None`` consumes one successful visit. Unplanned visits succeed. Use :meth:`assert_consumed` to catch a misspelled or unreached boundary. Adopter wrappers can use ``hit`` / ``hit_sync`` for additional destination, notification, and receipt boundaries. Queue access is thread-safe, but tests that need a particular ordering of concurrent calls must impose it with barriers. Errors are test inputs and must not contain real credentials or provider responses. """ def __init__(self) -> None: self._steps: dict[str, deque[BaseException | ReportingTestBarrier | None]] = defaultdict( deque ) self._hits: list[str] = [] self._lock = threading.Lock() def at(self, point: str, *steps: BaseException | ReportingTestBarrier | None) -> None: if not point or not steps: raise ValueError("a failure point needs a name and at least one step") if any( step is not None and not isinstance(step, (BaseException, ReportingTestBarrier)) for step in steps ): raise TypeError("failure steps must be exceptions, barriers, or None") with self._lock: self._steps[point].extend(steps) @property def hits(self) -> tuple[str, ...]: with self._lock: return tuple(self._hits) def _take(self, point: str) -> BaseException | ReportingTestBarrier | None: with self._lock: self._hits.append(point) steps = self._steps.get(point) return steps.popleft() if steps else None async def hit(self, point: str) -> None: step = self._take(point) if isinstance(step, BaseException): raise step if isinstance(step, ReportingTestBarrier): await step.pause() def hit_sync(self, point: str) -> None: step = self._take(point) if isinstance(step, BaseException): raise step if isinstance(step, ReportingTestBarrier): step.pause_sync() def assert_consumed(self) -> None: with self._lock: remaining = {point: len(steps) for point, steps in self._steps.items() if steps} if remaining: raise AssertionError(f"unreached reporting failure steps: {remaining}")FIFO faults and barriers at named operation boundaries.
at("seal.after", OSError("lost acknowledgement"))fails once after the seal store commits.Noneconsumes one successful visit. Unplanned visits succeed. Use :meth:assert_consumedto catch a misspelled or unreached boundary. Adopter wrappers can usehit/hit_syncfor additional destination, notification, and receipt boundaries.Queue access is thread-safe, but tests that need a particular ordering of concurrent calls must impose it with barriers. Errors are test inputs and must not contain real credentials or provider responses.
Instance variables
prop hits : tuple[str, ...]-
Expand source code
@property def hits(self) -> tuple[str, ...]: with self._lock: return tuple(self._hits)
Methods
def assert_consumed(self) ‑> None-
Expand source code
def assert_consumed(self) -> None: with self._lock: remaining = {point: len(steps) for point, steps in self._steps.items() if steps} if remaining: raise AssertionError(f"unreached reporting failure steps: {remaining}") def at(self,
point: str,
*steps: BaseException | ReportingTestBarrier | None) ‑> None-
Expand source code
def at(self, point: str, *steps: BaseException | ReportingTestBarrier | None) -> None: if not point or not steps: raise ValueError("a failure point needs a name and at least one step") if any( step is not None and not isinstance(step, (BaseException, ReportingTestBarrier)) for step in steps ): raise TypeError("failure steps must be exceptions, barriers, or None") with self._lock: self._steps[point].extend(steps) async def hit(self, point: str) ‑> None-
Expand source code
async def hit(self, point: str) -> None: step = self._take(point) if isinstance(step, BaseException): raise step if isinstance(step, ReportingTestBarrier): await step.pause() def hit_sync(self, point: str) ‑> None-
Expand source code
def hit_sync(self, point: str) -> None: step = self._take(point) if isinstance(step, BaseException): raise step if isinstance(step, ReportingTestBarrier): step.pause_sync()
class ReportingLifecycleSnapshot (obligation: ReportingObligationRecord,
revisions: tuple[ReportingRevisionRecord, ...],
rows: tuple[bytes, ...])-
Expand source code
@dataclass(frozen=True) class ReportingLifecycleSnapshot: """An obligation, its exact revisions, and detached canonical row bytes. Capture a quiescent obligation after publication, then compare it after closing/rebuilding the service and retrying the same work. The assertion checks every revision, so an extra publication is a failure too. This is evidence about the supplied store, not a durability or tier certification. """ obligation: ReportingObligationRecord revisions: tuple[ReportingRevisionRecord, ...] rows: tuple[bytes, ...]An obligation, its exact revisions, and detached canonical row bytes.
Capture a quiescent obligation after publication, then compare it after closing/rebuilding the service and retrying the same work. The assertion checks every revision, so an extra publication is a failure too. This is evidence about the supplied store, not a durability or tier certification.
Instance variables
var obligation : ReportingObligationRecordvar revisions : tuple[ReportingRevisionRecord, ...]var rows : tuple[bytes, ...]
class ReportingTestBarrier (*, timeout: float = 10)-
Expand source code
class ReportingTestBarrier: """Pause an async operation or a synchronous adapter running in a worker. Construct inside the test's event loop, await :meth:`wait` to observe the boundary, and always call :meth:`release` in test cleanup. The timeout is a real-time deadlock watchdog; it does not advance the reporting clock. """ def __init__(self, *, timeout: float = 10) -> None: if not math.isfinite(timeout) or timeout <= 0: raise ValueError("a reporting test barrier needs a finite positive timeout") self._loop = asyncio.get_running_loop() self._loop_thread = threading.get_ident() self._timeout = timeout self._entered = asyncio.Event() self._released = asyncio.Event() self._thread_released = threading.Event() async def wait(self) -> None: """Wait until the operation reaches this boundary, without polling.""" await asyncio.wait_for(self._entered.wait(), self._timeout) async def pause(self) -> None: self._entered.set() await asyncio.wait_for(self._released.wait(), self._timeout) def pause_sync(self) -> None: """Pause a worker thread; never block the event-loop thread.""" if threading.get_ident() == self._loop_thread: raise RuntimeError("pause_sync must run outside the barrier's event-loop thread") self._loop.call_soon_threadsafe(self._entered.set) if not self._thread_released.wait(self._timeout): raise TimeoutError("reporting test barrier was not released") def release(self) -> None: """Release all waiters, including a synchronous in-flight fetch.""" self._thread_released.set() self._loop.call_soon_threadsafe(self._released.set)Pause an async operation or a synchronous adapter running in a worker.
Construct inside the test's event loop, await :meth:
waitto observe the boundary, and always call :meth:releasein test cleanup. The timeout is a real-time deadlock watchdog; it does not advance the reporting clock.Methods
async def pause(self) ‑> None-
Expand source code
async def pause(self) -> None: self._entered.set() await asyncio.wait_for(self._released.wait(), self._timeout) def pause_sync(self) ‑> None-
Expand source code
def pause_sync(self) -> None: """Pause a worker thread; never block the event-loop thread.""" if threading.get_ident() == self._loop_thread: raise RuntimeError("pause_sync must run outside the barrier's event-loop thread") self._loop.call_soon_threadsafe(self._entered.set) if not self._thread_released.wait(self._timeout): raise TimeoutError("reporting test barrier was not released")Pause a worker thread; never block the event-loop thread.
def release(self) ‑> None-
Expand source code
def release(self) -> None: """Release all waiters, including a synchronous in-flight fetch.""" self._thread_released.set() self._loop.call_soon_threadsafe(self._released.set)Release all waiters, including a synchronous in-flight fetch.
async def wait(self) ‑> None-
Expand source code
async def wait(self) -> None: """Wait until the operation reaches this boundary, without polling.""" await asyncio.wait_for(self._entered.wait(), self._timeout)Wait until the operation reaches this boundary, without polling.
class ScriptedReportingAdapter (capabilities: ReportingSourceCapabilitiesV1,
results: Sequence[InlineFetchResult | Sequence[Mapping[str, Any]] | None | BaseException] = (),
*,
asynchronous: bool = False,
failures: ReportingFailurePlan | None = None)-
Expand source code
class ScriptedReportingAdapter: """Queue source answers or exceptions for deterministic service tests. Set ``asynchronous=True`` to exercise the async provider-SDK path. The default synchronous path is run in a worker thread by :class:`InlineReportingSource`, just like a synchronous production SDK. With ``failures``, each call visits ``fetch.before`` before consuming its queued answer and ``fetch.after`` after a successful answer. These points accept either an exception or a :class:`ReportingTestBarrier`. """ def __init__( self, capabilities: ReportingSourceCapabilitiesV1, results: Sequence[ InlineFetchResult | Sequence[Mapping[str, Any]] | None | BaseException ] = (), *, asynchronous: bool = False, failures: ReportingFailurePlan | None = None, ) -> None: self._capabilities = capabilities self._results = deque(results) self._lock = threading.Lock() self._failures = failures self.asynchronous = asynchronous self.calls: list[ReportingSourceSliceRequestV1] = [] @property def capabilities(self) -> ReportingSourceCapabilitiesV1: return self._capabilities def queue( self, result: InlineFetchResult | Sequence[Mapping[str, Any]] | None | BaseException, ) -> None: with self._lock: self._results.append(result) def _next(self) -> Any: with self._lock: if not self._results: raise AssertionError("ScriptedReportingAdapter has no queued source answer") result = self._results.popleft() if isinstance(result, BaseException): raise result return result def fetch_slice(self, request: ReportingSourceSliceRequestV1) -> Any: if not self.asynchronous: with self._lock: self.calls.append(request) if self._failures is not None: self._failures.hit_sync("fetch.before") result = self._next() if self._failures is not None: self._failures.hit_sync("fetch.after") return result async def fetch_async() -> Any: with self._lock: self.calls.append(request) if self._failures is not None: await self._failures.hit("fetch.before") result = self._next() if self._failures is not None: await self._failures.hit("fetch.after") return result return fetch_async()Queue source answers or exceptions for deterministic service tests.
Set
asynchronous=Trueto exercise the async provider-SDK path. The default synchronous path is run in a worker thread by :class:InlineReportingSource, just like a synchronous production SDK. Withfailures, each call visitsfetch.beforebefore consuming its queued answer andfetch.afterafter a successful answer. These points accept either an exception or a :class:ReportingTestBarrier.Instance variables
prop capabilities : ReportingSourceCapabilitiesV1-
Expand source code
@property def capabilities(self) -> ReportingSourceCapabilitiesV1: return self._capabilities
Methods
def fetch_slice(self, request: ReportingSourceSliceRequestV1) ‑> Any-
Expand source code
def fetch_slice(self, request: ReportingSourceSliceRequestV1) -> Any: if not self.asynchronous: with self._lock: self.calls.append(request) if self._failures is not None: self._failures.hit_sync("fetch.before") result = self._next() if self._failures is not None: self._failures.hit_sync("fetch.after") return result async def fetch_async() -> Any: with self._lock: self.calls.append(request) if self._failures is not None: await self._failures.hit("fetch.before") result = self._next() if self._failures is not None: await self._failures.hit("fetch.after") return result return fetch_async() def queue(self,
result: InlineFetchResult | Sequence[Mapping[str, Any]] | None | BaseException) ‑> None-
Expand source code
def queue( self, result: InlineFetchResult | Sequence[Mapping[str, Any]] | None | BaseException, ) -> None: with self._lock: self._results.append(result)