Module adcp.reporting.conformance
Executable conformance for reporting source executors.
An executor that type-checks is not necessarily a conforming executor.
These
validators are the difference: run them in your own test suite against your own
executor and they will tell you, before a buyer does, that you widened a date
range, published a partial result as complete, minted your own publication id,
or returned a different manifest the second time the same
source_execution_key was replayed.
Three levels, each strictly stronger:
:func:validate_reporting_source_execution()
One execution.
Checks the slice against the capability declaration, the
manifest against the slice, and every staged object's bytes against the
digest the manifest claims for it.
:func:run_reporting_source_replay_conformance()
Two executions of the same source_execution_key.
Both must validate
and both must produce byte-identical manifests.
This is the check that
catches an executor that embeds now() or a fresh UUID in its manifest.
:func:validate_reporting_revision_sequence()
A whole revision sequence.
Checks the rules no single manifest can:
freshness never regresses, availability never silently degrades, a snapshot
never follows an official publication, and a correction names the exact
publication it supersedes.
Failures raise :class:ReportingSourceConformanceError with a stable code, so
a harness can assert which rule was broken rather than string-matching a
message.
Functions
async def run_reporting_source_replay_conformance(*,
executor: ReportingSourceExecutor,
request: ReportingSourceSliceRequestV1,
object_reader: ReportingSourceStagedObjectReader,
cancel: asyncio.Event | None = None,
clock: Callable[[], datetime] | None = None) ‑> SourceBatchManifestV1-
Expand source code
async def run_reporting_source_replay_conformance( *, executor: ReportingSourceExecutor, request: ReportingSourceSliceRequestV1, object_reader: ReportingSourceStagedObjectReader, cancel: asyncio.Event | None = None, clock: Callable[[], datetime] | None = None, ) -> SourceBatchManifestV1: """Execute the same slice twice and require an identical immutable publication. Reusing a ``source_execution_key`` must return byte-identical manifest bytes and the same staged object set. This is the check that catches the two most common non-conformances: a ``now()`` timestamp baked into the manifest, and a fresh UUID minted per attempt. ``clock`` supplies the deadline-check instant for both executions and their staged-object reads. """ capabilities = executor.capabilities cancel = cancel or asyncio.Event() async def once() -> tuple[ReportingSourceExecutorResult, SourceBatchManifestV1]: result = await _with_deadline( lambda: executor.execute(request, cancel=cancel), deadline_at=request.deadline_at, cancel=cancel, clock=clock, ) manifest = await validate_reporting_source_execution( capabilities=capabilities, request=request, result=result, object_reader=object_reader, cancel=cancel, clock=clock, ) return result, manifest first_result, first_manifest = await once() second_result, second_manifest = await once() first_reference = first_result.response.manifest if first_result.response else None second_reference = second_result.response.manifest if second_result.response else None if ( first_reference != second_reference or first_result.manifest_bytes != second_result.manifest_bytes or first_manifest != second_manifest ): raise _fail( "REPLAY_MISMATCH", "repeating a source_execution_key must return the identical immutable publication", ) return first_manifestExecute the same slice twice and require an identical immutable publication.
Reusing a
source_execution_keymust return byte-identical manifest bytes and the same staged object set. This is the check that catches the two most common non-conformances: anow()timestamp baked into the manifest, and a fresh UUID minted per attempt.clocksupplies the deadline-check instant for both executions and their staged-object reads. def validate_reporting_revision_sequence(manifests: Sequence[SourceBatchManifestV1],
*,
allow_cross_finality: bool = False) ‑> None-
Expand source code
def validate_reporting_revision_sequence( manifests: Sequence[SourceBatchManifestV1], *, allow_cross_finality: bool = False, ) -> None: """Validate an ordered sequence of publications for one logical window. The rules a single manifest cannot enforce: * every publication describes the same account, config, source scope, and logical window; * ``observed_at``, ``acquired_at``, ``source_read_cutoff_at``, and ``data_through`` never regress; * a constituent or metric cell that was once ``present`` / ``explicit_zero`` never silently degrades, and its watermark never moves backwards -- a later fetch may report *new* unavailability, but it may not replace a fresh committed cell with a stale one; * a snapshot never follows the official publication; * exactly one initial ``authoritative`` publication exists, and every later ``correction`` names its immediate predecessor. Snapshot and authoritative offerings usually have different semantic scope (different metrics, grain, or windowing), so mixing them in one sequence is refused unless ``allow_cross_finality`` says the caller has already established that the two scopes address the same buyer-visible window. """ if not manifests: return first = manifests[0] baseline = _sequence_scope(first) for manifest in manifests: if _sequence_scope(manifest) != baseline: raise _fail( "REVISION_SEQUENCE_INVALID", "every revision in a sequence must describe the same account, source scope, " "and logical window", ) has_snapshot = any(item.revision_kind == "snapshot" for item in manifests) has_official = any(item.revision_kind != "snapshot" for item in manifests) if has_snapshot and has_official and not allow_cross_finality: raise _fail( "REVISION_SEQUENCE_INVALID", "a snapshot-to-official transition needs an explicit cross-finality decision; " "pass allow_cross_finality=True once the two semantic scopes are known to " "address the same buyer-visible window", ) prior_snapshot: SourceBatchManifestV1 | None = None official: SourceBatchManifestV1 | None = None for manifest in manifests: if manifest.revision_kind == "snapshot": if official is not None: raise _fail( "REVISION_SEQUENCE_INVALID", "a snapshot cannot replace or follow the official publication", ) if prior_snapshot is not None: if _semantic_scope(prior_snapshot) != _semantic_scope(manifest): raise _fail( "REVISION_SEQUENCE_INVALID", "a snapshot must retain the exact semantic scope of its predecessor", ) _assert_freshness_does_not_regress(prior_snapshot, manifest, "a snapshot") _assert_cells_do_not_regress(prior_snapshot, manifest, "a snapshot") prior_snapshot = manifest continue if manifest.revision_kind == "authoritative": if official is not None: raise _fail( "REVISION_SEQUENCE_INVALID", "only one initial authoritative publication is allowed; a later " "restatement is a correction", ) if prior_snapshot is not None: _assert_freshness_does_not_regress( prior_snapshot, manifest, "an authoritative publication" ) _assert_cells_do_not_regress( prior_snapshot, manifest, "an authoritative publication", shared_only=True ) official = manifest continue if official is None: raise _fail( "REVISION_SEQUENCE_INVALID", "a correction must follow an authoritative publication", ) if manifest.supersedes_publication_id != official.publication_id: raise _fail( "REVISION_SEQUENCE_INVALID", "a correction must name the exact publication it supersedes", ) if _semantic_scope(official) != _semantic_scope(manifest): raise _fail( "REVISION_SEQUENCE_INVALID", "a correction must retain the semantic scope of its predecessor", ) _assert_freshness_does_not_regress(official, manifest, "a correction") _assert_cells_do_not_regress(official, manifest, "a correction") official = manifestValidate an ordered sequence of publications for one logical window.
The rules a single manifest cannot enforce:
- every publication describes the same account, config, source scope, and logical window;
observed_at,acquired_at,source_read_cutoff_at, anddata_throughnever regress;- a constituent or metric cell that was once
present/explicit_zeronever silently degrades, and its watermark never moves backwards – a later fetch may report new unavailability, but it may not replace a fresh committed cell with a stale one; - a snapshot never follows the official publication;
- exactly one initial
authoritativepublication exists, and every latercorrectionnames its immediate predecessor.
Snapshot and authoritative offerings usually have different semantic scope (different metrics, grain, or windowing), so mixing them in one sequence is refused unless
allow_cross_finalitysays the caller has already established that the two scopes address the same buyer-visible window. async def validate_reporting_source_execution(*,
capabilities: ReportingSourceCapabilitiesV1,
request: ReportingSourceSliceRequestV1,
result: ReportingSourceExecutorResult,
object_reader: ReportingSourceStagedObjectReader,
cancel: asyncio.Event | None = None,
clock: Callable[[], datetime] | None = None) ‑> SourceBatchManifestV1-
Expand source code
async def validate_reporting_source_execution( *, capabilities: ReportingSourceCapabilitiesV1, request: ReportingSourceSliceRequestV1, result: ReportingSourceExecutorResult, object_reader: ReportingSourceStagedObjectReader, cancel: asyncio.Event | None = None, clock: Callable[[], datetime] | None = None, ) -> SourceBatchManifestV1: """Validate one execution end to end and return its verified manifest. Reads every staged object the manifest names and checks its bytes against the declared digest and size. A manifest whose objects cannot be read, or read differently than claimed, is not evidence of anything. ``clock`` permits deterministic deadline checks for retained test slices. """ offering = _validate_request_against_capabilities(capabilities, request) if not result.ok: error = result.error assert error is not None # guaranteed by ReportingSourceExecutorResult raise _fail("EXECUTION_FAILED", f"{error.code}: {error.safe_message}") response = result.response assert response is not None manifest_bytes = result.manifest_bytes assert manifest_bytes is not None if response.identity != request.identity: raise _fail( "IDENTITY_MISMATCH", "the response must echo the complete frozen request identity" ) manifest = parse_verified_source_batch_manifest_v1(response.manifest, manifest_bytes) _validate_manifest_against_request(offering, request, manifest) cancel = cancel or asyncio.Event() for item in manifest.objects: try: payload = await _with_deadline( lambda item=item: object_reader.read( # type: ignore[misc] object_ref=item.object_ref, object_generation=item.object_generation, account_id=request.identity.account_id, source_scope=request.identity.source_scope, cancel=cancel, ), deadline_at=request.deadline_at, cancel=cancel, clock=clock, ) except ReportingSourceConformanceError: raise except Exception as error: raise _fail( "OBJECT_MISMATCH", f"the generation-pinned object {item.ordinal} is not readable", ) from error if len(payload) != item.byte_count or hashlib.sha256(payload).hexdigest() != item.sha256: raise _fail( "OBJECT_MISMATCH", f"the generation-pinned object {item.ordinal} does not match its evidence", ) return manifestValidate one execution end to end and return its verified manifest.
Reads every staged object the manifest names and checks its bytes against the declared digest and size. A manifest whose objects cannot be read, or read differently than claimed, is not evidence of anything.
clockpermits deterministic deadline checks for retained test slices. def validate_reporting_source_failure(result: ReportingSourceExecutorResult,
expected_code: ReportingSourceErrorCode) ‑> ReportingSourceErrorV1-
Expand source code
def validate_reporting_source_failure( result: ReportingSourceExecutorResult, expected_code: ReportingSourceErrorCode ) -> ReportingSourceErrorV1: """Assert a slice failed with exactly ``expected_code`` and return the error. A failed upstream page or job must not publish ``complete: true``; this is how you assert that in a test. """ if result.ok: raise _fail( "EXECUTION_FAILED", "a failed or partial source fetch must not return a completed manifest", ) error = result.error assert error is not None if error.code != expected_code: raise _fail("EXECUTION_FAILED", f"expected {expected_code}, received {error.code}") return errorAssert a slice failed with exactly
expected_codeand return the error.A failed upstream page or job must not publish
complete: true; this is how you assert that in a test.
Classes
class ReportingSourceConformanceError (code: ReportingSourceConformanceCode, message: str)-
Expand source code
class ReportingSourceConformanceError(AssertionError): """An executor violated the source contract. Subclasses :class:`AssertionError` so an unhandled one reads like the test failure it almost always is, while still carrying a machine-checkable :attr:`code`. """ def __init__(self, code: ReportingSourceConformanceCode, message: str) -> None: super().__init__(f"{code}: {message}") self.code = codeAn executor violated the source contract.
Subclasses :class:
AssertionErrorso an unhandled one reads like the test failure it almost always is, while still carrying a machine-checkable :attr:code.Ancestors
- builtins.AssertionError
- builtins.Exception
- builtins.BaseException