Module adcp.reporting.inline_source

Wrap an existing delivery fetch into a conforming reporting source executor.

Most sellers adopting Reliable Reporting already have the numbers. They have a get_media_buy_delivery handler, or a nightly stats cache, or a report-API client someone wrote years ago. What they do not have is an immutable, replayable, evidence-bearing publication – and hand-rolling one means getting eight interlocking fingerprints, a zero-row taxonomy, and a replay seal right on the first try.

:class:InlineReportingSource is that adapter. Give it a callable that answers "what did this scope deliver in this window?" and it produces a conforming basic manifest, complete with staged objects, coverage evidence, and byte-identical replay.

For uneven metric support, return :class:InlineFetchResult with cell_availability={constituent_id: {metric_name: MetricEvidence(...)}}. Omitted cells retain the constituent defaults. Explicit evidence controls each cell independently. Unsupported cells are inapplicable: alongside available cells they do not make the constituent partial. A cell that withdraws a metric its own constituent's rows report withdraws that metric's control total, so no total ever sums a disclaimed value; a cell that staged no rows reaches no sum and leaves its neighbours' subtotal intact. Money has no such choice – see monetary_metrics on :class:InlineReportingSource. Metric semantics always come from the selected SDK offering.

The Three Answers A Fetch Can Give

None Not ready. The source does not have this period yet. This is not a zero and it is not a failure; it is "ask again later". For a full-coverage request that means a retryable PARTIAL_RESULT and no publication. For a partial-coverage request it means a completed publication whose cells are all delayed – which is useful, because the producer's health projection can then distinguish "late" from "dead".

[] (an empty row sequence) A real zero-row period. The source answered for the whole requested scope and the answer was nothing. That is data: it commits an explicit_zero revision and satisfies the obligation. Conflating this with None is the single most common Reliable Reporting bug, because it makes a dead feed look like a quiet campaign forever.

an exception A failure. Raise :class:~adcp.reporting.source.ReportingSourceError to classify it yourself. Anything else is treated as an unclassified transient fault (see :attr:InlineReportingSource.include_exception_detail for why the message is redacted by default).

Sync Callables

A synchronous fetch runs in a worker thread. Python cannot interrupt a running thread, so a sync fetch is not cancellable: on cancellation this executor waits for the thread rather than returning while work continues behind it, because a detached fetch would race the producer's next attempt for the same staged object names. Give your sync callable its own timeout.

Global variables

var InlineFetch : TypeAlias

The callable shape. Sync or async; a sync callable runs in a worker thread.

Functions

def metric_delayed_through(request: ReportingSourceSliceRequestV1,
metric_name: str,
watermark: datetime,
*,
reason: str = 'The reporting source has not finalized this metric yet') ‑> dict[str, dict[str, MetricEvidence]]
Expand source code
def metric_delayed_through(
    request: ReportingSourceSliceRequestV1,
    metric_name: str,
    watermark: datetime,
    *,
    reason: str = "The reporting source has not finalized this metric yet",
) -> dict[str, dict[str, MetricEvidence]]:
    """Build ``cell_availability`` retaining a delayed metric's known watermark."""
    return _metric_everywhere(
        request, metric_name, MetricEvidence.delayed(reason, data_through=watermark)
    )

Build cell_availability retaining a delayed metric's known watermark.

def metric_unsupported_everywhere(request: ReportingSourceSliceRequestV1, metric_name: str, reason: str) ‑> dict[str, dict[str, MetricEvidence]]
Expand source code
def metric_unsupported_everywhere(
    request: ReportingSourceSliceRequestV1, metric_name: str, reason: str
) -> dict[str, dict[str, MetricEvidence]]:
    """Build ``cell_availability`` for a metric unsupported across the request.

    Other metrics retain the fetch's constituent defaults. Pass the returned
    map to :class:`InlineFetchResult`, alongside rows, watermarks, or warnings.
    """
    return _metric_everywhere(request, metric_name, MetricEvidence.unavailable(reason))

Build cell_availability for a metric unsupported across the request.

Other metrics retain the fetch's constituent defaults. Pass the returned map to :class:InlineFetchResult, alongside rows, watermarks, or warnings.

Classes

class FileSystemStagingStore (root: Path | str)
Expand source code
class FileSystemStagingStore:
    """Account-scoped content-addressed staging on a local filesystem.

    Equal payload bytes in an account share their ref and SHA-256 generation
    across acquisitions. Legacy opaque refs remain readable at their existing
    paths. Observation identity belongs to the manifest/seal, not the payload.

    The filesystem must support atomic same-directory hard links and directory
    fsync. A flushed temporary file is linked without replacing a winner, then
    its temporary name is removed and the directory chain is synced. Readers
    never see an SDK writer's partial file. Existing bytes are verified before
    reuse; corruption fails closed rather than being overwritten. A failed
    publication may retain complete bytes, but never returns a successful pair.
    """

    def __init__(self, root: Path | str) -> None:
        self._root = Path(root)
        self._root_sync_lock = threading.Lock()
        self._root_synced = False

    def _path(self, account_id: str, object_ref: str, object_generation: str) -> Path:
        # The account is part of the path so an authorization bug cannot become
        # a cross-tenant read by way of a guessed object name.
        account = hashlib.sha256(account_id.encode("utf-8")).hexdigest()[:32]
        safe_ref = hashlib.sha256(object_ref.encode("utf-8")).hexdigest()[:32]
        return self._root / account / safe_ref / f"{object_generation}.bin"

    async def stage(
        self, *, account_id: str, source_execution_key: str, ordinal: int, payload: bytes
    ) -> tuple[str, str]:
        if not isinstance(payload, bytes):
            raise TypeError("staging payload must be immutable bytes")
        digest = hashlib.sha256(payload).hexdigest()
        account = hashlib.sha256(account_id.encode("utf-8")).hexdigest()
        object_ref = f"sha256.{account}.{digest}"
        target = self._path(account_id, object_ref, digest)
        await settle_task(asyncio.create_task(asyncio.to_thread(self._write, target, payload)))
        return object_ref, digest

    def _write(self, target: Path, payload: bytes) -> None:
        self._ensure_root_durable()
        try:
            self._verify_existing(target, payload)
        except FileNotFoundError:
            pass
        else:
            # Another writer may have linked the complete file but not yet
            # synced its name/ancestors. A reuse must complete that durability
            # barrier itself before returning a successful pair.
            self._sync_directory(target.parent)
            return
        target.parent.mkdir(parents=True, exist_ok=True)
        handle, temporary = tempfile.mkstemp(dir=str(target.parent))
        try:
            with os.fdopen(handle, "wb") as stream:
                stream.write(payload)
                stream.flush()
                os.fsync(stream.fileno())
            try:
                self._publish_file(Path(temporary), target)
            except FileExistsError:
                self._verify_existing(target, payload)
        finally:
            Path(temporary).unlink(missing_ok=True)
        # Sync after unlink: both the final link and removal of our temporary
        # link are committed together. Sync ancestors too, because this turn
        # (or a concurrent publisher) may have just created those directories.
        # If any sync fails, the complete file is retained for verified retry,
        # but no successful pair escapes to the source's seal operation.
        self._sync_directory(target.parent)

    def _ensure_root_durable(self) -> None:
        # The root itself may be created on the first stage. Commit its name
        # and any newly created parent names once; subsequent stages only need
        # to sync the directories inside this store's root.
        with self._root_sync_lock:
            if self._root_synced:
                return
            missing: list[Path] = []
            boundary = self._root
            while not boundary.is_dir():
                missing.append(boundary)
                parent = boundary.parent
                if parent == boundary:
                    raise FileNotFoundError("no existing parent for staging root")
                boundary = parent
            self._root.mkdir(parents=True, exist_ok=True)
            # A created directory's name is committed by syncing its parent.
            # The first pre-existing ancestor is sufficient; ancestors above
            # it may be traversable but unreadable to the staging process.
            for path in (*missing, boundary) if missing else (self._root,):
                descriptor = os.open(path, os.O_RDONLY | getattr(os, "O_DIRECTORY", 0))
                try:
                    os.fsync(descriptor)
                finally:
                    os.close(descriptor)
            if not missing and self._root.parent != self._root:
                # A fresh instance can repair a root another writer created
                # just before crashing. A pre-existing root is still usable
                # when its parent allows traversal but not directory reads.
                try:
                    descriptor = os.open(
                        self._root.parent, os.O_RDONLY | getattr(os, "O_DIRECTORY", 0)
                    )
                except PermissionError:
                    pass
                else:
                    try:
                        os.fsync(descriptor)
                    finally:
                        os.close(descriptor)
            self._root_synced = True

    @staticmethod
    def _verify_existing(target: Path, payload: bytes) -> None:
        with target.open("rb") as stream:
            if stream.read() != payload:
                raise OSError("staged object bytes no longer match their pinned generation")
            os.fsync(stream.fileno())

    @staticmethod
    def _publish_file(temporary: Path, target: Path) -> None:
        os.link(temporary, target)

    def _sync_directory(self, directory: Path) -> None:
        path = directory
        while True:
            descriptor = os.open(path, os.O_RDONLY | getattr(os, "O_DIRECTORY", 0))
            try:
                os.fsync(descriptor)
            finally:
                os.close(descriptor)
            if path == self._root:
                break
            path = path.parent

    async def read(
        self,
        *,
        object_ref: str,
        object_generation: str,
        account_id: str,
        source_scope: Mapping[str, Any],
        cancel: asyncio.Event,
    ) -> bytes:
        target = self._path(account_id, object_ref, object_generation)
        payload: bytes = await settle_task(
            asyncio.create_task(asyncio.to_thread(target.read_bytes))
        )
        if hashlib.sha256(payload).hexdigest() != object_generation:
            raise OSError("staged object bytes no longer match their pinned generation")
        return payload

