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_id means 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_key cannot 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 encoded

Seal a manifest into its canonical_json_utf8_v1 bytes.

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_sha256 itself 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 self

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.

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.

self is explicitly positional-only to allow self as 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 : str
var days_after_period_end : int
var expected_availability_lag : str
var model_config
var publication_class : Literal['AUTHORITATIVE']
var revision_semantics : Literal['official_with_declared_correction_policy']
var source_local_ready_time : str
var 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: ExternalId

A 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.

self is explicitly positional-only to allow self as a field name.

Ancestors

  • adcp.reporting.source._Frozen
  • pydantic.main.BaseModel

Class variables

var constituent_id : str
var constituent_kind : Literal['media_buy']
var media_buy_id : str
var model_config
var 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: ExternalId

A 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.

self is explicitly positional-only to allow self as a field name.

Ancestors

  • adcp.reporting.source._Frozen
  • pydantic.main.BaseModel

Class variables

var constituent_id : str
var constituent_kind : Literal['package_item']
var model_config
var package_id : str
var 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: ExternalId

A 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.

self is explicitly positional-only to allow self as a field name.

Ancestors

  • adcp.reporting.source._Frozen
  • pydantic.main.BaseModel

Class variables

var constituent_id : str
var constituent_kind : Literal['product']
var model_config
var 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 self

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.

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.

self is explicitly positional-only to allow self as 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 : str
var fastest_safe_cadence : str
var model_config
var official_close_lag : str | None
var publication_class : Literal['PROVISIONAL_SNAPSHOT']
var restatement_cadence : str | None
var restatement_window : str | None
var 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: Sha256

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.

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.

self is explicitly positional-only to allow self as a field name.

Ancestors

  • adcp.reporting.source._Frozen
  • pydantic.main.BaseModel

Class variables

var adapter_build_sha256 : str
var adapter_id : str
var adapter_version : str
var executor_id : str
var 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_id

One 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.

self is explicitly positional-only to allow self as a field name.

Ancestors

  • adcp.reporting.source._Frozen
  • pydantic.main.BaseModel

Class variables

var constituent : ProductConstituentV1 | PackageItemConstituentV1 | MediaBuyConstituentV1
var data_through : datetime.datetime | None
var model_config
var reason : str | None
var 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 self

The 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.

self is explicitly positional-only to allow self as a field name.

Ancestors

  • adcp.reporting.source._Frozen
  • pydantic.main.BaseModel

Class variables

var model_config
var report_definition_id : str
var report_definition_sha256 : str
var report_definition_uri : str
var reporting_profile : str
var schema_dialect : Literal['https://json-schema.org/draft/2020-12/schema']
var schema_ref_policy : Literal['local_fragment_only']
var schema_sha256 : str
var schema_uri : str
var 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 self

Why 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.

self is explicitly positional-only to allow self as 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 | None
var model_config
var observed_at : datetime.datetime
var 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 self

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".

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.

self is explicitly positional-only to allow self as a field name.

Ancestors

  • adcp.reporting.source._Frozen
  • pydantic.main.BaseModel

Class variables

var constituent_id : str
var data_through : datetime.datetime | None
var metric : str
var model_config
var reason : str | None
var semantic_contract_id : str
var semantic_contract_sha256 : str
var semantic_contract_version : str
var 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: Sha256

One 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.

self is explicitly positional-only to allow self as a field name.

Ancestors

  • adcp.reporting.source._Frozen
  • pydantic.main.BaseModel

Class variables

var job_id : str
var model_config
var observed_at : datetime.datetime
var ordinal : int
var request_id : str
var response_sha256 : str
var 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: Sha256

One 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.

self is explicitly positional-only to allow self as a field name.

Ancestors

  • adcp.reporting.source._Frozen
  • pydantic.main.BaseModel

Class variables

var job_id : str
var model_config
var submission_request_id : str
var 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.

self is explicitly positional-only to allow self as a field name.

Ancestors

  • adcp.reporting.source._Frozen
  • pydantic.main.BaseModel

Class variables

var attempt_id : str
var attempted_at : datetime.datetime
var classification : Literal['RATE_LIMITED', 'QUOTA_EXHAUSTED', 'PROVIDER_TRANSIENT', 'PROVIDER_JOB_FAILED']
var model_config
var 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.

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.

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.

self is explicitly positional-only to allow self as a field name.

Ancestors

  • adcp.reporting.source._Frozen
  • pydantic.main.BaseModel

Class variables

var adapter_build : ReportingAdapterBuildIdentityV1
var capabilities_sha256 : str
var capability_version : str
var contract_version : Literal['1.0']
var model_config
var 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_source and :class:~adcp.reporting.ledger.ReportingProducer both 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 self

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.

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.

self is explicitly positional-only to allow self as 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_config
var provider_code : str | None
var retry : Literal['retryable', 'terminal', 'cancelled']
var retry_after_seconds : float | None
var safe_message : str
var 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: SourceBatchManifestReferenceV1

The 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.

self is explicitly positional-only to allow self as a field name.

Ancestors

  • adcp.reporting.source._Frozen
  • pydantic.main.BaseModel

Class variables

