Module adcp.reporting.materializer.verification
SDK-owned verification of frozen source content and actual destination reads.
The registry contains immutable, explicitly installed bytes. It never fetches a URI or trusts a writer's digest as proof. B1 returns evidence only; B2 must reselect and fence under its final account transaction before publishing it.
Functions
def re_decimal(value: str, scale: int) ‑> bool-
Expand source code
def re_decimal(value: str, scale: int) -> bool: return bool( re.fullmatch(r"-?(?:0|[1-9][0-9]{0,37})" + (rf"\.[0-9]{{{scale}}}" if scale else ""), value) ) def validate_materialization_target(prepared: ReportingPreparedRevision,
*,
binding: ReportingDestinationBinding,
revisions: Sequence[ReportingRevisionRecord]) ‑> None-
Expand source code
def validate_materialization_target( prepared: ReportingPreparedRevision, *, binding: ReportingDestinationBinding, revisions: Sequence[ReportingRevisionRecord], ) -> None: """Pure B2 finish seam. Caller must supply its locked, complete current history.""" selection = select_reporting_revision( revisions, account_id=prepared.request.principal.account_id, reporting_obligation_id=prepared.request.reporting_obligation_id, required_finality=prepared.obligation.required_finality, ) if selection.kind == "corrupt": raise failure("HISTORY_CORRUPT") if ( selection.kind != "selected" or not selection.revision.readable or selection.revision != prepared.revision or binding_fingerprint(binding) != prepared.request.binding_fingerprint ): raise failure("CURRENT_REVISION_CHANGED")Pure B2 finish seam. Caller must supply its locked, complete current history.
def validate_verified_destination(value: ReportingVerifiedDestination,
request: ReportingDestinationRequest,
*,
token: str) ‑> None-
Expand source code
def validate_verified_destination( value: ReportingVerifiedDestination, request: ReportingDestinationRequest, *, token: str ) -> None: """Require unmodified evidence minted by this process's SDK readback. A public value constructor remains source compatible with B1. Constructed, copied or restored claims cannot authorize B2 finish: restart repeats the actual readback. The seal is never a persisted destination credential. """ if ( type(value) is not ReportingVerifiedDestination or value.request != request or getattr(value, "_materializer_fence", None) != token or not hmac.compare_digest(getattr(value, "_sdk_seal", b""), _verification_seal(value)) ): raise failure("DESTINATION_CORRUPT")Require unmodified evidence minted by this process's SDK readback.
A public value constructor remains source compatible with B1. Constructed, copied or restored claims cannot authorize B2 finish: restart repeats the actual readback. The seal is never a persisted destination credential.
Classes
class ReportingDestinationIO (registry: ReportingRevisionVerifierRegistry,
resolver: ReportingDestinationResolver)-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingDestinationIO: """Explicit one-shot write/readback operations; no coordinator or lease loop.""" registry: ReportingRevisionVerifierRegistry resolver: ReportingDestinationResolver = field(repr=False) @asynccontextmanager async def _session( self, prepared: ReportingPreparedRevision, phase: ReportingIOPhase, context: ReportingIOContext, ) -> AsyncIterator[ReportingDestinationSession]: self.registry.require(prepared.request.verification_key) problem = False canceled = False try: session = self.resolver.resolve(prepared.request, phase=phase, context=context) except asyncio.CancelledError: canceled = True except Exception: problem = True if canceled: raise asyncio.CancelledError if problem: raise ReportingWriterError( ReportingWriterFailure("RESOURCE_UNAVAILABLE", "same_identity") ) if not isinstance(session, ReportingDestinationSession): raise failure("BINDING_MISMATCH") try: _session_binding(session, prepared.request, phase, context) except ReportingWriterError: problem = True if problem: await session.aclose() raise failure("BINDING_MISMATCH") failure_record: ReportingWriterFailure | None = None canceled = False try: async with session: _session_binding(session, prepared.request, phase, context) yield session except asyncio.CancelledError: canceled = True except ReportingWriterError as exc: failure_record = exc.failure except Exception: failure_record = ReportingWriterFailure( "RESOURCE_UNAVAILABLE", "same_identity", "unknown" ) if canceled: raise asyncio.CancelledError if failure_record is not None: raise ReportingWriterError(failure_record) async def write( self, prepared: ReportingPreparedRevision, *, context: ReportingIOContext ) -> ReportingDestinationLocator: self._validate_prepared(prepared) problem: ReportingWriterFailure | None = None canceled = False try: async with self._session(prepared, "write", context) as session: _session_binding(session, prepared.request, "write", context) locator = await context.run(lambda: session.write(prepared), effect="unknown") invalid_locator = False try: _locator_binding(locator, prepared.request) except ReportingWriterError: invalid_locator = True if invalid_locator: raise ReportingWriterError( ReportingWriterFailure("BINDING_MISMATCH", "same_identity", "unknown") ) return locator except asyncio.CancelledError: canceled = True except ReportingWriterError as exc: problem = exc.failure if canceled: raise asyncio.CancelledError assert problem is not None raise ReportingWriterError(problem) def _validate_prepared(self, prepared: ReportingPreparedRevision) -> ReportingRevisionVerifier: verifier = self.registry.require(prepared.request.verification_key) request, obligation, binding = prepared.request, prepared.obligation, prepared.binding cap = request.verification_key.capability if ( binding_fingerprint(binding) != request.binding_fingerprint or request.principal != binding.principal or request.generation != binding.generation_key or request.destination_ref != binding.destination_ref or request.trusted_binding_ref != binding.trusted_binding_ref or request.generation != obligation.generation_key or request.reporting_obligation_id != obligation.reporting_obligation_id or request.reporting_revision_id != prepared.revision.reporting_revision_id or prepared.revision.account_id != obligation.account_id or prepared.revision.reporting_obligation_id != obligation.reporting_obligation_id or not prepared.revision.readable or prepared.delivery.scope.principal != request.principal or prepared.delivery.scope.generation_key != request.generation or prepared.delivery.scope.reporting_obligation_id != request.reporting_obligation_id or prepared.delivery.currency != obligation.currency or not _same_definition(verifier.key, obligation.definition) or (obligation.report_definition_id, obligation.reporting_profile) != (verifier.key.report_definition_id, verifier.key.reporting_profile) or (binding.method, binding.transport, binding.format, binding.verification_profile) != (cap.method, cap.transport, cap.format, cap.verification_profile) or (binding.success_status == "delivered" and cap.verification_path != "destination") ): raise failure("BINDING_MISMATCH") rows, totals = verifier.canonicalize( [parse_reporting_json(r, verifier.limits) for r in prepared.rows] ) if rows != prepared.rows: raise failure("SOURCE_INVALID") _verify_content(prepared.revision, verifier.key, rows, totals) return verifier async def verify( self, prepared: ReportingPreparedRevision, locator: ReportingDestinationLocator, *, context: ReportingIOContext, ) -> ReportingVerifiedDestination: verifier = self._validate_prepared(prepared) _locator_binding(locator, prepared.request) problem: ReportingWriterFailure | None = None canceled = False try: async with self._session(prepared, "readback", context) as session: _session_binding(session, prepared.request, "readback", context) return await _verify_destination(verifier, prepared, locator, session, context) except asyncio.CancelledError: canceled = True except ReportingWriterError as exc: problem = exc.failure if canceled: raise asyncio.CancelledError assert problem is not None if problem.code == "SOURCE_INVALID": problem = ReportingWriterFailure("DESTINATION_CORRUPT", "new_attempt", "applied") raise ReportingWriterError(problem)Explicit one-shot write/readback operations; no coordinator or lease loop.
Instance variables
var registry : ReportingRevisionVerifierRegistry-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingDestinationIO: """Explicit one-shot write/readback operations; no coordinator or lease loop.""" registry: ReportingRevisionVerifierRegistry resolver: ReportingDestinationResolver = field(repr=False) @asynccontextmanager async def _session( self, prepared: ReportingPreparedRevision, phase: ReportingIOPhase, context: ReportingIOContext, ) -> AsyncIterator[ReportingDestinationSession]: self.registry.require(prepared.request.verification_key) problem = False canceled = False try: session = self.resolver.resolve(prepared.request, phase=phase, context=context) except asyncio.CancelledError: canceled = True except Exception: problem = True if canceled: raise asyncio.CancelledError if problem: raise ReportingWriterError( ReportingWriterFailure("RESOURCE_UNAVAILABLE", "same_identity") ) if not isinstance(session, ReportingDestinationSession): raise failure("BINDING_MISMATCH") try: _session_binding(session, prepared.request, phase, context) except ReportingWriterError: problem = True if problem: await session.aclose() raise failure("BINDING_MISMATCH") failure_record: ReportingWriterFailure | None = None canceled = False try: async with session: _session_binding(session, prepared.request, phase, context) yield session except asyncio.CancelledError: canceled = True except ReportingWriterError as exc: failure_record = exc.failure except Exception: failure_record = ReportingWriterFailure( "RESOURCE_UNAVAILABLE", "same_identity", "unknown" ) if canceled: raise asyncio.CancelledError if failure_record is not None: raise ReportingWriterError(failure_record) async def write( self, prepared: ReportingPreparedRevision, *, context: ReportingIOContext ) -> ReportingDestinationLocator: self._validate_prepared(prepared) problem: ReportingWriterFailure | None = None canceled = False try: async with self._session(prepared, "write", context) as session: _session_binding(session, prepared.request, "write", context) locator = await context.run(lambda: session.write(prepared), effect="unknown") invalid_locator = False try: _locator_binding(locator, prepared.request) except ReportingWriterError: invalid_locator = True if invalid_locator: raise ReportingWriterError( ReportingWriterFailure("BINDING_MISMATCH", "same_identity", "unknown") ) return locator except asyncio.CancelledError: canceled = True except ReportingWriterError as exc: problem = exc.failure if canceled: raise asyncio.CancelledError assert problem is not None raise ReportingWriterError(problem) def _validate_prepared(self, prepared: ReportingPreparedRevision) -> ReportingRevisionVerifier: verifier = self.registry.require(prepared.request.verification_key) request, obligation, binding = prepared.request, prepared.obligation, prepared.binding cap = request.verification_key.capability if ( binding_fingerprint(binding) != request.binding_fingerprint or request.principal != binding.principal or request.generation != binding.generation_key or request.destination_ref != binding.destination_ref or request.trusted_binding_ref != binding.trusted_binding_ref or request.generation != obligation.generation_key or request.reporting_obligation_id != obligation.reporting_obligation_id or request.reporting_revision_id != prepared.revision.reporting_revision_id or prepared.revision.account_id != obligation.account_id or prepared.revision.reporting_obligation_id != obligation.reporting_obligation_id or not prepared.revision.readable or prepared.delivery.scope.principal != request.principal or prepared.delivery.scope.generation_key != request.generation or prepared.delivery.scope.reporting_obligation_id != request.reporting_obligation_id or prepared.delivery.currency != obligation.currency or not _same_definition(verifier.key, obligation.definition) or (obligation.report_definition_id, obligation.reporting_profile) != (verifier.key.report_definition_id, verifier.key.reporting_profile) or (binding.method, binding.transport, binding.format, binding.verification_profile) != (cap.method, cap.transport, cap.format, cap.verification_profile) or (binding.success_status == "delivered" and cap.verification_path != "destination") ): raise failure("BINDING_MISMATCH") rows, totals = verifier.canonicalize( [parse_reporting_json(r, verifier.limits) for r in prepared.rows] ) if rows != prepared.rows: raise failure("SOURCE_INVALID") _verify_content(prepared.revision, verifier.key, rows, totals) return verifier async def verify( self, prepared: ReportingPreparedRevision, locator: ReportingDestinationLocator, *, context: ReportingIOContext, ) -> ReportingVerifiedDestination: verifier = self._validate_prepared(prepared) _locator_binding(locator, prepared.request) problem: ReportingWriterFailure | None = None canceled = False try: async with self._session(prepared, "readback", context) as session: _session_binding(session, prepared.request, "readback", context) return await _verify_destination(verifier, prepared, locator, session, context) except asyncio.CancelledError: canceled = True except ReportingWriterError as exc: problem = exc.failure if canceled: raise asyncio.CancelledError assert problem is not None if problem.code == "SOURCE_INVALID": problem = ReportingWriterFailure("DESTINATION_CORRUPT", "new_attempt", "applied") raise ReportingWriterError(problem) var resolver : ReportingDestinationResolver-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingDestinationIO: """Explicit one-shot write/readback operations; no coordinator or lease loop.""" registry: ReportingRevisionVerifierRegistry resolver: ReportingDestinationResolver = field(repr=False) @asynccontextmanager async def _session( self, prepared: ReportingPreparedRevision, phase: ReportingIOPhase, context: ReportingIOContext, ) -> AsyncIterator[ReportingDestinationSession]: self.registry.require(prepared.request.verification_key) problem = False canceled = False try: session = self.resolver.resolve(prepared.request, phase=phase, context=context) except asyncio.CancelledError: canceled = True except Exception: problem = True if canceled: raise asyncio.CancelledError if problem: raise ReportingWriterError( ReportingWriterFailure("RESOURCE_UNAVAILABLE", "same_identity") ) if not isinstance(session, ReportingDestinationSession): raise failure("BINDING_MISMATCH") try: _session_binding(session, prepared.request, phase, context) except ReportingWriterError: problem = True if problem: await session.aclose() raise failure("BINDING_MISMATCH") failure_record: ReportingWriterFailure | None = None canceled = False try: async with session: _session_binding(session, prepared.request, phase, context) yield session except asyncio.CancelledError: canceled = True except ReportingWriterError as exc: failure_record = exc.failure except Exception: failure_record = ReportingWriterFailure( "RESOURCE_UNAVAILABLE", "same_identity", "unknown" ) if canceled: raise asyncio.CancelledError if failure_record is not None: raise ReportingWriterError(failure_record) async def write( self, prepared: ReportingPreparedRevision, *, context: ReportingIOContext ) -> ReportingDestinationLocator: self._validate_prepared(prepared) problem: ReportingWriterFailure | None = None canceled = False try: async with self._session(prepared, "write", context) as session: _session_binding(session, prepared.request, "write", context) locator = await context.run(lambda: session.write(prepared), effect="unknown") invalid_locator = False try: _locator_binding(locator, prepared.request) except ReportingWriterError: invalid_locator = True if invalid_locator: raise ReportingWriterError( ReportingWriterFailure("BINDING_MISMATCH", "same_identity", "unknown") ) return locator except asyncio.CancelledError: canceled = True except ReportingWriterError as exc: problem = exc.failure if canceled: raise asyncio.CancelledError assert problem is not None raise ReportingWriterError(problem) def _validate_prepared(self, prepared: ReportingPreparedRevision) -> ReportingRevisionVerifier: verifier = self.registry.require(prepared.request.verification_key) request, obligation, binding = prepared.request, prepared.obligation, prepared.binding cap = request.verification_key.capability if ( binding_fingerprint(binding) != request.binding_fingerprint or request.principal != binding.principal or request.generation != binding.generation_key or request.destination_ref != binding.destination_ref or request.trusted_binding_ref != binding.trusted_binding_ref or request.generation != obligation.generation_key or request.reporting_obligation_id != obligation.reporting_obligation_id or request.reporting_revision_id != prepared.revision.reporting_revision_id or prepared.revision.account_id != obligation.account_id or prepared.revision.reporting_obligation_id != obligation.reporting_obligation_id or not prepared.revision.readable or prepared.delivery.scope.principal != request.principal or prepared.delivery.scope.generation_key != request.generation or prepared.delivery.scope.reporting_obligation_id != request.reporting_obligation_id or prepared.delivery.currency != obligation.currency or not _same_definition(verifier.key, obligation.definition) or (obligation.report_definition_id, obligation.reporting_profile) != (verifier.key.report_definition_id, verifier.key.reporting_profile) or (binding.method, binding.transport, binding.format, binding.verification_profile) != (cap.method, cap.transport, cap.format, cap.verification_profile) or (binding.success_status == "delivered" and cap.verification_path != "destination") ): raise failure("BINDING_MISMATCH") rows, totals = verifier.canonicalize( [parse_reporting_json(r, verifier.limits) for r in prepared.rows] ) if rows != prepared.rows: raise failure("SOURCE_INVALID") _verify_content(prepared.revision, verifier.key, rows, totals) return verifier async def verify( self, prepared: ReportingPreparedRevision, locator: ReportingDestinationLocator, *, context: ReportingIOContext, ) -> ReportingVerifiedDestination: verifier = self._validate_prepared(prepared) _locator_binding(locator, prepared.request) problem: ReportingWriterFailure | None = None canceled = False try: async with self._session(prepared, "readback", context) as session: _session_binding(session, prepared.request, "readback", context) return await _verify_destination(verifier, prepared, locator, session, context) except asyncio.CancelledError: canceled = True except ReportingWriterError as exc: problem = exc.failure if canceled: raise asyncio.CancelledError assert problem is not None if problem.code == "SOURCE_INVALID": problem = ReportingWriterFailure("DESTINATION_CORRUPT", "new_attempt", "applied") raise ReportingWriterError(problem)
Methods
async def verify(self,
prepared: ReportingPreparedRevision,
locator: ReportingDestinationLocator,
*,
context: ReportingIOContext) ‑> ReportingVerifiedDestination-
Expand source code
async def verify( self, prepared: ReportingPreparedRevision, locator: ReportingDestinationLocator, *, context: ReportingIOContext, ) -> ReportingVerifiedDestination: verifier = self._validate_prepared(prepared) _locator_binding(locator, prepared.request) problem: ReportingWriterFailure | None = None canceled = False try: async with self._session(prepared, "readback", context) as session: _session_binding(session, prepared.request, "readback", context) return await _verify_destination(verifier, prepared, locator, session, context) except asyncio.CancelledError: canceled = True except ReportingWriterError as exc: problem = exc.failure if canceled: raise asyncio.CancelledError assert problem is not None if problem.code == "SOURCE_INVALID": problem = ReportingWriterFailure("DESTINATION_CORRUPT", "new_attempt", "applied") raise ReportingWriterError(problem) async def write(self, prepared: ReportingPreparedRevision, *, context: ReportingIOContext) ‑> ReportingDestinationLocator-
Expand source code
async def write( self, prepared: ReportingPreparedRevision, *, context: ReportingIOContext ) -> ReportingDestinationLocator: self._validate_prepared(prepared) problem: ReportingWriterFailure | None = None canceled = False try: async with self._session(prepared, "write", context) as session: _session_binding(session, prepared.request, "write", context) locator = await context.run(lambda: session.write(prepared), effect="unknown") invalid_locator = False try: _locator_binding(locator, prepared.request) except ReportingWriterError: invalid_locator = True if invalid_locator: raise ReportingWriterError( ReportingWriterFailure("BINDING_MISMATCH", "same_identity", "unknown") ) return locator except asyncio.CancelledError: canceled = True except ReportingWriterError as exc: problem = exc.failure if canceled: raise asyncio.CancelledError assert problem is not None raise ReportingWriterError(problem)
class ReportingRevisionRowReader (*args, **kwargs)-
Expand source code
class ReportingRevisionRowReader(Protocol): async def read_revision_rows( self, *, account_id: str, reporting_revision_id: str, cursor: str | None = None, limit: int = 500, ) -> ReportingRowPage: ...Base class for protocol classes.
Protocol classes are defined as::
class Proto(Protocol): def meth(self) -> int: ...Such classes are primarily used with static type checkers that recognize structural subtyping (static duck-typing).
For example::
class C: def meth(self) -> int: return 0 def func(x: Proto) -> int: return x.meth() func(C()) # Passes static type checkSee PEP 544 for details. Protocol classes decorated with @typing.runtime_checkable act as simple-minded runtime protocols that check only the presence of given attributes, ignoring their type signatures. Protocol classes can be generic, they are defined as::
class GenProto(Protocol[T]): def meth(self) -> T: ...Ancestors
- typing.Protocol
- typing.Generic
Subclasses
Methods
async def read_revision_rows(self,
*,
account_id: str,
reporting_revision_id: str,
cursor: str | None = None,
limit: int = 500) ‑> ReportingRowPage-
Expand source code
async def read_revision_rows( self, *, account_id: str, reporting_revision_id: str, cursor: str | None = None, limit: int = 500, ) -> ReportingRowPage: ...
class ReportingRevisionVerifier (key: ReportingVerificationKey,
definition_bytes: bytes,
schema_bytes: bytes,
canonicalization_bytes: bytes,
limits: ReportingVerificationLimits = <factory>)-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingRevisionVerifier(_ClosedValue): """One exact installed contract using the SDK's safe-integer JCS subset. Sum totals are declared on schema properties with ``x-adcp-control-total`` (value_type, optional unit, decimal scale). This explicit pinned subset has no expression evaluator, custom callbacks, mutable registration or network. Unsupported formats and definition semantics fail during construction. """ key: ReportingVerificationKey definition_bytes: bytes = field(repr=False) schema_bytes: bytes = field(repr=False) canonicalization_bytes: bytes = field(repr=False) limits: ReportingVerificationLimits = field(default_factory=ReportingVerificationLimits) def __post_init__(self) -> None: _freeze_fields(self) error = False try: self._runtime() except Exception: error = True if error: raise failure("UNSUPPORTED_VERIFICATION") def _runtime(self) -> _Canonicalizer: key, definition = self.key, self.key.definition if key.capability.format not in {"jsonl", None}: raise failure("UNSUPPORTED_VERIFICATION") if ( key.capability.method == "file_transfer" and key.capability.immutability != "immutable_location" ): raise failure("UNSUPPORTED_VERIFICATION") if ( definition.schema_dialect != "https://json-schema.org/draft/2020-12/schema" or definition.schema_ref_policy != "local_fragment_only" ): raise failure("UNSUPPORTED_VERIFICATION") for raw, expected in ( (self.definition_bytes, definition.report_definition_sha256), (self.schema_bytes, definition.schema_sha256), (self.canonicalization_bytes, key.canonicalization.canonicalization_sha256), ): if hashlib.sha256(raw).hexdigest() != expected.lower(): raise failure("UNSUPPORTED_VERIFICATION") report = _mapping(parse_reporting_json(self.definition_bytes, self.limits)) schema = _mapping(parse_reporting_json(self.schema_bytes, self.limits)) contract = _mapping(parse_reporting_json(self.canonicalization_bytes, self.limits)) if ( report.get("report_definition_id") != key.report_definition_id or report.get("reporting_profile") != key.reporting_profile or schema.get("$schema") != definition.schema_dialect or set(contract) != { "contract_version", "media_type", "algorithm", "schema_sha256", "primary_keys", "golden_vectors", } or contract["contract_version"] != "1.0" or contract["media_type"] != "application/vnd.adcp.reporting-canonicalization+json" or contract["algorithm"] != "adcp_jcs_rows_v1" or contract["schema_sha256"].lower() != definition.schema_sha256.lower() ): raise failure("UNSUPPORTED_VERIFICATION") _local_schema(schema) Draft202012Validator.check_schema(schema) keys = contract["primary_keys"] if ( type(keys) is not list or not keys or any(type(k) is not str for k in keys) or len(set(keys)) != len(keys) ): raise failure("UNSUPPORTED_VERIFICATION") totals: list[_Total] = [] properties = _mapping(schema.get("properties")) metrics = report.get("metrics") if type(metrics) is not list or not metrics: raise failure("UNSUPPORTED_VERIFICATION") for metric in metrics: metric = _mapping(metric) name = metric["name"] prop = _mapping(properties[name]) total = _mapping(prop["x-adcp-control-total"]) if ( metric.get("source_expression") != name or metric.get("aggregation") != "sum" or set(total) - {"value_type", "unit", "scale"} or total.get("unit") != metric.get("unit") or (total["value_type"], prop.get("type")) not in {("integer", "integer"), ("decimal", "string")} ): raise failure("UNSUPPORTED_VERIFICATION") scale = total.get("scale", 0) if ( type(scale) is not int or not 0 <= scale <= 18 or (total["value_type"] == "integer" and scale != 0) ): raise failure("UNSUPPORTED_VERIFICATION") totals.append(_Total(name, total["value_type"], total.get("unit"), scale)) if len({t.name for t in totals}) != len(totals): raise failure("UNSUPPORTED_VERIFICATION") if {t.name for t in totals} != { name for name, prop in properties.items() if type(prop) is dict and "x-adcp-control-total" in prop }: raise failure("UNSUPPORTED_VERIFICATION") if not set(keys).union(t.name for t in totals) <= set(schema.get("required", [])): raise failure("UNSUPPORTED_VERIFICATION") units = {t.name: t.unit for t in totals} if any( units.get(name) != unit for name, unit in ( *definition.monetary_metric_units, *definition.monetary_control_total_units, ) ): raise failure("UNSUPPORTED_VERIFICATION") runtime = _Canonicalizer( tuple(keys), tuple(totals), Draft202012Validator(schema), self.limits ) vectors = _mapping(contract["golden_vectors"]) if not {"empty_report", "ordering_encoding"} <= set(vectors) or set(vectors) - { "empty_report", "ordering_encoding", "additional", }: raise failure("UNSUPPORTED_VERIFICATION") seen: set[str] = set() for vector in [ vectors["empty_report"], vectors["ordering_encoding"], *vectors.get("additional", []), ]: vector = _mapping(vector) if ( set(vector) != {"name", "purpose", "input_rows", "canonical_utf8_base64", "sha256"} or vector["name"] in seen ): raise failure("UNSUPPORTED_VERIFICATION") seen.add(vector["name"]) rows, _ = runtime.canonicalize(vector["input_rows"]) golden_bytes = base64.b64decode(vector["canonical_utf8_base64"], validate=True) if ( b"[" + b",".join(rows) + b"]" != golden_bytes or hashlib.sha256(golden_bytes).hexdigest() != vector["sha256"].lower() ): raise failure("UNSUPPORTED_VERIFICATION") if ( vectors["empty_report"]["input_rows"] != [] or vectors["empty_report"]["purpose"] != "empty_report" ): raise failure("UNSUPPORTED_VERIFICATION") ordering = vectors["ordering_encoding"] order_rows = ordering["input_rows"] if ( ordering["purpose"] != "ordering_encoding" or len(order_rows) < 2 or [runtime.order_key(r) for r in order_rows] == sorted(runtime.order_key(r) for r in order_rows) or not _unordered_members(order_rows) ): raise failure("UNSUPPORTED_VERIFICATION") return runtime def canonicalize( self, rows: Sequence[object] ) -> tuple[tuple[bytes, ...], tuple[ReportingControlTotalRecord, ...]]: """Useful to trusted publishers preparing expected evidence before storage.""" return self._runtime().canonicalize(rows)One exact installed contract using the SDK's safe-integer JCS subset.
Sum totals are declared on schema properties with
x-adcp-control-total(value_type, optional unit, decimal scale). This explicit pinned subset has no expression evaluator, custom callbacks, mutable registration or network. Unsupported formats and definition semantics fail during construction.Ancestors
- adcp.reporting.ledger.delivery_models._ClosedValue
Instance variables
var canonicalization_bytes : bytes-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingRevisionVerifier(_ClosedValue): """One exact installed contract using the SDK's safe-integer JCS subset. Sum totals are declared on schema properties with ``x-adcp-control-total`` (value_type, optional unit, decimal scale). This explicit pinned subset has no expression evaluator, custom callbacks, mutable registration or network. Unsupported formats and definition semantics fail during construction. """ key: ReportingVerificationKey definition_bytes: bytes = field(repr=False) schema_bytes: bytes = field(repr=False) canonicalization_bytes: bytes = field(repr=False) limits: ReportingVerificationLimits = field(default_factory=ReportingVerificationLimits) def __post_init__(self) -> None: _freeze_fields(self) error = False try: self._runtime() except Exception: error = True if error: raise failure("UNSUPPORTED_VERIFICATION") def _runtime(self) -> _Canonicalizer: key, definition = self.key, self.key.definition if key.capability.format not in {"jsonl", None}: raise failure("UNSUPPORTED_VERIFICATION") if ( key.capability.method == "file_transfer" and key.capability.immutability != "immutable_location" ): raise failure("UNSUPPORTED_VERIFICATION") if ( definition.schema_dialect != "https://json-schema.org/draft/2020-12/schema" or definition.schema_ref_policy != "local_fragment_only" ): raise failure("UNSUPPORTED_VERIFICATION") for raw, expected in ( (self.definition_bytes, definition.report_definition_sha256), (self.schema_bytes, definition.schema_sha256), (self.canonicalization_bytes, key.canonicalization.canonicalization_sha256), ): if hashlib.sha256(raw).hexdigest() != expected.lower(): raise failure("UNSUPPORTED_VERIFICATION") report = _mapping(parse_reporting_json(self.definition_bytes, self.limits)) schema = _mapping(parse_reporting_json(self.schema_bytes, self.limits)) contract = _mapping(parse_reporting_json(self.canonicalization_bytes, self.limits)) if ( report.get("report_definition_id") != key.report_definition_id or report.get("reporting_profile") != key.reporting_profile or schema.get("$schema") != definition.schema_dialect or set(contract) != { "contract_version", "media_type", "algorithm", "schema_sha256", "primary_keys", "golden_vectors", } or contract["contract_version"] != "1.0" or contract["media_type"] != "application/vnd.adcp.reporting-canonicalization+json" or contract["algorithm"] != "adcp_jcs_rows_v1" or contract["schema_sha256"].lower() != definition.schema_sha256.lower() ): raise failure("UNSUPPORTED_VERIFICATION") _local_schema(schema) Draft202012Validator.check_schema(schema) keys = contract["primary_keys"] if ( type(keys) is not list or not keys or any(type(k) is not str for k in keys) or len(set(keys)) != len(keys) ): raise failure("UNSUPPORTED_VERIFICATION") totals: list[_Total] = [] properties = _mapping(schema.get("properties")) metrics = report.get("metrics") if type(metrics) is not list or not metrics: raise failure("UNSUPPORTED_VERIFICATION") for metric in metrics: metric = _mapping(metric) name = metric["name"] prop = _mapping(properties[name]) total = _mapping(prop["x-adcp-control-total"]) if ( metric.get("source_expression") != name or metric.get("aggregation") != "sum" or set(total) - {"value_type", "unit", "scale"} or total.get("unit") != metric.get("unit") or (total["value_type"], prop.get("type")) not in {("integer", "integer"), ("decimal", "string")} ): raise failure("UNSUPPORTED_VERIFICATION") scale = total.get("scale", 0) if ( type(scale) is not int or not 0 <= scale <= 18 or (total["value_type"] == "integer" and scale != 0) ): raise failure("UNSUPPORTED_VERIFICATION") totals.append(_Total(name, total["value_type"], total.get("unit"), scale)) if len({t.name for t in totals}) != len(totals): raise failure("UNSUPPORTED_VERIFICATION") if {t.name for t in totals} != { name for name, prop in properties.items() if type(prop) is dict and "x-adcp-control-total" in prop }: raise failure("UNSUPPORTED_VERIFICATION") if not set(keys).union(t.name for t in totals) <= set(schema.get("required", [])): raise failure("UNSUPPORTED_VERIFICATION") units = {t.name: t.unit for t in totals} if any( units.get(name) != unit for name, unit in ( *definition.monetary_metric_units, *definition.monetary_control_total_units, ) ): raise failure("UNSUPPORTED_VERIFICATION") runtime = _Canonicalizer( tuple(keys), tuple(totals), Draft202012Validator(schema), self.limits ) vectors = _mapping(contract["golden_vectors"]) if not {"empty_report", "ordering_encoding"} <= set(vectors) or set(vectors) - { "empty_report", "ordering_encoding", "additional", }: raise failure("UNSUPPORTED_VERIFICATION") seen: set[str] = set() for vector in [ vectors["empty_report"], vectors["ordering_encoding"], *vectors.get("additional", []), ]: vector = _mapping(vector) if ( set(vector) != {"name", "purpose", "input_rows", "canonical_utf8_base64", "sha256"} or vector["name"] in seen ): raise failure("UNSUPPORTED_VERIFICATION") seen.add(vector["name"]) rows, _ = runtime.canonicalize(vector["input_rows"]) golden_bytes = base64.b64decode(vector["canonical_utf8_base64"], validate=True) if ( b"[" + b",".join(rows) + b"]" != golden_bytes or hashlib.sha256(golden_bytes).hexdigest() != vector["sha256"].lower() ): raise failure("UNSUPPORTED_VERIFICATION") if ( vectors["empty_report"]["input_rows"] != [] or vectors["empty_report"]["purpose"] != "empty_report" ): raise failure("UNSUPPORTED_VERIFICATION") ordering = vectors["ordering_encoding"] order_rows = ordering["input_rows"] if ( ordering["purpose"] != "ordering_encoding" or len(order_rows) < 2 or [runtime.order_key(r) for r in order_rows] == sorted(runtime.order_key(r) for r in order_rows) or not _unordered_members(order_rows) ): raise failure("UNSUPPORTED_VERIFICATION") return runtime def canonicalize( self, rows: Sequence[object] ) -> tuple[tuple[bytes, ...], tuple[ReportingControlTotalRecord, ...]]: """Useful to trusted publishers preparing expected evidence before storage.""" return self._runtime().canonicalize(rows) var definition_bytes : bytes-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingRevisionVerifier(_ClosedValue): """One exact installed contract using the SDK's safe-integer JCS subset. Sum totals are declared on schema properties with ``x-adcp-control-total`` (value_type, optional unit, decimal scale). This explicit pinned subset has no expression evaluator, custom callbacks, mutable registration or network. Unsupported formats and definition semantics fail during construction. """ key: ReportingVerificationKey definition_bytes: bytes = field(repr=False) schema_bytes: bytes = field(repr=False) canonicalization_bytes: bytes = field(repr=False) limits: ReportingVerificationLimits = field(default_factory=ReportingVerificationLimits) def __post_init__(self) -> None: _freeze_fields(self) error = False try: self._runtime() except Exception: error = True if error: raise failure("UNSUPPORTED_VERIFICATION") def _runtime(self) -> _Canonicalizer: key, definition = self.key, self.key.definition if key.capability.format not in {"jsonl", None}: raise failure("UNSUPPORTED_VERIFICATION") if ( key.capability.method == "file_transfer" and key.capability.immutability != "immutable_location" ): raise failure("UNSUPPORTED_VERIFICATION") if ( definition.schema_dialect != "https://json-schema.org/draft/2020-12/schema" or definition.schema_ref_policy != "local_fragment_only" ): raise failure("UNSUPPORTED_VERIFICATION") for raw, expected in ( (self.definition_bytes, definition.report_definition_sha256), (self.schema_bytes, definition.schema_sha256), (self.canonicalization_bytes, key.canonicalization.canonicalization_sha256), ): if hashlib.sha256(raw).hexdigest() != expected.lower(): raise failure("UNSUPPORTED_VERIFICATION") report = _mapping(parse_reporting_json(self.definition_bytes, self.limits)) schema = _mapping(parse_reporting_json(self.schema_bytes, self.limits)) contract = _mapping(parse_reporting_json(self.canonicalization_bytes, self.limits)) if ( report.get("report_definition_id") != key.report_definition_id or report.get("reporting_profile") != key.reporting_profile or schema.get("$schema") != definition.schema_dialect or set(contract) != { "contract_version", "media_type", "algorithm", "schema_sha256", "primary_keys", "golden_vectors", } or contract["contract_version"] != "1.0" or contract["media_type"] != "application/vnd.adcp.reporting-canonicalization+json" or contract["algorithm"] != "adcp_jcs_rows_v1" or contract["schema_sha256"].lower() != definition.schema_sha256.lower() ): raise failure("UNSUPPORTED_VERIFICATION") _local_schema(schema) Draft202012Validator.check_schema(schema) keys = contract["primary_keys"] if ( type(keys) is not list or not keys or any(type(k) is not str for k in keys) or len(set(keys)) != len(keys) ): raise failure("UNSUPPORTED_VERIFICATION") totals: list[_Total] = [] properties = _mapping(schema.get("properties")) metrics = report.get("metrics") if type(metrics) is not list or not metrics: raise failure("UNSUPPORTED_VERIFICATION") for metric in metrics: metric = _mapping(metric) name = metric["name"] prop = _mapping(properties[name]) total = _mapping(prop["x-adcp-control-total"]) if ( metric.get("source_expression") != name or metric.get("aggregation") != "sum" or set(total) - {"value_type", "unit", "scale"} or total.get("unit") != metric.get("unit") or (total["value_type"], prop.get("type")) not in {("integer", "integer"), ("decimal", "string")} ): raise failure("UNSUPPORTED_VERIFICATION") scale = total.get("scale", 0) if ( type(scale) is not int or not 0 <= scale <= 18 or (total["value_type"] == "integer" and scale != 0) ): raise failure("UNSUPPORTED_VERIFICATION") totals.append(_Total(name, total["value_type"], total.get("unit"), scale)) if len({t.name for t in totals}) != len(totals): raise failure("UNSUPPORTED_VERIFICATION") if {t.name for t in totals} != { name for name, prop in properties.items() if type(prop) is dict and "x-adcp-control-total" in prop }: raise failure("UNSUPPORTED_VERIFICATION") if not set(keys).union(t.name for t in totals) <= set(schema.get("required", [])): raise failure("UNSUPPORTED_VERIFICATION") units = {t.name: t.unit for t in totals} if any( units.get(name) != unit for name, unit in ( *definition.monetary_metric_units, *definition.monetary_control_total_units, ) ): raise failure("UNSUPPORTED_VERIFICATION") runtime = _Canonicalizer( tuple(keys), tuple(totals), Draft202012Validator(schema), self.limits ) vectors = _mapping(contract["golden_vectors"]) if not {"empty_report", "ordering_encoding"} <= set(vectors) or set(vectors) - { "empty_report", "ordering_encoding", "additional", }: raise failure("UNSUPPORTED_VERIFICATION") seen: set[str] = set() for vector in [ vectors["empty_report"], vectors["ordering_encoding"], *vectors.get("additional", []), ]: vector = _mapping(vector) if ( set(vector) != {"name", "purpose", "input_rows", "canonical_utf8_base64", "sha256"} or vector["name"] in seen ): raise failure("UNSUPPORTED_VERIFICATION") seen.add(vector["name"]) rows, _ = runtime.canonicalize(vector["input_rows"]) golden_bytes = base64.b64decode(vector["canonical_utf8_base64"], validate=True) if ( b"[" + b",".join(rows) + b"]" != golden_bytes or hashlib.sha256(golden_bytes).hexdigest() != vector["sha256"].lower() ): raise failure("UNSUPPORTED_VERIFICATION") if ( vectors["empty_report"]["input_rows"] != [] or vectors["empty_report"]["purpose"] != "empty_report" ): raise failure("UNSUPPORTED_VERIFICATION") ordering = vectors["ordering_encoding"] order_rows = ordering["input_rows"] if ( ordering["purpose"] != "ordering_encoding" or len(order_rows) < 2 or [runtime.order_key(r) for r in order_rows] == sorted(runtime.order_key(r) for r in order_rows) or not _unordered_members(order_rows) ): raise failure("UNSUPPORTED_VERIFICATION") return runtime def canonicalize( self, rows: Sequence[object] ) -> tuple[tuple[bytes, ...], tuple[ReportingControlTotalRecord, ...]]: """Useful to trusted publishers preparing expected evidence before storage.""" return self._runtime().canonicalize(rows) var key : ReportingVerificationKey-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingRevisionVerifier(_ClosedValue): """One exact installed contract using the SDK's safe-integer JCS subset. Sum totals are declared on schema properties with ``x-adcp-control-total`` (value_type, optional unit, decimal scale). This explicit pinned subset has no expression evaluator, custom callbacks, mutable registration or network. Unsupported formats and definition semantics fail during construction. """ key: ReportingVerificationKey definition_bytes: bytes = field(repr=False) schema_bytes: bytes = field(repr=False) canonicalization_bytes: bytes = field(repr=False) limits: ReportingVerificationLimits = field(default_factory=ReportingVerificationLimits) def __post_init__(self) -> None: _freeze_fields(self) error = False try: self._runtime() except Exception: error = True if error: raise failure("UNSUPPORTED_VERIFICATION") def _runtime(self) -> _Canonicalizer: key, definition = self.key, self.key.definition if key.capability.format not in {"jsonl", None}: raise failure("UNSUPPORTED_VERIFICATION") if ( key.capability.method == "file_transfer" and key.capability.immutability != "immutable_location" ): raise failure("UNSUPPORTED_VERIFICATION") if ( definition.schema_dialect != "https://json-schema.org/draft/2020-12/schema" or definition.schema_ref_policy != "local_fragment_only" ): raise failure("UNSUPPORTED_VERIFICATION") for raw, expected in ( (self.definition_bytes, definition.report_definition_sha256), (self.schema_bytes, definition.schema_sha256), (self.canonicalization_bytes, key.canonicalization.canonicalization_sha256), ): if hashlib.sha256(raw).hexdigest() != expected.lower(): raise failure("UNSUPPORTED_VERIFICATION") report = _mapping(parse_reporting_json(self.definition_bytes, self.limits)) schema = _mapping(parse_reporting_json(self.schema_bytes, self.limits)) contract = _mapping(parse_reporting_json(self.canonicalization_bytes, self.limits)) if ( report.get("report_definition_id") != key.report_definition_id or report.get("reporting_profile") != key.reporting_profile or schema.get("$schema") != definition.schema_dialect or set(contract) != { "contract_version", "media_type", "algorithm", "schema_sha256", "primary_keys", "golden_vectors", } or contract["contract_version"] != "1.0" or contract["media_type"] != "application/vnd.adcp.reporting-canonicalization+json" or contract["algorithm"] != "adcp_jcs_rows_v1" or contract["schema_sha256"].lower() != definition.schema_sha256.lower() ): raise failure("UNSUPPORTED_VERIFICATION") _local_schema(schema) Draft202012Validator.check_schema(schema) keys = contract["primary_keys"] if ( type(keys) is not list or not keys or any(type(k) is not str for k in keys) or len(set(keys)) != len(keys) ): raise failure("UNSUPPORTED_VERIFICATION") totals: list[_Total] = [] properties = _mapping(schema.get("properties")) metrics = report.get("metrics") if type(metrics) is not list or not metrics: raise failure("UNSUPPORTED_VERIFICATION") for metric in metrics: metric = _mapping(metric) name = metric["name"] prop = _mapping(properties[name]) total = _mapping(prop["x-adcp-control-total"]) if ( metric.get("source_expression") != name or metric.get("aggregation") != "sum" or set(total) - {"value_type", "unit", "scale"} or total.get("unit") != metric.get("unit") or (total["value_type"], prop.get("type")) not in {("integer", "integer"), ("decimal", "string")} ): raise failure("UNSUPPORTED_VERIFICATION") scale = total.get("scale", 0) if ( type(scale) is not int or not 0 <= scale <= 18 or (total["value_type"] == "integer" and scale != 0) ): raise failure("UNSUPPORTED_VERIFICATION") totals.append(_Total(name, total["value_type"], total.get("unit"), scale)) if len({t.name for t in totals}) != len(totals): raise failure("UNSUPPORTED_VERIFICATION") if {t.name for t in totals} != { name for name, prop in properties.items() if type(prop) is dict and "x-adcp-control-total" in prop }: raise failure("UNSUPPORTED_VERIFICATION") if not set(keys).union(t.name for t in totals) <= set(schema.get("required", [])): raise failure("UNSUPPORTED_VERIFICATION") units = {t.name: t.unit for t in totals} if any( units.get(name) != unit for name, unit in ( *definition.monetary_metric_units, *definition.monetary_control_total_units, ) ): raise failure("UNSUPPORTED_VERIFICATION") runtime = _Canonicalizer( tuple(keys), tuple(totals), Draft202012Validator(schema), self.limits ) vectors = _mapping(contract["golden_vectors"]) if not {"empty_report", "ordering_encoding"} <= set(vectors) or set(vectors) - { "empty_report", "ordering_encoding", "additional", }: raise failure("UNSUPPORTED_VERIFICATION") seen: set[str] = set() for vector in [ vectors["empty_report"], vectors["ordering_encoding"], *vectors.get("additional", []), ]: vector = _mapping(vector) if ( set(vector) != {"name", "purpose", "input_rows", "canonical_utf8_base64", "sha256"} or vector["name"] in seen ): raise failure("UNSUPPORTED_VERIFICATION") seen.add(vector["name"]) rows, _ = runtime.canonicalize(vector["input_rows"]) golden_bytes = base64.b64decode(vector["canonical_utf8_base64"], validate=True) if ( b"[" + b",".join(rows) + b"]" != golden_bytes or hashlib.sha256(golden_bytes).hexdigest() != vector["sha256"].lower() ): raise failure("UNSUPPORTED_VERIFICATION") if ( vectors["empty_report"]["input_rows"] != [] or vectors["empty_report"]["purpose"] != "empty_report" ): raise failure("UNSUPPORTED_VERIFICATION") ordering = vectors["ordering_encoding"] order_rows = ordering["input_rows"] if ( ordering["purpose"] != "ordering_encoding" or len(order_rows) < 2 or [runtime.order_key(r) for r in order_rows] == sorted(runtime.order_key(r) for r in order_rows) or not _unordered_members(order_rows) ): raise failure("UNSUPPORTED_VERIFICATION") return runtime def canonicalize( self, rows: Sequence[object] ) -> tuple[tuple[bytes, ...], tuple[ReportingControlTotalRecord, ...]]: """Useful to trusted publishers preparing expected evidence before storage.""" return self._runtime().canonicalize(rows) var limits : adcp.reporting.materializer._json.ReportingVerificationLimits-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingRevisionVerifier(_ClosedValue): """One exact installed contract using the SDK's safe-integer JCS subset. Sum totals are declared on schema properties with ``x-adcp-control-total`` (value_type, optional unit, decimal scale). This explicit pinned subset has no expression evaluator, custom callbacks, mutable registration or network. Unsupported formats and definition semantics fail during construction. """ key: ReportingVerificationKey definition_bytes: bytes = field(repr=False) schema_bytes: bytes = field(repr=False) canonicalization_bytes: bytes = field(repr=False) limits: ReportingVerificationLimits = field(default_factory=ReportingVerificationLimits) def __post_init__(self) -> None: _freeze_fields(self) error = False try: self._runtime() except Exception: error = True if error: raise failure("UNSUPPORTED_VERIFICATION") def _runtime(self) -> _Canonicalizer: key, definition = self.key, self.key.definition if key.capability.format not in {"jsonl", None}: raise failure("UNSUPPORTED_VERIFICATION") if ( key.capability.method == "file_transfer" and key.capability.immutability != "immutable_location" ): raise failure("UNSUPPORTED_VERIFICATION") if ( definition.schema_dialect != "https://json-schema.org/draft/2020-12/schema" or definition.schema_ref_policy != "local_fragment_only" ): raise failure("UNSUPPORTED_VERIFICATION") for raw, expected in ( (self.definition_bytes, definition.report_definition_sha256), (self.schema_bytes, definition.schema_sha256), (self.canonicalization_bytes, key.canonicalization.canonicalization_sha256), ): if hashlib.sha256(raw).hexdigest() != expected.lower(): raise failure("UNSUPPORTED_VERIFICATION") report = _mapping(parse_reporting_json(self.definition_bytes, self.limits)) schema = _mapping(parse_reporting_json(self.schema_bytes, self.limits)) contract = _mapping(parse_reporting_json(self.canonicalization_bytes, self.limits)) if ( report.get("report_definition_id") != key.report_definition_id or report.get("reporting_profile") != key.reporting_profile or schema.get("$schema") != definition.schema_dialect or set(contract) != { "contract_version", "media_type", "algorithm", "schema_sha256", "primary_keys", "golden_vectors", } or contract["contract_version"] != "1.0" or contract["media_type"] != "application/vnd.adcp.reporting-canonicalization+json" or contract["algorithm"] != "adcp_jcs_rows_v1" or contract["schema_sha256"].lower() != definition.schema_sha256.lower() ): raise failure("UNSUPPORTED_VERIFICATION") _local_schema(schema) Draft202012Validator.check_schema(schema) keys = contract["primary_keys"] if ( type(keys) is not list or not keys or any(type(k) is not str for k in keys) or len(set(keys)) != len(keys) ): raise failure("UNSUPPORTED_VERIFICATION") totals: list[_Total] = [] properties = _mapping(schema.get("properties")) metrics = report.get("metrics") if type(metrics) is not list or not metrics: raise failure("UNSUPPORTED_VERIFICATION") for metric in metrics: metric = _mapping(metric) name = metric["name"] prop = _mapping(properties[name]) total = _mapping(prop["x-adcp-control-total"]) if ( metric.get("source_expression") != name or metric.get("aggregation") != "sum" or set(total) - {"value_type", "unit", "scale"} or total.get("unit") != metric.get("unit") or (total["value_type"], prop.get("type")) not in {("integer", "integer"), ("decimal", "string")} ): raise failure("UNSUPPORTED_VERIFICATION") scale = total.get("scale", 0) if ( type(scale) is not int or not 0 <= scale <= 18 or (total["value_type"] == "integer" and scale != 0) ): raise failure("UNSUPPORTED_VERIFICATION") totals.append(_Total(name, total["value_type"], total.get("unit"), scale)) if len({t.name for t in totals}) != len(totals): raise failure("UNSUPPORTED_VERIFICATION") if {t.name for t in totals} != { name for name, prop in properties.items() if type(prop) is dict and "x-adcp-control-total" in prop }: raise failure("UNSUPPORTED_VERIFICATION") if not set(keys).union(t.name for t in totals) <= set(schema.get("required", [])): raise failure("UNSUPPORTED_VERIFICATION") units = {t.name: t.unit for t in totals} if any( units.get(name) != unit for name, unit in ( *definition.monetary_metric_units, *definition.monetary_control_total_units, ) ): raise failure("UNSUPPORTED_VERIFICATION") runtime = _Canonicalizer( tuple(keys), tuple(totals), Draft202012Validator(schema), self.limits ) vectors = _mapping(contract["golden_vectors"]) if not {"empty_report", "ordering_encoding"} <= set(vectors) or set(vectors) - { "empty_report", "ordering_encoding", "additional", }: raise failure("UNSUPPORTED_VERIFICATION") seen: set[str] = set() for vector in [ vectors["empty_report"], vectors["ordering_encoding"], *vectors.get("additional", []), ]: vector = _mapping(vector) if ( set(vector) != {"name", "purpose", "input_rows", "canonical_utf8_base64", "sha256"} or vector["name"] in seen ): raise failure("UNSUPPORTED_VERIFICATION") seen.add(vector["name"]) rows, _ = runtime.canonicalize(vector["input_rows"]) golden_bytes = base64.b64decode(vector["canonical_utf8_base64"], validate=True) if ( b"[" + b",".join(rows) + b"]" != golden_bytes or hashlib.sha256(golden_bytes).hexdigest() != vector["sha256"].lower() ): raise failure("UNSUPPORTED_VERIFICATION") if ( vectors["empty_report"]["input_rows"] != [] or vectors["empty_report"]["purpose"] != "empty_report" ): raise failure("UNSUPPORTED_VERIFICATION") ordering = vectors["ordering_encoding"] order_rows = ordering["input_rows"] if ( ordering["purpose"] != "ordering_encoding" or len(order_rows) < 2 or [runtime.order_key(r) for r in order_rows] == sorted(runtime.order_key(r) for r in order_rows) or not _unordered_members(order_rows) ): raise failure("UNSUPPORTED_VERIFICATION") return runtime def canonicalize( self, rows: Sequence[object] ) -> tuple[tuple[bytes, ...], tuple[ReportingControlTotalRecord, ...]]: """Useful to trusted publishers preparing expected evidence before storage.""" return self._runtime().canonicalize(rows) var schema_bytes : bytes-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingRevisionVerifier(_ClosedValue): """One exact installed contract using the SDK's safe-integer JCS subset. Sum totals are declared on schema properties with ``x-adcp-control-total`` (value_type, optional unit, decimal scale). This explicit pinned subset has no expression evaluator, custom callbacks, mutable registration or network. Unsupported formats and definition semantics fail during construction. """ key: ReportingVerificationKey definition_bytes: bytes = field(repr=False) schema_bytes: bytes = field(repr=False) canonicalization_bytes: bytes = field(repr=False) limits: ReportingVerificationLimits = field(default_factory=ReportingVerificationLimits) def __post_init__(self) -> None: _freeze_fields(self) error = False try: self._runtime() except Exception: error = True if error: raise failure("UNSUPPORTED_VERIFICATION") def _runtime(self) -> _Canonicalizer: key, definition = self.key, self.key.definition if key.capability.format not in {"jsonl", None}: raise failure("UNSUPPORTED_VERIFICATION") if ( key.capability.method == "file_transfer" and key.capability.immutability != "immutable_location" ): raise failure("UNSUPPORTED_VERIFICATION") if ( definition.schema_dialect != "https://json-schema.org/draft/2020-12/schema" or definition.schema_ref_policy != "local_fragment_only" ): raise failure("UNSUPPORTED_VERIFICATION") for raw, expected in ( (self.definition_bytes, definition.report_definition_sha256), (self.schema_bytes, definition.schema_sha256), (self.canonicalization_bytes, key.canonicalization.canonicalization_sha256), ): if hashlib.sha256(raw).hexdigest() != expected.lower(): raise failure("UNSUPPORTED_VERIFICATION") report = _mapping(parse_reporting_json(self.definition_bytes, self.limits)) schema = _mapping(parse_reporting_json(self.schema_bytes, self.limits)) contract = _mapping(parse_reporting_json(self.canonicalization_bytes, self.limits)) if ( report.get("report_definition_id") != key.report_definition_id or report.get("reporting_profile") != key.reporting_profile or schema.get("$schema") != definition.schema_dialect or set(contract) != { "contract_version", "media_type", "algorithm", "schema_sha256", "primary_keys", "golden_vectors", } or contract["contract_version"] != "1.0" or contract["media_type"] != "application/vnd.adcp.reporting-canonicalization+json" or contract["algorithm"] != "adcp_jcs_rows_v1" or contract["schema_sha256"].lower() != definition.schema_sha256.lower() ): raise failure("UNSUPPORTED_VERIFICATION") _local_schema(schema) Draft202012Validator.check_schema(schema) keys = contract["primary_keys"] if ( type(keys) is not list or not keys or any(type(k) is not str for k in keys) or len(set(keys)) != len(keys) ): raise failure("UNSUPPORTED_VERIFICATION") totals: list[_Total] = [] properties = _mapping(schema.get("properties")) metrics = report.get("metrics") if type(metrics) is not list or not metrics: raise failure("UNSUPPORTED_VERIFICATION") for metric in metrics: metric = _mapping(metric) name = metric["name"] prop = _mapping(properties[name]) total = _mapping(prop["x-adcp-control-total"]) if ( metric.get("source_expression") != name or metric.get("aggregation") != "sum" or set(total) - {"value_type", "unit", "scale"} or total.get("unit") != metric.get("unit") or (total["value_type"], prop.get("type")) not in {("integer", "integer"), ("decimal", "string")} ): raise failure("UNSUPPORTED_VERIFICATION") scale = total.get("scale", 0) if ( type(scale) is not int or not 0 <= scale <= 18 or (total["value_type"] == "integer" and scale != 0) ): raise failure("UNSUPPORTED_VERIFICATION") totals.append(_Total(name, total["value_type"], total.get("unit"), scale)) if len({t.name for t in totals}) != len(totals): raise failure("UNSUPPORTED_VERIFICATION") if {t.name for t in totals} != { name for name, prop in properties.items() if type(prop) is dict and "x-adcp-control-total" in prop }: raise failure("UNSUPPORTED_VERIFICATION") if not set(keys).union(t.name for t in totals) <= set(schema.get("required", [])): raise failure("UNSUPPORTED_VERIFICATION") units = {t.name: t.unit for t in totals} if any( units.get(name) != unit for name, unit in ( *definition.monetary_metric_units, *definition.monetary_control_total_units, ) ): raise failure("UNSUPPORTED_VERIFICATION") runtime = _Canonicalizer( tuple(keys), tuple(totals), Draft202012Validator(schema), self.limits ) vectors = _mapping(contract["golden_vectors"]) if not {"empty_report", "ordering_encoding"} <= set(vectors) or set(vectors) - { "empty_report", "ordering_encoding", "additional", }: raise failure("UNSUPPORTED_VERIFICATION") seen: set[str] = set() for vector in [ vectors["empty_report"], vectors["ordering_encoding"], *vectors.get("additional", []), ]: vector = _mapping(vector) if ( set(vector) != {"name", "purpose", "input_rows", "canonical_utf8_base64", "sha256"} or vector["name"] in seen ): raise failure("UNSUPPORTED_VERIFICATION") seen.add(vector["name"]) rows, _ = runtime.canonicalize(vector["input_rows"]) golden_bytes = base64.b64decode(vector["canonical_utf8_base64"], validate=True) if ( b"[" + b",".join(rows) + b"]" != golden_bytes or hashlib.sha256(golden_bytes).hexdigest() != vector["sha256"].lower() ): raise failure("UNSUPPORTED_VERIFICATION") if ( vectors["empty_report"]["input_rows"] != [] or vectors["empty_report"]["purpose"] != "empty_report" ): raise failure("UNSUPPORTED_VERIFICATION") ordering = vectors["ordering_encoding"] order_rows = ordering["input_rows"] if ( ordering["purpose"] != "ordering_encoding" or len(order_rows) < 2 or [runtime.order_key(r) for r in order_rows] == sorted(runtime.order_key(r) for r in order_rows) or not _unordered_members(order_rows) ): raise failure("UNSUPPORTED_VERIFICATION") return runtime def canonicalize( self, rows: Sequence[object] ) -> tuple[tuple[bytes, ...], tuple[ReportingControlTotalRecord, ...]]: """Useful to trusted publishers preparing expected evidence before storage.""" return self._runtime().canonicalize(rows)
Methods
def canonicalize(self, rows: Sequence[object]) ‑> tuple[tuple[bytes, ...], tuple[ReportingControlTotalRecord, ...]]-
Expand source code
def canonicalize( self, rows: Sequence[object] ) -> tuple[tuple[bytes, ...], tuple[ReportingControlTotalRecord, ...]]: """Useful to trusted publishers preparing expected evidence before storage.""" return self._runtime().canonicalize(rows)Useful to trusted publishers preparing expected evidence before storage.
class ReportingRevisionVerifierRegistry (verifiers: tuple[ReportingRevisionVerifier, ...])-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingRevisionVerifierRegistry(_ClosedValue): verifiers: tuple[ReportingRevisionVerifier, ...] def __post_init__(self) -> None: _freeze_fields(self) if len({v.key for v in self.verifiers}) != len(self.verifiers): raise failure("UNSUPPORTED_VERIFICATION") def require(self, key: ReportingVerificationKey) -> ReportingRevisionVerifier: for verifier in self.verifiers: if verifier.key == key: return verifier raise failure("UNSUPPORTED_VERIFICATION") async def prepare( self, *, key: ReportingVerificationKey, binding: ReportingDestinationBinding, delivery: ReportingObligationDeliveryRecord, obligation: ReportingObligationRecord, revisions: Sequence[ReportingRevisionRecord], attempt: ReportingMaterializationAttempt, reader: ReportingRevisionRowReader, context: ReportingIOContext, ) -> ReportingPreparedRevision: verifier = self.require(key) # Unsupported tuples fail before any source/resolver I/O. selection = select_reporting_revision( revisions, account_id=obligation.account_id, reporting_obligation_id=obligation.reporting_obligation_id, required_finality=obligation.required_finality, ) if selection.kind == "corrupt": raise failure("HISTORY_CORRUPT") if selection.kind != "selected": raise failure("REVISION_NOT_READY") revision = selection.revision _revision_metadata(revision) if not revision.readable: raise failure("REVISION_NOT_READY") request = ReportingDestinationRequest.from_binding(binding, attempt, key) if ( not _same_definition(key, obligation.definition) or (obligation.report_definition_id, obligation.reporting_profile) != (key.report_definition_id, key.reporting_profile) or request.generation != obligation.generation_key or request.reporting_obligation_id != obligation.reporting_obligation_id or request.reporting_revision_id != revision.reporting_revision_id or delivery.scope != attempt.scope or delivery.currency != obligation.currency ): raise failure("BINDING_MISMATCH") _expected_digest(revision, key) rows: list[dict[str, Any]] = [] cursor: str | None = None seen: set[str] = set() budget = _Budget(verifier.limits) while True: budget.page() page = await context.run( lambda: reader.read_revision_rows( account_id=obligation.account_id, reporting_revision_id=revision.reporting_revision_id, cursor=cursor, limit=500, ) ) if ( type(page) is not ReportingRowPage or type(page.rows) is not tuple or page.reporting_revision_id != revision.reporting_revision_id ): raise failure("SOURCE_INVALID") _page( page.total_count, page.has_more, page.cursor, len(page.rows), len(rows), revision.row_count, seen, ) if page.cursor is not None: cursor_invalid = False try: cursor_invalid = revision_row_offset( page.cursor, revision.reporting_revision_id, 500 ) != len(rows) + len(page.rows) except (ValueError, LedgerConflictError): cursor_invalid = True if cursor_invalid: raise failure("SOURCE_INVALID") for row in page.rows: encoded = strict_reporting_json(row, verifier.limits) budget.add(encoded) rows.append(_mapping(parse_reporting_json(encoded, verifier.limits))) if not page.has_more: break cursor = page.cursor if ( revision_content_sha256( reporting_revision_id=revision.reporting_revision_id, row_count=revision.row_count, control_totals=revision.control_totals, reporting_rows=rows, control_total_evidence=revision.managed_control_totals, ) != revision.revision_content_sha256.lower() ): raise failure("SOURCE_INVALID") canonical_rows, totals = verifier.canonicalize(rows) _verify_content(revision, key, canonical_rows, totals) await context.run(_checkpoint) return ReportingPreparedRevision( request, obligation, revision, delivery, binding, canonical_rows )ReportingRevisionVerifierRegistry(verifiers: 'tuple[ReportingRevisionVerifier, …]')
Ancestors
- adcp.reporting.ledger.delivery_models._ClosedValue
Instance variables
var verifiers : tuple[ReportingRevisionVerifier, ...]-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingRevisionVerifierRegistry(_ClosedValue): verifiers: tuple[ReportingRevisionVerifier, ...] def __post_init__(self) -> None: _freeze_fields(self) if len({v.key for v in self.verifiers}) != len(self.verifiers): raise failure("UNSUPPORTED_VERIFICATION") def require(self, key: ReportingVerificationKey) -> ReportingRevisionVerifier: for verifier in self.verifiers: if verifier.key == key: return verifier raise failure("UNSUPPORTED_VERIFICATION") async def prepare( self, *, key: ReportingVerificationKey, binding: ReportingDestinationBinding, delivery: ReportingObligationDeliveryRecord, obligation: ReportingObligationRecord, revisions: Sequence[ReportingRevisionRecord], attempt: ReportingMaterializationAttempt, reader: ReportingRevisionRowReader, context: ReportingIOContext, ) -> ReportingPreparedRevision: verifier = self.require(key) # Unsupported tuples fail before any source/resolver I/O. selection = select_reporting_revision( revisions, account_id=obligation.account_id, reporting_obligation_id=obligation.reporting_obligation_id, required_finality=obligation.required_finality, ) if selection.kind == "corrupt": raise failure("HISTORY_CORRUPT") if selection.kind != "selected": raise failure("REVISION_NOT_READY") revision = selection.revision _revision_metadata(revision) if not revision.readable: raise failure("REVISION_NOT_READY") request = ReportingDestinationRequest.from_binding(binding, attempt, key) if ( not _same_definition(key, obligation.definition) or (obligation.report_definition_id, obligation.reporting_profile) != (key.report_definition_id, key.reporting_profile) or request.generation != obligation.generation_key or request.reporting_obligation_id != obligation.reporting_obligation_id or request.reporting_revision_id != revision.reporting_revision_id or delivery.scope != attempt.scope or delivery.currency != obligation.currency ): raise failure("BINDING_MISMATCH") _expected_digest(revision, key) rows: list[dict[str, Any]] = [] cursor: str | None = None seen: set[str] = set() budget = _Budget(verifier.limits) while True: budget.page() page = await context.run( lambda: reader.read_revision_rows( account_id=obligation.account_id, reporting_revision_id=revision.reporting_revision_id, cursor=cursor, limit=500, ) ) if ( type(page) is not ReportingRowPage or type(page.rows) is not tuple or page.reporting_revision_id != revision.reporting_revision_id ): raise failure("SOURCE_INVALID") _page( page.total_count, page.has_more, page.cursor, len(page.rows), len(rows), revision.row_count, seen, ) if page.cursor is not None: cursor_invalid = False try: cursor_invalid = revision_row_offset( page.cursor, revision.reporting_revision_id, 500 ) != len(rows) + len(page.rows) except (ValueError, LedgerConflictError): cursor_invalid = True if cursor_invalid: raise failure("SOURCE_INVALID") for row in page.rows: encoded = strict_reporting_json(row, verifier.limits) budget.add(encoded) rows.append(_mapping(parse_reporting_json(encoded, verifier.limits))) if not page.has_more: break cursor = page.cursor if ( revision_content_sha256( reporting_revision_id=revision.reporting_revision_id, row_count=revision.row_count, control_totals=revision.control_totals, reporting_rows=rows, control_total_evidence=revision.managed_control_totals, ) != revision.revision_content_sha256.lower() ): raise failure("SOURCE_INVALID") canonical_rows, totals = verifier.canonicalize(rows) _verify_content(revision, key, canonical_rows, totals) await context.run(_checkpoint) return ReportingPreparedRevision( request, obligation, revision, delivery, binding, canonical_rows )
Methods
async def prepare(self,
*,
key: ReportingVerificationKey,
binding: ReportingDestinationBinding,
delivery: ReportingObligationDeliveryRecord,
obligation: ReportingObligationRecord,
revisions: Sequence[ReportingRevisionRecord],
attempt: ReportingMaterializationAttempt,
reader: ReportingRevisionRowReader,
context: ReportingIOContext) ‑> ReportingPreparedRevision-
Expand source code
async def prepare( self, *, key: ReportingVerificationKey, binding: ReportingDestinationBinding, delivery: ReportingObligationDeliveryRecord, obligation: ReportingObligationRecord, revisions: Sequence[ReportingRevisionRecord], attempt: ReportingMaterializationAttempt, reader: ReportingRevisionRowReader, context: ReportingIOContext, ) -> ReportingPreparedRevision: verifier = self.require(key) # Unsupported tuples fail before any source/resolver I/O. selection = select_reporting_revision( revisions, account_id=obligation.account_id, reporting_obligation_id=obligation.reporting_obligation_id, required_finality=obligation.required_finality, ) if selection.kind == "corrupt": raise failure("HISTORY_CORRUPT") if selection.kind != "selected": raise failure("REVISION_NOT_READY") revision = selection.revision _revision_metadata(revision) if not revision.readable: raise failure("REVISION_NOT_READY") request = ReportingDestinationRequest.from_binding(binding, attempt, key) if ( not _same_definition(key, obligation.definition) or (obligation.report_definition_id, obligation.reporting_profile) != (key.report_definition_id, key.reporting_profile) or request.generation != obligation.generation_key or request.reporting_obligation_id != obligation.reporting_obligation_id or request.reporting_revision_id != revision.reporting_revision_id or delivery.scope != attempt.scope or delivery.currency != obligation.currency ): raise failure("BINDING_MISMATCH") _expected_digest(revision, key) rows: list[dict[str, Any]] = [] cursor: str | None = None seen: set[str] = set() budget = _Budget(verifier.limits) while True: budget.page() page = await context.run( lambda: reader.read_revision_rows( account_id=obligation.account_id, reporting_revision_id=revision.reporting_revision_id, cursor=cursor, limit=500, ) ) if ( type(page) is not ReportingRowPage or type(page.rows) is not tuple or page.reporting_revision_id != revision.reporting_revision_id ): raise failure("SOURCE_INVALID") _page( page.total_count, page.has_more, page.cursor, len(page.rows), len(rows), revision.row_count, seen, ) if page.cursor is not None: cursor_invalid = False try: cursor_invalid = revision_row_offset( page.cursor, revision.reporting_revision_id, 500 ) != len(rows) + len(page.rows) except (ValueError, LedgerConflictError): cursor_invalid = True if cursor_invalid: raise failure("SOURCE_INVALID") for row in page.rows: encoded = strict_reporting_json(row, verifier.limits) budget.add(encoded) rows.append(_mapping(parse_reporting_json(encoded, verifier.limits))) if not page.has_more: break cursor = page.cursor if ( revision_content_sha256( reporting_revision_id=revision.reporting_revision_id, row_count=revision.row_count, control_totals=revision.control_totals, reporting_rows=rows, control_total_evidence=revision.managed_control_totals, ) != revision.revision_content_sha256.lower() ): raise failure("SOURCE_INVALID") canonical_rows, totals = verifier.canonicalize(rows) _verify_content(revision, key, canonical_rows, totals) await context.run(_checkpoint) return ReportingPreparedRevision( request, obligation, revision, delivery, binding, canonical_rows ) def require(self, key: ReportingVerificationKey) ‑> ReportingRevisionVerifier-
Expand source code
def require(self, key: ReportingVerificationKey) -> ReportingRevisionVerifier: for verifier in self.verifiers: if verifier.key == key: return verifier raise failure("UNSUPPORTED_VERIFICATION")
class ReportingVerifiedDestination (request: ReportingDestinationRequest,
resource: ReportingResourceRecord,
verification: ReportingVerificationRecord)-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingVerifiedDestination(_VerifiedProvenance): """Immutable SDK observations. B2 alone owns atomic publication/readiness.""" request: ReportingDestinationRequest resource: ReportingResourceRecord verification: ReportingVerificationRecord def __post_init__(self) -> None: _freeze_fields(self) object.__setattr__(self, "_sdk_seal", b"") object.__setattr__(self, "_materializer_fence", None)Immutable SDK observations. B2 alone owns atomic publication/readiness.
Ancestors
- adcp.reporting.materializer.verification._VerifiedProvenance
- adcp.reporting.ledger.delivery_models._ClosedValue
Instance variables
var request : ReportingDestinationRequest-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingVerifiedDestination(_VerifiedProvenance): """Immutable SDK observations. B2 alone owns atomic publication/readiness.""" request: ReportingDestinationRequest resource: ReportingResourceRecord verification: ReportingVerificationRecord def __post_init__(self) -> None: _freeze_fields(self) object.__setattr__(self, "_sdk_seal", b"") object.__setattr__(self, "_materializer_fence", None) var resource : ReportingResourceRecord-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingVerifiedDestination(_VerifiedProvenance): """Immutable SDK observations. B2 alone owns atomic publication/readiness.""" request: ReportingDestinationRequest resource: ReportingResourceRecord verification: ReportingVerificationRecord def __post_init__(self) -> None: _freeze_fields(self) object.__setattr__(self, "_sdk_seal", b"") object.__setattr__(self, "_materializer_fence", None) var verification : ReportingVerificationRecord-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingVerifiedDestination(_VerifiedProvenance): """Immutable SDK observations. B2 alone owns atomic publication/readiness.""" request: ReportingDestinationRequest resource: ReportingResourceRecord verification: ReportingVerificationRecord def __post_init__(self) -> None: _freeze_fields(self) object.__setattr__(self, "_sdk_seal", b"") object.__setattr__(self, "_materializer_fence", None)