Account-scoped content-addressed staging on a local filesystem.

Equal payload bytes in an account share their ref and SHA-256 generation across acquisitions. Legacy opaque refs remain readable at their existing paths. Observation identity belongs to the manifest/seal, not the payload.

The filesystem must support atomic same-directory hard links and directory fsync. A flushed temporary file is linked without replacing a winner, then its temporary name is removed and the directory chain is synced. Readers never see an SDK writer's partial file. Existing bytes are verified before reuse; corruption fails closed rather than being overwritten. A failed publication may retain complete bytes, but never returns a successful pair.

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:
    target = self._path(account_id, object_ref, object_generation)
    payload: bytes = await settle_task(
        asyncio.create_task(asyncio.to_thread(target.read_bytes))
    )
    if hashlib.sha256(payload).hexdigest() != object_generation:
        raise OSError("staged object bytes no longer match their pinned generation")
    return payload
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]:
    if not isinstance(payload, bytes):
        raise TypeError("staging payload must be immutable bytes")
    digest = hashlib.sha256(payload).hexdigest()
    account = hashlib.sha256(account_id.encode("utf-8")).hexdigest()
    object_ref = f"sha256.{account}.{digest}"
    target = self._path(account_id, object_ref, digest)
    await settle_task(asyncio.create_task(asyncio.to_thread(self._write, target, payload)))
    return object_ref, digest
class InMemorySealStore
Expand source code
class InMemorySealStore:
    """Process-local replay seals.  Durable stores belong in your database."""

    def __init__(self) -> None:
        self._seals: dict[tuple[str, str], SealedSlice] = {}

    async def get(self, *, account_id: str, source_execution_key: str) -> SealedSlice | None:
        return self._seals.get((account_id, source_execution_key))

    async def put(
        self, *, account_id: str, source_execution_key: str, sealed: SealedSlice
    ) -> SealedSlice:
        return self._seals.setdefault((account_id, source_execution_key), sealed)

Process-local replay seals. Durable stores belong in your database.

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:
    return self._seals.get((account_id, source_execution_key))
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:
    return self._seals.setdefault((account_id, source_execution_key), sealed)
class InMemoryStagingStore
Expand source code
class InMemoryStagingStore:
    """Process-local, account-scoped payload reuse for tests and pilots.

    Not durable: a restart loses payloads needed for uncommitted acquisitions
    or replay of retained source manifests. Committed ledger rows use the
    ledger's independent exact-revision read path. Use
    :class:`FileSystemStagingStore` or your own durable staging store when
    source payloads must survive a restart.
    """

    def __init__(self) -> None:
        self._objects: dict[tuple[str, str, str], bytes] = {}

    async def stage(
        self, *, account_id: str, source_execution_key: str, ordinal: int, payload: bytes
    ) -> tuple[str, str]:
        if not isinstance(payload, bytes):
            raise TypeError("staging payload must be immutable bytes")
        digest = hashlib.sha256(payload).hexdigest()
        account = hashlib.sha256(account_id.encode("utf-8")).hexdigest()
        object_ref = f"sha256.{account}.{digest}"
        key = (account_id, object_ref, digest)
        if key in self._objects:
            if self._objects[key] != payload:
                raise OSError("staged object bytes no longer match their pinned generation")
        else:
            self._objects[key] = payload
        return object_ref, digest

    async def read(
        self,
        *,
        object_ref: str,
        object_generation: str,
        account_id: str,
        source_scope: Mapping[str, Any],
        cancel: asyncio.Event,
    ) -> bytes:
        key = (account_id, object_ref, object_generation)
        if key not in self._objects:
            raise PermissionError("staged object is outside the requested account scope")
        payload = self._objects[key]
        if hashlib.sha256(payload).hexdigest() != object_generation:
            raise OSError("staged object bytes no longer match their pinned generation")
        return payload

Process-local, account-scoped payload reuse for tests and pilots.

Not durable: a restart loses payloads needed for uncommitted acquisitions or replay of retained source manifests. Committed ledger rows use the ledger's independent exact-revision read path. Use :class:FileSystemStagingStore or your own durable staging store when source payloads must survive a restart.

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:
    key = (account_id, object_ref, object_generation)
    if key not in self._objects:
        raise PermissionError("staged object is outside the requested account scope")
    payload = self._objects[key]
    if hashlib.sha256(payload).hexdigest() != object_generation:
        raise OSError("staged object bytes no longer match their pinned generation")
    return payload
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]:
    if not isinstance(payload, bytes):
        raise TypeError("staging payload must be immutable bytes")
    digest = hashlib.sha256(payload).hexdigest()
    account = hashlib.sha256(account_id.encode("utf-8")).hexdigest()
    object_ref = f"sha256.{account}.{digest}"
    key = (account_id, object_ref, digest)
    if key in self._objects:
        if self._objects[key] != payload:
            raise OSError("staged object bytes no longer match their pinned generation")
    else:
        self._objects[key] = payload
    return object_ref, digest
