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
5 changes: 3 additions & 2 deletions .github/workflows/ci.yml
Original file line number Diff line number Diff line change
Expand Up @@ -764,7 +764,7 @@ jobs:
runs-on: ubuntu-latest
permissions:
contents: read
timeout-minutes: 35
timeout-minutes: 45
services:
postgres:
image: postgres:16
Expand All @@ -786,7 +786,8 @@ jobs:
run: pip install -e ".[dev,pg]"
- name: Run complete production and projection conformance
shell: bash
timeout-minutes: 30
# A passing suite plus harness cleanup can exceed 30 minutes.
timeout-minutes: 40
env:
ADCP_PG_TEST_URL: postgresql://postgres@localhost:5432/adcp_production_test
run: |
Expand Down
12 changes: 12 additions & 0 deletions docs/reporting-production.md
Original file line number Diff line number Diff line change
Expand Up @@ -398,3 +398,15 @@ pre-`e16eb8cf` hardening, and pre-`34c8f6d9` production. In particular, old A wh
is not supported, and historical false notification-readiness results do not
become healthy through later source integration. These notes do not qualify
simultaneous old autonomous writers or release/activation acceptance.

### Source revocation

`ReportingProductionSource.configuration_binding(configuration)` is also the
live authorization callback for each dispatch and publication. The SDK checks
again under the account lock before sealing or committing a fetched result.
`None` discards the result without a success checkpoint and stops further source
work for that account in the current turn; restoration resumes on the next turn.
Replay and restart perform fresh checks. The SDK does not cancel an in-flight
fetch for revocation: revocation takes effect at the next dispatch or publish.
Adopters own any caching or latency inside the callback. See the
[release guide](reporting-release-notes.md#source-authorization-d5-option-a).
19 changes: 19 additions & 0 deletions docs/reporting-release-notes.md
Original file line number Diff line number Diff line change
Expand Up @@ -46,6 +46,25 @@ writes do not authorize concurrent incompatible autonomous materializers,
status projectors or clock sweepers. Follow the component-specific migration
procedures before starting those workers.

## Source authorization (D5 Option A)

The adopter's existing `ReportingProductionSource.configuration_binding(configuration)`
is the authority for source work. Returning `None` withdraws authorization.
The SDK calls it immediately before each source dispatch and again under the
account lock before sealing or publishing the result. If authorization has been
withdrawn, the result is discarded without publication or a success checkpoint,
and the account gets no further source work during that turn. A restored
binding resumes work on the next turn. Recovery and replay after a restart make
fresh checks; an earlier successful check grants no authority to publish later.

An in-flight fetch is allowed to finish: revocation takes effect at the next dispatch or publish.
Adopters own any caching or latency inside their callback. There is no additional
authorization adapter, persisted authorization grant, freshness budget, epoch,
compare-and-swap protocol, or settlement-only mode; the proposed F/L/P values
are not part of this contract. Existing frozen generation mappings still protect
historical scope. Buyer exposure retains the existing per-request feed
reauthorization and per-session destination authorization.

## Current and historical protocol versions

Live reporting mounts and callers use AdCP `3.2-rc.6`, whose bundle spelling is
Expand Down
62 changes: 62 additions & 0 deletions examples/reporting_service_production.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,62 @@
"""Own an existing B2 graph through the service lifecycle.

This explicit bridge still requires production providers, fixed source profiles
and the typed account task. It does not implement the separate adapter-first
factory, durable acquisition-envelope, or database-time fencing contracts.
"""

from collections.abc import Mapping
from typing import Any

from adcp.reporting.ledger import ReportingProducer
from adcp.reporting.materializer import ReportingMaterializerService
from adcp.reporting.production import (
ReportingProductionConfigurationTask,
ReportingProductionHandler,
ReportingProductionOffering,
ReportingProductionSourceRegistry,
ReportingProductionSupport,
)
from adcp.reporting.projection import InMemoryReportingStatusProjection, PgReportingStatusProjection
from adcp.reporting.receipts import ReceiptAccountResolver
from adcp.reporting.service import ReliableReportingService, ReportingContextResolver
from adcp.server import ADCPHandler


def compose_service(
*,
materializer: ReportingMaterializerService,
projection: InMemoryReportingStatusProjection | PgReportingStatusProjection,
offerings: tuple[ReportingProductionOffering, ...],
producers: Mapping[str, ReportingProducer],
account_context: ReportingContextResolver,
configuration_task: ReportingProductionConfigurationTask,
resolve_account: ReceiptAccountResolver,
application: ADCPHandler[Any],
) -> tuple[ReliableReportingService, ReportingProductionHandler]:
"""Register before mounting; start/close own workers, while pools stay borrowed.

Each stable name identifies one fixed execution profile. Account resolution
must return that profile's currency, scope, metric set and source offerings.
Recovery uses persisted generation facts and the provider's live account
binding; it does not call account_context or enumerate accounts.

Mount the returned exact handler on MCP and A2A, then use service.start and
service.close as lifespan hooks. Retain the same event loop until shutdown
settles; configure owned pools with ReportingServiceResource when needed.
Reporting admission comes from the typed sync_accounts task, never configure.
"""
registry = ReportingProductionSourceRegistry(account_context=account_context)
for name, producer in producers.items():
registry.register(name, producer)
support = ReportingProductionSupport(
materializer,
projection,
offerings=offerings,
configuration_task=configuration_task,
resolve_account=resolve_account,
source_registry=registry,
)
service = ReliableReportingService.from_production(support)
service.install(application)
return service, support.handler
1 change: 1 addition & 0 deletions pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -202,6 +202,7 @@ adcp = [
"reporting/feed/*.json",
"reporting/projection/*.json",
"reporting/production/*.json",
"reporting/production/*.sql",
# PREVIEW: vendored sync_reporting_status schemas. They are the runtime
# validator for the wire conditionals codegen cannot express, so the wheel
# must carry them. Removed with the rest of _preview/ at rc.2.
Expand Down
131 changes: 131 additions & 0 deletions src/adcp/reporting/_inline_storage_schema.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,131 @@
"""Required catalog objects from reporting_inline_storage.sql on PostgreSQL 16.

Generated from a clean migration; unrelated adopter objects are permitted.
"""

REQUIRED_OBJECTS: dict[str, dict[str, str | bool]] = {
"column:reporting_inline_objects.account_id": {
"enabled": True,
"fingerprint": "fdc4b549138bf90a5af5ed00b87a86ff9a0cbea553e7f017b4d95510708f6f88",
},
"column:reporting_inline_objects.payload": {
"enabled": True,
"fingerprint": "554c34e416bd4546469b42d5773bc7c59b2104366ffa5259a2838c64cd826e57",
},
"column:reporting_inline_objects.payload_sha256": {
"enabled": True,
"fingerprint": "fdc4b549138bf90a5af5ed00b87a86ff9a0cbea553e7f017b4d95510708f6f88",
},
"column:reporting_inline_seals.account_id": {
"enabled": True,
"fingerprint": "fdc4b549138bf90a5af5ed00b87a86ff9a0cbea553e7f017b4d95510708f6f88",
},
"column:reporting_inline_seals.byte_count": {
"enabled": True,
"fingerprint": "64be57437fdc0a07a97985c2aa058031f8082db7251bdb4d5afa1a9b088de97a",
},
"column:reporting_inline_seals.manifest": {
"enabled": True,
"fingerprint": "554c34e416bd4546469b42d5773bc7c59b2104366ffa5259a2838c64cd826e57",
},
"column:reporting_inline_seals.manifest_sha256": {
"enabled": True,
"fingerprint": "fdc4b549138bf90a5af5ed00b87a86ff9a0cbea553e7f017b4d95510708f6f88",
},
"column:reporting_inline_seals.source_execution_key": {
"enabled": True,
"fingerprint": "fdc4b549138bf90a5af5ed00b87a86ff9a0cbea553e7f017b4d95510708f6f88",
},
"column:reporting_inline_seals.staged_commit_ref": {
"enabled": True,
"fingerprint": "fdc4b549138bf90a5af5ed00b87a86ff9a0cbea553e7f017b4d95510708f6f88",
},
"constraint:reporting_inline_objects.reporting_inline_objects_account": {
"enabled": True,
"fingerprint": "ff3e5aec84611ad080f664147be02d92ab106b3e2c9407bc63b5396bc0c8006f",
},
"constraint:reporting_inline_objects.reporting_inline_objects_bytes": {
"enabled": True,
"fingerprint": "5cf3d75f67d8c5c6631ba87b9d6d301a4f807f9ce619bc77e3b81f43707d6669",
},
"constraint:reporting_inline_objects.reporting_inline_objects_digest": {
"enabled": True,
"fingerprint": "3a303f02341c258a572f9eaf1e6869cdfd7547bc2932c44c41a8d1a7eed280a2",
},
"constraint:reporting_inline_objects.reporting_inline_objects_pk": {
"enabled": True,
"fingerprint": "52c73c955e80fe2c8de482afc62851429635aa24f3a10cd4e15925eaa5a05bb3",
},
"constraint:reporting_inline_objects.reporting_inline_objects_size": {
"enabled": True,
"fingerprint": "7c4fb6f6c0ff3c05be6ae3986c8a849e886e06e7088f5d09fef81b9f3b326ed8",
},
"constraint:reporting_inline_seals.reporting_inline_seals_account": {
"enabled": True,
"fingerprint": "ff3e5aec84611ad080f664147be02d92ab106b3e2c9407bc63b5396bc0c8006f",
},
"constraint:reporting_inline_seals.reporting_inline_seals_account_binding": {
"enabled": True,
"fingerprint": "ae8a06d397d61baba67cc4e77001596682a09ebecac0babaa0268c79ac7f242d",
},
"constraint:reporting_inline_seals.reporting_inline_seals_bytes": {
"enabled": True,
"fingerprint": "2cf54d2016b1d93e6cc2e8f6538ed87fc5b04b980048275210661926eeb49296",
},
"constraint:reporting_inline_seals.reporting_inline_seals_count": {
"enabled": True,
"fingerprint": "1e87e7574019dc01b1b39584ef6a420e82330e7089490b928f04f5bce4482f7c",
},
"constraint:reporting_inline_seals.reporting_inline_seals_digest": {
"enabled": True,
"fingerprint": "6d150c5fa94598b18c9adf65800d6ddd58479f97cb9b845b3d0f7a791aafced4",
},
"constraint:reporting_inline_seals.reporting_inline_seals_key": {
"enabled": True,
"fingerprint": "fadbcdf37ecabe857d9111775f689255637a0cef3028301593c285b69a0551a2",
},
"constraint:reporting_inline_seals.reporting_inline_seals_key_binding": {
"enabled": True,
"fingerprint": "897155db2cfdde3b4ed88ea3066c43b8938b9961e9bcca08eac4acbe71e5d50b",
},
"constraint:reporting_inline_seals.reporting_inline_seals_pk": {
"enabled": True,
"fingerprint": "123e2c394bd92bfab2862432ed6191484dd5438ca9b9f6372108b59c901b4248",
},
"constraint:reporting_inline_seals.reporting_inline_seals_ref": {
"enabled": True,
"fingerprint": "6ef63ffe9db91b8fa1fe7f75817e958e343cfbc0f8d9134a188fc65ce7d39a88",
},
"constraint:reporting_inline_seals.reporting_inline_seals_size": {
"enabled": True,
"fingerprint": "b6c20825c79f2f4cbee21b439842ace6958549b83913f82d958de73980754ab2",
},
"function:reporting_inline_immutable()": {
"enabled": True,
"fingerprint": "f2baed2e6158fc173155a71c5dedd3423b32d5923582600ae3c7062086c8ba41",
},
"index:reporting_inline_objects.reporting_inline_objects_pk": {
"enabled": True,
"fingerprint": "ba41f90f12a63ad244a676e7d7ae37aacca579d32f6a1ced16a5676c420083d0",
},
"index:reporting_inline_seals.reporting_inline_seals_pk": {
"enabled": True,
"fingerprint": "0816d86030df87c7d77ed0df4cec81b2a30e61d959fb53063750e4c2934623aa",
},
"table:reporting_inline_objects": {
"enabled": True,
"fingerprint": "1f824779ff80f110344420b019786663d8c9beaad230da90e0439795e734ccda",
},
"table:reporting_inline_seals": {
"enabled": True,
"fingerprint": "1f824779ff80f110344420b019786663d8c9beaad230da90e0439795e734ccda",
},
"trigger:reporting_inline_objects.reporting_inline_objects_immutable": {
"enabled": True,
"fingerprint": "da630163871019e8457c2d54b41f23e7f182ad231e539d4f57360ea45c39a1c1",
},
"trigger:reporting_inline_seals.reporting_inline_seals_immutable": {
"enabled": True,
"fingerprint": "201a920cf5c5343295735a14da5eee8def6fd80bd72826628ba3706ebc75d9d2",
},
}
89 changes: 89 additions & 0 deletions src/adcp/reporting/_source_authorization.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,89 @@
"""Turn-local source revocation and the SDK inline publication boundary.

Only denials are remembered, for the remainder of a scheduling turn. Every
dispatch and publication calls the adopter again; no authorization is cached.
"""

from __future__ import annotations

from collections.abc import AsyncIterator, Awaitable, Callable, Iterator
from contextlib import AbstractAsyncContextManager, asynccontextmanager, contextmanager
from contextvars import ContextVar
from typing import TYPE_CHECKING, NoReturn, TypeAlias

if TYPE_CHECKING:
from adcp.reporting.inline_source import ReportingSealStore, SealedSlice

InlineSealPublisher: TypeAlias = Callable[[str, "SealedSlice"], Awaitable["SealedSlice"]]

_REVOKED_ACCOUNTS: ContextVar[set[str] | None] = ContextVar(
"reporting_source_revoked_accounts", default=None
)
_INLINE_PUBLICATION: ContextVar[
tuple[
str,
Callable[[ReportingSealStore], AbstractAsyncContextManager[InlineSealPublisher | None]],
]
| None
] = ContextVar("reporting_inline_publication", default=None)


@contextmanager
def source_turn() -> Iterator[None]:
"""Share denials across configurations in this turn, never across turns."""
if _REVOKED_ACCOUNTS.get() is not None:
yield
return
token = _REVOKED_ACCOUNTS.set(set())
try:
yield
finally:
_REVOKED_ACCOUNTS.reset(token)


def account_revoked(account_id: str) -> bool:
return account_id in (_REVOKED_ACCOUNTS.get() or ())


def source_revoked(account_id: str) -> NoReturn:
from adcp.reporting.production.contracts import _SourceAuthorizationRevokedError

revoked = _REVOKED_ACCOUNTS.get()
if revoked is not None:
revoked.add(account_id)
raise _SourceAuthorizationRevokedError()


def require_account_work(account_id: str) -> None:
if account_revoked(account_id):
source_revoked(account_id)


@contextmanager
def bind_inline_publication(
account_id: str,
guard: Callable[[ReportingSealStore], AbstractAsyncContextManager[InlineSealPublisher | None]],
) -> Iterator[None]:
"""Carry the producer's lock and live check into its inline executor task."""
token = _INLINE_PUBLICATION.set((account_id, guard))
try:
yield
finally:
_INLINE_PUBLICATION.reset(token)


@asynccontextmanager
async def inline_publication(
account_id: str, seals: ReportingSealStore
) -> AsyncIterator[InlineSealPublisher]:
async def publish(key: str, sealed: SealedSlice) -> SealedSlice:
return await seals.put(account_id=account_id, source_execution_key=key, sealed=sealed)

bound = _INLINE_PUBLICATION.get()
if bound is None:
yield publish
return
if bound[0] != account_id:
raise ValueError("source publication must belong to the dispatched account")
async with bound[1](seals) as owned_publication:
yield owned_publication or publish
Loading
Loading