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_availabilityretaining 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_availabilityfor 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 payloadAccount-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 payloadProcess-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:
FileSystemStagingStoreor 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_availabilityfor 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
unsupportedand cannot satisfy full coverage. Other mixed availability promotes it topartial. 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-formatmetric_availabilitylist.Unknown keys, duplicate mapping entries, invalid evidence, contradictory explicit zeros, and a withdrawn monetary cell whose own constituent's rows report that metric raise
ValueErrorbefore 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 reportedmissingrather 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.
-
constituent_idto 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.
-
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.ReportingSourceSliceRequestV1and returning rows, an :class:InlineFetchResult, an AdCPGetMediaBuyDeliveryResponse, orNonefor "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 theconstituent_idit belongs to. The default matches a row'smedia_buy_id/package_id/product_idagainst 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-inspend. SupplyReportingDefinitionBinding.monetary_metric_unitskeys (and anymonetary_control_total_unitskeys) when the obligation declares them: the frozen slice request cannot carry those declarations, becauseto_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'sstr()may enter the retained safe message. Off by default: HTTP client exceptions routinely stringify the full request URL, and a signed URL or anaccess_tokenquery 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._stagingThe 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
ValueErrorfor invalid evidence.Watermarks must be timezone-aware. Reasons use the manifest's bounded ASCII evidence format; use stable, redacted explanations such as
not_video_inventoryrather 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.
-
The source cannot measure this cell (wire status
unsupported).
Instance variables
var data_through : datetime.datetime | Nonevar reason : str | Nonevar 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_keyalready 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.
stagereturns 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: bytesThe immutable result of one execution, kept for replay.
Subclasses
- adcp.reporting.inline_storage._StoredSeal
Instance variables
var manifest_bytes : bytesvar reference : SourceBatchManifestReferenceV1