class InlineFetchResult (rows: Sequence[Mapping[str, Any]],
data_through: datetime | None = None,
covered_constituent_ids: Collection[str] | None = None,
unavailable_constituents: Mapping[str, str] = <factory>,
unavailable_status: ReportingAvailabilityStatus = 'unsupported',
warnings: Sequence[str] = (),
*,
currency: str | None = None,
provisional_until: datetime | None = None,
cell_availability: Mapping[str, Mapping[str, MetricEvidence]] | None = None)
Expand source code
@dataclass(frozen=True)
class InlineFetchResult:
    """A fetch answer richer than a bare row list.

    Return this when the honest answer is nuanced: the cache covers some
    constituents but not others, the numbers are only good through a
    freshness watermark earlier than the period end, or a particular scope is
    known-unsupported rather than known-empty.
    """

    rows: Sequence[Mapping[str, Any]]
    """Normalized rows.  Empty means an observed zero for every covered scope."""

    data_through: datetime | None = None
    """How far into the period these numbers are actually good.

    This is the freshness gate.  An ad server that finalizes yesterday at
    07:00 local is not telling you about the last four hours at 05:00, and
    publishing those hours as complete is how an official close ends up
    understating delivery.  Defaults to ``min(period end, source read
    cutoff)``; supply the real watermark when you know it.
    """

    covered_constituent_ids: Collection[str] | None = None
    """Which constituents this answer actually covers.

    ``None`` (the default) means "all of them" -- an answer for the whole
    requested scope, so a constituent with no rows is an observed zero.  Name
    a subset when your cache genuinely only holds part of the denominator;
    everything outside it is reported ``missing`` rather than zero.
    """

    unavailable_constituents: Mapping[str, str] = field(default_factory=dict)
    """``constituent_id`` to a stable reason, for scopes this source cannot serve.

    Use for known-unsupported scopes (a media buy on a different ad server, a
    package the report definition does not model).  Distinct from "covered but
    empty", which is a zero.
    """

    unavailable_status: ReportingAvailabilityStatus = "unsupported"
    """The status to attach to :attr:`unavailable_constituents`."""

    warnings: Sequence[str] = ()
    """Safe, redacted operator notes retained on the publication."""

    currency: str | None = field(default=None, kw_only=True)
    """Optional source corroboration; must match the already frozen request."""

    provisional_until: datetime | None = field(default=None, kw_only=True)
    """Source evidence that this snapshot may still change until this instant.

    When present it overrides the offering's default restatement window for
    this slice. It is retained in the manifest so a durable producer can keep
    scheduling correctly after a restart.
    """

    cell_availability: Mapping[str, Mapping[str, MetricEvidence]] | None = field(
        default=None, kw_only=True
    )
    """Optional ``{constituent_id: {metric_name: evidence}}`` overrides.

    Keys are the frozen request's constituent IDs, which need not equal media
    buy IDs. Omitted cells retain the derived constituent status and watermark.
    Explicit cells override even missing/unsupported constituent defaults;
    unsupported cells do not demote otherwise available coverage. With no
    applicable cells, the constituent stays ``unsupported`` and cannot satisfy
    full coverage. Other mixed availability promotes it to ``partial``. The
    SDK still applies the source cutoff, observation ceiling, and authoritative
    freshness gate. A cell watermark may advance the batch watermark without
    advancing other cells. This is an adapter surface, not the manifest's
    wire-format ``metric_availability`` list.

    Unknown keys, duplicate mapping entries, invalid evidence, contradictory
    explicit zeros, and a withdrawn *monetary* cell whose own constituent's rows
    report that metric raise ``ValueError`` before staging. A zero-row batch must
    be wholly explicit-zero or wholly unavailable under the existing contract.
    """

    @classmethod
    def all_present(
        cls, rows: Sequence[Mapping[str, Any]], *, data_through: datetime
    ) -> InlineFetchResult:
        """Cover the whole request through a watermark, using the legacy defaults.

        Constituents with rows are present; covered constituents without rows
        are observed zeros. Use ``cell_availability`` for exceptions.
        """
        return cls(rows=rows, data_through=data_through)

A fetch answer richer than a bare row list.

Return this when the honest answer is nuanced: the cache covers some constituents but not others, the numbers are only good through a freshness watermark earlier than the period end, or a particular scope is known-unsupported rather than known-empty.

Static methods

def all_present(rows: Sequence[Mapping[str, Any]], *, data_through: datetime) ‑> InlineFetchResult

Cover the whole request through a watermark, using the legacy defaults.

Constituents with rows are present; covered constituents without rows are observed zeros. Use cell_availability for exceptions.

Instance variables

var cell_availability : collections.abc.Mapping[str, collections.abc.Mapping[str, MetricEvidence]] | None

Optional {constituent_id: {metric_name: evidence}} overrides.

Keys are the frozen request's constituent IDs, which need not equal media buy IDs. Omitted cells retain the derived constituent status and watermark. Explicit cells override even missing/unsupported constituent defaults; unsupported cells do not demote otherwise available coverage. With no applicable cells, the constituent stays unsupported and cannot satisfy full coverage. Other mixed availability promotes it to partial. The SDK still applies the source cutoff, observation ceiling, and authoritative freshness gate. A cell watermark may advance the batch watermark without advancing other cells. This is an adapter surface, not the manifest's wire-format metric_availability list.

Unknown keys, duplicate mapping entries, invalid evidence, contradictory explicit zeros, and a withdrawn monetary cell whose own constituent's rows report that metric raise ValueError before staging. A zero-row batch must be wholly explicit-zero or wholly unavailable under the existing contract.

var covered_constituent_ids : collections.abc.Collection[str] | None

Which constituents this answer actually covers.

None (the default) means "all of them" – an answer for the whole requested scope, so a constituent with no rows is an observed zero. Name a subset when your cache genuinely only holds part of the denominator; everything outside it is reported missing rather than zero.

var currency : str | None

Optional source corroboration; must match the already frozen request.

var data_through : datetime.datetime | None

How far into the period these numbers are actually good.

This is the freshness gate. An ad server that finalizes yesterday at 07:00 local is not telling you about the last four hours at 05:00, and publishing those hours as complete is how an official close ends up understating delivery. Defaults to min(period end, source read cutoff); supply the real watermark when you know it.

var provisional_until : datetime.datetime | None

Source evidence that this snapshot may still change until this instant.

When present it overrides the offering's default restatement window for this slice. It is retained in the manifest so a durable producer can keep scheduling correctly after a restart.

var rows : Sequence[Mapping[str, typing.Any]]

Normalized rows. Empty means an observed zero for every covered scope.

var unavailable_constituents : Mapping[str, str]

constituent_id to a stable reason, for scopes this source cannot serve.

Use for known-unsupported scopes (a media buy on a different ad server, a package the report definition does not model). Distinct from "covered but empty", which is a zero.

var unavailable_status : Literal['present', 'explicit_zero', 'unsupported', 'delayed', 'partial', 'stale', 'missing']

The status to attach to :attr:unavailable_constituents.

var warnings : Sequence[str]

Safe, redacted operator notes retained on the publication.