var contract_version : Literal['1.0']
var identity : ReportingSourceIdentityV1
var manifest : SourceBatchManifestReferenceV1
var model_config
var 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 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.

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.

self is explicitly positional-only to allow self as a field name.

Ancestors

  • adcp.reporting.source._Frozen
  • pydantic.main.BaseModel

Class variables

var error : ReportingSourceErrorV1 | None
var manifest_bytes : bytes | None
var model_config
var response : ReportingSourceExecutionResponseV1 | None

Static methods

def completed(*,
request: ReportingSourceSliceRequestV1,
manifest: SourceBatchManifestReferenceV1,
manifest_bytes: bytes) ‑> ReportingSourceExecutorResult
def 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 self

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.

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.

self is explicitly positional-only to allow self as a field name.

Ancestors

  • adcp.reporting.source._Frozen
  • pydantic.main.BaseModel

Class variables

var account_id : str
var delivery_config_id : str
var delivery_config_version : int
var logical_slice_fingerprint : str
var model_config
var period_key : str
var report_definition_id : str
var reporting_obligation_id : str
var run_id : str
var source_execution_key : str
var 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 self

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.

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.

self is explicitly positional-only to allow self as a field name.

Ancestors

  • adcp.reporting.source._Frozen
  • pydantic.main.BaseModel

Class variables

var end : datetime.datetime
var grain : str
var model_config
var period_key : str
var source_local_date : str
var source_read_cutoff_at : datetime.datetime
var source_timezone : str
var start : datetime.datetime
var 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.

self is explicitly positional-only to allow self as a field name.

Ancestors

  • adcp.reporting.source._Frozen
  • pydantic.main.BaseModel

Class variables

var adapter_build : ReportingAdapterBuildIdentityV1
var contract : ReportingContractIdentityV1
var contract_version : Literal['1.0']
var coverage : adcp.reporting.source.ReportingSourceCoverageRequestV1
var currency : str
var deadline_at : datetime.datetime
var identity : ReportingSourceIdentityV1
var model_config
var offering_id : str
var period : ReportingSourcePeriodV1
var prior_checkpoint : adcp.reporting.source.ReportingSourceCheckpointV1 | None
var publication_class : Literal['PROVISIONAL_SNAPSHOT', 'AUTHORITATIVE']
var publication_namespace : str
var requested_dimensions : list[str]
var requested_metrics : list[str]
var revision_kind : Literal['snapshot', 'authoritative', 'correction']
var supersedes_publication_id : str | None
var 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_scope pair 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 self

How 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.

self is explicitly positional-only to allow self as 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 : str
var minimum_window : str
var model_config
var 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.

self is explicitly positional-only to allow self as a field name.

Ancestors

  • adcp.reporting.source._Frozen
  • pydantic.main.BaseModel

Class variables

var byte_count : int
var encoding : Literal['canonical_json_utf8_v1']
var manifest_sha256 : str
var model_config
var 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_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.

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.

self is explicitly positional-only to allow self as a field name.

Ancestors

  • adcp.reporting.source._Frozen
  • pydantic.main.BaseModel

Class variables

var acquired_at : datetime.datetime
var adapter_build : ReportingAdapterBuildIdentityV1
var byte_count : int
var complete : Literal[True]
var content_fingerprint : str
var contract : ReportingContractIdentityV1
var control_totals : list[SourceControlTotalV1]
var coverage : adcp.reporting.source.SourceBatchCoverageV1
var currency : str
var data_through : datetime.datetime
var event_time_range : tuple[datetime.datetime, datetime.datetime] | None
var evidence : adcp.reporting.source.SourceBatchEvidenceV1 | None
var explicit_zero : bool
var finality_evidence : ReportingFinalityEvidenceV1
var identity : ReportingSourceIdentityV1
var manifest_version : Literal['1.0']
var metric_availability : list[ReportingMetricAvailabilityV1]
var model_config
var object_set_sha256 : str
var objects : list[SourceBatchObjectV1]
var observed_at : datetime.datetime
var offering_id : str
var period : ReportingSourcePeriodV1
var prior_checkpoint : adcp.reporting.source.ReportingSourceCheckpointV1 | None
var publication_class : Literal['PROVISIONAL_SNAPSHOT', 'AUTHORITATIVE']
var publication_id : str
var publication_namespace : str
var requested_dimensions : list[str]
var revision_kind : Literal['snapshot', 'authoritative', 'correction']
var row_count : int
var strictness : Literal['basic', 'evidenced']
var supersedes_publication_id : str | None
var 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_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.

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.

self is explicitly positional-only to allow self as a field name.

Ancestors

  • adcp.reporting.source._Frozen
  • pydantic.main.BaseModel

Class variables

var byte_count : int
var compression : Literal['none', 'gzip', 'zstd']
var media_type : Literal['application/json', 'application/x-ndjson', 'text/csv', 'application/vnd.apache.parquet']
var model_config
var object_generation : str
var object_ref : str
var ordinal : int
var row_count : int
var 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 None

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.

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.

self is explicitly positional-only to allow self as a field name.

Ancestors

  • adcp.reporting.source._Frozen
  • pydantic.main.BaseModel

Class variables

var model_config
var name : str
var unit : str | None
var value : str
var 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