diff --git a/src/adcp/reporting/ledger/producer.py b/src/adcp/reporting/ledger/producer.py index b35b4e45d..66f6d8fb8 100644 --- a/src/adcp/reporting/ledger/producer.py +++ b/src/adcp/reporting/ledger/producer.py @@ -579,6 +579,9 @@ async def acquire_obligation( snapshot obligation. Without it a satisfied obligation is left alone, because re-reading a settled period on every worker turn would burn upstream quota to republish bytes nobody asked for. + + ``now`` freezes dispatch and the source read cutoff. Revision creation + uses a fresh producer clock sample after the staged objects are read. """ if configuration.generation_key != obligation.generation_key: raise LedgerConflictError( @@ -664,12 +667,21 @@ async def acquire_obligation( manifest = self._verified_manifest(result) self._validate_manifest_currency(obligation, manifest) + rows = await self._read_rows(request, manifest) + # ``now`` freezes dispatch/lease/cutoff decisions, not publication. + # A conforming source can observe finality while acquisition is running. + published_at = self._clock() + if _utc(published_at) < _utc(now): + raise LedgerConflictError( + "PUBLICATION_TIME_INVALID", + "producer clock regressed during acquisition; correct the clock before retrying", + ) return await self.commit_revision_from_manifest( obligation, manifest, - rows=await self._read_rows(request, manifest), + rows=rows, finality=finality, - now=now, + now=published_at, turn=turn, ) @@ -734,7 +746,9 @@ async def commit_revision_from_manifest( A snapshot restatement supersedes the current snapshot leaf; there is no edit path. An official close is terminal, so a later source correction - must arrive as an adjustment instead. + must arrive as an adjustment instead. ``now`` is a trusted publication + instant, unlike the dispatch instant accepted by ``acquire_obligation``. + Replaying a publication retains its original creation time and parent. """ obligation = await self._stored_obligation(obligation) self._validate_manifest_currency(obligation, manifest) @@ -771,6 +785,26 @@ async def commit_revision_from_manifest( control_totals = tuple((total.name, total.value) for total in manifest.control_totals) revision_id = f"rpr_{manifest.publication_id[4:44]}" + prior = next((item for item in existing if item.reporting_revision_id == revision_id), None) + created_at = prior.created_at if prior is not None else now + if prior is not None: + # Still reconstruct and verify the supplied content below. Merely + # finding the ID must not bypass immutable-content validation. + supersedes = prior.supersedes_reporting_revision_id + if ( + _utc(manifest.acquired_at) > _utc(now) + or _utc(manifest.observed_at) > _utc(created_at) + or _utc(manifest.finality_evidence.observed_at) > _utc(created_at) + or ( + finality == "official" + and _utc(manifest.finality_evidence.observed_at) < _utc(obligation.period.end) + ) + ): + raise LedgerConflictError( + "PUBLICATION_TIME_INVALID", + "source observation or finality is outside the publication time bounds; " + "check source evidence and the producer clock before retrying", + ) revision = ReportingRevisionRecord( reporting_revision_id=revision_id, account_id=obligation.account_id, @@ -786,7 +820,7 @@ async def commit_revision_from_manifest( control_totals=control_totals, observed_at=manifest.observed_at, data_through=manifest.data_through, - created_at=now, + created_at=created_at, supersedes_reporting_revision_id=supersedes, finality_basis="source_final" if finality == "official" else None, finality_policy_id=( diff --git a/tests/conformance/reporting/_late_account_server.py b/tests/conformance/reporting/_late_account_server.py index 97bdd9482..8bcd484d3 100644 --- a/tests/conformance/reporting/_late_account_server.py +++ b/tests/conformance/reporting/_late_account_server.py @@ -33,6 +33,7 @@ ) from adcp.reporting.projection import PgReportingStatusProjection from adcp.reporting.receipts import ReportingReceiptError +from adcp.reporting.source import parse_verified_source_batch_manifest_v1 from adcp.server import serve from adcp.server.auth import BearerTokenAuth, Principal, auth_context_factory from adcp.types import ReportingDeliveryOffering @@ -43,6 +44,39 @@ CONSUMERS = {"usd": "urn:buyer:usd", "eur": "urn:buyer:eur"} +class LiveAccountSource(AccountSource): + """Record the public source boundary, without changing its returned evidence.""" + + async def execute(self, request, *, cancel, heartbeat=None): + result = await super().execute(request, cancel=cancel, heartbeat=heartbeat) + if result.ok: + manifest = parse_verified_source_batch_manifest_v1( + result.response.manifest, result.manifest_bytes + ) + with self.observation_log.open("a") as output: + output.write( + json.dumps( + { + "account": request.identity.account_id, + "period": { + "start": request.period.start.isoformat(), + "end": request.period.end.isoformat(), + }, + "source_read_cutoff_at": ( + request.period.source_read_cutoff_at.isoformat() + ), + "observed_at": manifest.observed_at.isoformat(), + "acquired_at": manifest.acquired_at.isoformat(), + "finalized_at": manifest.finality_evidence.observed_at.isoformat(), + "data_through": manifest.data_through.isoformat(), + "publication_id": manifest.publication_id, + } + ) + + "\n" + ) + return result + + class WireCapture: def __init__(self, app, path): self.app, self.path = app, Path(path) @@ -91,6 +125,7 @@ def main(): parser.add_argument("--root", type=Path, required=True) parser.add_argument("--schema", required=True) parser.add_argument("--notifications", action="store_true") + parser.add_argument("--live-source-observation", action="store_true") args = parser.parse_args() root = args.root pool = AsyncConnectionPool( @@ -121,18 +156,25 @@ def main(): ) registry = ReportingRevisionVerifierRegistry(tuple(verifiers.values())) sources, producers, offerings = {}, {}, {} - # This adopter exposes an immutable historical dataset observed before - # startup. Its source observation remains that timestamp on every fetch; - # production turns and PostgreSQL lease scheduling use their real clocks. + # The runtime fairness test keeps a historical source observation. The + # separate publication-time test selects a real UTC sample after fetch. + # Production turns and PostgreSQL lease scheduling always use real clocks. source_observed_at = datetime.now(timezone.utc) for currency, verifier in verifiers.items(): key = verifier.key - source = AccountSource( + source_class = LiveAccountSource if args.live_source_observation else AccountSource + source = source_class( key, root / ("source-" + currency), - clock=lambda: source_observed_at, + clock=( + (lambda: datetime.now(timezone.utc)) + if args.live_source_observation + else (lambda: source_observed_at) + ), official=True, ) + if args.live_source_observation: + source.observation_log = root / "source-observations.jsonl" producer = ReportingProducer( source=source, store=store, @@ -342,7 +384,10 @@ async def startup(): "templates": templates, "initial_configurations": initial[0], "pool_size": 1, - "source_observed_at": source_observed_at.isoformat(), + "source_observed_at": ( + None if args.live_source_observation else source_observed_at.isoformat() + ), + "live_source_observation": args.live_source_observation, "python": __import__("sys").version, "adcp_file": adcp.__file__, "pydantic": importlib.metadata.version("pydantic"), diff --git a/tests/conformance/reporting/_production_packaging.py b/tests/conformance/reporting/_production_packaging.py index d0316d13e..cfc1e639a 100644 --- a/tests/conformance/reporting/_production_packaging.py +++ b/tests/conformance/reporting/_production_packaging.py @@ -258,6 +258,7 @@ def installed_production(root, python, wheel, source, *, label, driver_absent): "tests/conformance/reporting/test_reporting_tier_projection.py", "tests/conformance/reporting/test_reporting_schedule_schema.py", "tests/conformance/reporting/test_reporting_materializer_progress.py", + "tests/conformance/reporting/test_reporting_publication_time.py", "tests/test_reporting_revision_ownership.py", "tests/test_reporting_capability_models.py", "tests/test_reporting_scope_models.py", diff --git a/tests/conformance/reporting/test_reporting_production_late_accounts.py b/tests/conformance/reporting/test_reporting_production_late_accounts.py index c68246c10..1b760157b 100644 --- a/tests/conformance/reporting/test_reporting_production_late_accounts.py +++ b/tests/conformance/reporting/test_reporting_production_late_accounts.py @@ -40,7 +40,7 @@ @asynccontextmanager -async def running_server(root, schema, index, *, notifications): +async def running_server(root, schema, index, *, notifications, live_source_observation=False): fixture_root = Path(__file__).resolve().parents[3] launcher = ( "import sys; sys.path.insert(0,sys.argv.pop(1)); " @@ -66,6 +66,8 @@ async def running_server(root, schema, index, *, notifications): ] if notifications: command.append("--notifications") + if live_source_observation: + command.append("--live-source-observation") with ( (root / f"server-{index}.stdout").open("xb") as output, (root / f"server-{index}.stderr").open("xb") as errors, diff --git a/tests/conformance/reporting/test_reporting_production_publication_time.py b/tests/conformance/reporting/test_reporting_production_publication_time.py new file mode 100644 index 000000000..c1642e2f9 --- /dev/null +++ b/tests/conformance/reporting/test_reporting_production_publication_time.py @@ -0,0 +1,121 @@ +"""Actual MCP/PG publication after real source observation, including restart.""" + +import json +from datetime import datetime + +import pytest + +from adcp.reporting import load_reporting_ledger +from adcp.types import GetReportingStatusRequest, SyncAccountsRequest + +from .test_reporting_materializer_progress import progress_pool +from .test_reporting_production_late_accounts import ( + client_for, + data, + delivered_and_reconciled, + running_server, +) + + +@pytest.mark.parametrize("account", ["usd", "eur"]) +async def test_real_source_observation_before_publication_reconciles_and_replays(account, tmp_path): + async with progress_pool(autocommit=True) as (pool, schema): + runs = [] + for index in range(2): + async with running_server( + tmp_path, + schema, + index, + notifications=False, + live_source_observation=True, + ) as (uri, ready): + assert ready["pool_size"] == 1 and ready["live_source_observation"] + assert ready["initial_configurations"] == (0 if index == 0 else 1) + async with client_for(uri, account) as client: + if index == 0: + admitted = data( + await client.sync_accounts( + SyncAccountsRequest.model_validate( + { + "idempotency_key": "publication-time-" + account, + "accounts": [ + { + "account": {"account_id": account}, + "reporting_delivery_configs": [ + ready["templates"][account] + ], + } + ], + } + ) + ) + ) + assert ( + admitted["accounts"][0]["reporting_delivery_configs"][0]["state"] + == "ready" + ) + result = await delivered_and_reconciled( + client, + tmp_path, + ready, + account, + replay=bool(index), + period=runs[0]["period"] if index else None, + ) + ledger = await load_reporting_ledger( + client, + GetReportingStatusRequest.model_validate( + { + "account": {"account_id": account}, + "view": "periods", + "period": result["period"], + } + ), + ) + assert len(ledger.revisions) == 1 + revision = ledger.revisions[0] + observations = [ + json.loads(line) + for line in (tmp_path / "source-observations.jsonl") + .read_text() + .splitlines() + ] + matches = [ + item + for item in observations + if item["account"] == account and item["period"] == result["period"] + ] + assert len(matches) == 1 + source = matches[0] + # These values come from the actual source's public request + # and result, and the revision comes from a typed HTTP read. + assert datetime.fromisoformat(source["source_read_cutoff_at"]) < ( + revision.observed_at + ) + assert revision.observed_at == datetime.fromisoformat(source["observed_at"]) + assert revision.finalized_at == datetime.fromisoformat(source["finalized_at"]) + assert revision.data_through == datetime.fromisoformat(source["data_through"]) + assert datetime.fromisoformat(source["acquired_at"]) <= revision.created_at + assert revision.finalized_at <= revision.created_at + result["source"] = source + result["revision_evidence"] = revision.model_dump( + mode="json", exclude_none=True + ) + runs.append(result) + async with pool.connection() as connection: + assert await ( + await connection.execute( + "SELECT count(*) FROM pg_stat_activity WHERE application_name=%s", (schema,) + ) + ).fetchone() == (0,) + for field in ( + "revision", + "digest", + "receipt_ids", + "materialization_id", + "source", + "revision_evidence", + ): + assert runs[0][field] == runs[1][field] + assert "-test-token" not in (tmp_path / "wire.jsonl").read_text() + print(json.dumps({"public_real_clock_publication": {"account": account, "runs": runs}})) diff --git a/tests/conformance/reporting/test_reporting_publication_time.py b/tests/conformance/reporting/test_reporting_publication_time.py new file mode 100644 index 000000000..59b8c75a4 --- /dev/null +++ b/tests/conformance/reporting/test_reporting_publication_time.py @@ -0,0 +1,409 @@ +"""Publication clocks are distinct from dispatch and immutable source evidence.""" + +from __future__ import annotations + +import asyncio +from dataclasses import dataclass +from datetime import datetime, timedelta, timezone + +import pytest +from pydantic import ValidationError + +from adcp.reporting import ( + ExpectedReportingPeriod, + ReportingLedger, + evaluate_reporting_ledger, +) +from adcp.reporting.conformance import validate_reporting_source_execution +from adcp.reporting.ledger import ( + InMemoryReportingLedgerStore, + LedgerConflictError, + PgReportingLedgerStore, + ProducerOfferings, + ReportingConfiguration, + ReportingProducer, + ReportingScheduleSpec, + ReportingStatusCaller, + ReportingStatusHandler, +) +from adcp.reporting.materializer import ReportingWriterCapability, reference_verifier +from adcp.reporting.source import SourceBatchManifestV1, parse_verified_source_batch_manifest_v1 +from adcp.types import GetReportingStatusResponse +from adcp.validation.schema_loader import get_named_validator + +from ._generation_support import END, START, isolated_reporting_pool +from ._production_support import Source + +TURN = END + timedelta(hours=1) +PUBLISHED = TURN + timedelta(seconds=2) +ROWS = [{"row_id": "1", "impressions": 5, "spend": "1.25", "currency": "USD"}] + + +@dataclass +class Clock: + now: datetime = TURN + + def __call__(self): + return self.now + + +class RecordedSource(Source): + async def execute(self, request, *, cancel, heartbeat=None): + await asyncio.sleep(0) # The source can complete after the dispatch instant. + result = await super().execute(request, cancel=cancel, heartbeat=heartbeat) + self.executions.append((request, result)) + return result + + +class DelayedReader: + def __init__(self, reader, complete): + self.reader, self.complete = reader, complete + self.calls = 0 + + async def read(self, **kwargs): + result = await self.reader.read(**kwargs) + await asyncio.sleep(0) + self.complete() + self.calls += 1 + return result + + +@pytest.fixture(params=["memory", "postgres"]) +async def store(request): + if request.param == "memory": + yield InMemoryReportingLedgerStore(clock=lambda: datetime.now(timezone.utc)) + else: + async with isolated_reporting_pool() as pool: + ledger = PgReportingLedgerStore(pool=pool) + await ledger.create_schema() + yield ledger + + +async def setup( + store, path, *, observed=TURN, completion=PUBLISHED, real=False, finality="official" +): + verifier = reference_verifier( + ReportingWriterCapability( + "warehouse_materialization", + "fixture-sql", + "jsonl", + "canonical_digest", + "destination", + "immutable_location", + "sha256", + "conditional_create", + ) + ) + key = verifier.key + clock = Clock() + source = RecordedSource( + key, + path, + ROWS, + official=finality == "official", + product_ids=(key.report_definition_id,), + clock=(lambda: datetime.now(timezone.utc)) if real else (lambda: observed), + ) + source.executions = [] + config = ReportingConfiguration( + "publication-clock", + 1, + "account-clock", + key.report_definition_id, + key.reporting_profile, + "analytics", + ReportingScheduleSpec("PT1H", "PT1H", period_anchor=START), + finality, + activated_at=START, + deactivated_at=END, + media_buy_ids=("mb-clock",), + definition=key.definition, + ) + await store.put_configuration(config) + + def complete_read(): + if not real: + clock.now = completion + + reader = DelayedReader(source.reader, complete_read) + + def producer(): + return ReportingProducer( + source=source, + store=store, + object_reader=reader, + offerings=ProducerOfferings( + official_offering_id=source.source_id if finality == "official" else None, + snapshot_offering_id=source.source_id if finality == "snapshot" else None, + publication_namespace=source.capabilities.offerings[0].publication_namespace, + source_scope=source.capabilities.source_scope, + ), + revision_verifier=verifier, + max_periods_per_turn=1, + **({} if real else {"clock": clock}), + ) + + return config, source, reader, clock, producer + + +async def public_outcome(store, config): + payload = await ReportingStatusHandler(store).handle( + { + "adcp_version": "3.2-rc.4", + "account": {"account_id": config.account_id}, + "view": "periods", + "period": {"start": START.isoformat(), "end": END.isoformat()}, + }, + caller=ReportingStatusCaller(account_id=config.account_id, consumer_id="clock-buyer"), + ) + response = GetReportingStatusResponse.model_validate(payload) + validator = get_named_validator("core/reporting-revision.json", version="3.2.0-rc.4") + assert validator is not None + for revision in payload["revisions"]: + validator.validate(revision) + ledger = ReportingLedger( + ledger_snapshot_id=response.ledger_snapshot_id, + ledger_as_of=response.ledger_as_of, + account_id=response.account_id, + scope=response.scope, + obligations=response.periods, + revisions=response.revisions, + materializations=response.materializations, + receipts=response.receipts, + ) + expected = [ + ExpectedReportingPeriod( + config.delivery_config_id, + 1, + config.report_definition_id, + config.feed_purpose, + config.reporting_profile, + config.media_buy_ids, + START.isoformat(), + END.isoformat(), + ) + ] + return response, evaluate_reporting_ledger(ledger, expected_periods=expected) + + +@pytest.mark.parametrize("observed", [END, TURN, TURN + timedelta(seconds=1)]) +async def test_creation_follows_acquisition_and_staged_read(store, tmp_path, observed): + config, source, reader, clock, factory = await setup( + store, tmp_path / "source", observed=observed + ) + turn = await factory().run_worker() + assert len(turn.revisions_committed) == len(source.executions) == reader.calls == 1 + request, result = source.executions[0] + manifest = await validate_reporting_source_execution( + capabilities=source.capabilities, + request=request, + result=result, + object_reader=source.reader, + clock=clock, + ) + response, outcome = await public_outcome(store, config) + revision = response.revisions[0] + assert outcome.definitive, [o.reasons for o in outcome.obligations] + assert request.period.source_read_cutoff_at == TURN + assert revision.created_at == PUBLISHED + assert revision.observed_at == revision.finalized_at == manifest.observed_at == observed + assert revision.data_through == END + + +async def test_real_clock_observation_after_dispatch_is_definitive(store, tmp_path): + config, source, reader, _, factory = await setup(store, tmp_path / "source", real=True) + turn = await factory().run_worker() + assert len(turn.revisions_committed) == reader.calls == 1 + request, result = source.executions[0] + manifest = await validate_reporting_source_execution( + capabilities=source.capabilities, + request=request, + result=result, + object_reader=source.reader, + ) + response, outcome = await public_outcome(store, config) + revision = response.revisions[0] + assert outcome.definitive, [o.reasons for o in outcome.obligations] + assert request.period.source_read_cutoff_at < manifest.observed_at <= revision.created_at + assert revision.observed_at == revision.finalized_at == manifest.observed_at + + +@pytest.mark.parametrize( + "observed,completion", + [(TURN + timedelta(seconds=30), PUBLISHED), (END, TURN - timedelta(seconds=1))], + ids=["source-clock-in-future", "trusted-clock-regresses-after-dispatch"], +) +async def test_invalid_publication_clock_stops_before_immutable_write( + store, tmp_path, observed, completion, monkeypatch +): + config, source, _, _, factory = await setup( + store, + tmp_path / "source", + observed=observed, + completion=completion, + ) + writes = [] + original = store.commit_revision + + async def commit(revision, rows): + writes.append(revision) + return await original(revision, rows) + + monkeypatch.setattr(store, "commit_revision", commit) + with pytest.raises(LedgerConflictError) as raised: + await factory().run_worker() + assert raised.value.code == "PUBLICATION_TIME_INVALID" + assert not writes + request, result = source.executions[0] + manifest = parse_verified_source_batch_manifest_v1( + result.response.manifest, result.manifest_bytes + ) + assert manifest.observed_at == observed # Original source evidence was not backdated. + assert ( + await store.list_revisions( + account_id=config.account_id, + reporting_obligation_id=request.identity.reporting_obligation_id, + ) + == () + ) + + +@pytest.mark.parametrize("finality", ["official", "snapshot"]) +async def test_same_publication_replay_preserves_committed_time_and_source_evidence( + store, tmp_path, finality +): + config, source, reader, clock, factory = await setup( + store, tmp_path / "source", finality=finality + ) + producer = factory() + turn = await producer.run_worker() + assert len(turn.revisions_committed) == 1 + request, result = source.executions[0] + manifest = parse_verified_source_batch_manifest_v1( + result.response.manifest, result.manifest_bytes + ) + obligation = await store.get_obligation( + account_id=config.account_id, + reporting_obligation_id=request.identity.reporting_obligation_id, + ) + first = await store.get_revision( + account_id=config.account_id, reporting_revision_id=turn.revisions_committed[0] + ) + assert obligation is not None and first is not None + expected_revisions = (first,) + if finality == "snapshot": + clock.now = PUBLISHED + timedelta(minutes=1) + reader.complete = lambda: None + restated = await producer.acquire_obligation(config, obligation, restate=True) + assert restated is not None + assert restated.supersedes_reporting_revision_id == first.reporting_revision_id + expected_revisions = (first, restated) + clock.now = PUBLISHED + timedelta(minutes=10) + for caller in (producer, factory()): + replay = await caller.commit_revision_from_manifest( + obligation, + manifest, + rows=ROWS, + finality=finality, + now=clock.now, + ) + assert replay == first + assert ( + await store.list_revisions( + account_id=config.account_id, + reporting_obligation_id=obligation.reporting_obligation_id, + ) + == expected_revisions + ) + changed = [dict(ROWS[0], row_id="changed")] + with pytest.raises(LedgerConflictError) as conflict: + await factory().commit_revision_from_manifest( + obligation, + manifest, + rows=changed, + finality=finality, + now=clock.now, + ) + assert conflict.value.code == "REVISION_IMMUTABLE" + + +async def test_explicit_dispatch_time_and_source_temporal_negatives(store, tmp_path): + config, source, _, _, factory = await setup(store, tmp_path / "source") + producer = factory() + obligations = await producer.close_elapsed_periods(config, now=TURN) + assert len(obligations) == 1 + first = await producer.acquire_obligation(config, obligations[0], now=TURN) + request, result = source.executions[0] + manifest = parse_verified_source_batch_manifest_v1( + result.response.manifest, result.manifest_bytes + ) + assert first is not None and request.period.source_read_cutoff_at == TURN + assert first.created_at == PUBLISHED + for updates in ( + {"observed_at": END - timedelta(seconds=1)}, + { + "finality_evidence": { + **manifest.finality_evidence.model_dump(), + "observed_at": manifest.acquired_at + timedelta(seconds=1), + } + }, + ): + with pytest.raises(ValidationError): + SourceBatchManifestV1.model_validate({**manifest.model_dump(), **updates}) + + +@pytest.mark.parametrize("future_acquisition", [False, True]) +async def test_direct_commit_uses_explicit_publication_time_before_any_write( + store, tmp_path, monkeypatch, future_acquisition +): + config, source, reader, clock, factory = await setup(store, tmp_path / "source") + producer = factory() + obligations = await producer.close_elapsed_periods(config, now=TURN) + + async def unavailable(**kwargs): + raise OSError("fixture staged read interrupted") + + # Retain a real source publication after an interrupted read, before any + # revision exists. The public commit primitive can then retry that manifest. + with monkeypatch.context() as patch: + patch.setattr(reader, "read", unavailable) + with pytest.raises(OSError, match="fixture staged read interrupted"): + await producer.acquire_obligation(config, obligations[0], now=TURN) + _, result = source.executions[0] + manifest = parse_verified_source_batch_manifest_v1( + result.response.manifest, result.manifest_bytes + ) + clock.now = PUBLISHED + timedelta(minutes=10) + if future_acquisition: + # Source-contract-valid evidence may still be in the future relative to + # the trusted publication clock. It must fail at the producer boundary. + manifest = SourceBatchManifestV1.model_validate( + {**manifest.model_dump(), "acquired_at": PUBLISHED + timedelta(seconds=1)} + ) + writes = [] + commit = store.commit_revision + + async def record_write(*args, **kwargs): + writes.append(args) + return await commit(*args, **kwargs) + + monkeypatch.setattr(store, "commit_revision", record_write) + with pytest.raises(LedgerConflictError) as raised: + await producer.commit_revision_from_manifest( + obligations[0], manifest, rows=ROWS, finality="official", now=PUBLISHED + ) + assert raised.value.code == "PUBLICATION_TIME_INVALID" + assert not writes + assert ( + await store.list_revisions( + account_id=config.account_id, + reporting_obligation_id=obligations[0].reporting_obligation_id, + ) + == () + ) + else: + committed = await producer.commit_revision_from_manifest( + obligations[0], manifest, rows=ROWS, finality="official", now=PUBLISHED + ) + assert committed.created_at == PUBLISHED + assert committed.observed_at == committed.finalized_at == TURN