class InlineReportingSource (*,
capabilities: ReportingSourceCapabilitiesV1,
fetch: InlineFetch,
staging: ReportingStagingStore | None = None,
seals: ReportingSealStore | None = None,
constituent_of: Callable[[Mapping[str, Any], ReportingSourceSliceRequestV1], str | None] | None = None,
staged_commit_prefix: str = 'inline',
include_exception_detail: bool = False,
monetary_metrics: Collection[str] = (),
clock: Callable[[], datetime] | None = None)
Expand source code
class InlineReportingSource:
    """A conforming executor over an ordinary delivery fetch.

    Parameters
    ----------
    capabilities:
        Your truthful, scoped capability declaration.  Every offering it
        declares must use ``manifest_strictness='basic'`` -- this adapter
        cannot honestly produce per-call provider evidence for a fetch it did
        not make itself.
    fetch:
        Sync or async callable taking the frozen
        :class:`~adcp.reporting.source.ReportingSourceSliceRequestV1` and
        returning rows, an :class:`InlineFetchResult`, an AdCP
        ``GetMediaBuyDeliveryResponse``, or ``None`` for "not ready".
    staging / seals:
        Where rows and replay seals live.  The in-memory defaults are for tests
        and single-process pilots; they are not durable.
    constituent_of:
        Maps a row to the ``constituent_id`` it belongs to.  The default
        matches a row's ``media_buy_id`` / ``package_id`` / ``product_id``
        against the requested denominator, which covers the common shapes.
        Rows that match nothing are retained in the staged object but excluded
        from coverage, and a warning names how many were dropped.
    monetary_metrics:
        The metric names the trusted report definition froze as money, beyond
        the built-in ``spend``.  Supply
        ``ReportingDefinitionBinding.monetary_metric_units`` keys (and any
        ``monetary_control_total_units`` keys) when the obligation declares
        them: the frozen slice request cannot carry those declarations, because
        ``to_wire()`` keeps them off the wire so retained contract hashes do
        not move, so the adapter would otherwise not know which columns the
        ledger will reconcile.  A custom money column that is not declared here
        behaves like a non-monetary one.
    include_exception_detail:
        Whether an unclassified exception's ``str()`` may enter the retained
        safe message.  **Off by default**: HTTP client exceptions routinely
        stringify the full request URL, and a signed URL or an ``access_token``
        query parameter in a manifest is a credential leak into permanent,
        replicated evidence.  The exception *type* is always recorded.
    clock:
        Supplies the observation instant.  Override it to make a test's
        publications deterministic; the default is the wall clock, truncated
        to the millisecond precision the manifest binds.
    """

    def __init__(
        self,
        *,
        capabilities: ReportingSourceCapabilitiesV1,
        fetch: InlineFetch,
        staging: ReportingStagingStore | None = None,
        seals: ReportingSealStore | None = None,
        constituent_of: (
            Callable[[Mapping[str, Any], ReportingSourceSliceRequestV1], str | None] | None
        ) = None,
        staged_commit_prefix: str = "inline",
        include_exception_detail: bool = False,
        monetary_metrics: Collection[str] = (),
        clock: Callable[[], datetime] | None = None,
    ) -> None:
        unsupported = [
            offering.offering_id
            for offering in capabilities.offerings
            if offering.manifest_strictness != "basic"
        ]
        if unsupported:
            raise ValueError(
                "InlineReportingSource emits basic manifests only; offerings "
                f"{sorted(unsupported)} declare evidenced manifests, which would require "
                "per-call provider evidence this adapter cannot truthfully produce"
            )
        self._capabilities = capabilities
        self._fetch = fetch
        self._staging = staging or InMemoryStagingStore()
        self._seals = seals or InMemorySealStore()
        self._constituent_of = constituent_of or _default_constituent_of
        self._staged_commit_prefix = staged_commit_prefix
        self.include_exception_detail = include_exception_detail
        # Mirrors validate_monetary_content exactly: ``spend`` is the built-in
        # money column of the single-currency API, plus whatever the trusted
        # definition froze.  The frozen slice request cannot supply these --
        # ReportingDefinitionBinding.to_wire() deliberately keeps the unit
        # declarations off the wire so retained contract hashes do not move.
        self._monetary = frozenset(monetary_metrics) | {_MONEY}
        self._clock = clock or _now

    @property
    def capabilities(self) -> ReportingSourceCapabilitiesV1:
        return self._capabilities

    @property
    def staging(self) -> ReportingStagingStore:
        """The staging store, which also serves as the staged-object reader."""
        return self._staging

    async def execute(
        self,
        request: ReportingSourceSliceRequestV1,
        *,
        cancel: asyncio.Event,
        heartbeat: Callable[[], None] | None = None,
    ) -> ReportingSourceExecutorResult:
        account_id = request.identity.account_id
        key = request.identity.source_execution_key

        sealed = await self._seals.get(account_id=account_id, source_execution_key=key)
        if sealed is not None:
            # Replay, including replay after cancellation: the publication is
            # already immutable, so returning anything else would be a lie.
            return ReportingSourceExecutorResult.completed(
                request=request,
                manifest=sealed.reference,
                manifest_bytes=sealed.manifest_bytes,
            )

        if cancel.is_set():
            return ReportingSourceExecutorResult.failed(
                ReportingSourceErrorV1(
                    code="CANCELLED",
                    retry="cancelled",
                    safe_message="the slice was cancelled before any source read began",
                )
            )

        try:
            answer = await self._invoke(request, heartbeat)
            if answer is not None:
                # Coerce inside the guard: a fetch that returns something
                # unrecognizable is a source failure like any other, not an
                # exception escaping the executor.
                answer = _coerce(answer, currency=request.currency)
                # Retain the rows we checked across staging awaits. A fetch
                # may return a shared cache whose mappings later change.
                answer = replace(answer, rows=tuple(deepcopy(dict(row)) for row in answer.rows))
                if (
                    answer.currency is not None
                    and validate_currency(answer.currency) != request.currency
                ):
                    raise ReportingCurrencyError(
                        "CURRENCY_MISMATCH",
                        "inline source currency disagrees with the frozen request",
                    )
                # Before staging and, crucially, before _control_totals aggregates.
                validate_row_currencies(request.currency, answer.rows)
        except ReportingCurrencyError as error:
            return ReportingSourceExecutorResult.failed(
                ReportingSourceErrorV1(
                    code="INTEGRITY_FAILED", retry="terminal", safe_message=str(error)
                )
            )
        except ReportingSourceError as error:
            return ReportingSourceExecutorResult.failed(error.error)
        except _CellEvidenceError:
            # Constructors/helpers may run inside the fetch. Keep the same
            # actionable ValueError as matrix validation after the fetch,
            # rather than redacting an adapter bug as a retryable provider fault.
            raise
        except asyncio.CancelledError:
            raise
        except Exception as error:
            return ReportingSourceExecutorResult.failed(self._unclassified(error))

        if answer is None:
            not_ready = self._not_ready(request)
            if not_ready is not None:
                return not_ready
            answer = InlineFetchResult(
                rows=[],
                covered_constituent_ids=(),
                unavailable_constituents={
                    item.constituent_id: "Source has not produced this period yet"
                    for item in request.coverage.constituents
                },
                unavailable_status="delayed",
            )

        return await self._publish(request, answer, cancel=cancel)

    # -- fetch ----------------------------------------------------------

    async def _invoke(
        self, request: ReportingSourceSliceRequestV1, heartbeat: Callable[[], None] | None
    ) -> Any:
        """Call the fetch, keeping a synchronous one off the event loop.

        A blocking report API call on the loop thread stalls every other slice
        the producer is running, so a sync callable goes to a worker thread.
        The thread is then awaited to completion even under cancellation:
        Python cannot interrupt it, and returning while it still runs would
        leave it racing the producer's next attempt for the same staged object
        names.
        """
        if heartbeat is not None:
            heartbeat()
        answer: Any
        if _is_async_callable(self._fetch):
            answer = await self._fetch(request)
        else:
            # Keep both the thread and a dynamically returned awaitable owned
            # until they settle. Shield alone abandons the wait on cancellation.
            async def run_sync() -> Any:
                result = await asyncio.to_thread(self._fetch, request)
                return await result if isinstance(result, Awaitable) else result

            answer = await settle_task(asyncio.create_task(run_sync()))
        if isinstance(answer, Awaitable):
            # A plain callable that hands back a coroutine; await it here
            # rather than letting it reach the coercion step unfinished.
            answer = await answer
        return answer  # narrowed by _coerce

    def _unclassified(self, error: Exception) -> ReportingSourceErrorV1:
        """Classify an exception nobody classified for us.

        Retryable, deliberately.  A transient network blip that we call
        terminal permanently fails a period until a human replays it; a
        genuinely broken source that we call retryable still escalates on its
        own when the producer's recovery window elapses.  The second failure
        mode is strictly cheaper.
        """
        detail = f": {error}" if self.include_exception_detail else ""
        return ReportingSourceErrorV1(
            code="PROVIDER_TRANSIENT",
            retry="retryable",
            safe_message=f"the reporting source raised {type(error).__name__}{detail}"[:1_024],
        )

    def _not_ready(
        self, request: ReportingSourceSliceRequestV1
    ) -> ReportingSourceExecutorResult | None:
        if request.coverage.expected != "full":
            return None
        return ReportingSourceExecutorResult.failed(
            ReportingSourceErrorV1(
                code="PARTIAL_RESULT",
                retry="retryable",
                safe_message=(
                    "the reporting source has not produced this period yet; a full-coverage "
                    "slice cannot complete partially"
                ),
            )
        )

    # -- publication ----------------------------------------------------

    async def _publish(
        self,
        request: ReportingSourceSliceRequestV1,
        result: InlineFetchResult,
        *,
        cancel: asyncio.Event,
    ) -> ReportingSourceExecutorResult:
        observed_at = self._clock()
        data_through = _clamp_watermark(request, result.data_through, observed_at)
        overrides = _validate_cell_availability(request, result.cell_availability)
        covered = (
            {item.constituent_id for item in request.coverage.constituents}
            if result.covered_constituent_ids is None
            else set(result.covered_constituent_ids)
        )
        covered -= set(result.unavailable_constituents)

        rows_by_constituent: dict[str, list[Mapping[str, Any]]] = {
            item.constituent_id: [] for item in request.coverage.constituents
        }
        unmatched = 0
        for row in result.rows:
            constituent_id = self._constituent_of(row, request)
            if constituent_id is None or constituent_id not in rows_by_constituent:
                unmatched += 1
                continue
            rows_by_constituent[constituent_id].append(row)
            if constituent_id not in covered and not overrides.get(constituent_id):
                unmatched += 1

        warnings = list(result.warnings)
        if unmatched:
            warnings.append(
                f"{unmatched} source row(s) matched no requested constituent and are "
                "retained in staging but excluded from coverage"
            )

        statuses = self._derive_statuses(
            request, result, covered, rows_by_constituent, data_through
        )
        constituents, cells = self._resolve_availability(
            request,
            result,
            overrides=overrides,
            statuses=statuses,
            rows_by_constituent=rows_by_constituent,
            data_through=data_through,
            observed_at=observed_at,
        )
        # Granular watermarks are bounded by the batch watermark on the wire.
        # Advancing it for an explicit cell must not advance fallback cells.
        data_through = max(data_through, *(cell.data_through or data_through for cell in cells))
        coverage_status = _roll_up([item.status for item in constituents])

        if request.coverage.expected == "full" and coverage_status != "full":
            # Never publish a partial answer to a full-coverage slice; the
            # producer needs a truthful failure so it can retry or escalate.
            return ReportingSourceExecutorResult.failed(
                ReportingSourceErrorV1(
                    code="PARTIAL_RESULT",
                    retry="retryable",
                    safe_message=(
                        "the reporting source covered only part of the requested denominator; "
                        "a full-coverage slice cannot complete partially"
                    ),
                )
            )

        # A partial constituent is deliberately excluded: a known zero for
        # one metric must never turn unavailable neighbours into a batch zero.
        explicit_zero = not result.rows and all(cell.status == "explicit_zero" for cell in cells)
        if (
            not result.rows
            and not explicit_zero
            and any(cell.status in _AVAILABLE for cell in cells)
        ):
            # Reachable from a derived result too, so the message names the
            # shape rather than the optional field an adopter may never have set.
            raise _CellEvidenceError(
                "a zero-row batch cannot mix available and unavailable cells; supply rows "
                "for the observed metrics or withdraw their availability"
            )
        # Only a *declared* exception withdraws a column.  A control total is
        # a checksum over the staged rows -- consumers recompute it from the
        # revision's rows -- so an unmatched row or a legacy-derived
        # unavailable constituent leaves it verifiable and unchanged.
        # A withdrawal only reaches the sum through the rows its own constituent
        # staged.  A cell that staged none cannot put a disclaimed value in the
        # total, so the checksum over the rows that *were* measured is retained
        # -- which is also what keeps a monetary column reconcilable when one
        # constituent delivered nothing and its money is unavailable.  A batch
        # with no rows withdraws either way: a "0" there would publish
        # unavailability as an observed zero.
        declared_incomplete = {
            cell.metric
            for cell in cells
            if cell.status not in _AVAILABLE
            and cell.metric in overrides.get(cell.constituent_id, {})
            and (not result.rows or rows_by_constituent[cell.constituent_id])
        }
        control_totals = _control_totals(request, result.rows, declared_incomplete)
        payload = _encode_rows(result.rows)
        object_ref, object_generation = await self._staging.stage(
            account_id=request.identity.account_id,
            source_execution_key=request.identity.source_execution_key,
            ordinal=0,
            payload=payload,
        )
        staged = SourceBatchObjectV1(
            ordinal=0,
            object_ref=object_ref,
            object_generation=object_generation,
            media_type="application/x-ndjson",
            compression="none",
            sha256=hashlib.sha256(payload).hexdigest(),
            byte_count=len(payload),
            row_count=len(result.rows),
        )

        async with inline_publication(request.identity.account_id, self._seals) as publish_seal:
            manifest = self._seal_manifest(
                request,
                observed_at=observed_at,
                data_through=data_through,
                staged=staged,
                constituents=constituents,
                cells=cells,
                coverage_status=coverage_status,
                explicit_zero=explicit_zero,
                row_count=len(result.rows),
                control_totals=control_totals,
                warnings=warnings,
                provisional_until=result.provisional_until,
            )

            manifest_bytes = encode_source_batch_manifest_v1(manifest)
            reference = source_batch_manifest_reference_v1(
                f"{self._staged_commit_prefix}.{manifest.publication_id}", manifest_bytes
            )
            winner = await publish_seal(
                request.identity.source_execution_key,
                SealedSlice(reference=reference, manifest_bytes=manifest_bytes),
            )
            # A concurrent worker may have sealed first; its bytes are the
            # publication, and returning ours instead would make the same key
            # resolve two ways.
            return ReportingSourceExecutorResult.completed(
                request=request, manifest=winner.reference, manifest_bytes=winner.manifest_bytes
            )

    def _derive_statuses(
        self,
        request: ReportingSourceSliceRequestV1,
        result: InlineFetchResult,
        covered: set[str],
        rows_by_constituent: Mapping[str, Sequence[Mapping[str, Any]]],
        data_through: datetime,
    ) -> dict[str, ReportingAvailabilityStatus]:
        """Turn "what came back" into per-constituent publication evidence.

        The freshness gate lives here.  An ``AUTHORITATIVE`` close whose
        watermark has not reached the period end cannot claim ``present``: the
        numbers are real but the window is not finished, which is ``delayed``,
        not complete.
        """
        short_of_the_window = request.publication_class == "AUTHORITATIVE" and _utc(
            data_through
        ) < _utc(request.period.end)
        statuses: dict[str, ReportingAvailabilityStatus] = {}
        for item in request.coverage.constituents:
            constituent_id = item.constituent_id
            if constituent_id in result.unavailable_constituents:
                statuses[constituent_id] = result.unavailable_status
            elif constituent_id not in covered:
                statuses[constituent_id] = "missing"
            elif short_of_the_window:
                statuses[constituent_id] = "delayed"
            elif rows_by_constituent.get(constituent_id):
                statuses[constituent_id] = "present"
            else:
                statuses[constituent_id] = "explicit_zero"
        return statuses

    def _resolve_availability(
        self,
        request: ReportingSourceSliceRequestV1,
        result: InlineFetchResult,
        *,
        overrides: Mapping[str, Mapping[str, MetricEvidence]],
        statuses: Mapping[str, ReportingAvailabilityStatus],
        rows_by_constituent: Mapping[str, Sequence[Mapping[str, Any]]],
        data_through: datetime,
        observed_at: datetime,
    ) -> tuple[list[ReportingConstituentCoverageV1], list[ReportingMetricAvailabilityV1]]:
        """Resolve cells first, then derive coverage without widening their claims."""
        offering = self._capabilities.offering(request.offering_id)
        declared = {metric.name: metric for metric in offering.metrics}
        constituents: list[ReportingConstituentCoverageV1] = []
        cells: list[ReportingMetricAvailabilityV1] = []
        for item in request.coverage.constituents:
            constituent_id = item.constituent_id
            fallback = statuses[constituent_id]
            fallback_reason = (
                None
                if fallback in _AVAILABLE
                else result.unavailable_constituents.get(constituent_id)
                or _DEFAULT_REASONS[fallback]
            )
            constituent_cells: list[ReportingMetricAvailabilityV1] = []
            for metric in request.requested_metrics:
                status = fallback
                watermark = data_through if status in _AVAILABLE else None
                reason = fallback_reason
                explicit = overrides.get(constituent_id, {}).get(metric)
                if explicit is not None:
                    _validate_cell_rows(
                        constituent_id,
                        metric,
                        explicit,
                        rows_by_constituent[constituent_id],
                        monetary=metric in self._monetary,
                    )
                    status, reason = explicit.status, explicit.reason
                    watermark = explicit.data_through
                    if watermark is not None:
                        if _utc(watermark) < _utc(request.period.start):
                            raise _CellEvidenceError(
                                f"cell_availability ({constituent_id!r}, {metric!r}) "
                                "data_through must not precede the reporting period"
                            )
                        watermark = _clamp_watermark(request, watermark, observed_at)
                    elif status == "explicit_zero":
                        watermark = data_through
                    if (
                        status in _AVAILABLE
                        and request.publication_class == "AUTHORITATIVE"
                        and watermark is not None
                        and _utc(watermark) < _utc(request.period.end)
                    ):
                        status, reason = "delayed", _DEFAULT_REASONS["delayed"]
                # Evidence follows the cell status, never its constituent's
                # roll-up. Only the selected SDK offering supplies semantics.
                constituent_cells.append(
                    ReportingMetricAvailabilityV1(
                        constituent_id=constituent_id,
                        metric=metric,
                        semantic_contract_id=declared[metric].semantic_contract_id,
                        semantic_contract_version=declared[metric].semantic_contract_version,
                        semantic_contract_sha256=declared[metric].semantic_contract_sha256,
                        status=status,
                        data_through=watermark,
                        reason=reason,
                    )
                )
            cells.extend(constituent_cells)
            status = fallback
            watermark = data_through if status in _AVAILABLE else None
            reason = fallback_reason
            if overrides.get(constituent_id):
                cell_statuses = {cell.status for cell in constituent_cells}
                if cell_statuses == {"unsupported"}:
                    # Zero applicable metrics is no coverage evidence, not a
                    # vacuously complete measurement or an observed zero.
                    status = "unsupported"
                elif len(cell_statuses) == 1:
                    status = constituent_cells[0].status
                elif cell_statuses <= _AVAILABLE | {"unsupported"}:
                    status = "present"
                else:
                    status = "partial"
                if status in _AVAILABLE:
                    watermark = min(
                        cell.data_through
                        for cell in constituent_cells
                        if cell.data_through is not None
                    )
                    reason = None
                else:
                    watermark = None
                    reasons = {cell.reason for cell in constituent_cells}
                    reason = (
                        constituent_cells[0].reason
                        if len(reasons) == 1 and constituent_cells[0].reason is not None
                        else _DEFAULT_REASONS[status]
                    )
            constituents.append(
                ReportingConstituentCoverageV1(
                    constituent=item, status=status, data_through=watermark, reason=reason
                )
            )
        _validate_metric_applicability(offering.metrics, cells)
        return constituents, cells

    def _seal_manifest(
        self,
        request: ReportingSourceSliceRequestV1,
        *,
        observed_at: datetime,
        data_through: datetime,
        staged: SourceBatchObjectV1,
        constituents: list[ReportingConstituentCoverageV1],
        cells: list[ReportingMetricAvailabilityV1],
        coverage_status: str,
        explicit_zero: bool,
        row_count: int,
        control_totals: list[SourceControlTotalV1],
        warnings: list[str],
        provisional_until: datetime | None,
    ) -> SourceBatchManifestV1:
        coverage = SourceBatchCoverageV1(
            denominator_fingerprint=request.coverage.denominator_fingerprint,
            status=coverage_status,  # type: ignore[arg-type]
            constituents=constituents,
        )
        finality_at = (
            max(_utc(observed_at), _utc(request.period.end))
            if request.publication_class == "AUTHORITATIVE"
            else _utc(observed_at)
        )
        draft: dict[str, Any] = {
            "strictness": "basic",
            "identity": request.identity,
            "publication_namespace": request.publication_namespace,
            "publication_id": deterministic_source_publication_id_v1(
                source_execution_key=request.identity.source_execution_key,
                logical_slice_fingerprint=request.identity.logical_slice_fingerprint,
                publication_class=request.publication_class,
                revision_kind=request.revision_kind,
                supersedes_publication_id=request.supersedes_publication_id,
            ),
            "publication_class": request.publication_class,
            "content_fingerprint": f"sha256:{'0' * 64}",
            "adapter_build": request.adapter_build,
            "offering_id": request.offering_id,
            "contract": request.contract,
            "period": request.period,
            "revision_kind": request.revision_kind,
            "supersedes_publication_id": request.supersedes_publication_id,
            "finality_evidence": ReportingFinalityEvidenceV1(
                basis=(
                    "provisional_observation"
                    if request.publication_class == "PROVISIONAL_SNAPSHOT"
                    else "elapsed_settlement_window"
                ),
                observed_at=finality_at,
                provisional_until=(
                    provisional_until
                    if request.publication_class == "PROVISIONAL_SNAPSHOT"
                    else None
                ),
            ),
            "observed_at": observed_at,
            "data_through": data_through,
            "acquired_at": max(_utc(observed_at), finality_at),
            "currency": request.currency,
            "requested_dimensions": list(request.requested_dimensions),
            "objects": [staged],
            "object_set_sha256": source_batch_object_set_sha256_v1([staged]),
            "row_count": row_count,
            "byte_count": staged.byte_count,
            "control_totals": control_totals,
            "metric_availability": cells,
            "coverage": coverage,
            "explicit_zero": explicit_zero,
            "event_time_range": None if row_count == 0 else (request.period.start, data_through),
            "prior_checkpoint": request.prior_checkpoint,
            "warnings": warnings,
        }
        draft["content_fingerprint"] = publication_content_fingerprint_v1(
            SourceBatchManifestV1.model_construct(**draft)
        )
        return SourceBatchManifestV1(**draft)

