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 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.

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 result

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.

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 result

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.

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. 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.

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 : ReportingObligationRecord
var 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: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.

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=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.

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)