Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
42 changes: 38 additions & 4 deletions src/adcp/reporting/ledger/producer.py
Original file line number Diff line number Diff line change
Expand Up @@ -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(
Expand Down Expand Up @@ -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,
)

Expand Down Expand Up @@ -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)
Expand Down Expand Up @@ -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,
Expand All @@ -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=(
Expand Down
57 changes: 51 additions & 6 deletions tests/conformance/reporting/_late_account_server.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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)
Expand Down Expand Up @@ -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(
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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"),
Expand Down
1 change: 1 addition & 0 deletions tests/conformance/reporting/_production_packaging.py
Original file line number Diff line number Diff line change
Expand Up @@ -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",
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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)); "
Expand All @@ -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,
Expand Down
Original file line number Diff line number Diff line change
@@ -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}}))
Loading
Loading