A conforming executor over an ordinary delivery fetch.

Parameters

capabilities: Your truthful, scoped capability declaration. Every offering it declares must use manifest_strictness='basic' – this adapter cannot honestly produce per-call provider evidence for a fetch it did not make itself. fetch: Sync or async callable taking the frozen :class:~adcp.reporting.source.ReportingSourceSliceRequestV1 and returning rows, an :class:InlineFetchResult, an AdCP GetMediaBuyDeliveryResponse, or None for "not ready". staging / seals: Where rows and replay seals live. The in-memory defaults are for tests and single-process pilots; they are not durable. constituent_of: Maps a row to the constituent_id it belongs to. The default matches a row's media_buy_id / package_id / product_id against the requested denominator, which covers the common shapes. Rows that match nothing are retained in the staged object but excluded from coverage, and a warning names how many were dropped. monetary_metrics: The metric names the trusted report definition froze as money, beyond the built-in spend. Supply ReportingDefinitionBinding.monetary_metric_units keys (and any monetary_control_total_units keys) when the obligation declares them: the frozen slice request cannot carry those declarations, because to_wire() keeps them off the wire so retained contract hashes do not move, so the adapter would otherwise not know which columns the ledger will reconcile. A custom money column that is not declared here behaves like a non-monetary one. include_exception_detail: Whether an unclassified exception's str() may enter the retained safe message. Off by default: HTTP client exceptions routinely stringify the full request URL, and a signed URL or an access_token query parameter in a manifest is a credential leak into permanent, replicated evidence. The exception type is always recorded. clock: Supplies the observation instant. Override it to make a test's publications deterministic; the default is the wall clock, truncated to the millisecond precision the manifest binds.

