Module adcp.reporting.source
The seller-side reporting source contract: one frozen slice in, one immutable manifest out.
This is the transport-independent seam between the part of a seller that
schedules reporting (obligations, cadence, retry, gap detection – the
producer in :mod:adcp.reporting.ledger) and the part that fetches it from
whatever actually holds the numbers: an ad server report API, a stats cache, a
warehouse, a partner feed.
Who Owns What
The producer owns obligations, schedules, cadence, deduplication, gap detection, retry and backfill policy, durable checkpoints, manifest acceptance, revision identity, and status. It is the scheduler.
An executor owns exactly one thing: fulfilling one bounded, frozen slice.
It owns its own pagination, its own async-job polling, its own retry
classification, and it returns one conforming immutable manifest.
It never
advances a durable checkpoint and never commits a revision.
An executor that
loops past its deadline, ignores cancellation, or publishes a partial result as
complete is not conforming – see :mod:adcp.reporting.conformance, which says
so executably.
Two Publication Classes
PROVISIONAL_SNAPSHOT and AUTHORITATIVE are different offerings, even
when they cover the same media buy and window. A declared settling policy may
use them sequentially for one ledger obligation: a snapshot is replaceable
current-state evidence; an authoritative publication is official-close
evidence whose later corrections are explicit, immutable adjustments. Never
relabel a later authoritative fetch as a fresh first official publication.
Two Manifest Strictness Levels
basic
Identity, staged objects with SHA-256, data_through, coverage,
finality evidence, and completeness.
This is what a seller wrapping an
existing delivery read can honestly produce, and it is enough to commit a
revision.
:mod:adcp.reporting.inline_source emits exactly this.
evidenced
Everything in basic plus per-call provider evidence: every result page,
every async job submission and poll, every retry attempt, and a request
count that must exactly equal that evidence.
Required when a manifest
has to prove to a third party that no provider page was silently dropped.
A basic manifest is not a weaker claim about the data; it is a weaker
claim about the acquisition.
Both bind their bytes the same way.
Generic Identity
Unlike the internal platform contract this is ported from, identity here is
AdCP-shaped and vendor-neutral: account, delivery config, report definition,
period, and obligation.
Everything a particular source needs to identify
itself – provider name, credential binding generation, ad-server network code
– goes in one opaque :attr:ReportingSourceIdentityV1.source_scope mapping
that the SDK fingerprints but never interprets.
Global variables
var REPORTING_EVIDENCE_REASON_MAX_LENGTH_V1-
Evidence reasons are bounded ASCII so the budget is identical in Python, TypeScript, and generated JSON Schema. Provider-native or multilingual diagnostic detail belongs in redacted observability, not in an immutable manifest cell.
var SOURCE_BATCH_MANIFEST_MAX_BYTES_V1-
A retained manifest is evidence, not a payload. Rows belong in staged objects; the manifest that binds them stays small enough to keep forever.
var SOURCE_BATCH_MANIFEST_MAX_METRIC_AVAILABILITY_V1-
One availability record per accepted
(constituent, requested metric)cell, so a slice's cell matrix is bounded by the same number. var SOURCE_BATCH_MANIFEST_MAX_PROVIDER_CALLS_V1-
Successful calls plus retries, per slice, before publication.
Functions
def completed_reporting_source_response_v1(*,
request: ReportingSourceSliceRequestV1,
manifest: SourceBatchManifestReferenceV1) ‑> ReportingSourceExecutionResponseV1-
Expand source code
def completed_reporting_source_response_v1( *, request: ReportingSourceSliceRequestV1, manifest: SourceBatchManifestReferenceV1 ) -> ReportingSourceExecutionResponseV1: """Build the success arm, echoing the frozen request identity verbatim.""" return ReportingSourceExecutionResponseV1(identity=request.identity, manifest=manifest)Build the success arm, echoing the frozen request identity verbatim.
def coverage_denominator_fingerprint_v1(constituents: Sequence[Any]) ‑> str-
Expand source code
def coverage_denominator_fingerprint_v1(constituents: Sequence[Any]) -> str: """Fingerprint the *requested* constituent set, excluding publication status. This is what makes "did I ask for the same thing?" answerable across a request and its manifest without comparing two differently shaped documents. Sorting by ``constituent_id`` means request order is not part of identity; adding, removing, or re-binding a constituent changes it. """ ordered = sorted( ( { key: value for key, value in _constituent_identity(constituent).items() if value is not None } for constituent in constituents ), key=lambda item: str(item["constituent_id"]).encode("utf-16-be", errors="surrogatepass"), ) return reporting_fingerprint_v1( {"kind": "reporting_coverage_denominator_v1", "constituents": ordered} )Fingerprint the requested constituent set, excluding publication status.
This is what makes "did I ask for the same thing?" answerable across a request and its manifest without comparing two differently shaped documents. Sorting by
constituent_idmeans request order is not part of identity; adding, removing, or re-binding a constituent changes it. def deterministic_source_publication_id_v1(*,
source_execution_key: str,
logical_slice_fingerprint: str,
publication_class: str,
revision_kind: str,
supersedes_publication_id: str | None = None) ‑> str-
Expand source code
def deterministic_source_publication_id_v1( *, source_execution_key: str, logical_slice_fingerprint: str, publication_class: str, revision_kind: str, supersedes_publication_id: str | None = None, ) -> str: """The publication id a conforming executor MUST mint for this slice. Derived, not chosen, so replaying a ``source_execution_key`` cannot produce a second publication id for the same logical work. """ identity: dict[str, Any] = { "source_execution_key": source_execution_key, "logical_slice_fingerprint": logical_slice_fingerprint, "publication_class": publication_class, "revision_kind": revision_kind, } if supersedes_publication_id: identity["supersedes_publication_id"] = supersedes_publication_id return f"src-{canonical_json_sha256_v1(identity)}"The publication id a conforming executor MUST mint for this slice.
Derived, not chosen, so replaying a
source_execution_keycannot produce a second publication id for the same logical work. def encode_source_batch_manifest_v1(manifest: SourceBatchManifestV1) ‑> bytes-
Expand source code
def encode_source_batch_manifest_v1(manifest: SourceBatchManifestV1) -> bytes: """Seal a manifest into its ``canonical_json_utf8_v1`` bytes.""" payload = manifest.model_dump(mode="json", exclude_none=True) encoded = canonical_json_utf8_v1(payload) if len(encoded) > SOURCE_BATCH_MANIFEST_MAX_BYTES_V1: raise SourceBatchManifestError( "SOURCE_MANIFEST_TOO_LARGE", "the source manifest exceeds the 1 MiB retained-evidence limit; " "large data belongs in staged objects", ) return encodedSeal a manifest into its
canonical_json_utf8_v1bytes. def publication_content_fingerprint_v1(manifest: SourceBatchManifestV1 | Mapping[str, Any]) ‑> str-
Expand source code
def publication_content_fingerprint_v1(manifest: SourceBatchManifestV1 | Mapping[str, Any]) -> str: """Fingerprint the source-neutral *content* of a publication. Deliberately excludes acquisition timing, staged object references, and run identity: two fetches of an unchanged quiet campaign produce equal content fingerprints, which lets a producer record the successful refresh without committing a duplicate revision. A later changed observation still receives its own immutable publication and supersedes the prior snapshot. """ raw = ( manifest.model_dump(mode="json", exclude_none=True) if isinstance(manifest, BaseModel) else dict(manifest) ) payload: dict[str, Any] = dict(strip_absent(raw)) content = { "publication_class": payload["publication_class"], "publication_namespace": payload["publication_namespace"], "offering_id": payload["offering_id"], "contract": payload["contract"], "currency": payload["currency"], "requested_dimensions": payload.get("requested_dimensions", []), "period": { key: payload["period"][key] for key in ("period_key", "start", "end", "source_timezone", "grain", "windowing") }, "account_id": payload["identity"]["account_id"], "report_definition_id": payload["identity"]["report_definition_id"], "source_scope": payload["identity"].get("source_scope", {}), "objects": [ { key: item[key] for key in ( "ordinal", "media_type", "compression", "sha256", "byte_count", "row_count", ) } for item in payload["objects"] ], "row_count": payload["row_count"], "control_totals": payload.get("control_totals", []), "metric_availability": [ {key: value for key, value in cell.items() if key != "data_through"} for cell in payload["metric_availability"] ], "coverage": { "denominator_fingerprint": payload["coverage"]["denominator_fingerprint"], "status": payload["coverage"]["status"], "constituents": [ {key: value for key, value in item.items() if key != "data_through"} for item in payload["coverage"]["constituents"] ], }, "explicit_zero": payload.get("explicit_zero", False), } return reporting_fingerprint_v1(content)Fingerprint the source-neutral content of a publication.
Deliberately excludes acquisition timing, staged object references, and run identity: two fetches of an unchanged quiet campaign produce equal content fingerprints, which lets a producer record the successful refresh without committing a duplicate revision. A later changed observation still receives its own immutable publication and supersedes the prior snapshot.
def reporting_source_capabilities_sha256_v1(capabilities: ReportingSourceCapabilitiesV1 | Mapping[str, Any]) ‑> str-
Expand source code
def reporting_source_capabilities_sha256_v1( capabilities: ReportingSourceCapabilitiesV1 | Mapping[str, Any], ) -> str: """SHA-256 over the declaration with ``capabilities_sha256`` itself removed.""" raw = ( capabilities.model_dump(mode="json", exclude_none=True) if isinstance(capabilities, BaseModel) else dict(capabilities) ) payload = dict(strip_absent(raw)) payload.pop("capabilities_sha256", None) return canonical_json_sha256_v1(payload)SHA-256 over the declaration with
capabilities_sha256itself removed. def source_batch_coverage_fingerprint_v1(constituents: Sequence[ReportingConstituentCoverageV1]) ‑> str-
Expand source code
def source_batch_coverage_fingerprint_v1( constituents: Sequence[ReportingConstituentCoverageV1], ) -> str: """Fingerprint the complete constituent *publication evidence*, status included.""" ordered = sorted( (item.canonical_payload() for item in constituents), key=lambda item: str(item["constituent"]["constituent_id"]).encode( "utf-16-be", errors="surrogatepass" ), ) return reporting_fingerprint_v1( {"kind": "reporting_coverage_evidence_v1", "constituents": ordered} )Fingerprint the complete constituent publication evidence, status included.
def source_batch_manifest_reference_v1(staged_commit_ref: str, manifest_bytes: bytes) ‑> SourceBatchManifestReferenceV1-
Expand source code
def source_batch_manifest_reference_v1( staged_commit_ref: str, manifest_bytes: bytes ) -> SourceBatchManifestReferenceV1: """Bind sealed manifest bytes to a reference a consumer can verify.""" import hashlib if len(manifest_bytes) > SOURCE_BATCH_MANIFEST_MAX_BYTES_V1: raise SourceBatchManifestError( "SOURCE_MANIFEST_TOO_LARGE", "the source manifest exceeds the 1 MiB retained-evidence limit", ) return SourceBatchManifestReferenceV1( staged_commit_ref=staged_commit_ref, manifest_sha256=hashlib.sha256(manifest_bytes).hexdigest(), byte_count=len(manifest_bytes), )Bind sealed manifest bytes to a reference a consumer can verify.
def source_batch_object_set_sha256_v1(objects: Sequence[SourceBatchObjectV1]) ‑> str-
Expand source code
def source_batch_object_set_sha256_v1(objects: Sequence[SourceBatchObjectV1]) -> str: """Length-prefixed digest over the ordered, generation-pinned object set. Length prefixing rather than concatenation: without it, two different object sets whose fields happen to concatenate to the same string would collide. """ import hashlib digest = hashlib.sha256() def append(value: str) -> None: encoded = value.encode("utf-8") digest.update(len(encoded).to_bytes(4, "big")) digest.update(encoded) append("ordered_generation_pinned_object_set_v1") for item in objects: append(str(item.ordinal)) append(item.object_ref) append(item.object_generation) append(item.media_type) append(item.compression) append(item.sha256) append(str(item.byte_count)) append(str(item.row_count)) return digest.hexdigest()Length-prefixed digest over the ordered, generation-pinned object set.
Length prefixing rather than concatenation: without it, two different object sets whose fields happen to concatenate to the same string would collide.
Classes
class AuthoritativeOfferingV1 (**data: Any)-
Expand source code
class AuthoritativeOfferingV1(_OfferingBase): """Official-close evidence, with a declared correction policy. The first publication is ``authoritative``. A later restatement is a ``correction`` that names the exact publication it adjusts -- never a second "first official" publication. """ publication_class: Literal["AUTHORITATIVE"] = "AUTHORITATIVE" source_local_ready_time: Annotated[ str, StringConstraints(pattern=r"^(?:[01][0-9]|2[0-3]):[0-5][0-9]$") ] days_after_period_end: int = Field(ge=0) expected_availability_lag: IsoDuration worst_case_availability_lag: IsoDuration correction_window: IsoDuration correction_policy: Literal["none", "immutable_correction"] = "none" revision_semantics: Literal["official_with_declared_correction_policy"] = ( "official_with_declared_correction_policy" ) @model_validator(mode="after") def _lag_range_is_ordered(self) -> AuthoritativeOfferingV1: _validate_lag_range(self.expected_availability_lag, self.worst_case_availability_lag) return selfOfficial-close evidence, with a declared correction policy.
The first publication is
authoritative. A later restatement is acorrectionthat names the exact publication it adjusts – never a second "first official" publication.Create a new model by parsing and validating input data from keyword arguments.
Raises [
ValidationError][pydantic_core.ValidationError] if the input data cannot be validated to form a valid model.selfis explicitly positional-only to allowselfas a field name.Ancestors
- adcp.reporting.source._OfferingBase
- adcp.reporting.source._Frozen
- pydantic.main.BaseModel
Class variables
var correction_policy : Literal['none', 'immutable_correction']var correction_window : strvar days_after_period_end : intvar expected_availability_lag : strvar model_configvar publication_class : Literal['AUTHORITATIVE']var revision_semantics : Literal['official_with_declared_correction_policy']var source_local_ready_time : strvar worst_case_availability_lag : str
class MediaBuyConstituentV1 (**data: Any)-
Expand source code
class MediaBuyConstituentV1(_Frozen): """A media-buy constituent: a product inside a named media buy.""" constituent_kind: Literal["media_buy"] = "media_buy" constituent_id: ExternalId product_id: ExternalId media_buy_id: ExternalIdA media-buy constituent: a product inside a named media buy.
Create a new model by parsing and validating input data from keyword arguments.
Raises [
ValidationError][pydantic_core.ValidationError] if the input data cannot be validated to form a valid model.selfis explicitly positional-only to allowselfas a field name.Ancestors
- adcp.reporting.source._Frozen
- pydantic.main.BaseModel
Class variables
var constituent_id : strvar constituent_kind : Literal['media_buy']var media_buy_id : strvar model_configvar product_id : str
class PackageItemConstituentV1 (**data: Any)-
Expand source code
class PackageItemConstituentV1(_Frozen): """A package-item constituent: a product inside a named package.""" constituent_kind: Literal["package_item"] = "package_item" constituent_id: ExternalId product_id: ExternalId package_id: ExternalIdA package-item constituent: a product inside a named package.
Create a new model by parsing and validating input data from keyword arguments.
Raises [
ValidationError][pydantic_core.ValidationError] if the input data cannot be validated to form a valid model.selfis explicitly positional-only to allowselfas a field name.Ancestors
- adcp.reporting.source._Frozen
- pydantic.main.BaseModel
Class variables
var constituent_id : strvar constituent_kind : Literal['package_item']var model_configvar package_id : strvar product_id : str
class ProductConstituentV1 (**data: Any)-
Expand source code
class ProductConstituentV1(_Frozen): """A product-level constituent: no package or media-buy parent.""" constituent_kind: Literal["product"] = "product" constituent_id: ExternalId product_id: ExternalIdA product-level constituent: no package or media-buy parent.
Create a new model by parsing and validating input data from keyword arguments.
Raises [
ValidationError][pydantic_core.ValidationError] if the input data cannot be validated to form a valid model.selfis explicitly positional-only to allowselfas a field name.Ancestors
- adcp.reporting.source._Frozen
- pydantic.main.BaseModel
Class variables
var constituent_id : strvar constituent_kind : Literal['product']var model_configvar product_id : str
class ProvisionalSnapshotOfferingV1 (**data: Any)-
Expand source code
class ProvisionalSnapshotOfferingV1(_OfferingBase): """Fast, replaceable current-state evidence. ``fastest_safe_cadence`` is a *lower bound the producer may choose*, not an SLA and not an invitation to poll faster. Declaring a 15-minute cadence without upstream evidence that the account tolerates it is how a source gets itself rate-limited into ``action_required``. ``restatement_window`` overrides the SDK's default three-day provisional re-read window after period close. ``restatement_cadence`` defaults to the configured reporting period (and is never allowed to beat ``fastest_safe_cadence``). Every successful scheduled read is a new immutable observation, including unchanged content. A per-slice ``provisional_until`` is authoritative over this offering default. ``official_close_lag`` still requires an explicitly declared window and asks a producer with an authoritative offering to publish an official revision after settlement; the SDK fallback does not enable official close. """ publication_class: Literal["PROVISIONAL_SNAPSHOT"] = "PROVISIONAL_SNAPSHOT" fastest_safe_cadence: IsoDuration alignment: Literal["source_timezone", "provider_defined", "unaligned"] = "provider_defined" expected_availability_lag: IsoDuration worst_case_availability_lag: IsoDuration revision_semantics: Literal["provisional_replaceable"] = "provisional_replaceable" restatement_window: IsoDuration | None = None restatement_cadence: IsoDuration | None = None official_close_lag: IsoDuration | None = None @model_validator(mode="after") def _lag_range_is_ordered(self) -> ProvisionalSnapshotOfferingV1: _validate_lag_range(self.expected_availability_lag, self.worst_case_availability_lag) if self.restatement_window is None: if self.restatement_cadence is not None or self.official_close_lag is not None: raise ValueError( "restatement_cadence and official_close_lag require restatement_window" ) return self iso_duration_milliseconds_v1(self.restatement_window) if self.restatement_cadence is not None: cadence_ms = iso_duration_milliseconds_v1(self.restatement_cadence) if cadence_ms <= 0: raise ValueError("restatement_cadence must be greater than zero") if cadence_ms < iso_duration_milliseconds_v1(self.fastest_safe_cadence): raise ValueError("restatement_cadence must not be faster than fastest_safe_cadence") if self.official_close_lag is not None: iso_duration_milliseconds_v1(self.official_close_lag) return selfFast, replaceable current-state evidence.
fastest_safe_cadenceis a lower bound the producer may choose, not an SLA and not an invitation to poll faster. Declaring a 15-minute cadence without upstream evidence that the account tolerates it is how a source gets itself rate-limited intoaction_required.restatement_windowoverrides the SDK's default three-day provisional re-read window after period close.restatement_cadencedefaults to the configured reporting period (and is never allowed to beatfastest_safe_cadence). Every successful scheduled read is a new immutable observation, including unchanged content. A per-sliceprovisional_untilis authoritative over this offering default.official_close_lagstill requires an explicitly declared window and asks a producer with an authoritative offering to publish an official revision after settlement; the SDK fallback does not enable official close.Create a new model by parsing and validating input data from keyword arguments.
Raises [
ValidationError][pydantic_core.ValidationError] if the input data cannot be validated to form a valid model.selfis explicitly positional-only to allowselfas a field name.Ancestors
- adcp.reporting.source._OfferingBase
- adcp.reporting.source._Frozen
- pydantic.main.BaseModel
Class variables
var alignment : Literal['source_timezone', 'provider_defined', 'unaligned']var expected_availability_lag : strvar fastest_safe_cadence : strvar model_configvar official_close_lag : str | Nonevar publication_class : Literal['PROVISIONAL_SNAPSHOT']var restatement_cadence : str | Nonevar restatement_window : str | Nonevar revision_semantics : Literal['provisional_replaceable']var worst_case_availability_lag : str
class ReportingAdapterBuildIdentityV1 (**data: Any)-
Expand source code
class ReportingAdapterBuildIdentityV1(_Frozen): """Which build of which executor produced a manifest. Pinned in the manifest so a bug traced to one adapter release can be scoped to exactly the publications that release produced. """ executor_id: Identifier adapter_id: Identifier adapter_version: Annotated[str, StringConstraints(min_length=1, max_length=128)] adapter_build_sha256: Sha256Which build of which executor produced a manifest.
Pinned in the manifest so a bug traced to one adapter release can be scoped to exactly the publications that release produced.
Create a new model by parsing and validating input data from keyword arguments.
Raises [
ValidationError][pydantic_core.ValidationError] if the input data cannot be validated to form a valid model.selfis explicitly positional-only to allowselfas a field name.Ancestors
- adcp.reporting.source._Frozen
- pydantic.main.BaseModel
Class variables
var adapter_build_sha256 : strvar adapter_id : strvar adapter_version : strvar executor_id : strvar model_config
class ReportingConstituentCoverageV1 (**data: Any)-
Expand source code
class ReportingConstituentCoverageV1(_Frozen): """One constituent's publication evidence for this slice.""" constituent: ReportingConstituent status: ReportingAvailabilityStatus data_through: datetime | None = None reason: EvidenceReason | None = None @model_validator(mode="after") def _unavailable_coverage_is_explained(self) -> ReportingConstituentCoverageV1: if self.status not in _AVAILABLE_STATUSES and not self.reason: raise ValueError(f"constituent status {self.status!r} needs a stable reason") return self @property def constituent_id(self) -> str: return self.constituent.constituent_idOne constituent's publication evidence for this slice.
Create a new model by parsing and validating input data from keyword arguments.
Raises [
ValidationError][pydantic_core.ValidationError] if the input data cannot be validated to form a valid model.selfis explicitly positional-only to allowselfas a field name.Ancestors
- adcp.reporting.source._Frozen
- pydantic.main.BaseModel
Class variables
var constituent : ProductConstituentV1 | PackageItemConstituentV1 | MediaBuyConstituentV1var data_through : datetime.datetime | Nonevar model_configvar reason : str | Nonevar status : Literal['present', 'explicit_zero', 'unsupported', 'delayed', 'partial', 'stale', 'missing']
Instance variables
prop constituent_id : str-
Expand source code
@property def constituent_id(self) -> str: return self.constituent.constituent_id
class ReportingContractIdentityV1 (**data: Any)-
Expand source code
class ReportingContractIdentityV1(_Frozen): """The content-addressed report definition this slice was fetched against.""" report_definition_id: Identifier report_definition_uri: Annotated[str, StringConstraints(max_length=2_048)] report_definition_sha256: Sha256 reporting_profile: Annotated[str, StringConstraints(pattern=r"^[A-Za-z0-9_.:-]{1,128}$")] schema_version: Annotated[str, StringConstraints(min_length=1, max_length=64)] schema_uri: Annotated[str, StringConstraints(max_length=2_048)] schema_sha256: Sha256 schema_dialect: Literal["https://json-schema.org/draft/2020-12/schema"] = ( "https://json-schema.org/draft/2020-12/schema" ) schema_ref_policy: Literal["local_fragment_only"] = "local_fragment_only" @model_validator(mode="after") def _uris_are_public_https(self) -> ReportingContractIdentityV1: for field, value in ( ("report_definition_uri", self.report_definition_uri), ("schema_uri", self.schema_uri), ): _require_public_https(field, value) return selfThe content-addressed report definition this slice was fetched against.
Create a new model by parsing and validating input data from keyword arguments.
Raises [
ValidationError][pydantic_core.ValidationError] if the input data cannot be validated to form a valid model.selfis explicitly positional-only to allowselfas a field name.Ancestors
- adcp.reporting.source._Frozen
- pydantic.main.BaseModel
Class variables
var model_configvar report_definition_id : strvar report_definition_sha256 : strvar report_definition_uri : strvar reporting_profile : strvar schema_dialect : Literal['https://json-schema.org/draft/2020-12/schema']var schema_ref_policy : Literal['local_fragment_only']var schema_sha256 : strvar schema_uri : strvar schema_version : str
class ReportingFinalityEvidenceV1 (**data: Any)-
Expand source code
class ReportingFinalityEvidenceV1(_Frozen): """Why the executor believes the data is as final as it claims.""" basis: Literal[ "provisional_observation", "provider_declared", "provider_job_terminal", "elapsed_settlement_window", ] observed_at: datetime evidence_ref: ExternalId | None = None provisional_until: datetime | None = None @model_validator(mode="after") def _provisional_signal_matches_basis(self) -> ReportingFinalityEvidenceV1: if self.provisional_until is None: return self if self.basis != "provisional_observation": raise ValueError("provisional_until belongs only to provisional observations") if _aware(self.provisional_until, "provisional_until") < _aware( self.observed_at, "observed_at" ): raise ValueError("provisional_until must not predate the source observation") return selfWhy the executor believes the data is as final as it claims.
Create a new model by parsing and validating input data from keyword arguments.
Raises [
ValidationError][pydantic_core.ValidationError] if the input data cannot be validated to form a valid model.selfis explicitly positional-only to allowselfas a field name.Ancestors
- adcp.reporting.source._Frozen
- pydantic.main.BaseModel
Class variables
var basis : Literal['provisional_observation', 'provider_declared', 'provider_job_terminal', 'elapsed_settlement_window']var evidence_ref : str | Nonevar model_configvar observed_at : datetime.datetimevar provisional_until : datetime.datetime | None
class ReportingMetricAvailabilityV1 (**data: Any)-
Expand source code
class ReportingMetricAvailabilityV1(_Frozen): """One ``(constituent, metric)`` cell's publication evidence. Every accepted cell gets exactly one of these. ``missing`` and ``explicit_zero`` are never interchangeable: the first says the source failed to answer, the second says it answered "nothing happened". """ constituent_id: ExternalId metric: MetricName semantic_contract_id: Annotated[ str, StringConstraints(pattern=r"^[A-Za-z][A-Za-z0-9_.:-]{0,255}$") ] semantic_contract_version: Annotated[str, StringConstraints(min_length=1, max_length=128)] semantic_contract_sha256: Sha256 status: ReportingAvailabilityStatus data_through: datetime | None = None reason: EvidenceReason | None = None @model_validator(mode="after") def _evidence_matches_status(self) -> ReportingMetricAvailabilityV1: if self.status in _AVAILABLE_STATUSES and self.data_through is None: raise ValueError(f"metric status {self.status!r} needs its source watermark") if self.status not in _AVAILABLE_STATUSES and not self.reason: raise ValueError(f"metric status {self.status!r} needs a stable reason") return selfOne
(constituent, metric)cell's publication evidence.Every accepted cell gets exactly one of these.
missingandexplicit_zeroare never interchangeable: the first says the source failed to answer, the second says it answered "nothing happened".Create a new model by parsing and validating input data from keyword arguments.
Raises [
ValidationError][pydantic_core.ValidationError] if the input data cannot be validated to form a valid model.selfis explicitly positional-only to allowselfas a field name.Ancestors
- adcp.reporting.source._Frozen
- pydantic.main.BaseModel
Class variables
var constituent_id : strvar data_through : datetime.datetime | Nonevar metric : strvar model_configvar reason : str | Nonevar semantic_contract_id : strvar semantic_contract_sha256 : strvar semantic_contract_version : strvar status : Literal['present', 'explicit_zero', 'unsupported', 'delayed', 'partial', 'stale', 'missing']
class ReportingProviderJobPollV1 (**data: Any)-
Expand source code
class ReportingProviderJobPollV1(_Frozen): """One poll of an upstream async job.""" job_id: ExternalId ordinal: int = Field(ge=0) request_id: ExternalId observed_at: datetime state: Literal["queued", "running", "succeeded"] response_sha256: Sha256One poll of an upstream async job.
Create a new model by parsing and validating input data from keyword arguments.
Raises [
ValidationError][pydantic_core.ValidationError] if the input data cannot be validated to form a valid model.selfis explicitly positional-only to allowselfas a field name.Ancestors
- adcp.reporting.source._Frozen
- pydantic.main.BaseModel
Class variables
var job_id : strvar model_configvar observed_at : datetime.datetimevar ordinal : intvar request_id : strvar response_sha256 : strvar state : Literal['queued', 'running', 'succeeded']
class ReportingProviderJobV1 (**data: Any)-
Expand source code
class ReportingProviderJobV1(_Frozen): """One upstream async job submitted for this slice.""" job_id: ExternalId submission_request_id: ExternalId submission_response_sha256: Sha256One upstream async job submitted for this slice.
Create a new model by parsing and validating input data from keyword arguments.
Raises [
ValidationError][pydantic_core.ValidationError] if the input data cannot be validated to form a valid model.selfis explicitly positional-only to allowselfas a field name.Ancestors
- adcp.reporting.source._Frozen
- pydantic.main.BaseModel
Class variables
var job_id : strvar model_configvar submission_request_id : strvar submission_response_sha256 : str
class ReportingProviderRetryAttemptV1 (**data: Any)-
Expand source code
class ReportingProviderRetryAttemptV1(_Frozen): """One retried upstream call, with the classification that justified it.""" attempt_id: OpaqueReference request_id: ExternalId | None = None attempted_at: datetime classification: Literal[ "RATE_LIMITED", "QUOTA_EXHAUSTED", "PROVIDER_TRANSIENT", "PROVIDER_JOB_FAILED" ]One retried upstream call, with the classification that justified it.
Create a new model by parsing and validating input data from keyword arguments.
Raises [
ValidationError][pydantic_core.ValidationError] if the input data cannot be validated to form a valid model.selfis explicitly positional-only to allowselfas a field name.Ancestors
- adcp.reporting.source._Frozen
- pydantic.main.BaseModel
Class variables
var attempt_id : strvar attempted_at : datetime.datetimevar classification : Literal['RATE_LIMITED', 'QUOTA_EXHAUSTED', 'PROVIDER_TRANSIENT', 'PROVIDER_JOB_FAILED']var model_configvar request_id : str | None
class ReportingSourceCapabilitiesV1 (**data: Any)-
Expand source code
class ReportingSourceCapabilitiesV1(_Frozen): """What one executor build can truthfully do, for one bound account scope. ``scope`` distinguishes a *product declaration* (what this adapter can do in general) from an *effective account* result (what these credentials can do for this account right now). Only the latter is executable; the distinction exists so an integration-wide brochure cannot masquerade as account eligibility. ``capabilities_sha256`` binds the whole scoped declaration. It is computed for you by :func:`reporting_source_capabilities_sha256_v1` and verified on every construction, so a declaration edited in flight fails loudly. """ contract_version: Literal["1.0"] = REPORTING_SOURCE_CONTRACT_VERSION_V1 scope: Literal["product_declaration", "effective_account"] adapter_build: ReportingAdapterBuildIdentityV1 source_scope: dict[str, Any] = Field(default_factory=dict) capability_version: Annotated[str, StringConstraints(min_length=1, max_length=128)] capabilities_sha256: Sha256 offerings: list[ReportingSourceOffering] = Field(min_length=1, max_length=100) @model_validator(mode="after") def _checksum_binds_the_declaration(self) -> ReportingSourceCapabilitiesV1: offering_ids = [offering.offering_id for offering in self.offerings] if len(set(offering_ids)) != len(offering_ids): raise ValueError("offering_id must be unique within a capability declaration") expected = reporting_source_capabilities_sha256_v1(self) if expected != self.capabilities_sha256: raise ValueError( "capabilities_sha256 must bind the complete scoped declaration; " f"expected {expected}" ) return self def offering(self, offering_id: str) -> ProvisionalSnapshotOfferingV1 | AuthoritativeOfferingV1: """The offering with this id, or ``KeyError``.""" for candidate in self.offerings: if candidate.offering_id == offering_id: return candidate raise KeyError(offering_id)What one executor build can truthfully do, for one bound account scope.
scopedistinguishes a product declaration (what this adapter can do in general) from an effective account result (what these credentials can do for this account right now). Only the latter is executable; the distinction exists so an integration-wide brochure cannot masquerade as account eligibility.capabilities_sha256binds the whole scoped declaration. It is computed for you by :func:reporting_source_capabilities_sha256_v1()and verified on every construction, so a declaration edited in flight fails loudly.Create a new model by parsing and validating input data from keyword arguments.
Raises [
ValidationError][pydantic_core.ValidationError] if the input data cannot be validated to form a valid model.selfis explicitly positional-only to allowselfas a field name.Ancestors
- adcp.reporting.source._Frozen
- pydantic.main.BaseModel
Class variables
var adapter_build : ReportingAdapterBuildIdentityV1var capabilities_sha256 : strvar capability_version : strvar contract_version : Literal['1.0']var model_configvar offerings : list[ProvisionalSnapshotOfferingV1 | AuthoritativeOfferingV1]var scope : Literal['product_declaration', 'effective_account']var source_scope : dict[str, typing.Any]
Methods
def offering(self, offering_id: str) ‑> ProvisionalSnapshotOfferingV1 | AuthoritativeOfferingV1-
Expand source code
def offering(self, offering_id: str) -> ProvisionalSnapshotOfferingV1 | AuthoritativeOfferingV1: """The offering with this id, or ``KeyError``.""" for candidate in self.offerings: if candidate.offering_id == offering_id: return candidate raise KeyError(offering_id)The offering with this id, or
KeyError.
class ReportingSourceError (code: ReportingSourceErrorCode,
safe_message: str,
*,
retry: ReportingSourceRetryClass | None = None,
scope: "Literal['slice', 'account', 'source']" = 'slice',
provider_code: str | None = None,
retry_after_seconds: float | None = None)-
Expand source code
class ReportingSourceError(Exception): """Raise this from an executor to settle the slice as a typed failure. Convenience for the common case where the failure is discovered deep in a call stack. :mod:`adcp.reporting.inline_source` and :class:`~adcp.reporting.ledger.ReportingProducer` both convert it into the corresponding :class:`ReportingSourceErrorV1`. """ def __init__( self, code: ReportingSourceErrorCode, safe_message: str, *, retry: ReportingSourceRetryClass | None = None, scope: Literal["slice", "account", "source"] = "slice", provider_code: str | None = None, retry_after_seconds: float | None = None, ) -> None: super().__init__(safe_message) permitted = _PERMITTED_RETRY_CLASSES[code] self.error = ReportingSourceErrorV1( code=code, retry=retry if retry is not None else _default_retry_class(code), scope=scope, safe_message=safe_message, provider_code=provider_code, retry_after_seconds=retry_after_seconds, ) if self.error.retry not in permitted: # pragma: no cover - guarded by the model raise ValueError(f"error {code} cannot be classified {self.error.retry}")Raise this from an executor to settle the slice as a typed failure.
Convenience for the common case where the failure is discovered deep in a call stack. :mod:
adcp.reporting.inline_sourceand :class:~adcp.reporting.ledger.ReportingProducerboth convert it into the corresponding :class:ReportingSourceErrorV1.Ancestors
- builtins.Exception
- builtins.BaseException
class ReportingSourceErrorV1 (**data: Any)-
Expand source code
class ReportingSourceErrorV1(_Frozen): """The failure arm: a stable code, a retry classification, and a safe message. ``safe_message`` is retained and may be shown to a counterparty. It must not carry tokens, signed URLs, raw upstream bodies, or customer secrets -- put those in redacted observability instead. """ contract_version: Literal["1.0"] = REPORTING_SOURCE_CONTRACT_VERSION_V1 code: ReportingSourceErrorCode retry: ReportingSourceRetryClass scope: Literal["slice", "account", "source"] = "slice" safe_message: Annotated[str, StringConstraints(min_length=1, max_length=1_024)] provider_code: Annotated[str, StringConstraints(min_length=1, max_length=128)] | None = None retry_after_seconds: float | None = Field(default=None, ge=0) @model_validator(mode="after") def _classification_is_permitted(self) -> ReportingSourceErrorV1: if self.retry not in _PERMITTED_RETRY_CLASSES[self.code]: raise ValueError(f"error {self.code} cannot be classified {self.retry}") if self.retry_after_seconds is not None and self.code not in _BACKOFF_EVIDENCE_CODES: raise ValueError("retry_after_seconds is reserved for upstream backoff evidence") return selfThe failure arm: a stable code, a retry classification, and a safe message.
safe_messageis retained and may be shown to a counterparty. It must not carry tokens, signed URLs, raw upstream bodies, or customer secrets – put those in redacted observability instead.Create a new model by parsing and validating input data from keyword arguments.
Raises [
ValidationError][pydantic_core.ValidationError] if the input data cannot be validated to form a valid model.selfis explicitly positional-only to allowselfas a field name.Ancestors
- adcp.reporting.source._Frozen
- pydantic.main.BaseModel
Class variables
var code : Literal['AUTHENTICATION_FAILED', 'AUTHORIZATION_FAILED', 'INVALID_REQUEST', 'UNSUPPORTED_OFFERING', 'RATE_LIMITED', 'QUOTA_EXHAUSTED', 'PROVIDER_TRANSIENT', 'PROVIDER_PERMANENT', 'PROVIDER_JOB_FAILED', 'PARTIAL_RESULT', 'CANCELLED', 'DEADLINE_EXCEEDED', 'STAGING_FAILED', 'INTEGRITY_FAILED']var contract_version : Literal['1.0']var model_configvar provider_code : str | Nonevar retry : Literal['retryable', 'terminal', 'cancelled']var retry_after_seconds : float | Nonevar safe_message : strvar scope : Literal['slice', 'account', 'source']
class ReportingSourceExecutionResponseV1 (**data: Any)-
Expand source code
class ReportingSourceExecutionResponseV1(_Frozen): """The success arm: the frozen identity echoed back, plus a manifest reference.""" contract_version: Literal["1.0"] = REPORTING_SOURCE_CONTRACT_VERSION_V1 identity: ReportingSourceIdentityV1 outcome: Literal["completed"] = "completed" manifest: SourceBatchManifestReferenceV1The success arm: the frozen identity echoed back, plus a manifest reference.
Create a new model by parsing and validating input data from keyword arguments.
Raises [
ValidationError][pydantic_core.ValidationError] if the input data cannot be validated to form a valid model.selfis explicitly positional-only to allowselfas a field name.Ancestors
- adcp.reporting.source._Frozen
- pydantic.main.BaseModel
Class variables
var contract_version : Literal['1.0']var identity : ReportingSourceIdentityV1var manifest : SourceBatchManifestReferenceV1var model_configvar outcome : Literal['completed']
class ReportingSourceExecutor (*args, **kwargs)-
Expand source code
@runtime_checkable class ReportingSourceExecutor(Protocol): """Fetches exactly one frozen slice. Cancellation is cooperative and *binding*: when ``cancel`` is set, an executor must stop its upstream requests, its job polling, and its staging I/O, and only then return. Returning early while leaving work running in the background is not conforming -- the producer's next attempt would race the abandoned one for the same staged object names. ``heartbeat`` is optional and advisory. Call it during long polls so the producer can distinguish "still working" from "hung" and avoid re-leasing the job out from under you. """ @property def capabilities(self) -> ReportingSourceCapabilitiesV1: """This build's truthful, scoped capability declaration.""" ... async def execute( self, request: ReportingSourceSliceRequestV1, *, cancel: asyncio.Event, heartbeat: Callable[[], None] | None = None, ) -> ReportingSourceExecutorResult: ...Fetches exactly one frozen slice.
Cancellation is cooperative and binding: when
cancelis set, an executor must stop its upstream requests, its job polling, and its staging I/O, and only then return. Returning early while leaving work running in the background is not conforming – the producer's next attempt would race the abandoned one for the same staged object names.heartbeatis optional and advisory. Call it during long polls so the producer can distinguish "still working" from "hung" and avoid re-leasing the job out from under you.Ancestors
- typing.Protocol
- typing.Generic
Subclasses
Instance variables
prop capabilities : ReportingSourceCapabilitiesV1-
Expand source code
@property def capabilities(self) -> ReportingSourceCapabilitiesV1: """This build's truthful, scoped capability declaration.""" ...This build's truthful, scoped capability declaration.
Methods
async def execute(self,
request: ReportingSourceSliceRequestV1,
*,
cancel: asyncio.Event,
heartbeat: Callable[[], None] | None = None) ‑> ReportingSourceExecutorResult-
Expand source code
async def execute( self, request: ReportingSourceSliceRequestV1, *, cancel: asyncio.Event, heartbeat: Callable[[], None] | None = None, ) -> ReportingSourceExecutorResult: ...
class ReportingSourceExecutorResult (**data: Any)-
Expand source code
class ReportingSourceExecutorResult(_Frozen): """Exactly one of: a completed publication, or a typed failure.""" response: ReportingSourceExecutionResponseV1 | None = None manifest_bytes: bytes | None = None error: ReportingSourceErrorV1 | None = None @model_validator(mode="after") def _exactly_one_arm(self) -> ReportingSourceExecutorResult: completed = self.response is not None if completed == (self.error is not None): raise ValueError("an executor result must be exactly one of completed or failed") if completed and self.manifest_bytes is None: raise ValueError("a completed result must carry its sealed manifest bytes") if not completed and self.manifest_bytes is not None: raise ValueError("a failed result must not carry manifest bytes") return self @property def ok(self) -> bool: return self.response is not None @classmethod def completed( cls, *, request: ReportingSourceSliceRequestV1, manifest: SourceBatchManifestReferenceV1, manifest_bytes: bytes, ) -> ReportingSourceExecutorResult: return cls( response=completed_reporting_source_response_v1(request=request, manifest=manifest), manifest_bytes=manifest_bytes, ) @classmethod def failed(cls, error: ReportingSourceErrorV1) -> ReportingSourceExecutorResult: return cls(error=error)Exactly one of: a completed publication, or a typed failure.
Create a new model by parsing and validating input data from keyword arguments.
Raises [
ValidationError][pydantic_core.ValidationError] if the input data cannot be validated to form a valid model.selfis explicitly positional-only to allowselfas a field name.Ancestors
- adcp.reporting.source._Frozen
- pydantic.main.BaseModel
Class variables
var error : ReportingSourceErrorV1 | Nonevar manifest_bytes : bytes | Nonevar model_configvar response : ReportingSourceExecutionResponseV1 | None
Static methods
def completed(*,
request: ReportingSourceSliceRequestV1,
manifest: SourceBatchManifestReferenceV1,
manifest_bytes: bytes) ‑> ReportingSourceExecutorResultdef failed(error: ReportingSourceErrorV1) ‑> ReportingSourceExecutorResult
Instance variables
prop ok : bool-
Expand source code
@property def ok(self) -> bool: return self.response is not None
class ReportingSourceIdentityV1 (**data: Any)-
Expand source code
class ReportingSourceIdentityV1(_Frozen): """The frozen identity of one slice execution. ``source_execution_key`` is the idempotency key: reusing it MUST return the byte-identical manifest and the same staged object set. ``run_id`` may change across retries of the same logical work; ``source_execution_key`` may not. ``source_scope`` is the deliberate escape hatch. Everything a particular reporting source needs to pin itself -- provider name, ad-server network code, credential binding generation, region -- lives there as opaque JSON. The SDK fingerprints it into the logical slice identity (so a re-bound account produces a different slice, not a silent continuation of the old one) but never reads a key out of it. """ account_id: ExternalId delivery_config_id: ExternalId delivery_config_version: int = Field(ge=1) report_definition_id: Identifier reporting_obligation_id: ExternalId period_key: Identifier source_execution_key: Annotated[str, StringConstraints(pattern=_EXECUTION_KEY_PATTERN)] run_id: ExternalId logical_slice_fingerprint: Fingerprint source_scope: dict[str, Any] = Field(default_factory=dict) @model_validator(mode="after") def _source_scope_is_canonical(self) -> ReportingSourceIdentityV1: # Fail at construction rather than at publication: a source_scope # carrying a float or a non-string key cannot be fingerprinted, and # discovering that while sealing a manifest loses the fetch. canonical_json_utf8_v1(self.source_scope) return selfThe frozen identity of one slice execution.
source_execution_keyis the idempotency key: reusing it MUST return the byte-identical manifest and the same staged object set.run_idmay change across retries of the same logical work;source_execution_keymay not.source_scopeis the deliberate escape hatch. Everything a particular reporting source needs to pin itself – provider name, ad-server network code, credential binding generation, region – lives there as opaque JSON. The SDK fingerprints it into the logical slice identity (so a re-bound account produces a different slice, not a silent continuation of the old one) but never reads a key out of it.Create a new model by parsing and validating input data from keyword arguments.
Raises [
ValidationError][pydantic_core.ValidationError] if the input data cannot be validated to form a valid model.selfis explicitly positional-only to allowselfas a field name.Ancestors
- adcp.reporting.source._Frozen
- pydantic.main.BaseModel
Class variables
var account_id : strvar delivery_config_id : strvar delivery_config_version : intvar logical_slice_fingerprint : strvar model_configvar period_key : strvar report_definition_id : strvar reporting_obligation_id : strvar run_id : strvar source_execution_key : strvar source_scope : dict[str, typing.Any]
class ReportingSourcePeriodV1 (**data: Any)-
Expand source code
class ReportingSourcePeriodV1(_Frozen): """The half-open reporting window, frozen before dispatch. Every field here is part of admission identity and MUST be echoed exactly by the manifest. The request deliberately does *not* predict ``observed_at`` or ``data_through`` -- those are results the source reports, not instructions the producer gives. """ period_key: Identifier source_local_date: Annotated[str, StringConstraints(pattern=_ISO_DATE_PATTERN)] start: datetime end: datetime source_timezone: Annotated[str, StringConstraints(min_length=1, max_length=255)] source_read_cutoff_at: datetime grain: Grain windowing: ReportingSourceWindowingV1 @model_validator(mode="after") def _period_is_coherent(self) -> ReportingSourcePeriodV1: start = _aware(self.start, "start") end = _aware(self.end, "end") cutoff = _aware(self.source_read_cutoff_at, "source_read_cutoff_at") if end <= start: raise ValueError("period must be nonempty (end must follow start)") if cutoff < start: raise ValueError("source_read_cutoff_at must not precede the bounded period") _require_iana_timezone(self.source_timezone) if _source_local_date(start, self.source_timezone) != self.source_local_date: raise ValueError("source_local_date must be the source-local date at period start") return selfThe half-open reporting window, frozen before dispatch.
Every field here is part of admission identity and MUST be echoed exactly by the manifest. The request deliberately does not predict
observed_atordata_through– those are results the source reports, not instructions the producer gives.Create a new model by parsing and validating input data from keyword arguments.
Raises [
ValidationError][pydantic_core.ValidationError] if the input data cannot be validated to form a valid model.selfis explicitly positional-only to allowselfas a field name.Ancestors
- adcp.reporting.source._Frozen
- pydantic.main.BaseModel
Class variables
var end : datetime.datetimevar grain : strvar model_configvar period_key : strvar source_local_date : strvar source_read_cutoff_at : datetime.datetimevar source_timezone : strvar start : datetime.datetimevar windowing : ReportingSourceWindowingV1
class ReportingSourceSliceRequestV1 (**data: Any)-
Expand source code
class ReportingSourceSliceRequestV1(_Frozen): """One bounded, frozen unit of work handed to an executor. Everything here is immutable for the life of the execution. The executor echoes it into the manifest; the conformance validator compares the two field by field, so a source that "helpfully" widens a date range or swaps a currency fails admission rather than silently publishing something the producer never asked for. """ contract_version: Literal["1.0"] = REPORTING_SOURCE_CONTRACT_VERSION_V1 identity: ReportingSourceIdentityV1 adapter_build: ReportingAdapterBuildIdentityV1 offering_id: Identifier publication_namespace: Annotated[str, StringConstraints(pattern=_NAMESPACE_PATTERN)] publication_class: ReportingPublicationClass contract: ReportingContractIdentityV1 period: ReportingSourcePeriodV1 revision_kind: ReportingRevisionKind supersedes_publication_id: OpaqueReference | None = None trigger: ReportingTriggerKind = "scheduled_poll" coverage: ReportingSourceCoverageRequestV1 requested_metrics: list[MetricName] = Field(min_length=1, max_length=1_000) requested_dimensions: list[MetricName] = Field(default_factory=list, max_length=1_000) currency: Annotated[str, StringConstraints(pattern=r"^[A-Z]{3}$")] prior_checkpoint: ReportingSourceCheckpointV1 | None = None deadline_at: datetime @model_validator(mode="after") def _request_is_coherent(self) -> ReportingSourceSliceRequestV1: snapshot = self.publication_class == "PROVISIONAL_SNAPSHOT" if snapshot != (self.revision_kind == "snapshot"): raise ValueError( "snapshot offerings produce snapshot revisions; authoritative offerings do not" ) if (self.revision_kind == "correction") != (self.supersedes_publication_id is not None): raise ValueError( "supersedes_publication_id is required for a correction and forbidden otherwise" ) for label, values in ( ("requested_metrics", self.requested_metrics), ("requested_dimensions", self.requested_dimensions), ): if len(set(values)) != len(values): raise ValueError(f"{label} must be unique") cells = len(self.coverage.constituents) * len(self.requested_metrics) if cells > SOURCE_BATCH_MANIFEST_MAX_METRIC_AVAILABILITY_V1: raise ValueError( "a source slice may request at most " f"{SOURCE_BATCH_MANIFEST_MAX_METRIC_AVAILABILITY_V1} constituent-metric cells, " f"got {cells}" ) _aware(self.deadline_at, "deadline_at") return self @property def product_ids(self) -> tuple[str, ...]: """The unique products represented by the constituent denominator.""" return tuple( dict.fromkeys(constituent.product_id for constituent in self.coverage.constituents) ) @property def media_buy_ids(self) -> tuple[str, ...]: """The unique media buys represented by the constituent denominator.""" return tuple( dict.fromkeys( constituent.media_buy_id for constituent in self.coverage.constituents if isinstance(constituent, MediaBuyConstituentV1) ) ) @property def package_ids(self) -> tuple[str, ...]: """The unique packages represented by the constituent denominator.""" return tuple( dict.fromkeys( constituent.package_id for constituent in self.coverage.constituents if isinstance(constituent, PackageItemConstituentV1) ) )One bounded, frozen unit of work handed to an executor.
Everything here is immutable for the life of the execution. The executor echoes it into the manifest; the conformance validator compares the two field by field, so a source that "helpfully" widens a date range or swaps a currency fails admission rather than silently publishing something the producer never asked for.
Create a new model by parsing and validating input data from keyword arguments.
Raises [
ValidationError][pydantic_core.ValidationError] if the input data cannot be validated to form a valid model.selfis explicitly positional-only to allowselfas a field name.Ancestors
- adcp.reporting.source._Frozen
- pydantic.main.BaseModel
Class variables
var adapter_build : ReportingAdapterBuildIdentityV1var contract : ReportingContractIdentityV1var contract_version : Literal['1.0']var coverage : adcp.reporting.source.ReportingSourceCoverageRequestV1var currency : strvar deadline_at : datetime.datetimevar identity : ReportingSourceIdentityV1var model_configvar offering_id : strvar period : ReportingSourcePeriodV1var prior_checkpoint : adcp.reporting.source.ReportingSourceCheckpointV1 | Nonevar publication_class : Literal['PROVISIONAL_SNAPSHOT', 'AUTHORITATIVE']var publication_namespace : strvar requested_dimensions : list[str]var requested_metrics : list[str]var revision_kind : Literal['snapshot', 'authoritative', 'correction']var supersedes_publication_id : str | Nonevar trigger : Literal['scheduled_poll', 'webhook_hint', 'retry', 'backfill', 'gap_recovery', 'manual_replay']
Instance variables
prop media_buy_ids : tuple[str, ...]-
Expand source code
@property def media_buy_ids(self) -> tuple[str, ...]: """The unique media buys represented by the constituent denominator.""" return tuple( dict.fromkeys( constituent.media_buy_id for constituent in self.coverage.constituents if isinstance(constituent, MediaBuyConstituentV1) ) )The unique media buys represented by the constituent denominator.
prop package_ids : tuple[str, ...]-
Expand source code
@property def package_ids(self) -> tuple[str, ...]: """The unique packages represented by the constituent denominator.""" return tuple( dict.fromkeys( constituent.package_id for constituent in self.coverage.constituents if isinstance(constituent, PackageItemConstituentV1) ) )The unique packages represented by the constituent denominator.
prop product_ids : tuple[str, ...]-
Expand source code
@property def product_ids(self) -> tuple[str, ...]: """The unique products represented by the constituent denominator.""" return tuple( dict.fromkeys(constituent.product_id for constituent in self.coverage.constituents) )The unique products represented by the constituent denominator.
class ReportingSourceStagedObjectReader (*args, **kwargs)-
Expand source code
@runtime_checkable class ReportingSourceStagedObjectReader(Protocol): """Reads one immutable, generation-pinned staged object. The reader MUST authorize the ``account_id`` / ``source_scope`` pair before returning bytes. Possession of an opaque object reference is not authorization: references leak into logs, manifests, and support tickets. """ async def read( self, *, object_ref: str, object_generation: str, account_id: str, source_scope: Mapping[str, Any], cancel: asyncio.Event, ) -> bytes: ...Reads one immutable, generation-pinned staged object.
The reader MUST authorize the
account_id/source_scopepair before returning bytes. Possession of an opaque object reference is not authorization: references leak into logs, manifests, and support tickets.Ancestors
- typing.Protocol
- typing.Generic
Methods
async def read(self,
*,
object_ref: str,
object_generation: str,
account_id: str,
source_scope: Mapping[str, Any],
cancel: asyncio.Event) ‑> bytes-
Expand source code
async def read( self, *, object_ref: str, object_generation: str, account_id: str, source_scope: Mapping[str, Any], cancel: asyncio.Event, ) -> bytes: ...
class ReportingSourceWindowingV1 (**data: Any)-
Expand source code
class ReportingSourceWindowingV1(_Frozen): """How a source's reporting window behaves.""" kind: Literal["cumulative_current_window", "fixed_closed_window", "provider_defined"] minimum_window: IsoDuration maximum_window: IsoDuration overlapping_windows_supported: bool = False @model_validator(mode="after") def _window_range_is_ordered(self) -> ReportingSourceWindowingV1: if iso_duration_milliseconds_v1(self.minimum_window) > iso_duration_milliseconds_v1( self.maximum_window ): raise ValueError("maximum_window must be at least minimum_window") return selfHow a source's reporting window behaves.
Create a new model by parsing and validating input data from keyword arguments.
Raises [
ValidationError][pydantic_core.ValidationError] if the input data cannot be validated to form a valid model.selfis explicitly positional-only to allowselfas a field name.Ancestors
- adcp.reporting.source._Frozen
- pydantic.main.BaseModel
Class variables
var kind : Literal['cumulative_current_window', 'fixed_closed_window', 'provider_defined']var maximum_window : strvar minimum_window : strvar model_configvar overlapping_windows_supported : bool
class SourceBatchManifestReferenceV1 (**data: Any)-
Expand source code
class SourceBatchManifestReferenceV1(_Frozen): """A pointer to sealed manifest bytes, plus what those bytes must be. Possessing this reference is not authorization to read the bytes; the reader still enforces scope. It is only enough to *verify* them. """ staged_commit_ref: OpaqueReference manifest_sha256: Sha256 byte_count: int = Field(ge=1, le=SOURCE_BATCH_MANIFEST_MAX_BYTES_V1) encoding: Literal["canonical_json_utf8_v1"] = "canonical_json_utf8_v1"A pointer to sealed manifest bytes, plus what those bytes must be.
Possessing this reference is not authorization to read the bytes; the reader still enforces scope. It is only enough to verify them.
Create a new model by parsing and validating input data from keyword arguments.
Raises [
ValidationError][pydantic_core.ValidationError] if the input data cannot be validated to form a valid model.selfis explicitly positional-only to allowselfas a field name.Ancestors
- adcp.reporting.source._Frozen
- pydantic.main.BaseModel
Class variables
var byte_count : intvar encoding : Literal['canonical_json_utf8_v1']var manifest_sha256 : strvar model_configvar staged_commit_ref : str
class SourceBatchManifestV1 (**data: Any)-
Expand source code
class SourceBatchManifestV1(_Frozen): """One immutable publication: what was fetched, from when, and how completely. The invariants enforced here are *single-publication* admission rules. Cross-revision rules (a snapshot may not regress freshness, a correction must follow an authoritative publication) live in :func:`adcp.reporting.conformance.validate_reporting_revision_sequence` because they need more than one manifest to evaluate. Time evidence is ordered: ``data_through`` cannot follow the observation, acquisition cannot precede the observation, and finality evidence must be observed between ``data_through`` and acquisition. Authoritative finality evidence additionally cannot predate the period end -- you cannot declare a period official before it has finished. """ manifest_version: Literal["1.0"] = REPORTING_SOURCE_CONTRACT_VERSION_V1 strictness: ReportingManifestStrictness = "basic" complete: Literal[True] = True identity: ReportingSourceIdentityV1 publication_namespace: Annotated[str, StringConstraints(pattern=_NAMESPACE_PATTERN)] publication_id: OpaqueReference publication_class: ReportingPublicationClass content_fingerprint: Fingerprint adapter_build: ReportingAdapterBuildIdentityV1 offering_id: Identifier contract: ReportingContractIdentityV1 period: ReportingSourcePeriodV1 revision_kind: ReportingRevisionKind supersedes_publication_id: OpaqueReference | None = None finality_evidence: ReportingFinalityEvidenceV1 observed_at: datetime data_through: datetime acquired_at: datetime currency: Annotated[str, StringConstraints(pattern=r"^[A-Z]{3}$")] requested_dimensions: list[MetricName] = Field(default_factory=list, max_length=1_000) objects: list[SourceBatchObjectV1] = Field(min_length=1, max_length=100_000) object_set_sha256: Sha256 row_count: int = Field(ge=0) byte_count: int = Field(ge=0) control_totals: list[SourceControlTotalV1] = Field(default_factory=list, max_length=1_000) metric_availability: list[ReportingMetricAvailabilityV1] = Field( min_length=1, max_length=SOURCE_BATCH_MANIFEST_MAX_METRIC_AVAILABILITY_V1 ) coverage: SourceBatchCoverageV1 explicit_zero: bool = False event_time_range: tuple[datetime, datetime] | None = None evidence: SourceBatchEvidenceV1 | None = None prior_checkpoint: ReportingSourceCheckpointV1 | None = None warnings: list[Annotated[str, StringConstraints(min_length=1, max_length=1_024)]] = Field( default_factory=list, max_length=1_000 ) @model_validator(mode="after") def _manifest_is_internally_consistent(self) -> SourceBatchManifestV1: self._validate_objects() self._validate_time_evidence() self._validate_finality() self._validate_cells() self._validate_zero_rows() self._validate_event_range() self._validate_strictness() expected = publication_content_fingerprint_v1(self) if expected != self.content_fingerprint: raise ValueError( f"content_fingerprint must bind the publication content; expected {expected}" ) return self # -- object set ------------------------------------------------------ def _validate_objects(self) -> None: seen_refs: set[str] = set() seen_generations: set[tuple[str, str]] = set() byte_count = 0 row_count = 0 for index, item in enumerate(self.objects): if item.ordinal != index: raise ValueError("object ordinals must be contiguous and ordered from zero") if item.object_ref in seen_refs: raise ValueError(f"object_ref {item.object_ref!r} is not unique") if (item.object_ref, item.object_generation) in seen_generations: raise ValueError("object reference and generation pairs must be unique") seen_refs.add(item.object_ref) seen_generations.add((item.object_ref, item.object_generation)) byte_count += item.byte_count row_count += item.row_count if byte_count != self.byte_count: raise ValueError(f"byte_count must equal the object sum, expected {byte_count}") if row_count != self.row_count: raise ValueError(f"row_count must equal the object sum, expected {row_count}") expected = source_batch_object_set_sha256_v1(self.objects) if expected != self.object_set_sha256: raise ValueError( f"object_set_sha256 must bind the ordered pinned objects, expected {expected}" ) # -- time ------------------------------------------------------------ def _validate_time_evidence(self) -> None: start = _aware(self.period.start, "period.start") end = _aware(self.period.end, "period.end") cutoff = _aware(self.period.source_read_cutoff_at, "period.source_read_cutoff_at") observed_at = _aware(self.observed_at, "observed_at") data_through = _aware(self.data_through, "data_through") acquired_at = _aware(self.acquired_at, "acquired_at") finality_at = _aware(self.finality_evidence.observed_at, "finality_evidence.observed_at") if not (start <= data_through <= end): raise ValueError("data_through must be inside the reporting period") if data_through > observed_at or data_through > cutoff: raise ValueError("data_through must not follow observed_at or the source read cutoff") if acquired_at < observed_at: raise ValueError("acquired_at must not precede the reported source observation") if finality_at < data_through or finality_at > acquired_at: raise ValueError( "finality evidence must be observed no earlier than data_through and " "no later than acquisition completion" ) if self.publication_class == "AUTHORITATIVE" and finality_at < end: raise ValueError("authoritative finality evidence must not predate the period end") for item in self.coverage.constituents: if item.data_through is not None: self._require_bounded_watermark(item.data_through, "constituent") for cell in self.metric_availability: if cell.data_through is not None: self._require_bounded_watermark(cell.data_through, "metric") self._validate_authoritative_full_window(end) def _require_bounded_watermark(self, value: datetime, label: str) -> None: moment = _aware(value, f"{label} data_through") if not ( _aware(self.period.start, "period.start") <= moment <= min( _aware(self.period.end, "period.end"), _aware(self.data_through, "data_through"), _aware(self.period.source_read_cutoff_at, "period.source_read_cutoff_at"), ) ): raise ValueError( f"granular {label} data_through must be inside the period and no later than " "the batch watermark or the requested source cutoff" ) def _validate_authoritative_full_window(self, end: datetime) -> None: # Full official coverage is the strongest claim in the contract: it # says the period is closed and complete. Data that never reached the # period boundary stays delayed/stale/partial instead. if self.publication_class != "AUTHORITATIVE" or self.coverage.status != "full": return if _aware(self.data_through, "data_through") != end: raise ValueError("full authoritative coverage must prove the complete reporting window") for item in self.coverage.constituents: if item.status in _AVAILABLE_STATUSES and ( item.data_through is None or _aware(item.data_through, "data_through") != end ): raise ValueError( "available constituents in full authoritative coverage must prove the " "complete reporting window" ) for cell in self.metric_availability: if cell.status in _AVAILABLE_STATUSES and ( cell.data_through is None or _aware(cell.data_through, "data_through") != end ): raise ValueError( "available metrics in full authoritative coverage must prove the " "complete reporting window" ) # -- finality -------------------------------------------------------- def _validate_finality(self) -> None: snapshot = self.publication_class == "PROVISIONAL_SNAPSHOT" if snapshot != (self.revision_kind == "snapshot"): raise ValueError("snapshot and authoritative publications use distinct revision kinds") if (self.revision_kind == "correction") != (self.supersedes_publication_id is not None): raise ValueError( "supersedes_publication_id is required for a correction and forbidden otherwise" ) basis = self.finality_evidence.basis if snapshot and basis != "provisional_observation": raise ValueError("a snapshot must be identified as a provisional observation") if not snapshot and basis == "provisional_observation": raise ValueError("authoritative data needs finality evidence beyond an observation") # -- cells ----------------------------------------------------------- def _validate_cells(self) -> None: cells = [(cell.constituent_id, cell.metric) for cell in self.metric_availability] if len(set(cells)) != len(cells): raise ValueError( "metric availability must carry one record per constituent-metric cell" ) by_constituent: dict[str, set[str]] = {} applicable_constituents: set[str] = set() for cell in self.metric_availability: by_constituent.setdefault(cell.constituent_id, set()).add(cell.metric) if cell.status != "unsupported": applicable_constituents.add(cell.constituent_id) published_metrics = {metric for _, metric in cells} statuses = {item.constituent_id: item.status for item in self.coverage.constituents} for constituent_id, status in statuses.items(): if by_constituent.get(constituent_id, set()) != published_metrics: raise ValueError( f"constituent {constituent_id!r} needs one availability record for every " "published metric" ) if status != "unsupported" and constituent_id not in applicable_constituents: raise ValueError("a constituent without applicable metrics must be unsupported") for cell in self.metric_availability: constituent_status = statuses.get(cell.constituent_id) if constituent_status is None: raise ValueError( f"metric evidence for {cell.constituent_id!r} is outside the coverage " "denominator" ) if cell.status not in _CONSTITUENT_ALLOWS_METRIC[constituent_status]: raise ValueError( f"metric status {cell.status!r} contradicts constituent status " f"{constituent_status!r}" ) # -- zero rows ------------------------------------------------------- def _validate_zero_rows(self) -> None: """A zero-row period has exactly two honest shapes, and they differ. ``explicit_zero``: the source answered for every requested cell and the answer was zero. That is data, and it satisfies an obligation. Unavailable: the source did not answer. That is not data, and it must not be dressed up as a zero -- otherwise a dead feed looks like a quiet campaign forever. """ if self.explicit_zero: if self.event_time_range is not None: raise ValueError("an explicit-zero batch must not declare an event-time range") if self.row_count != 0 or any(item.row_count for item in self.objects): raise ValueError("every explicit-zero object must have zero rows") if any(not total.is_zero for total in self.control_totals): raise ValueError("every explicit-zero control total must be zero") if any(item.status != "explicit_zero" for item in self.coverage.constituents) or any( cell.status != "explicit_zero" for cell in self.metric_availability ): raise ValueError( "a batch-wide explicit zero requires explicit-zero evidence for every " "requested constituent and metric" ) return if self.row_count == 0: unavailable = ( self.coverage.status != "full" and all( item.status not in _AVAILABLE_STATUSES for item in self.coverage.constituents ) and all(cell.status not in _AVAILABLE_STATUSES for cell in self.metric_availability) ) if not unavailable: raise ValueError("a completed zero-row available batch must be explicitly zero") if self.event_time_range is not None: raise ValueError( "a zero-row unavailable batch must not declare an event-time range" ) if any(not total.is_zero for total in self.control_totals): raise ValueError("a zero-row unavailable batch cannot declare nonzero totals") return if self.event_time_range is None: raise ValueError("a nonzero batch must declare its event-time range") def _validate_event_range(self) -> None: if self.event_time_range is None: return # Row-level evidence, not a second freshness clock: rows beyond the # watermark or the frozen cutoff make the manifest self-contradictory. event_start = _aware(self.event_time_range[0], "event_time_range start") event_end = _aware(self.event_time_range[1], "event_time_range end") if event_end < event_start: raise ValueError("event-time range end must not precede its start") if event_start < _aware(self.period.start, "period.start"): raise ValueError("event-time range must start within the period") if event_end > _aware(self.period.end, "period.end"): raise ValueError("event-time range must end within the period") if event_end > _aware(self.data_through, "data_through") or event_end > _aware( self.period.source_read_cutoff_at, "period.source_read_cutoff_at" ): raise ValueError( "event-time range must end no later than data_through or the source cutoff" ) def _validate_strictness(self) -> None: if self.strictness == "evidenced": if self.evidence is None: raise ValueError("an evidenced manifest must retain per-call provider evidence") evidenced_rows = sum(page.item_count for page in self.evidence.pages) if evidenced_rows != self.row_count: raise ValueError( "evidenced page item counts must equal the normalized row count, " f"expected {self.row_count}" ) ordinals = [ ordinal for page in self.evidence.pages for ordinal in page.output_object_ordinals ] if sorted(ordinals) != [item.ordinal for item in self.objects]: raise ValueError("page evidence must exactly partition the staged object set") acquired_at = _aware(self.acquired_at, "acquired_at") for poll in self.evidence.polls: if _aware(poll.observed_at, "observed_at") > acquired_at: raise ValueError("job poll observations must not postdate acquisition") for attempt in self.evidence.retry_attempts: if _aware(attempt.attempted_at, "attempted_at") > acquired_at: raise ValueError("retry attempts must not postdate acquisition") elif self.evidence is not None: raise ValueError( "a basic manifest must omit per-call provider evidence; " "declare strictness='evidenced' to retain it" )One immutable publication: what was fetched, from when, and how completely.
The invariants enforced here are single-publication admission rules. Cross-revision rules (a snapshot may not regress freshness, a correction must follow an authoritative publication) live in :func:
validate_reporting_revision_sequence()because they need more than one manifest to evaluate.Time evidence is ordered:
data_throughcannot follow the observation, acquisition cannot precede the observation, and finality evidence must be observed betweendata_throughand acquisition. Authoritative finality evidence additionally cannot predate the period end – you cannot declare a period official before it has finished.Create a new model by parsing and validating input data from keyword arguments.
Raises [
ValidationError][pydantic_core.ValidationError] if the input data cannot be validated to form a valid model.selfis explicitly positional-only to allowselfas a field name.Ancestors
- adcp.reporting.source._Frozen
- pydantic.main.BaseModel
Class variables
var acquired_at : datetime.datetimevar adapter_build : ReportingAdapterBuildIdentityV1var byte_count : intvar complete : Literal[True]var content_fingerprint : strvar contract : ReportingContractIdentityV1var control_totals : list[SourceControlTotalV1]var coverage : adcp.reporting.source.SourceBatchCoverageV1var currency : strvar data_through : datetime.datetimevar event_time_range : tuple[datetime.datetime, datetime.datetime] | Nonevar evidence : adcp.reporting.source.SourceBatchEvidenceV1 | Nonevar explicit_zero : boolvar finality_evidence : ReportingFinalityEvidenceV1var identity : ReportingSourceIdentityV1var manifest_version : Literal['1.0']var metric_availability : list[ReportingMetricAvailabilityV1]var model_configvar object_set_sha256 : strvar objects : list[SourceBatchObjectV1]var observed_at : datetime.datetimevar offering_id : strvar period : ReportingSourcePeriodV1var prior_checkpoint : adcp.reporting.source.ReportingSourceCheckpointV1 | Nonevar publication_class : Literal['PROVISIONAL_SNAPSHOT', 'AUTHORITATIVE']var publication_id : strvar publication_namespace : strvar requested_dimensions : list[str]var revision_kind : Literal['snapshot', 'authoritative', 'correction']var row_count : intvar strictness : Literal['basic', 'evidenced']var supersedes_publication_id : str | Nonevar warnings : list[str]
class SourceBatchObjectV1 (**data: Any)-
Expand source code
class SourceBatchObjectV1(_Frozen): """One immutable, generation-pinned staged object holding normalized rows. ``object_generation`` is what makes the reference immutable: an object store key can be overwritten, a key plus generation cannot. Re-reading the pair years later either returns the same bytes or fails -- it never returns different bytes. """ ordinal: int = Field(ge=0) object_ref: OpaqueReference object_generation: OpaqueReference media_type: Literal[ "application/json", "application/x-ndjson", "text/csv", "application/vnd.apache.parquet", ] compression: Literal["none", "gzip", "zstd"] = "none" sha256: Sha256 byte_count: int = Field(ge=0) row_count: int = Field(ge=0)One immutable, generation-pinned staged object holding normalized rows.
object_generationis what makes the reference immutable: an object store key can be overwritten, a key plus generation cannot. Re-reading the pair years later either returns the same bytes or fails – it never returns different bytes.Create a new model by parsing and validating input data from keyword arguments.
Raises [
ValidationError][pydantic_core.ValidationError] if the input data cannot be validated to form a valid model.selfis explicitly positional-only to allowselfas a field name.Ancestors
- adcp.reporting.source._Frozen
- pydantic.main.BaseModel
Class variables
var byte_count : intvar compression : Literal['none', 'gzip', 'zstd']var media_type : Literal['application/json', 'application/x-ndjson', 'text/csv', 'application/vnd.apache.parquet']var model_configvar object_generation : strvar object_ref : strvar ordinal : intvar row_count : intvar sha256 : str
class SourceControlTotalV1 (**data: Any)-
Expand source code
class SourceControlTotalV1(_Frozen): """A named total a consumer can check its own ingest against. Decimals are strings. A JSON float here would round differently in two languages and turn an equality check into a flake. """ name: MetricName value: str value_type: Literal["integer", "decimal"] = "integer" unit: Annotated[str, StringConstraints(min_length=1, max_length=32)] | None = None @model_validator(mode="after") def _value_matches_its_type(self) -> SourceControlTotalV1: pattern = _INTEGER_TOTAL_PATTERN if self.value_type == "integer" else _DECIMAL_TOTAL_PATTERN if not re.match(pattern, self.value): raise ValueError(f"{self.value!r} is not a valid {self.value_type} control total") return self @property def is_zero(self) -> bool: return re.fullmatch(r"-?0(?:\.0+)?", self.value) is not NoneA named total a consumer can check its own ingest against.
Decimals are strings. A JSON float here would round differently in two languages and turn an equality check into a flake.
Create a new model by parsing and validating input data from keyword arguments.
Raises [
ValidationError][pydantic_core.ValidationError] if the input data cannot be validated to form a valid model.selfis explicitly positional-only to allowselfas a field name.Ancestors
- adcp.reporting.source._Frozen
- pydantic.main.BaseModel
Class variables
var model_configvar name : strvar unit : str | Nonevar value : strvar value_type : Literal['integer', 'decimal']
Instance variables
prop is_zero : bool-
Expand source code
@property def is_zero(self) -> bool: return re.fullmatch(r"-?0(?:\.0+)?", self.value) is not None