From e5ab6816a0c0fe59d742d48b94fd9638ec75d9b6 Mon Sep 17 00:00:00 2001 From: Brian O'Kelley Date: Tue, 29 Sep 2026 14:57:04 +0000 Subject: [PATCH 1/3] fix(reporting): allow conditional metrics under full coverage --- docs/reporting-source-adapters.md | 26 +- src/adcp/reporting/conformance.py | 9 +- src/adcp/reporting/inline_source.py | 27 +- src/adcp/reporting/source.py | 16 +- .../test_inline_cell_availability.py | 66 +++- .../test_reporting_conditional_metrics.py | 296 ++++++++++++++++++ ...reporting_evidence_currency_integration.py | 4 +- tests/test_reporting_source_contract.py | 55 ++++ 8 files changed, 469 insertions(+), 30 deletions(-) create mode 100644 tests/conformance/reporting/test_reporting_conditional_metrics.py diff --git a/docs/reporting-source-adapters.md b/docs/reporting-source-adapters.md index 5ffe41794..a09393714 100644 --- a/docs/reporting-source-adapters.md +++ b/docs/reporting-source-adapters.md @@ -26,8 +26,16 @@ return InlineFetchResult( This publishes one `partial` constituent with four independent metric cells. It can complete a request whose `coverage.expected` is `partial`. A `full` -request receives retryable `PARTIAL_RESULT` and no publication until every -requested cell is available. +request receives retryable `PARTIAL_RESULT` while `viewability` is delayed. +When every applicable cell is available, the constituent is `present` and can +satisfy full coverage even though `completed_views` remains `unsupported`. + +Offering metrics with `support: partial` are requestable alongside `exact` +metrics in both the service and source conformance. Every requested cell must +still have manifest evidence with the offering's semantic contract. The inline +adapter builds that evidence from the defaults and `cell_availability` +overrides; use an override for each inapplicable or unready cell. An offering's +`unavailable` metrics cannot pass source conformance. | Constructor | Meaning and required evidence | | --- | --- | @@ -35,7 +43,7 @@ requested cell is available. | `MetricEvidence.explicit_zero(data_through=...)` | An observed zero. The watermark defaults to the fetch watermark. Any supplied values must be finite numeric zeros. | | `MetricEvidence.missing(reason)` | No answer for this metric. A stable reason is required; a watermark is forbidden. | | `MetricEvidence.delayed(reason, data_through=...)` | Not ready yet. A stable reason is required; retain a known watermark when available. | -| `MetricEvidence.unavailable(reason)` | The source cannot measure this metric for this constituent. Emits the existing wire status `unsupported`; a reason is required and a watermark is forbidden. | +| `MetricEvidence.unavailable(reason)` | This metric is inapplicable to this constituent. Emits `unsupported`; a reason is required and a watermark is forbidden. Use `missing` or `delayed` for an applicable measurement that has not arrived. | Evidence is immutable. Direct `MetricEvidence(...)` construction enforces the same invariants. Reasons use the existing manifest format: bounded ASCII, @@ -56,8 +64,16 @@ behavior. Omitted cells inherit their constituent's status, reason, and watermark. Explicit cells take precedence over `covered_constituent_ids` and `unavailable_constituents`, including explicit measurements for an otherwise missing constituent. Coverage is then reconciled from the resolved cells: -uniform statuses stay uniform, present plus explicit-zero is `present`, and -other mixtures are `partial`. +uniform statuses stay uniform. A mixture of `present`, `explicit_zero`, and +`unsupported` cells is `present` if at least one cell is available. Other +mixtures are `partial`; missing, delayed, and stale cells still prevent full +coverage. + +An all-`unsupported` constituent has zero applicable metrics. It remains +`unsupported`, with a reason and no watermark, and cannot satisfy full +coverage. A result containing only such constituents has coverage `none`, +never an observed zero. A partial-coverage request can retain that diagnostic +result; a full-coverage request returns `PARTIAL_RESULT` without publishing. Explicit watermarks are bounded by period end, source read cutoff, and the observation instant. A watermark before the period is rejected. An explicit diff --git a/src/adcp/reporting/conformance.py b/src/adcp/reporting/conformance.py index 86792f4cd..1b2e70545 100644 --- a/src/adcp/reporting/conformance.py +++ b/src/adcp/reporting/conformance.py @@ -221,11 +221,14 @@ def _validate_request_against_capabilities( "the selected offering does not support immutable corrections", ) - exact_metrics = {metric.name for metric in offering.metrics if metric.support == "exact"} + requestable_metrics = { + metric.name for metric in offering.metrics if metric.support in {"exact", "partial"} + } for metric in request.requested_metrics: - if metric not in exact_metrics: + if metric not in requestable_metrics: raise _fail( - "CAPABILITY_MISMATCH", f"requested metric {metric!r} is not exact in the offering" + "CAPABILITY_MISMATCH", + f"requested metric {metric!r} is not supported in the offering", ) exact_dimensions = { dimension.name for dimension in offering.dimensions if dimension.support == "exact" diff --git a/src/adcp/reporting/inline_source.py b/src/adcp/reporting/inline_source.py index 861711a79..b8d297bfd 100644 --- a/src/adcp/reporting/inline_source.py +++ b/src/adcp/reporting/inline_source.py @@ -15,9 +15,10 @@ For uneven metric support, return :class:`InlineFetchResult` with ``cell_availability={constituent_id: {metric_name: MetricEvidence(...)}}``. Omitted cells retain the constituent defaults. Explicit evidence controls each -cell independently and mixed cells make the constituent partial. A cell that -withdraws a metric its own constituent's rows report withdraws that metric's -control total, so no total ever sums a disclaimed value; a cell that staged no +cell independently. Unsupported cells are inapplicable: alongside available +cells they do not make the constituent partial. A cell that withdraws a metric +its own constituent's rows report withdraws that metric's control total, so no +total ever sums a disclaimed value; a cell that staged no rows reaches no sum and leaves its neighbours' subtotal intact. Money has no such choice -- see ``monetary_metrics`` on :class:`InlineReportingSource`. Metric semantics always come from the selected SDK offering. @@ -285,11 +286,13 @@ class InlineFetchResult: Keys are the frozen request's constituent IDs, which need not equal media buy IDs. Omitted cells retain the derived constituent status and watermark. Explicit cells override even missing/unsupported constituent defaults; - mixed availability promotes the constituent to ``partial``. The SDK still - applies the source cutoff, observation ceiling, and authoritative freshness - gate. A cell watermark may advance the batch watermark without advancing - other cells. This is an adapter surface, not the manifest's wire-format - ``metric_availability`` list. + unsupported cells do not demote otherwise available coverage. With no + applicable cells, the constituent stays ``unsupported`` and cannot satisfy + full coverage. Other mixed availability promotes it to ``partial``. The + SDK still applies the source cutoff, observation ceiling, and authoritative + freshness gate. A cell watermark may advance the batch watermark without + advancing other cells. This is an adapter surface, not the manifest's + wire-format ``metric_availability`` list. Unknown keys, duplicate mapping entries, invalid evidence, contradictory explicit zeros, and a withdrawn *monetary* cell whose own constituent's rows @@ -1139,9 +1142,13 @@ def _resolve_availability( reason = fallback_reason if overrides.get(constituent_id): cell_statuses = {cell.status for cell in constituent_cells} - if len(cell_statuses) == 1: + if cell_statuses == {"unsupported"}: + # Zero applicable metrics is no coverage evidence, not a + # vacuously complete measurement or an observed zero. + status = "unsupported" + elif len(cell_statuses) == 1: status = constituent_cells[0].status - elif cell_statuses <= _AVAILABLE: + elif cell_statuses <= _AVAILABLE | {"unsupported"}: status = "present" else: status = "partial" diff --git a/src/adcp/reporting/source.py b/src/adcp/reporting/source.py index eb0b89c70..cf81ebd66 100644 --- a/src/adcp/reporting/source.py +++ b/src/adcp/reporting/source.py @@ -235,10 +235,11 @@ _AVAILABLE_STATUSES = frozenset({"present", "explicit_zero"}) #: A constituent's roll-up status constrains what its metric cells may claim. -#: ``partial`` is the only genuinely mixed roll-up: its cells keep independent -#: statuses precisely so a consumer can select a compatible subset. +#: Unsupported cells are inapplicable and may accompany available cells in a +#: present constituent. Other unavailable cells still require a degraded +#: roll-up, retaining their independent evidence for consumers. _CONSTITUENT_ALLOWS_METRIC: Mapping[str, frozenset[str]] = { - "present": frozenset({"present", "explicit_zero"}), + "present": frozenset({"present", "explicit_zero", "unsupported"}), "explicit_zero": frozenset({"explicit_zero"}), "unsupported": frozenset({"unsupported"}), "delayed": frozenset({"unsupported", "delayed", "missing"}), @@ -1502,16 +1503,23 @@ def _validate_cells(self) -> None: "metric availability must carry one record per constituent-metric cell" ) by_constituent: dict[str, set[str]] = {} + available_constituents: set[str] = set() for cell in self.metric_availability: by_constituent.setdefault(cell.constituent_id, set()).add(cell.metric) + if cell.status in _AVAILABLE_STATUSES: + available_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 in statuses: + 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 in _AVAILABLE_STATUSES and constituent_id not in available_constituents: + raise ValueError( + "available constituent requires at least one present or explicit-zero metric" + ) for cell in self.metric_availability: constituent_status = statuses.get(cell.constituent_id) if constituent_status is None: diff --git a/tests/conformance/reporting/test_inline_cell_availability.py b/tests/conformance/reporting/test_inline_cell_availability.py index c078cd582..1683dd609 100644 --- a/tests/conformance/reporting/test_inline_cell_availability.py +++ b/tests/conformance/reporting/test_inline_cell_availability.py @@ -218,7 +218,9 @@ async def test_explicit_cells_override_constituent_defaults(fallback: str) -> No assert cells["impressions"].status == explicit.status assert cells["impressions"].reason == explicit.reason assert {cells[metric].status for metric in METRICS[1:]} == {fallback} - assert manifest.coverage.constituents[0].status == "partial" + assert manifest.coverage.constituents[0].status == ( + "present" if fallback in {"present", "unsupported"} else "partial" + ) async def test_explicit_cells_can_replace_all_missing_constituent_defaults() -> None: @@ -291,7 +293,58 @@ async def test_zero_row_mixed_unavailability_has_no_manufactured_measurements() } -async def test_full_request_with_mixed_cells_fails_without_staging() -> None: +@pytest.mark.parametrize("authoritative", [False, True], ids=["snapshot", "official"]) +@pytest.mark.parametrize("available", ["present", "explicit_zero"]) +async def test_unsupported_cells_do_not_demote_full_coverage( + authoritative: bool, available: str +) -> None: + request = _request(authoritative=authoritative) + request = request.model_copy( + update={"coverage": request.coverage.model_copy(update={"expected": "full"})} + ) + rows = [ + { + "media_buy_id": "media-buy-redacted", + **{metric: 10 if available == "present" else 0 for metric in METRICS[:-1]}, + } + ] + watermark = min(request.period.end, request.period.source_read_cutoff_at) + manifest = await _seal( + InlineFetchResult( + rows=rows, + cell_availability={ + CID: { + **{ + metric: MetricEvidence(status=available, data_through=watermark) + for metric in METRICS[:-1] + }, + "completed_views": MetricEvidence.unavailable("not_video_inventory"), + } + }, + ), + request, + ) + assert manifest.coverage.status == "full" + constituent = manifest.coverage.constituents[0] + assert (constituent.status, constituent.data_through, constituent.reason) == ( + "present", + watermark, + None, + ) + assert not manifest.explicit_zero # A normalized measurement row, not a batch-wide zero. + assert {total.name for total in manifest.control_totals} == set(METRICS[:-1]) + unsupported = manifest.metric_availability[-1] + assert (unsupported.status, unsupported.reason, unsupported.data_through) == ( + "unsupported", + "not_video_inventory", + None, + ) + + +@pytest.mark.parametrize("status", ["missing", "delayed"]) +async def test_full_request_with_missing_or_delayed_cells_fails_without_staging( + status: str, +) -> None: request = _request() request = request.model_copy( update={"coverage": request.coverage.model_copy(update={"expected": "full"})} @@ -299,9 +352,9 @@ async def test_full_request_with_mixed_cells_fails_without_staging() -> None: source = _source( lambda req: InlineFetchResult( rows=[ROW], - cell_availability=metric_unsupported_everywhere( - req, "completed_views", "not_video_inventory" - ), + cell_availability={ + CID: {"completed_views": MetricEvidence(status=status, reason="not_ready")} + }, ), staging=_NoStaging(), ) @@ -615,7 +668,8 @@ async def test_a_covered_zero_row_constituent_keeps_zeros_only_for_its_available ), request, ) - assert manifest.coverage.constituents[1].status == "partial" + assert manifest.coverage.status == "full" + assert manifest.coverage.constituents[1].status == "present" assert not manifest.explicit_zero for cell in manifest.metric_availability: if cell.constituent_id == SECOND_CID: diff --git a/tests/conformance/reporting/test_reporting_conditional_metrics.py b/tests/conformance/reporting/test_reporting_conditional_metrics.py new file mode 100644 index 000000000..a3dd72643 --- /dev/null +++ b/tests/conformance/reporting/test_reporting_conditional_metrics.py @@ -0,0 +1,296 @@ +"""Conditional metrics have the same meaning in Core and source conformance.""" + +from __future__ import annotations + +import asyncio +from collections.abc import AsyncIterator +from contextlib import AsyncExitStack +from dataclasses import replace +from typing import Any + +import pytest + +from adcp.reporting.conformance import ( + run_reporting_source_replay_conformance, + validate_reporting_source_failure, +) +from adcp.reporting.fixtures import SNAPSHOT_OFFERING_ID +from adcp.reporting.inline_source import InlineFetchResult, MetricEvidence +from adcp.reporting.ledger import ReportingConfiguration +from adcp.reporting.service import ReliableReportingService, ReportingAccountContext +from adcp.reporting.source import ( + MetricOfferingV1, + ReportingSourceCapabilitiesV1, + ReportingSourceSliceRequestV1, + reporting_source_capabilities_sha256_v1, +) +from adcp.reporting.testing import DeterministicReportingClock +from tests.test_reliable_reporting_service import _capability_offering + +from ._generation_support import NOW, isolated_reporting_pool +from ._reliable_support import capabilities, configuration + +METRICS = ("impressions", "clicks", "spend", "completed_views") + + +def _conditional_capabilities() -> ReportingSourceCapabilitiesV1: + base = capabilities() + draft = base.model_copy( + update={ + "offerings": [ + offering.model_copy( + update={ + "metrics": [ + MetricOfferingV1( + name=name, + support="exact" if name in {"impressions", "spend"} else "partial", + reason=( + None + if name in {"impressions", "spend"} + else "inventory_dependent" + ), + semantic_contract_id=f"fixture.{name}", + semantic_contract_version="1", + semantic_contract_sha256="a" * 64, + ) + for name in METRICS + ] + } + ) + for offering in base.offerings + ] + } + ) + return ReportingSourceCapabilitiesV1.model_validate( + { + **draft.model_dump(), + "capabilities_sha256": reporting_source_capabilities_sha256_v1(draft), + } + ) + + +class _ConditionalAdapter: + def __init__(self, display_status: str, *, all_unsupported: bool = False) -> None: + self.capabilities = _conditional_capabilities() + self.calls: list[ReportingSourceSliceRequestV1] = [] + self.display_status = display_status + self.all_unsupported = all_unsupported + + def fetch_slice(self, request: ReportingSourceSliceRequestV1) -> InlineFetchResult: + self.calls.append(request) + rows: list[dict[str, Any]] = [] + evidence: dict[str, dict[str, MetricEvidence]] = {} + for item in request.coverage.constituents: + if self.all_unsupported and item.media_buy_id == "display-buy": + evidence[item.constituent_id] = { + name: MetricEvidence.unavailable("not_applicable") for name in METRICS + } + continue + status = self.display_status if item.media_buy_id == "display-buy" else "present" + cell = { + "present": MetricEvidence.present(request.period.end), + "explicit_zero": MetricEvidence.explicit_zero(), + "unsupported": MetricEvidence.unavailable("not_video_inventory"), + "missing": MetricEvidence.missing("measurement_missing"), + "delayed": MetricEvidence.delayed("measurement_pending"), + }[status] + row: dict[str, Any] = { + "media_buy_id": item.media_buy_id, + "impressions": 10, + "clicks": 0, + "spend": "1.25", + } + if status in {"present", "explicit_zero"}: + row["completed_views"] = 0 if status == "explicit_zero" else 5 + rows.append(row) + evidence[item.constituent_id] = { + "clicks": MetricEvidence.explicit_zero(), + "completed_views": cell, + } + return InlineFetchResult(rows=rows, currency=request.currency, cell_availability=evidence) + + +@pytest.fixture(params=["memory", "postgres"]) +async def conditional_service( + request: pytest.FixtureRequest, +) -> AsyncIterator[ReliableReportingService]: + source_capabilities = _conditional_capabilities() + + def context(config: ReportingConfiguration) -> ReportingAccountContext: + return ReportingAccountContext( + account_id=config.account_id, + adapter="conditional", + currency="EUR", + source_scope=source_capabilities.source_scope, + snapshot_offering_id=SNAPSHOT_OFFERING_ID, + requested_metrics=METRICS, + publication_namespace="reporting-source:fixture", + capability_offering=_capability_offering("conditional"), + ) + + clock = DeterministicReportingClock(NOW) + async with AsyncExitStack() as stack: + if request.param == "postgres": + pool = await stack.enter_async_context(isolated_reporting_pool()) + service = ReliableReportingService.postgres( + pool=pool, account_context=context, clock=clock + ) + else: + service = ReliableReportingService.memory(account_context=context, clock=clock) + stack.push_async_callback(service.close) + yield service + + +@pytest.mark.parametrize("display_status", ["present", "explicit_zero", "unsupported"]) +async def test_full_core_publication_and_conformance_accept_conditional_metrics( + conditional_service: ReliableReportingService, display_status: str +) -> None: + service = conditional_service + adapter = _ConditionalAdapter(display_status) + registration = service.sources.register("conditional", adapter) + config = replace(configuration("eur"), media_buy_ids=("display-buy", "video-buy")) + await service.configure(config) + turn = await service.run_worker(now=NOW) + + assert not turn.configuration_errors + worker = turn.configurations[config.generation_key] + assert not worker.slices_failed + assert len(worker.revisions_committed) == 1 + (request,) = adapter.calls + assert request.coverage.expected == "full" + assert tuple(request.requested_metrics) == METRICS + assert registration.object_reader is not None + manifest = await run_reporting_source_replay_conformance( + executor=registration.executor, + request=request, + object_reader=registration.object_reader, + clock=lambda: NOW, + ) + # Conformance reads the same sealed execution without calling the adapter again. + assert len(adapter.calls) == 1 + assert manifest.coverage.status == "full" + assert {item.status for item in manifest.coverage.constituents} == {"present"} + owners = { + item.constituent_id: item.constituent.media_buy_id + for item in manifest.coverage.constituents + } + conditional_cells = { + owners[cell.constituent_id]: cell + for cell in manifest.metric_availability + if cell.metric == "completed_views" + } + assert conditional_cells["display-buy"].status == display_status + assert conditional_cells["video-buy"].status == "present" + if display_status == "unsupported": + assert conditional_cells["display-buy"].reason == "not_video_inventory" + assert conditional_cells["display-buy"].data_through is None + totals = {total.name: total.value for total in manifest.control_totals} + assert totals == { + "impressions": "20", + "clicks": "0", + "spend": "2.50", + **( + {} + if display_status == "unsupported" + else {"completed_views": "5" if display_status == "explicit_zero" else "10"} + ), + } + (revision,) = await service.store.list_revisions( + account_id="eur", reporting_obligation_id=request.identity.reporting_obligation_id + ) + assert revision.source_manifest_sha256 == manifest.content_fingerprint.removeprefix("sha256:") + assert dict(revision.control_totals) == totals + page = await service.store.read_revision_rows( + account_id="eur", reporting_revision_id=revision.reporting_revision_id, limit=10 + ) + assert page.total_count == 2 and not page.has_more + display_row = next(row for row in page.rows if row["media_buy_id"] == "display-buy") + assert ("completed_views" in display_row) == (display_status != "unsupported") + + +@pytest.mark.parametrize("display_status", ["missing", "delayed"]) +async def test_full_core_publication_still_refuses_missing_or_delayed_applicable_cells( + conditional_service: ReliableReportingService, display_status: str +) -> None: + service = conditional_service + adapter = _ConditionalAdapter(display_status) + registration = service.sources.register("conditional", adapter) + config = replace(configuration("eur"), media_buy_ids=("display-buy", "video-buy")) + await service.configure(config) + turn = await service.run_worker(now=NOW) + + assert not turn.configuration_errors + worker = turn.configurations[config.generation_key] + assert len(worker.slices_failed) == 1 and not worker.revisions_committed + (request,) = adapter.calls + error = validate_reporting_source_failure( + await registration.executor.execute(request, cancel=asyncio.Event()), "PARTIAL_RESULT" + ) + assert error.retry == "retryable" + assert ( + await service.store.list_revisions( + account_id="eur", reporting_obligation_id=request.identity.reporting_obligation_id + ) + == () + ) + + +@pytest.mark.parametrize("second_available", [False, True]) +async def test_full_core_publication_does_not_turn_zero_applicable_metrics_into_coverage( + conditional_service: ReliableReportingService, + second_available: bool, +) -> None: + service = conditional_service + adapter = _ConditionalAdapter("unsupported", all_unsupported=True) + registration = service.sources.register("conditional", adapter) + config = replace( + configuration("eur"), + media_buy_ids=("display-buy", "video-buy") if second_available else ("display-buy",), + ) + await service.configure(config) + turn = await service.run_worker(now=NOW) + + assert not turn.configuration_errors + worker = turn.configurations[config.generation_key] + assert len(worker.slices_failed) == 1 and not worker.revisions_committed + (request,) = adapter.calls + error = validate_reporting_source_failure( + await registration.executor.execute(request, cancel=asyncio.Event()), "PARTIAL_RESULT" + ) + assert error.retry == "retryable" + assert registration.object_reader is not None + partial = request.model_copy( + update={ + "coverage": request.coverage.model_copy(update={"expected": "partial"}), + "identity": request.identity.model_copy( + update={"source_execution_key": request.identity.source_execution_key + "-partial"} + ), + } + ) + manifest = await run_reporting_source_replay_conformance( + executor=registration.executor, + request=partial, + object_reader=registration.object_reader, + clock=lambda: NOW, + ) + assert manifest.coverage.status == ("partial" if second_available else "none") + display = next( + item + for item in manifest.coverage.constituents + if item.constituent.media_buy_id == "display-buy" + ) + assert display.status == "unsupported" + assert display.data_through is None and display.reason == "not_applicable" + assert { + cell.status + for cell in manifest.metric_availability + if cell.constituent_id == display.constituent_id + } == {"unsupported"} + assert not manifest.explicit_zero + assert bool(manifest.control_totals) is second_available + assert ( + await service.store.list_revisions( + account_id="eur", reporting_obligation_id=request.identity.reporting_obligation_id + ) + == () + ) diff --git a/tests/conformance/reporting/test_reporting_evidence_currency_integration.py b/tests/conformance/reporting/test_reporting_evidence_currency_integration.py index 4942a69fd..bc2795e1a 100644 --- a/tests/conformance/reporting/test_reporting_evidence_currency_integration.py +++ b/tests/conformance/reporting/test_reporting_evidence_currency_integration.py @@ -161,7 +161,7 @@ async def test_colliding_accounts_keep_sparse_metric_evidence_and_staging_isolat clock=h.clock, ) assert manifest.currency == obligation.currency == h.currencies[account] - assert manifest.row_count == 1 and manifest.coverage.status == "partial" + assert manifest.row_count == 1 and manifest.coverage.status == "full" cells = {cell.metric: cell for cell in manifest.metric_availability} assert {name: cell.status for name, cell in cells.items()} == { "impressions": "present", @@ -307,7 +307,7 @@ async def test_one_unavailable_spend_cell_suppresses_only_its_metric_total( ("impressions", "10"), ("clicks", "0"), ] - assert [item.status for item in manifest.coverage.constituents] == ["present", "partial"] + assert [item.status for item in manifest.coverage.constituents] == ["present", "present"] def _with_second_constituent( diff --git a/tests/test_reporting_source_contract.py b/tests/test_reporting_source_contract.py index 047c0ef39..fb92dc986 100644 --- a/tests/test_reporting_source_contract.py +++ b/tests/test_reporting_source_contract.py @@ -42,6 +42,7 @@ encode_source_batch_manifest_v1, parse_verified_source_batch_manifest_v1, publication_content_fingerprint_v1, + reporting_source_capabilities_sha256_v1, source_batch_manifest_reference_v1, ) @@ -145,6 +146,33 @@ async def test_metric_outside_the_offering_is_refused() -> None: assert error.value.code == "CAPABILITY_MISMATCH" +@pytest.mark.parametrize("support", ["partial", "unavailable"]) +async def test_conditional_metric_support_is_requestable_but_unavailable_is_not( + support: str, +) -> None: + payload = redacted_capabilities().model_dump(mode="json") + for offering in payload["offerings"]: + for metric in offering["metrics"]: + metric.update(support=support, reason="inventory_dependent") + payload["capabilities_sha256"] = reporting_source_capabilities_sha256_v1(payload) + capabilities = ReportingSourceCapabilitiesV1.model_validate(payload) + request = redacted_snapshot_request() + result, manifest, reader = redacted_completed_result(request) + if support == "unavailable": + with pytest.raises(ReportingSourceConformanceError) as error: + await validate_reporting_source_execution( + capabilities=capabilities, request=request, result=result, object_reader=reader + ) + assert error.value.code == "CAPABILITY_MISMATCH" + else: + assert ( + await validate_reporting_source_execution( + capabilities=capabilities, request=request, result=result, object_reader=reader + ) + == manifest + ) + + async def test_slice_wider_than_the_offering_window_is_refused() -> None: base = redacted_snapshot_request() wide = base.model_copy( @@ -391,6 +419,33 @@ async def test_metric_status_cannot_contradict_its_constituent() -> None: _rebind(manifest, metric_availability=contradictory) +@pytest.mark.parametrize("status", ["missing", "delayed", "stale", "partial"]) +def test_full_constituent_cannot_hide_unavailable_applicable_cells(status: str) -> None: + _, manifest, _ = redacted_completed_result(redacted_snapshot_request()) + cells = [ + manifest.metric_availability[0].model_copy( + update={"status": status, "reason": "not_ready", "data_through": None} + ), + *manifest.metric_availability[1:], + ] + with pytest.raises(ValidationError, match="contradicts constituent status"): + _rebind(manifest, metric_availability=cells) + + +def test_full_constituent_needs_at_least_one_applicable_metric() -> None: + _, manifest, _ = redacted_completed_result(redacted_snapshot_request()) + cells = [ + cell.model_copy( + update={"status": "unsupported", "reason": "not_applicable", "data_through": None} + ) + for cell in manifest.metric_availability + ] + # Without a zero-denominator guard, allowing unsupported cells under a + # present constituent would let an entirely unsupported result claim full. + with pytest.raises(ValidationError): + _rebind(manifest, metric_availability=cells) + + # -- errors ----------------------------------------------------------------- From 76cfc1f0f3d49bf9e73142f1a593eb9e0528e156 Mon Sep 17 00:00:00 2001 From: Brian O'Kelley Date: Tue, 29 Sep 2026 15:26:52 +0000 Subject: [PATCH 2/3] fix(reporting): bind conditional evidence to declared support --- docs/reporting-source-adapters.md | 21 +- src/adcp/reporting/conformance.py | 7 + src/adcp/reporting/inline_source.py | 2 + src/adcp/reporting/ledger/producer.py | 10 + src/adcp/reporting/source.py | 23 ++- .../test_inline_cell_availability.py | 56 +++++ .../test_reporting_conditional_metrics.py | 191 +++++++++++++----- ...reporting_evidence_currency_integration.py | 18 +- tests/test_reporting_inline_source.py | 18 +- tests/test_reporting_source_contract.py | 101 ++++++++- ...reporting_evidence_currency_integration.py | 2 +- tests/type_checks/reporting_inline_source.py | 3 +- 12 files changed, 376 insertions(+), 76 deletions(-) diff --git a/docs/reporting-source-adapters.md b/docs/reporting-source-adapters.md index a09393714..f8ecde9ee 100644 --- a/docs/reporting-source-adapters.md +++ b/docs/reporting-source-adapters.md @@ -37,6 +37,19 @@ adapter builds that evidence from the defaults and `cell_availability` overrides; use an override for each inapplicable or unready cell. An offering's `unavailable` metrics cannot pass source conformance. +Only a metric declared with `support: partial` and a stable offering reason may +emit `unsupported` cells. Declaring `exact` support promises applicability; +use `missing` or `delayed` when that measurement has not arrived. The inline +adapter rejects contradictory evidence before staging, and both producer +admission and source conformance enforce the same rule for custom executors. +For example, an exact-support spend metric awaiting billing must use +`MetricEvidence.delayed("billing_pending")` and retain partial coverage. + +If a conditional metric is `unsupported` in every constituent, the batch can +still have full coverage when every constituent has other applicable, complete +measurements. Every cell of the inapplicable metric retains its `unsupported` +evidence, including its reason and lack of watermark. + | Constructor | Meaning and required evidence | | --- | --- | | `MetricEvidence.present(data_through)` | Measured values through a timezone-aware watermark. Every matched row for the constituent must carry a non-null value for this metric. | @@ -71,9 +84,11 @@ coverage. An all-`unsupported` constituent has zero applicable metrics. It remains `unsupported`, with a reason and no watermark, and cannot satisfy full -coverage. A result containing only such constituents has coverage `none`, -never an observed zero. A partial-coverage request can retain that diagnostic -result; a full-coverage request returns `PARTIAL_RESULT` without publishing. +coverage. A manifest cannot label that constituent `partial` either. A result +containing only such constituents has coverage `none`, never an observed zero. +A partial-coverage request can retain that diagnostic result; a full-coverage +request returns `PARTIAL_RESULT` without publishing. Conformance rejects a +completed `none` result for a full-coverage request. Explicit watermarks are bounded by period end, source read cutoff, and the observation instant. A watermark before the period is rejected. An explicit diff --git a/src/adcp/reporting/conformance.py b/src/adcp/reporting/conformance.py index 1b2e70545..cb147c05d 100644 --- a/src/adcp/reporting/conformance.py +++ b/src/adcp/reporting/conformance.py @@ -49,6 +49,7 @@ ReportingSourceSliceRequestV1, ReportingSourceStagedObjectReader, SourceBatchManifestV1, + _validate_metric_applicability, deterministic_source_publication_id_v1, iso_duration_milliseconds_v1, parse_verified_source_batch_manifest_v1, @@ -394,6 +395,12 @@ def _validate_manifest_against_request( requested_metrics = set(request.requested_metrics) requested_constituents = {item.constituent_id for item in request.coverage.constituents} declared = {metric.name: metric for metric in offering.metrics} + try: + _validate_metric_applicability(offering.metrics, manifest.metric_availability) + except ValueError: + raise _fail( + "MANIFEST_MISMATCH", "unsupported cells require partial metric support with a reason" + ) from None published_cells = { (cell.constituent_id, cell.metric): cell for cell in manifest.metric_availability } diff --git a/src/adcp/reporting/inline_source.py b/src/adcp/reporting/inline_source.py index b8d297bfd..9abeeb3ed 100644 --- a/src/adcp/reporting/inline_source.py +++ b/src/adcp/reporting/inline_source.py @@ -101,6 +101,7 @@ SourceBatchManifestV1, SourceBatchObjectV1, SourceControlTotalV1, + _validate_metric_applicability, deterministic_source_publication_id_v1, encode_source_batch_manifest_v1, publication_content_fingerprint_v1, @@ -1172,6 +1173,7 @@ def _resolve_availability( constituent=item, status=status, data_through=watermark, reason=reason ) ) + _validate_metric_applicability(offering.metrics, cells) return constituents, cells def _seal_manifest( diff --git a/src/adcp/reporting/ledger/producer.py b/src/adcp/reporting/ledger/producer.py index 86c7dee23..75f978332 100644 --- a/src/adcp/reporting/ledger/producer.py +++ b/src/adcp/reporting/ledger/producer.py @@ -87,6 +87,7 @@ ReportingSourceSliceRequestV1, ReportingSourceStagedObjectReader, SourceBatchManifestV1, + _validate_metric_applicability, coverage_denominator_fingerprint_v1, iso_duration_milliseconds_v1, parse_verified_source_batch_manifest_v1, @@ -1213,6 +1214,15 @@ async def commit_revision_from_manifest( """ obligation = await self._stored_obligation(obligation) self._validate_manifest_currency(obligation, manifest) + if any(cell.status == "unsupported" for cell in manifest.metric_availability): + try: + offering = self._source.capabilities.offering(manifest.offering_id) + _validate_metric_applicability(offering.metrics, manifest.metric_availability) + except (KeyError, ValueError): + raise LedgerConflictError( + "MANIFEST_MISMATCH", + "unsupported cells require partial metric support with a reason", + ) from None now = now or self._clock() turn = turn or WorkerTurn() existing = await self._store.list_revisions( diff --git a/src/adcp/reporting/source.py b/src/adcp/reporting/source.py index cf81ebd66..9cfa17956 100644 --- a/src/adcp/reporting/source.py +++ b/src/adcp/reporting/source.py @@ -1057,6 +1057,17 @@ def _evidence_matches_status(self) -> ReportingMetricAvailabilityV1: return self +def _validate_metric_applicability( + metrics: Sequence[MetricOfferingV1], cells: Sequence[ReportingMetricAvailabilityV1] +) -> None: + """Only declared conditional support permits an inapplicable cell.""" + conditional = { + metric.name for metric in metrics if metric.support == "partial" and metric.reason + } + if any(cell.status == "unsupported" and cell.metric not in conditional for cell in cells): + raise ValueError("unsupported cells require partial metric support with a reason") + + class SourceBatchCoverageV1(_Frozen): """The slice's coverage roll-up plus its per-constituent evidence.""" @@ -1503,11 +1514,11 @@ def _validate_cells(self) -> None: "metric availability must carry one record per constituent-metric cell" ) by_constituent: dict[str, set[str]] = {} - available_constituents: set[str] = set() + applicable_constituents: set[str] = set() for cell in self.metric_availability: by_constituent.setdefault(cell.constituent_id, set()).add(cell.metric) - if cell.status in _AVAILABLE_STATUSES: - available_constituents.add(cell.constituent_id) + 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(): @@ -1516,10 +1527,8 @@ def _validate_cells(self) -> None: f"constituent {constituent_id!r} needs one availability record for every " "published metric" ) - if status in _AVAILABLE_STATUSES and constituent_id not in available_constituents: - raise ValueError( - "available constituent requires at least one present or explicit-zero 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: diff --git a/tests/conformance/reporting/test_inline_cell_availability.py b/tests/conformance/reporting/test_inline_cell_availability.py index 1683dd609..d5cb5786c 100644 --- a/tests/conformance/reporting/test_inline_cell_availability.py +++ b/tests/conformance/reporting/test_inline_cell_availability.py @@ -60,6 +60,8 @@ def _capabilities() -> ReportingSourceCapabilitiesV1: "metrics": [ MetricOfferingV1( name=metric, + support="partial", + reason="inventory_dependent", semantic_contract_id=f"fixture.{offering.offering_id}.{metric}", semantic_contract_version=str(index + 1), semantic_contract_sha256=str(index + 1) * 64, @@ -369,6 +371,51 @@ async def stage(self, **kwargs: Any) -> tuple[str, str]: pytest.fail("invalid evidence must be rejected before staging") +@pytest.mark.parametrize("explicit", [False, True], ids=["constituent-default", "cell-override"]) +@pytest.mark.parametrize("expected", ["full", "partial"]) +async def test_exact_support_rejects_unsupported_before_staging_or_sealing( + explicit: bool, expected: str +) -> None: + payload = _capabilities().model_dump(mode="json") + for offering in payload["offerings"]: + for metric in offering["metrics"]: + metric.update(support="exact", reason=None) + payload["capabilities_sha256"] = reporting_source_capabilities_sha256_v1(payload) + request = _request() + request = request.model_copy( + update={"coverage": request.coverage.model_copy(update={"expected": expected})} + ) + answer = ( + InlineFetchResult( + rows=[ROW], + cell_availability={ + CID: {"completed_views": MetricEvidence.unavailable("not_video_inventory")} + }, + ) + if explicit + else InlineFetchResult(rows=[], unavailable_constituents={CID: "not_applicable"}) + ) + seals = InMemorySealStore() + source = InlineReportingSource( + capabilities=ReportingSourceCapabilitiesV1.model_validate(payload), + fetch=lambda req: answer, + staging=_NoStaging(), + seals=seals, + clock=lambda: OBSERVED_AT, + ) + with pytest.raises( + ValueError, match="unsupported cells require partial metric support with a reason" + ): + await source.execute(request, cancel=asyncio.Event()) + assert ( + await seals.get( + account_id=request.identity.account_id, + source_execution_key=request.identity.source_execution_key, + ) + is None + ) + + class _RepeatedMapping(Mapping[str, Any]): """A mapping that exposes repeated entries before a dict would discard them.""" @@ -741,6 +788,10 @@ async def test_declared_zero_without_a_row_value_does_not_manufacture_a_total() @pytest.mark.parametrize("pattern", ["unsupported", "delayed"]) async def test_bulk_helpers_cover_only_the_requested_metric(pattern: str) -> None: request = _request(second=True) + if pattern == "unsupported": + request = request.model_copy( + update={"coverage": request.coverage.model_copy(update={"expected": "full"})} + ) through = request.period.source_read_cutoff_at - timedelta(hours=4) overrides = ( metric_unsupported_everywhere(request, "completed_views", "not_video_inventory") @@ -753,6 +804,11 @@ async def test_bulk_helpers_cover_only_the_requested_metric(pattern: str) -> Non ), request, ) + assert manifest.coverage.status == ("full" if pattern == "unsupported" else "partial") + assert {item.status for item in manifest.coverage.constituents} == { + "present" if pattern == "unsupported" else "partial" + } + assert "completed_views" not in {total.name for total in manifest.control_totals} for cell in manifest.metric_availability: if cell.metric == "completed_views": assert cell.status == pattern diff --git a/tests/conformance/reporting/test_reporting_conditional_metrics.py b/tests/conformance/reporting/test_reporting_conditional_metrics.py index a3dd72643..a91457be0 100644 --- a/tests/conformance/reporting/test_reporting_conditional_metrics.py +++ b/tests/conformance/reporting/test_reporting_conditional_metrics.py @@ -11,12 +11,14 @@ import pytest from adcp.reporting.conformance import ( + ReportingSourceConformanceError, run_reporting_source_replay_conformance, + validate_reporting_source_execution, validate_reporting_source_failure, ) from adcp.reporting.fixtures import SNAPSHOT_OFFERING_ID from adcp.reporting.inline_source import InlineFetchResult, MetricEvidence -from adcp.reporting.ledger import ReportingConfiguration +from adcp.reporting.ledger import LedgerConflictError, ReportingConfiguration from adcp.reporting.service import ReliableReportingService, ReportingAccountContext from adcp.reporting.source import ( MetricOfferingV1, @@ -33,7 +35,9 @@ METRICS = ("impressions", "clicks", "spend", "completed_views") -def _conditional_capabilities() -> ReportingSourceCapabilitiesV1: +def _conditional_capabilities( + *, exact_metrics: tuple[str, ...] = ("impressions", "spend") +) -> ReportingSourceCapabilitiesV1: base = capabilities() draft = base.model_copy( update={ @@ -43,12 +47,8 @@ def _conditional_capabilities() -> ReportingSourceCapabilitiesV1: "metrics": [ MetricOfferingV1( name=name, - support="exact" if name in {"impressions", "spend"} else "partial", - reason=( - None - if name in {"impressions", "spend"} - else "inventory_dependent" - ), + support="exact" if name in exact_metrics else "partial", + reason=None if name in exact_metrics else "inventory_dependent", semantic_contract_id=f"fixture.{name}", semantic_contract_version="1", semantic_contract_sha256="a" * 64, @@ -70,23 +70,33 @@ def _conditional_capabilities() -> ReportingSourceCapabilitiesV1: class _ConditionalAdapter: - def __init__(self, display_status: str, *, all_unsupported: bool = False) -> None: - self.capabilities = _conditional_capabilities() + def __init__( + self, + display_status: str, + *, + video_status: str = "present", + unsupported_buys: tuple[str, ...] = (), + exact_metrics: tuple[str, ...] = ("impressions", "spend"), + ) -> None: + self.capabilities = _conditional_capabilities(exact_metrics=exact_metrics) self.calls: list[ReportingSourceSliceRequestV1] = [] self.display_status = display_status - self.all_unsupported = all_unsupported + self.video_status = video_status + self.unsupported_buys = unsupported_buys def fetch_slice(self, request: ReportingSourceSliceRequestV1) -> InlineFetchResult: self.calls.append(request) rows: list[dict[str, Any]] = [] evidence: dict[str, dict[str, MetricEvidence]] = {} for item in request.coverage.constituents: - if self.all_unsupported and item.media_buy_id == "display-buy": + if item.media_buy_id in self.unsupported_buys: evidence[item.constituent_id] = { name: MetricEvidence.unavailable("not_applicable") for name in METRICS } continue - status = self.display_status if item.media_buy_id == "display-buy" else "present" + status = ( + self.display_status if item.media_buy_id == "display-buy" else self.video_status + ) cell = { "present": MetricEvidence.present(request.period.end), "explicit_zero": MetricEvidence.explicit_zero(), @@ -141,12 +151,23 @@ def context(config: ReportingConfiguration) -> ReportingAccountContext: yield service -@pytest.mark.parametrize("display_status", ["present", "explicit_zero", "unsupported"]) +@pytest.mark.parametrize( + ("display_status", "video_status", "completed_views_total"), + [ + pytest.param("present", "present", "10", id="measured-everywhere"), + pytest.param("explicit_zero", "present", "5", id="measured-zero"), + pytest.param("unsupported", "present", None, id="conditional-inventory"), + pytest.param("unsupported", "unsupported", None, id="unsupported-everywhere"), + ], +) async def test_full_core_publication_and_conformance_accept_conditional_metrics( - conditional_service: ReliableReportingService, display_status: str + conditional_service: ReliableReportingService, + display_status: str, + video_status: str, + completed_views_total: str | None, ) -> None: service = conditional_service - adapter = _ConditionalAdapter(display_status) + adapter = _ConditionalAdapter(display_status, video_status=video_status) registration = service.sources.register("conditional", adapter) config = replace(configuration("eur"), media_buy_ids=("display-buy", "video-buy")) await service.configure(config) @@ -180,20 +201,17 @@ async def test_full_core_publication_and_conformance_accept_conditional_metrics( if cell.metric == "completed_views" } assert conditional_cells["display-buy"].status == display_status - assert conditional_cells["video-buy"].status == "present" - if display_status == "unsupported": - assert conditional_cells["display-buy"].reason == "not_video_inventory" - assert conditional_cells["display-buy"].data_through is None + assert conditional_cells["video-buy"].status == video_status + for cell in conditional_cells.values(): + if cell.status == "unsupported": + assert cell.reason == "not_video_inventory" + assert cell.data_through is None totals = {total.name: total.value for total in manifest.control_totals} assert totals == { "impressions": "20", "clicks": "0", "spend": "2.50", - **( - {} - if display_status == "unsupported" - else {"completed_views": "5" if display_status == "explicit_zero" else "10"} - ), + **({} if completed_views_total is None else {"completed_views": completed_views_total}), } (revision,) = await service.store.list_revisions( account_id="eur", reporting_obligation_id=request.identity.reporting_obligation_id @@ -235,18 +253,91 @@ async def test_full_core_publication_still_refuses_missing_or_delayed_applicable ) -@pytest.mark.parametrize("second_available", [False, True]) +@pytest.mark.parametrize("custom_executor", [False, True], ids=["inline", "custom-executor"]) +async def test_exact_support_cannot_withdraw_a_cell_as_unsupported( + conditional_service: ReliableReportingService, + custom_executor: bool, +) -> None: + service = conditional_service + if custom_executor: + adapter = _ConditionalAdapter("unsupported") + inner = service.sources.register("partial", adapter) + + class ExactExecutor: + capabilities = _conditional_capabilities(exact_metrics=METRICS) + + async def execute(self, request, *, cancel, heartbeat=None): + return await inner.executor.execute(request, cancel=cancel, heartbeat=heartbeat) + + registration = service.sources.register_executor( + "conditional", ExactExecutor(), object_reader=inner.object_reader + ) + else: + adapter = _ConditionalAdapter("unsupported", exact_metrics=METRICS) + registration = service.sources.register("conditional", adapter) + config = replace(configuration("eur"), media_buy_ids=("display-buy", "video-buy")) + await service.configure(config) + turn = await service.run_worker(now=NOW) + + assert not turn.configurations + assert set(turn.configuration_errors) == {config.generation_key} + error = turn.configuration_errors[config.generation_key] + assert isinstance(error, LedgerConflictError if custom_executor else ValueError) + assert "unsupported cells require partial metric support with a reason" in str(error) + (request,) = adapter.calls + assert ( + await service.store.list_revisions( + account_id="eur", reporting_obligation_id=request.identity.reporting_obligation_id + ) + == () + ) + if custom_executor: + assert isinstance(error, LedgerConflictError) and error.code == "MANIFEST_MISMATCH" + assert registration.object_reader is not None + result = await registration.executor.execute(request, cancel=asyncio.Event()) + with pytest.raises( + ReportingSourceConformanceError, match="partial metric support" + ) as caught: + await validate_reporting_source_execution( + capabilities=registration.executor.capabilities, + request=request, + result=result, + object_reader=registration.object_reader, + clock=lambda: NOW, + ) + assert caught.value.code == "MANIFEST_MISMATCH" + + +@pytest.mark.parametrize( + ("media_buy_ids", "unsupported_buys", "expected_coverage"), + [ + pytest.param(("display-buy",), ("display-buy",), "none", id="one-unsupported"), + pytest.param( + ("display-buy", "video-buy"), + ("display-buy",), + "partial", + id="unsupported-and-available", + ), + pytest.param( + ("display-buy", "video-buy"), + ("display-buy", "video-buy"), + "none", + id="all-constituents-all-metrics-unsupported", + ), + ], +) async def test_full_core_publication_does_not_turn_zero_applicable_metrics_into_coverage( conditional_service: ReliableReportingService, - second_available: bool, + media_buy_ids: tuple[str, ...], + unsupported_buys: tuple[str, ...], + expected_coverage: str, ) -> None: service = conditional_service - adapter = _ConditionalAdapter("unsupported", all_unsupported=True) - registration = service.sources.register("conditional", adapter) - config = replace( - configuration("eur"), - media_buy_ids=("display-buy", "video-buy") if second_available else ("display-buy",), + adapter = _ConditionalAdapter( + "unsupported", unsupported_buys=unsupported_buys, exact_metrics=() ) + registration = service.sources.register("conditional", adapter) + config = replace(configuration("eur"), media_buy_ids=media_buy_ids) await service.configure(config) turn = await service.run_worker(now=NOW) @@ -273,21 +364,29 @@ async def test_full_core_publication_does_not_turn_zero_applicable_metrics_into_ object_reader=registration.object_reader, clock=lambda: NOW, ) - assert manifest.coverage.status == ("partial" if second_available else "none") - display = next( - item - for item in manifest.coverage.constituents - if item.constituent.media_buy_id == "display-buy" - ) - assert display.status == "unsupported" - assert display.data_through is None and display.reason == "not_applicable" - assert { - cell.status - for cell in manifest.metric_availability - if cell.constituent_id == display.constituent_id - } == {"unsupported"} + assert manifest.coverage.status == expected_coverage + for item in manifest.coverage.constituents: + if item.constituent.media_buy_id in unsupported_buys: + assert item.status == "unsupported" + assert item.data_through is None and item.reason == "not_applicable" + assert { + cell.status + for cell in manifest.metric_availability + if cell.constituent_id == item.constituent_id + } == {"unsupported"} assert not manifest.explicit_zero - assert bool(manifest.control_totals) is second_available + assert bool(manifest.control_totals) is (expected_coverage == "partial") + if expected_coverage == "none": + assert manifest.row_count == 0 + assert {cell.status for cell in manifest.metric_availability} == {"unsupported"} + with pytest.raises(ReportingSourceConformanceError, match="full-coverage request"): + await validate_reporting_source_execution( + capabilities=adapter.capabilities, + request=partial.model_copy(update={"coverage": request.coverage}), + result=await registration.executor.execute(partial, cancel=asyncio.Event()), + object_reader=registration.object_reader, + clock=lambda: NOW, + ) assert ( await service.store.list_revisions( account_id="eur", reporting_obligation_id=request.identity.reporting_obligation_id diff --git a/tests/conformance/reporting/test_reporting_evidence_currency_integration.py b/tests/conformance/reporting/test_reporting_evidence_currency_integration.py index bc2795e1a..d9e2686d3 100644 --- a/tests/conformance/reporting/test_reporting_evidence_currency_integration.py +++ b/tests/conformance/reporting/test_reporting_evidence_currency_integration.py @@ -90,7 +90,7 @@ def sparse_fetch(request: ReportingSourceSliceRequestV1) -> InlineFetchResult: item.constituent_id: { "impressions": MetricEvidence.present(request.period.end), "clicks": MetricEvidence.explicit_zero(data_through=request.period.end), - "spend": MetricEvidence.unavailable("billing_pending"), + "spend": MetricEvidence.delayed("billing_pending"), } for item in request.coverage.constituents }, @@ -161,12 +161,12 @@ async def test_colliding_accounts_keep_sparse_metric_evidence_and_staging_isolat clock=h.clock, ) assert manifest.currency == obligation.currency == h.currencies[account] - assert manifest.row_count == 1 and manifest.coverage.status == "full" + assert manifest.row_count == 1 and manifest.coverage.status == "partial" cells = {cell.metric: cell for cell in manifest.metric_availability} assert {name: cell.status for name, cell in cells.items()} == { "impressions": "present", "clicks": "explicit_zero", - "spend": "unsupported", + "spend": "delayed", } assert cells["spend"].reason == "billing_pending" and cells["spend"].data_through is None assert [total.model_dump(exclude_none=True) for total in manifest.control_totals] == [ @@ -302,12 +302,12 @@ async def test_one_unavailable_spend_cell_suppresses_only_its_metric_total( (cell.constituent_id, cell.status) for cell in manifest.metric_availability if cell.metric == "spend" - } == {(constituents[0].constituent_id, "explicit_zero"), ("second", "unsupported")} + } == {(constituents[0].constituent_id, "explicit_zero"), ("second", "delayed")} assert [(t.name, t.value) for t in manifest.control_totals] == [ ("impressions", "10"), ("clicks", "0"), ] - assert [item.status for item in manifest.coverage.constituents] == ["present", "present"] + assert [item.status for item in manifest.coverage.constituents] == ["present", "partial"] def _with_second_constituent( @@ -379,7 +379,7 @@ async def test_a_no_row_unavailable_spend_cell_still_commits_its_neighbour_subto "spend": ( MetricEvidence.present(request.period.end) if item is first - else MetricEvidence.unavailable("billing_pending") + else MetricEvidence.delayed("billing_pending") ), } for item in (first, second) @@ -392,7 +392,7 @@ async def test_a_no_row_unavailable_spend_cell_still_commits_its_neighbour_subto (cell.constituent_id, cell.status) for cell in manifest.metric_availability if cell.metric == "spend" - } == {(first.constituent_id, "present"), (second.constituent_id, "unsupported")} + } == {(first.constituent_id, "present"), (second.constituent_id, "delayed")} spend = [ (total.value, total.unit, total.value_type) for total in manifest.control_totals @@ -437,7 +437,7 @@ async def test_a_withdrawn_money_cell_its_own_rows_contradict_fails_before_stagi ], currency=request.currency, cell_availability={ - first.constituent_id: {"spend": MetricEvidence(status=status, reason="billing_pending")} + first.constituent_id: {"spend": MetricEvidence(status=status, reason="not_measured")} }, ) with pytest.raises(ValueError, match=f"{status} contradicts a monetary value"): @@ -512,7 +512,7 @@ def build(*, contradict: bool) -> InlineFetchResult: "clicks": ( available if item is first - else MetricEvidence.unavailable("measurement_pending") + else MetricEvidence.delayed("measurement_pending") ), } for item in (first, second) diff --git a/tests/test_reporting_inline_source.py b/tests/test_reporting_inline_source.py index 11287e7ef..c92719917 100644 --- a/tests/test_reporting_inline_source.py +++ b/tests/test_reporting_inline_source.py @@ -30,7 +30,12 @@ InMemorySealStore, InMemoryStagingStore, ) -from adcp.reporting.source import ReportingSourceError, ReportingSourceSliceRequestV1 +from adcp.reporting.source import ( + ReportingSourceCapabilitiesV1, + ReportingSourceError, + ReportingSourceSliceRequestV1, + reporting_source_capabilities_sha256_v1, +) ROW = { "media_buy_id": "media-buy-redacted", @@ -48,7 +53,8 @@ def _source(fetch: Any, **kwargs: Any) -> InlineReportingSource: kwargs.setdefault("clock", lambda: OBSERVED_AT) - return InlineReportingSource(capabilities=redacted_capabilities(), fetch=fetch, **kwargs) + kwargs.setdefault("capabilities", redacted_capabilities()) + return InlineReportingSource(fetch=fetch, **kwargs) def _partial(request: ReportingSourceSliceRequestV1) -> ReportingSourceSliceRequestV1: @@ -175,13 +181,19 @@ async def test_an_uncovered_constituent_is_missing_not_zero() -> None: async def test_a_known_unsupported_scope_is_reported_as_unsupported() -> None: request = _partial(redacted_snapshot_request()) + payload = redacted_capabilities().model_dump(mode="json") + for offering in payload["offerings"]: + for metric in offering["metrics"]: + metric.update(support="partial", reason="inventory_dependent") + payload["capabilities_sha256"] = reporting_source_capabilities_sha256_v1(payload) source = _source( lambda _request: InlineFetchResult( rows=[], unavailable_constituents={ "campaign-redacted-1": "Media buy lives on another ad server" }, - ) + ), + capabilities=ReportingSourceCapabilitiesV1.model_validate(payload), ) manifest = await validate_reporting_source_execution( capabilities=source.capabilities, diff --git a/tests/test_reporting_source_contract.py b/tests/test_reporting_source_contract.py index fb92dc986..57e6d886f 100644 --- a/tests/test_reporting_source_contract.py +++ b/tests/test_reporting_source_contract.py @@ -173,6 +173,59 @@ async def test_conditional_metric_support_is_requestable_but_unavailable_is_not( ) +@pytest.mark.parametrize("support", ["exact", "partial"]) +async def test_unsupported_cell_must_bind_declared_conditional_support(support: str) -> None: + payload = redacted_capabilities().model_dump(mode="json") + for offering in payload["offerings"]: + for metric in offering["metrics"]: + if metric["name"] == "spend": + metric.update( + support=support, reason="inventory_dependent" if support == "partial" else None + ) + payload["capabilities_sha256"] = reporting_source_capabilities_sha256_v1(payload) + capabilities = ReportingSourceCapabilitiesV1.model_validate(payload) + request = redacted_snapshot_request() + _, manifest, reader = redacted_completed_result( + request, page_bodies=['{"campaign_id":"campaign-redacted-1","impressions":10}\n'] + ) + manifest = _rebind( + manifest, + metric_availability=[ + ( + cell.model_copy( + update={ + "status": "unsupported", + "reason": "not_applicable", + "data_through": None, + } + ) + if cell.metric == "spend" + else cell + ) + for cell in manifest.metric_availability + ], + ) + manifest_bytes = encode_source_batch_manifest_v1(manifest) + result = ReportingSourceExecutorResult.completed( + request=request, + manifest=source_batch_manifest_reference_v1("conditional-manifest", manifest_bytes), + manifest_bytes=manifest_bytes, + ) + if support == "exact": + with pytest.raises( + ReportingSourceConformanceError, match="partial metric support" + ) as caught: + await validate_reporting_source_execution( + capabilities=capabilities, request=request, result=result, object_reader=reader + ) + assert caught.value.code == "MANIFEST_MISMATCH" + else: + validated = await validate_reporting_source_execution( + capabilities=capabilities, request=request, result=result, object_reader=reader + ) + assert validated.coverage.status == "full" + + async def test_slice_wider_than_the_offering_window_is_refused() -> None: base = redacted_snapshot_request() wide = base.model_copy( @@ -432,18 +485,54 @@ def test_full_constituent_cannot_hide_unavailable_applicable_cells(status: str) _rebind(manifest, metric_availability=cells) -def test_full_constituent_needs_at_least_one_applicable_metric() -> None: - _, manifest, _ = redacted_completed_result(redacted_snapshot_request()) +@pytest.mark.parametrize( + ("constituent_status", "coverage_status"), + [("present", "full"), ("explicit_zero", "full"), ("partial", "partial"), ("delayed", "none")], +) +async def test_zero_applicable_metrics_cannot_claim_available_or_partial_coverage( + constituent_status: str, coverage_status: str +) -> None: + request = redacted_snapshot_request() + request = request.model_copy( + update={"coverage": request.coverage.model_copy(update={"expected": "partial"})} + ) + _, manifest, reader = redacted_completed_result(request) cells = [ cell.model_copy( update={"status": "unsupported", "reason": "not_applicable", "data_through": None} ) for cell in manifest.metric_availability ] - # Without a zero-denominator guard, allowing unsupported cells under a - # present constituent would let an entirely unsupported result claim full. - with pytest.raises(ValidationError): - _rebind(manifest, metric_availability=cells) + coverage = manifest.coverage.model_copy( + update={ + "status": coverage_status, + "constituents": [ + item.model_copy(update={"status": constituent_status, "reason": "not_applicable"}) + for item in manifest.coverage.constituents + ], + } + ) + with pytest.raises(ValidationError, match="without applicable metrics must be unsupported"): + _rebind(manifest, metric_availability=cells, coverage=coverage) + # A custom executor can supply canonical, correctly hashed bytes without + # using the model constructor. Public conformance must reject those too. + forged = manifest.model_copy(update={"metric_availability": cells, "coverage": coverage}) + forged = forged.model_copy( + update={"content_fingerprint": publication_content_fingerprint_v1(forged)} + ) + manifest_bytes = encode_source_batch_manifest_v1(forged) + result = ReportingSourceExecutorResult.completed( + request=request, + manifest=source_batch_manifest_reference_v1("forged-coverage", manifest_bytes), + manifest_bytes=manifest_bytes, + ) + with pytest.raises(ValidationError, match="without applicable metrics must be unsupported"): + await validate_reporting_source_execution( + capabilities=redacted_capabilities(), + request=request, + result=result, + object_reader=reader, + ) # -- errors ----------------------------------------------------------------- diff --git a/tests/type_checks/reporting_evidence_currency_integration.py b/tests/type_checks/reporting_evidence_currency_integration.py index 40aed8c77..3869479ef 100644 --- a/tests/type_checks/reporting_evidence_currency_integration.py +++ b/tests/type_checks/reporting_evidence_currency_integration.py @@ -21,7 +21,7 @@ def fetch(request: ReportingSourceSliceRequestV1) -> InlineFetchResult: item.constituent_id: { "impressions": MetricEvidence.present(through), "clicks": MetricEvidence.explicit_zero(data_through=through), - "spend": MetricEvidence.unavailable("billing_pending"), + "spend": MetricEvidence.delayed("billing_pending"), } for item in request.coverage.constituents } diff --git a/tests/type_checks/reporting_inline_source.py b/tests/type_checks/reporting_inline_source.py index 57f2d8dc0..ca702b294 100644 --- a/tests/type_checks/reporting_inline_source.py +++ b/tests/type_checks/reporting_inline_source.py @@ -1,7 +1,8 @@ """Adopter patterns: typed synchronous/asynchronous fetches and metric evidence. These adapters assume an offering requesting impressions, clicks, viewability, -and completed_views. The SDK supplies each metric's semantic-contract identity. +and completed_views, with partial support and a reason for completed_views. +The SDK supplies each metric's semantic-contract identity. """ from collections.abc import Awaitable, Callable, Mapping, Sequence From 07a53d1fb266ef77f47ad33ca60dd5d0c00d70e5 Mon Sep 17 00:00:00 2001 From: Brian O'Kelley Date: Tue, 29 Sep 2026 15:45:32 +0000 Subject: [PATCH 3/3] fix(reporting): enforce full coverage on custom executors --- src/adcp/reporting/ledger/producer.py | 5 ++ .../test_reporting_conditional_metrics.py | 67 +++++++++++++++++++ 2 files changed, 72 insertions(+) diff --git a/src/adcp/reporting/ledger/producer.py b/src/adcp/reporting/ledger/producer.py index 75f978332..641beb6de 100644 --- a/src/adcp/reporting/ledger/producer.py +++ b/src/adcp/reporting/ledger/producer.py @@ -968,6 +968,11 @@ async def acquire_obligation( return None manifest = self._verified_manifest(result) + if request.coverage.expected == "full" and manifest.coverage.status != "full": + raise LedgerConflictError( + "MANIFEST_MISMATCH", + "a full-coverage request cannot complete with partial or missing coverage", + ) self._validate_manifest_currency(obligation, manifest) rows = await self._read_rows(request, manifest) # ``now`` freezes dispatch/lease/cutoff decisions, not publication. diff --git a/tests/conformance/reporting/test_reporting_conditional_metrics.py b/tests/conformance/reporting/test_reporting_conditional_metrics.py index a91457be0..2f561ce37 100644 --- a/tests/conformance/reporting/test_reporting_conditional_metrics.py +++ b/tests/conformance/reporting/test_reporting_conditional_metrics.py @@ -24,6 +24,7 @@ MetricOfferingV1, ReportingSourceCapabilitiesV1, ReportingSourceSliceRequestV1, + parse_verified_source_batch_manifest_v1, reporting_source_capabilities_sha256_v1, ) from adcp.reporting.testing import DeterministicReportingClock @@ -393,3 +394,69 @@ async def test_full_core_publication_does_not_turn_zero_applicable_metrics_into_ ) == () ) + + +@pytest.mark.parametrize("diagnostic_coverage", ["none", "partial"]) +async def test_custom_executor_cannot_complete_a_full_request_with_diagnostic_coverage( + conditional_service: ReliableReportingService, + diagnostic_coverage: str, +) -> None: + service = conditional_service + adapter = _ConditionalAdapter( + "unsupported", + unsupported_buys=( + ("display-buy", "video-buy") if diagnostic_coverage == "none" else ("display-buy",) + ), + exact_metrics=(), + ) + inner = service.sources.register("partial", adapter) + + class DiagnosticExecutor: + capabilities = adapter.capabilities + + async def execute(self, request, *, cancel, heartbeat=None): + assert request.coverage.expected == "full" + diagnostic_request = request.model_copy( + update={"coverage": request.coverage.model_copy(update={"expected": "partial"})} + ) + return await inner.executor.execute( + diagnostic_request, cancel=cancel, heartbeat=heartbeat + ) + + registration = service.sources.register_executor( + "conditional", DiagnosticExecutor(), object_reader=inner.object_reader + ) + config = replace(configuration("eur"), media_buy_ids=("display-buy", "video-buy")) + await service.configure(config) + turn = await service.run_worker(now=NOW) + + assert not turn.configurations + assert set(turn.configuration_errors) == {config.generation_key} + error = turn.configuration_errors[config.generation_key] + assert isinstance(error, LedgerConflictError) and error.code == "MANIFEST_MISMATCH" + (diagnostic_request,) = adapter.calls + assert ( + await service.store.list_revisions( + account_id="eur", + reporting_obligation_id=diagnostic_request.identity.reporting_obligation_id, + ) + == () + ) + request = diagnostic_request.model_copy( + update={"coverage": diagnostic_request.coverage.model_copy(update={"expected": "full"})} + ) + assert registration.object_reader is not None + result = await registration.executor.execute(request, cancel=asyncio.Event()) + assert result.ok and result.response is not None and result.manifest_bytes is not None + manifest = parse_verified_source_batch_manifest_v1( + result.response.manifest, result.manifest_bytes + ) + assert manifest.coverage.status == diagnostic_coverage + with pytest.raises(ReportingSourceConformanceError, match="full-coverage request"): + await validate_reporting_source_execution( + capabilities=registration.executor.capabilities, + request=request, + result=result, + object_reader=registration.object_reader, + clock=lambda: NOW, + )