Subclasses

  • adcp.reporting.service._AuthorizedInlineSource

Instance variables

prop capabilities : ReportingSourceCapabilitiesV1
Expand source code
@property
def capabilities(self) -> ReportingSourceCapabilitiesV1:
    return self._capabilities
prop staging : ReportingStagingStore
Expand source code
@property
def staging(self) -> ReportingStagingStore:
    """The staging store, which also serves as the staged-object reader."""
    return self._staging

The staging store, which also serves as the staged-object reader.

Methods

async def execute(self,
request: ReportingSourceSliceRequestV1,
*,
cancel: asyncio.Event,
heartbeat: Callable[[], None] | None = None) ‑> ReportingSourceExecutorResult
Expand source code
async def execute(
    self,
    request: ReportingSourceSliceRequestV1,
    *,
    cancel: asyncio.Event,
    heartbeat: Callable[[], None] | None = None,
) -> ReportingSourceExecutorResult:
    account_id = request.identity.account_id
    key = request.identity.source_execution_key

    sealed = await self._seals.get(account_id=account_id, source_execution_key=key)
    if sealed is not None:
        # Replay, including replay after cancellation: the publication is
        # already immutable, so returning anything else would be a lie.
        return ReportingSourceExecutorResult.completed(
            request=request,
            manifest=sealed.reference,
            manifest_bytes=sealed.manifest_bytes,
        )

    if cancel.is_set():
        return ReportingSourceExecutorResult.failed(
            ReportingSourceErrorV1(
                code="CANCELLED",
                retry="cancelled",
                safe_message="the slice was cancelled before any source read began",
            )
        )

    try:
        answer = await self._invoke(request, heartbeat)
        if answer is not None:
            # Coerce inside the guard: a fetch that returns something
            # unrecognizable is a source failure like any other, not an
            # exception escaping the executor.
            answer = _coerce(answer, currency=request.currency)
            # Retain the rows we checked across staging awaits. A fetch
            # may return a shared cache whose mappings later change.
            answer = replace(answer, rows=tuple(deepcopy(dict(row)) for row in answer.rows))
            if (
                answer.currency is not None
                and validate_currency(answer.currency) != request.currency
            ):
                raise ReportingCurrencyError(
                    "CURRENCY_MISMATCH",
                    "inline source currency disagrees with the frozen request",
                )
            # Before staging and, crucially, before _control_totals aggregates.
            validate_row_currencies(request.currency, answer.rows)
    except ReportingCurrencyError as error:
        return ReportingSourceExecutorResult.failed(
            ReportingSourceErrorV1(
                code="INTEGRITY_FAILED", retry="terminal", safe_message=str(error)
            )
        )
    except ReportingSourceError as error:
        return ReportingSourceExecutorResult.failed(error.error)
    except _CellEvidenceError:
        # Constructors/helpers may run inside the fetch. Keep the same
        # actionable ValueError as matrix validation after the fetch,
        # rather than redacting an adapter bug as a retryable provider fault.
        raise
    except asyncio.CancelledError:
        raise
    except Exception as error:
        return ReportingSourceExecutorResult.failed(self._unclassified(error))

    if answer is None:
        not_ready = self._not_ready(request)
        if not_ready is not None:
            return not_ready
        answer = InlineFetchResult(
            rows=[],
            covered_constituent_ids=(),
            unavailable_constituents={
                item.constituent_id: "Source has not produced this period yet"
                for item in request.coverage.constituents
            },
            unavailable_status="delayed",
        )

    return await self._publish(request, answer, cancel=cancel)
