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 check

See 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)