Module adcp.reporting_inspection
Built-in inspection for immutable reporting file manifests.
Classes
class HttpsReportingResourceReader (credential_provider: ReportingCredentialProvider | None = None,
*,
timeout_seconds: float = 10.0,
trusted_origins: Iterable[str] | None = None,
trusted_origin_policy: ReportingTrustedOriginPolicy | None = None)-
Expand source code
class HttpsReportingResourceReader: """Redirect-free, DNS-pinned reader for explicitly trusted HTTPS origins. ``trusted_origins`` must contain the seller, provider, and/or registry origins authorized by the selected reporting offering. This deliberately has no permissive default: ledger-controlled contract URLs are not an authorization decision. A policy is useful when the trusted origin set is resolved from an authenticated principal at inspection time. """ def __init__( self, credential_provider: ReportingCredentialProvider | None = None, *, timeout_seconds: float = 10.0, trusted_origins: Iterable[str] | None = None, trusted_origin_policy: ReportingTrustedOriginPolicy | None = None, ) -> None: if timeout_seconds <= 0: raise ValueError("timeout_seconds must be positive") if trusted_origins is not None and trusted_origin_policy is not None: raise ValueError("pass trusted_origins or trusted_origin_policy, not both") self._credential_provider = credential_provider self._timeout = timeout_seconds self._trusted_origins = ( frozenset(_origin(value) for value in trusted_origins) if trusted_origins is not None else None ) self._trusted_origin_policy = trusted_origin_policy async def read( self, locator: str, *, base: str | None = None, max_bytes: int, expected_content_types: frozenset[str] | None = None, ) -> bytes: url = urljoin(base, locator) if base else locator origin = _origin(url) if base: if origin != _origin(base): raise ReportingInspectionError( ReportingInspectionCode.UNSAFE_RESOURCE, "reporting object reference crossed the manifest origin", ) origin_text = _origin_text(origin) if self._trusted_origins is not None: trusted = origin in self._trusted_origins elif self._trusted_origin_policy is not None: trusted = self._trusted_origin_policy(origin_text) else: trusted = False if not trusted: raise ReportingInspectionError( ReportingInspectionCode.UNSAFE_RESOURCE, f"reporting resource origin {origin_text!r} is not trusted", ) try: # Host resolution in the pin factory uses socket.getaddrinfo. Keep # it off the event loop so the reconciler's inspection deadline can # still cancel the awaiting task during a slow DNS lookup. transport = await asyncio.to_thread( build_async_ip_pinned_transport, url, allowed_ports=frozenset({443}) ) headers = ( dict(await self._credential_provider(url)) if self._credential_provider else {} ) async with httpx.AsyncClient( transport=transport, timeout=self._timeout, follow_redirects=False, trust_env=False, ) as client: async with client.stream("GET", url, headers=headers) as response: if 300 <= response.status_code < 400: raise ReportingInspectionError( ReportingInspectionCode.UNSAFE_RESOURCE, "redirects are not allowed for reporting resources", ) if response.status_code != 200: raise ReportingInspectionError( ReportingInspectionCode.RESOURCE_UNAVAILABLE, f"reporting resource returned HTTP {response.status_code}", retryable=response.status_code >= 500 or response.status_code in {408, 429}, ) if expected_content_types is not None: content_type = response.headers.get("content-type", "").split(";", 1)[0] if content_type.strip().lower() not in expected_content_types: raise ReportingInspectionError( ReportingInspectionCode.UNEXPECTED_CONTENT_TYPE, "reporting resource returned an unexpected content type", ) return await async_read_limited_bytes(response, limit=max_bytes) except ReportingInspectionError: raise except ResponseTooLargeError as error: raise ReportingInspectionError( ReportingInspectionCode.RESOURCE_TOO_LARGE, str(error) ) from error except SSRFValidationError as error: raise ReportingInspectionError( ReportingInspectionCode.UNSAFE_RESOURCE, "reporting resource failed public-network validation", ) from error except (httpx.HTTPError, OSError, ValueError) as error: raise ReportingInspectionError( ReportingInspectionCode.RESOURCE_UNAVAILABLE, "reporting resource fetch failed", retryable=True, ) from errorRedirect-free, DNS-pinned reader for explicitly trusted HTTPS origins.
trusted_originsmust contain the seller, provider, and/or registry origins authorized by the selected reporting offering. This deliberately has no permissive default: ledger-controlled contract URLs are not an authorization decision. A policy is useful when the trusted origin set is resolved from an authenticated principal at inspection time.Methods
async def read(self,
locator: str,
*,
base: str | None = None,
max_bytes: int,
expected_content_types: frozenset[str] | None = None) ‑> bytes-
Expand source code
async def read( self, locator: str, *, base: str | None = None, max_bytes: int, expected_content_types: frozenset[str] | None = None, ) -> bytes: url = urljoin(base, locator) if base else locator origin = _origin(url) if base: if origin != _origin(base): raise ReportingInspectionError( ReportingInspectionCode.UNSAFE_RESOURCE, "reporting object reference crossed the manifest origin", ) origin_text = _origin_text(origin) if self._trusted_origins is not None: trusted = origin in self._trusted_origins elif self._trusted_origin_policy is not None: trusted = self._trusted_origin_policy(origin_text) else: trusted = False if not trusted: raise ReportingInspectionError( ReportingInspectionCode.UNSAFE_RESOURCE, f"reporting resource origin {origin_text!r} is not trusted", ) try: # Host resolution in the pin factory uses socket.getaddrinfo. Keep # it off the event loop so the reconciler's inspection deadline can # still cancel the awaiting task during a slow DNS lookup. transport = await asyncio.to_thread( build_async_ip_pinned_transport, url, allowed_ports=frozenset({443}) ) headers = ( dict(await self._credential_provider(url)) if self._credential_provider else {} ) async with httpx.AsyncClient( transport=transport, timeout=self._timeout, follow_redirects=False, trust_env=False, ) as client: async with client.stream("GET", url, headers=headers) as response: if 300 <= response.status_code < 400: raise ReportingInspectionError( ReportingInspectionCode.UNSAFE_RESOURCE, "redirects are not allowed for reporting resources", ) if response.status_code != 200: raise ReportingInspectionError( ReportingInspectionCode.RESOURCE_UNAVAILABLE, f"reporting resource returned HTTP {response.status_code}", retryable=response.status_code >= 500 or response.status_code in {408, 429}, ) if expected_content_types is not None: content_type = response.headers.get("content-type", "").split(";", 1)[0] if content_type.strip().lower() not in expected_content_types: raise ReportingInspectionError( ReportingInspectionCode.UNEXPECTED_CONTENT_TYPE, "reporting resource returned an unexpected content type", ) return await async_read_limited_bytes(response, limit=max_bytes) except ReportingInspectionError: raise except ResponseTooLargeError as error: raise ReportingInspectionError( ReportingInspectionCode.RESOURCE_TOO_LARGE, str(error) ) from error except SSRFValidationError as error: raise ReportingInspectionError( ReportingInspectionCode.UNSAFE_RESOURCE, "reporting resource failed public-network validation", ) from error except (httpx.HTTPError, OSError, ValueError) as error: raise ReportingInspectionError( ReportingInspectionCode.RESOURCE_UNAVAILABLE, "reporting resource fetch failed", retryable=True, ) from error
class ManifestReportingInspector (reader: ReportingResourceReader,
*,
max_manifest_bytes: int = 1048576,
max_object_bytes: int = 134217728,
max_contract_bytes: int = 4194304,
max_total_object_bytes: int = 268435456,
max_total_decoded_bytes: int = 268435456,
max_total_rows: int = 500000,
max_files: int = 4096,
control_total_calculator: ControlTotalCalculator | None = None)-
Expand source code
class ManifestReportingInspector: """Retrieve and verify a manifest materialization into an observation.""" def __init__( self, reader: ReportingResourceReader, *, max_manifest_bytes: int = 1024 * 1024, max_object_bytes: int = 128 * 1024 * 1024, max_contract_bytes: int = 4 * 1024 * 1024, max_total_object_bytes: int = 256 * 1024 * 1024, max_total_decoded_bytes: int = 256 * 1024 * 1024, max_total_rows: int = 500_000, max_files: int = 4_096, control_total_calculator: ControlTotalCalculator | None = None, ) -> None: if ( min( max_manifest_bytes, max_object_bytes, max_contract_bytes, max_total_object_bytes, max_total_decoded_bytes, max_total_rows, max_files, ) <= 0 ): raise ValueError("inspection byte limits must be positive") self._reader = reader self._max_manifest_bytes = max_manifest_bytes self._max_object_bytes = max_object_bytes self._max_contract_bytes = max_contract_bytes self._max_total_object_bytes = max_total_object_bytes self._max_total_decoded_bytes = max_total_decoded_bytes self._max_total_rows = max_total_rows self._max_files = max_files self._control_total_calculator = control_total_calculator async def __call__(self, context: ReportingInspectionContext) -> ReportingObservation: materialization = context.materialization resource = materialization.resource revision = context.revision obligation = context.obligation if resource is None or str(resource.kind) != "manifest" or not resource.manifest_sha256: raise ReportingInspectionError( ReportingInspectionCode.INVALID_MANIFEST, "materialization does not expose a digest-pinned manifest", ) manifest_bytes = await self._read( resource.location, max_bytes=self._max_manifest_bytes, expected_content_types=frozenset( {"application/json", "application/vnd.adcp.reporting-file-manifest+json"} ), ) manifest_digest = _digest(manifest_bytes) if manifest_digest.lower() != resource.manifest_sha256.lower(): raise ReportingInspectionError( ReportingInspectionCode.MANIFEST_DIGEST_MISMATCH, "manifest digest does not match the resource descriptor", ) try: manifest_value = _json_loads_no_duplicates(manifest_bytes) manifest = ReportingFileManifest.model_validate(manifest_value) except (ValidationError, _InvalidJsonError) as error: raise ReportingInspectionError( ReportingInspectionCode.INVALID_MANIFEST, "manifest is not valid reporting-file-manifest 1.0", ) from error if not ( manifest.reporting_revision_id == revision.reporting_revision_id and manifest.reporting_obligation_id == obligation.reporting_obligation_id and manifest.reporting_materialization_id == materialization.reporting_materialization_id and _same(manifest.period, revision.period) and manifest.row_count == revision.row_count and _totals_same(manifest.control_totals, revision.control_totals) ): raise ReportingInspectionError( ReportingInspectionCode.MANIFEST_IDENTITY_MISMATCH, "manifest does not match the reporting ledger records", ) refs = [str(entry.object_ref) for entry in manifest.files] if len(refs) != len(set(refs)): raise ReportingInspectionError( ReportingInspectionCode.DUPLICATE_OBJECT, "manifest contains duplicate object references", ) if len(manifest.files) > self._max_files: raise ReportingInspectionError( ReportingInspectionCode.RESOURCE_TOO_LARGE, "manifest exceeds the configured file limit", ) verification = materialization.verification physical = { str(item.object_ref): item.value.lower() for item in (verification.physical_checksums if verification else None) or [] if str(item.algorithm) == "sha256" } if any( physical.get(str(entry.object_ref)) != entry.sha256.lower() for entry in manifest.files ): raise ReportingInspectionError( ReportingInspectionCode.MANIFEST_IDENTITY_MISMATCH, "manifest checksums do not match producer verification evidence", ) if manifest.total_size_bytes != sum(entry.size_bytes for entry in manifest.files): raise ReportingInspectionError( ReportingInspectionCode.OBJECT_SIZE_MISMATCH, "manifest total_size_bytes does not equal its file entries", ) if manifest.row_count != sum(entry.row_count for entry in manifest.files): raise ReportingInspectionError( ReportingInspectionCode.ROW_COUNT_MISMATCH, "manifest row_count does not equal its file entries", ) rows: list[dict[str, Any]] = [] total_object_bytes = 0 total_decoded_bytes = 0 total_rows = 0 for entry in manifest.files: object_ref = str(entry.object_ref) if entry.size_bytes > self._max_object_bytes: raise ReportingInspectionError( ReportingInspectionCode.RESOURCE_TOO_LARGE, "manifest object exceeds the configured byte limit", ) total_object_bytes += entry.size_bytes total_rows += entry.row_count if ( total_object_bytes > self._max_total_object_bytes or total_rows > self._max_total_rows ): raise ReportingInspectionError( ReportingInspectionCode.RESOURCE_TOO_LARGE, "manifest exceeds the configured aggregate inspection limits", ) body = await self._read( object_ref, base=resource.location, max_bytes=min(self._max_object_bytes, entry.size_bytes + 1), expected_content_types=_object_content_types(str(manifest.format)), ) if len(body) != entry.size_bytes: raise ReportingInspectionError( ReportingInspectionCode.OBJECT_SIZE_MISMATCH, f"reporting object {object_ref!r} has the wrong size", ) if _digest(body).lower() != entry.sha256.lower(): raise ReportingInspectionError( ReportingInspectionCode.OBJECT_DIGEST_MISMATCH, f"reporting object {object_ref!r} failed checksum validation", ) file_rows, decoded_bytes = _decode_rows( body, str(manifest.format), str(manifest.compression), max_decoded_bytes=self._max_object_bytes, max_rows=self._max_total_rows - len(rows), ) total_decoded_bytes += decoded_bytes if total_decoded_bytes > self._max_total_decoded_bytes: raise ReportingInspectionError( ReportingInspectionCode.RESOURCE_TOO_LARGE, "manifest exceeds the configured aggregate decoded-byte limit", ) if len(file_rows) != entry.row_count: raise ReportingInspectionError( ReportingInspectionCode.ROW_COUNT_MISMATCH, f"reporting object {object_ref!r} has the wrong row count", ) rows.extend(file_rows) schema = await self._read_json_contract( str(revision.schema_uri), revision.schema_sha256, "row schema", frozenset({"application/schema+json", "application/json"}), ) _validate_row_schema(schema, str(revision.schema_dialect)) try: jsonschema.Draft202012Validator.check_schema(schema) validator = jsonschema.Draft202012Validator(schema) first_error = next( (error for row in rows for error in validator.iter_errors(row)), None ) except (jsonschema.SchemaError, RecursionError, TypeError, ValueError) as error: raise ReportingInspectionError( ReportingInspectionCode.INVALID_CONTRACT, "reporting row schema is invalid" ) from error if first_error is not None: raise ReportingInspectionError( ReportingInspectionCode.ROW_SCHEMA_INVALID, "a reporting row failed its pinned schema", ) definition_json = await self._read_json_contract( str(revision.report_definition_uri), revision.report_definition_sha256, "report definition", frozenset({"application/vnd.adcp.reporting-definition+json"}), ) try: definition = ReportingReportDefinition.model_validate(definition_json) except ValidationError as error: raise ReportingInspectionError( ReportingInspectionCode.INVALID_CONTRACT, "report definition is invalid", ) from error if ( definition.report_definition_id != revision.report_definition_id or definition.reporting_profile != revision.reporting_profile ): raise ReportingInspectionError( ReportingInspectionCode.INVALID_CONTRACT, "report definition identity does not match the revision", ) policy_keys = [ (policy.finality_policy_id, str(policy.basis)) for policy in definition.finality_policies ] if len({policy.finality_policy_id for policy in definition.finality_policies}) != len( definition.finality_policies ): raise ReportingInspectionError( ReportingInspectionCode.INVALID_CONTRACT, "report definition has duplicate finality policy identifiers", ) if ( str(revision.finality) == "official" and ( revision.finality_policy_id, str(revision.finality_basis), ) not in policy_keys ): raise ReportingInspectionError( ReportingInspectionCode.INVALID_CONTRACT, "official revision finality does not match the report definition", ) observed_totals = ( self._control_total_calculator(rows, definition) if self._control_total_calculator else _default_control_totals(rows, definition, revision.control_totals) ) if not _totals_same(observed_totals, revision.control_totals): raise ReportingInspectionError( ReportingInspectionCode.CONTROL_TOTAL_MISMATCH, "recomputed control totals do not match the revision", ) canonical_digest = None if revision.canonical_content_digest: canonical_digest = await self._canonical_digest( rows, revision.canonical_content_digest, revision.schema_sha256 ) return ReportingObservation( row_count=len(rows), control_totals=list(observed_totals), canonical_content_digest=canonical_digest, manifest_sha256=manifest_digest, ) async def _read_json_contract( self, uri: str, expected_digest: str, description: str, expected_content_types: frozenset[str], ) -> dict[str, Any]: body = await self._read( uri, max_bytes=self._max_contract_bytes, expected_content_types=expected_content_types, ) if _digest(body).lower() != expected_digest.lower(): raise ReportingInspectionError( ReportingInspectionCode.INVALID_CONTRACT, f"{description} digest mismatch", ) try: value = _json_loads_no_duplicates(body) except _InvalidJsonError as error: raise ReportingInspectionError( ReportingInspectionCode.INVALID_CONTRACT, f"{description} is not valid UTF-8 JSON", ) from error if not isinstance(value, dict): raise ReportingInspectionError( ReportingInspectionCode.INVALID_CONTRACT, f"{description} root must be an object", ) return value async def _read( self, locator: str, *, max_bytes: int, base: str | None = None, expected_content_types: frozenset[str] | None = None, ) -> bytes: """Use typed content checks for the built-in HTTPS reader. The small resource-reader protocol intentionally remains byte-only so existing destination adapters stay compatible; those adapters own equivalent media-type validation for their native transports. """ if isinstance(self._reader, HttpsReportingResourceReader): return await self._reader.read( locator, base=base, max_bytes=max_bytes, expected_content_types=expected_content_types, ) return await self._reader.read(locator, base=base, max_bytes=max_bytes) async def _canonical_digest( self, rows: list[dict[str, Any]], expected: ReportingCanonicalContentDigest, schema_sha256: str, ) -> ReportingCanonicalContentDigest: contract_json = await self._read_json_contract( str(expected.canonicalization_uri), expected.canonicalization_sha256, "canonicalization contract", frozenset({"application/vnd.adcp.reporting-canonicalization+json"}), ) try: contract = ReportingCanonicalizationContract.model_validate(contract_json) except ValidationError as error: raise ReportingInspectionError( ReportingInspectionCode.INVALID_CONTRACT, "canonicalization contract is invalid", ) from error if contract.schema_sha256.lower() != schema_sha256.lower(): raise ReportingInspectionError( ReportingInspectionCode.INVALID_CONTRACT, "canonicalization contract targets a different row schema", ) keys = [str(item) for item in contract.primary_keys] try: _validate_golden_vectors(contract, keys) value = hashlib.sha256(_canonical_rows_bytes(rows, keys)).hexdigest() except ReportingInspectionError: raise except (KeyError, TypeError, ValueError, rfc8785.CanonicalizationError) as error: raise ReportingInspectionError( ReportingInspectionCode.INVALID_ROWS, "reporting rows cannot be canonicalized with the declared primary keys", ) from error if value.lower() != expected.value.lower(): raise ReportingInspectionError( ReportingInspectionCode.CANONICAL_DIGEST_MISMATCH, "canonical logical-content digest does not match the revision", ) return expectedRetrieve and verify a manifest materialization into an observation.
class ReportingInspectionCode (*args, **kwds)-
Expand source code
class ReportingInspectionCode(str, Enum): RESOURCE_UNAVAILABLE = "RESOURCE_UNAVAILABLE" UNSAFE_RESOURCE = "UNSAFE_RESOURCE" RESOURCE_TOO_LARGE = "RESOURCE_TOO_LARGE" UNEXPECTED_CONTENT_TYPE = "UNEXPECTED_CONTENT_TYPE" MANIFEST_DIGEST_MISMATCH = "MANIFEST_DIGEST_MISMATCH" INVALID_MANIFEST = "INVALID_MANIFEST" MANIFEST_IDENTITY_MISMATCH = "MANIFEST_IDENTITY_MISMATCH" DUPLICATE_OBJECT = "DUPLICATE_OBJECT" OBJECT_SIZE_MISMATCH = "OBJECT_SIZE_MISMATCH" OBJECT_DIGEST_MISMATCH = "OBJECT_DIGEST_MISMATCH" UNSUPPORTED_FORMAT = "UNSUPPORTED_FORMAT" UNSUPPORTED_COMPRESSION = "UNSUPPORTED_COMPRESSION" INVALID_ROWS = "INVALID_ROWS" ROW_SCHEMA_INVALID = "ROW_SCHEMA_INVALID" ROW_COUNT_MISMATCH = "ROW_COUNT_MISMATCH" CONTROL_TOTAL_MISMATCH = "CONTROL_TOTAL_MISMATCH" UNSUPPORTED_CONTROL_TOTAL = "UNSUPPORTED_CONTROL_TOTAL" CANONICAL_DIGEST_MISMATCH = "CANONICAL_DIGEST_MISMATCH" INVALID_CONTRACT = "INVALID_CONTRACT"str(object='') -> str str(bytes_or_buffer[, encoding[, errors]]) -> str
Create a new string object from the given object. If encoding or errors is specified, then the object must expose a data buffer that will be decoded using the given encoding and error handler. Otherwise, returns the result of object.str() (if defined) or repr(object). encoding defaults to sys.getdefaultencoding(). errors defaults to 'strict'.
Ancestors
- builtins.str
- enum.Enum
Class variables
var CANONICAL_DIGEST_MISMATCHvar CONTROL_TOTAL_MISMATCHvar DUPLICATE_OBJECTvar INVALID_CONTRACTvar INVALID_MANIFESTvar INVALID_ROWSvar MANIFEST_DIGEST_MISMATCHvar MANIFEST_IDENTITY_MISMATCHvar OBJECT_DIGEST_MISMATCHvar OBJECT_SIZE_MISMATCHvar RESOURCE_TOO_LARGEvar RESOURCE_UNAVAILABLEvar ROW_COUNT_MISMATCHvar ROW_SCHEMA_INVALIDvar UNEXPECTED_CONTENT_TYPEvar UNSAFE_RESOURCEvar UNSUPPORTED_COMPRESSIONvar UNSUPPORTED_CONTROL_TOTALvar UNSUPPORTED_FORMAT
class ReportingInspectionError (code: ReportingInspectionCode,
message: str,
*,
retryable: bool = False)-
Expand source code
class ReportingInspectionError(RuntimeError): def __init__( self, code: ReportingInspectionCode, message: str, *, retryable: bool = False ) -> None: super().__init__(message) self.code = code self.retryable = retryableUnspecified run-time error.
Ancestors
- builtins.RuntimeError
- builtins.Exception
- builtins.BaseException
class ReportingResourceReader (*args, **kwargs)-
Expand source code
class ReportingResourceReader(Protocol): """Read a bounded locator; adapters own credentials and provider routing.""" async def read(self, locator: str, *, base: str | None = None, max_bytes: int) -> bytes: raise NotImplementedErrorRead a bounded locator; adapters own credentials and provider routing.
Ancestors
- typing.Protocol
- typing.Generic
Methods
async def read(self, locator: str, *, base: str | None = None, max_bytes: int) ‑> bytes-
Expand source code
async def read(self, locator: str, *, base: str | None = None, max_bytes: int) -> bytes: raise NotImplementedError