class MetricEvidence (status: "Literal['present', 'explicit_zero', 'missing', 'delayed', 'unsupported']",
data_through: datetime | None = None,
reason: str | None = None)
Expand source code
@dataclass(frozen=True)
class MetricEvidence:
    """Source evidence for one requested constituent/metric cell.

    This carries availability only. The SDK binds metric semantics from the
    selected offering. Prefer the named constructors; direct construction
    enforces the same invariants and raises ``ValueError`` for invalid evidence.

    Watermarks must be timezone-aware. Reasons use the manifest's bounded ASCII
    evidence format; use stable, redacted explanations such as
    ``not_video_inventory`` rather than provider diagnostics or credentials.
    """

    status: Literal["present", "explicit_zero", "missing", "delayed", "unsupported"]
    data_through: datetime | None = None
    reason: str | None = None

    def __post_init__(self) -> None:
        if self.status not in ("present", "explicit_zero", "missing", "delayed", "unsupported"):
            raise _CellEvidenceError("MetricEvidence.status is not a supported cell status")
        if self.data_through is not None:
            if (
                not isinstance(self.data_through, datetime)
                or self.data_through.tzinfo is None
                or self.data_through.utcoffset() is None
            ):
                raise _CellEvidenceError(
                    "MetricEvidence.data_through must be a timezone-aware datetime"
                )
        if self.status == "present" and self.data_through is None:
            raise _CellEvidenceError("MetricEvidence present requires data_through")
        if self.status in _AVAILABLE:
            if self.reason is not None:
                raise _CellEvidenceError(f"MetricEvidence {self.status} must not carry a reason")
        else:
            if self.reason is None:
                raise _CellEvidenceError(f"MetricEvidence {self.status} requires a stable reason")
            try:
                _reason_adapter().validate_python(self.reason, strict=True)
            except ValidationError as error:
                raise _CellEvidenceError(
                    "MetricEvidence.reason must be a wire-valid stable reason"
                ) from error
        if self.status in {"missing", "unsupported"} and self.data_through is not None:
            raise _CellEvidenceError(f"MetricEvidence {self.status} must not carry data_through")

    @classmethod
    def present(cls, data_through: datetime) -> MetricEvidence:
        """The source measured this metric through the supplied watermark."""
        return cls("present", data_through=data_through)

    @classmethod
    def explicit_zero(cls, *, data_through: datetime | None = None) -> MetricEvidence:
        """An observed zero; omit the watermark to inherit the fetch watermark.

        Any supplied row values for this cell must be numeric zeros. An
        omitted row value is not converted into a control-total contribution.
        """
        return cls("explicit_zero", data_through=data_through)

    @classmethod
    def missing(cls, reason: str) -> MetricEvidence:
        """The source returned no answer for this cell."""
        return cls("missing", reason=reason)

    @classmethod
    def delayed(cls, reason: str, *, data_through: datetime | None = None) -> MetricEvidence:
        """The metric is not ready; retain its watermark when one is known."""
        return cls("delayed", data_through=data_through, reason=reason)

    @classmethod
    def unavailable(cls, reason: str) -> MetricEvidence:
        """The source cannot measure this cell (wire status ``unsupported``)."""
        return cls("unsupported", reason=reason)

Source evidence for one requested constituent/metric cell.

This carries availability only. The SDK binds metric semantics from the selected offering. Prefer the named constructors; direct construction enforces the same invariants and raises ValueError for invalid evidence.

Watermarks must be timezone-aware. Reasons use the manifest's bounded ASCII evidence format; use stable, redacted explanations such as not_video_inventory rather than provider diagnostics or credentials.

Static methods

def delayed(reason: str, *, data_through: datetime | None = None) ‑> MetricEvidence

The metric is not ready; retain its watermark when one is known.

def explicit_zero(*, data_through: datetime | None = None) ‑> MetricEvidence

An observed zero; omit the watermark to inherit the fetch watermark.

Any supplied row values for this cell must be numeric zeros. An omitted row value is not converted into a control-total contribution.

def missing(reason: str) ‑> MetricEvidence

The source returned no answer for this cell.

def present(data_through: datetime) ‑> MetricEvidence

The source measured this metric through the supplied watermark.

def unavailable(reason: str) ‑> MetricEvidence

The source cannot measure this cell (wire status unsupported).

Instance variables

var data_through : datetime.datetime | None
var reason : str | None
var status : Literal['present', 'explicit_zero', 'missing', 'delayed', 'unsupported']
class ReportingSealStore (*args, **kwargs)
Expand source code
@runtime_checkable
class ReportingSealStore(Protocol):
    """Remembers what a ``source_execution_key`` already published.

    This is what makes replay identity possible at all.  A manifest carries
    real observation and acquisition times; recomputing them on a second
    execution would produce different bytes for the same logical work.  Sealing
    the first result and returning it verbatim is the only correct answer, and
    it is also what cancellation semantics require: a cancelled execution
    publishes nothing *unless the exact idempotent result was already sealed*.
    """

    async def get(self, *, account_id: str, source_execution_key: str) -> SealedSlice | None: ...

    async def put(
        self, *, account_id: str, source_execution_key: str, sealed: SealedSlice
    ) -> SealedSlice:
        """Persist the seal, returning whichever seal won a concurrent race.

        Two workers may execute the same key at once.  The store decides the
        winner; both callers then return the winning bytes, so replay identity
        holds even under a double dispatch.
        """
        ...

Remembers what a source_execution_key already published.

This is what makes replay identity possible at all. A manifest carries real observation and acquisition times; recomputing them on a second execution would produce different bytes for the same logical work. Sealing the first result and returning it verbatim is the only correct answer, and it is also what cancellation semantics require: a cancelled execution publishes nothing unless the exact idempotent result was already sealed.

Ancestors

  • typing.Protocol
  • typing.Generic

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: ...
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:
    """Persist the seal, returning whichever seal won a concurrent race.

    Two workers may execute the same key at once.  The store decides the
    winner; both callers then return the winning bytes, so replay identity
    holds even under a double dispatch.
    """
    ...

Persist the seal, returning whichever seal won a concurrent race.

Two workers may execute the same key at once. The store decides the winner; both callers then return the winning bytes, so replay identity holds even under a double dispatch.

class ReportingStagingStore (*args, **kwargs)
Expand source code
@runtime_checkable
class ReportingStagingStore(Protocol):
    """Where normalized rows live, immutably, for as long as the manifest does.

    ``stage`` returns an ``(object_ref, object_generation)`` pair.  The pair,
    not the ref alone, is the immutable identity: an object-store key can be
    overwritten, a key plus generation cannot.  Re-reading the pair must either
    return the same bytes or fail -- never different bytes.
    """

    async def stage(
        self, *, account_id: str, source_execution_key: str, ordinal: int, payload: bytes
    ) -> tuple[str, str]: ...

    async def read(
        self,
        *,
        object_ref: str,
        object_generation: str,
        account_id: str,
        source_scope: Mapping[str, Any],
        cancel: asyncio.Event,
    ) -> bytes: ...

Where normalized rows live, immutably, for as long as the manifest does.

stage returns an (object_ref, object_generation) pair. The pair, not the ref alone, is the immutable identity: an object-store key can be overwritten, a key plus generation cannot. Re-reading the pair must either return the same bytes or fail – never different bytes.

Ancestors

  • typing.Protocol
  • typing.Generic

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: ...
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]: ...
class SealedSlice (reference: SourceBatchManifestReferenceV1, manifest_bytes: bytes)
Expand source code
@dataclass(frozen=True)
class SealedSlice:
    """The immutable result of one execution, kept for replay."""

    reference: SourceBatchManifestReferenceV1
    manifest_bytes: bytes

The immutable result of one execution, kept for replay.

Subclasses

  • adcp.reporting.inline_storage._StoredSeal

Instance variables

var manifest_bytes : bytes
var reference : SourceBatchManifestReferenceV1