Skip to content

Commit 65be431

Browse files
authored
feat(reporting): compose PostgreSQL services from adapters (#1231)
* feat(reporting): compose PostgreSQL services from adapters * fix(reporting): expose production options and document source schema upgrade
1 parent ac5e432 commit 65be431

7 files changed

Lines changed: 875 additions & 66 deletions

File tree

‎docs/reliable-reporting-service.md‎

Lines changed: 86 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -108,33 +108,104 @@ asserted only in the request body.
108108

109109
## PostgreSQL production wiring
110110

111+
The PostgreSQL factory installs ledger, source staging, and replay-seal schemas
112+
in the caller's pool. Ordinary adapter registrations use those durable stores;
113+
restart replay reads the sealed bytes without refetching the source. Open the
114+
pool before starting the service and close it after service shutdown settles.
115+
This is a change for existing `postgres()` callers, including those that do not
116+
enable `production`: `initialize()` now bootstraps `reporting_inline_objects`
117+
and `reporting_inline_seals` in addition to the ledger schema. Before upgrading,
118+
ensure the pool's startup role can run that source-schema DDL. A DDL-restricted
119+
runtime role cannot use this factory's schema bootstrap until its deployment
120+
grants or startup sequence are updated.
121+
122+
Pass `ReportingProductionOptions` to compose the managed materializer, status
123+
projection, exact reads, consumer/receipt handlers, and optional signed
124+
notification workers. Adopters supply domain declarations and their actual
125+
destination, verifier registry, trusted account task and authorization callbacks.
126+
The SDK constructs the component graph:
127+
111128
```python
112129
from psycopg_pool import AsyncConnectionPool
113-
from adcp.reporting import ReliableReportingService
130+
from adcp.reporting import (
131+
ReliableReportingService,
132+
ReportingProductionOptions,
133+
ReportingServiceOffering,
134+
)
114135

115136
pool = AsyncConnectionPool(DATABASE_URL, open=False)
116137
await pool.open()
117138

118139
reporting = ReliableReportingService.postgres(
119140
pool=pool,
120141
account_context=resolve_reporting_context,
121-
caller_resolver=resolve_authenticated_caller,
122-
worker_interval=timedelta(minutes=1),
123-
materialization_worker=managed_delivery_worker, # optional tier
124-
notification_worker=notification_worker, # optional extension
125-
notification_attempt_store=notification_attempts,
126-
receipt_handler=receipt_handler, # optional tier
127-
reconciled_billing=True,
128-
worker_error_handler=report_reporting_worker_error,
142+
consumer_status_enabled=True,
143+
production=ReportingProductionOptions(
144+
offerings=(
145+
ReportingServiceOffering(
146+
adapter="gam",
147+
offering=gam_public_offering, # ReportingDeliveryOffering
148+
profile=gam_execution_profile, # ProducerOfferings
149+
source_offering_id=gam_source_offering_id,
150+
verification_key=gam_verifier.key,
151+
),
152+
),
153+
destination=warehouse_provider, # ReportingProductionDestination
154+
registry=verifier_registry, # ReportingRevisionVerifierRegistry
155+
configuration_task=account_task, # ReportingProductionConfigurationTask
156+
resolve_account=authorize_reporting_account,
157+
),
129158
)
159+
reporting.sources.register("gam", gam_adapter)
160+
platform = reporting.install(platform)
161+
# Mount the returned handler on MCP/A2A before reporting.start().
162+
# Use reporting.start / reporting.close as the application's lifespan hooks.
130163
```
131164

132-
The PostgreSQL factory makes the ledger durable; it does not make the default
133-
adapter staging or replay-seal stores durable. Production adapters should pass
134-
durable `staging=` and `seals=` implementations to `sources.register`, or use
135-
`sources.register_executor` for a custom executor and object reader. Committed revision rows are retained in the ledger for exact reads. Durable
136-
staging and seals are needed to recover interrupted acquisitions and replay
137-
previously sealed source results across a restart.
165+
Each registered production adapter needs at least one public offering. Multiple
166+
public offerings for an adapter share its fixed execution profile and verifier;
167+
different adapters may reuse a local source offering ID. Profiles include the
168+
source scope, currency, metric/dimension sets, and snapshot/official selections.
169+
Admission checks trusted account context against the selected profile and
170+
persists it once per generation. Restart recovery uses the stored route.
171+
172+
Managed adapters also implement the existing
173+
`ReportingProductionSource.configuration_binding(configuration)` method. Return
174+
the current source binding, including the exact media-buy/product mapping, or
175+
`None` when unauthorized. The wrapper forwards this callback before dispatch and
176+
again under the account lock before seal/publication. Revocation or remapping
177+
discards in-flight results; restoring the admitted mapping allows the next turn
178+
to retry. **Revocation takes effect at the next dispatch or publish.** Adopters
179+
own any caching and latency inside the callback. A fetch already in progress is
180+
allowed to return. Buyer-facing feed and destination authorization still run on
181+
each request and session.
182+
183+
To enable signed push, set `notifications=True` and provide `subscriptions`,
184+
`cipher` (`ReportingEnvelopeCipher`), and `signing` (`ReportingProductionSigning`)
185+
on the options. All three are required together. The factory owns the outboxes,
186+
attempt storage and workers for source, materialization and status events.
187+
With only `notifications=True`, events are retained for polling without claiming
188+
push delivery. Consumer status is controlled by the service's
189+
`consumer_status_enabled` argument. Reconciled Billing follows the validated
190+
official offering, receipt method, and destination contracts; a flag cannot
191+
promote an unsupported provider. The memory factory accepts the same options
192+
for conformance, but never advertises durable Managed/Reconciled guarantees.
193+
194+
Use the typed `sync_accounts` admission callback for this composition, rather
195+
than `configure()`. Production workers belong to `start()`/`close()`; the Core
196+
`run_worker()` entry point is not used. See the
197+
[production guide](reporting-production.md) for account admission and provider
198+
contracts, and [the wiring example](../examples/reporting_service_production.py).
199+
200+
## Advanced composition
201+
202+
`postgres()` without production options remains the Core factory. Advanced
203+
adopters can still inject workers and receipt handlers, pass explicit
204+
`staging=`/`seals=` to `sources.register`, use `constituent_of=` for a custom row
205+
identity, or register a complete executor and object reader. Explicit stores are
206+
borrowed; initialize their schemas yourself. `from_production()` continues to
207+
own an already composed production graph. Keep injected components separate
208+
from `production=` options, which own that graph themselves.
138209

139210
Managed delivery, notification, and receipt components are replaceable
140211
extensions. Startup rejects combinations that cannot be advertised honestly,
Lines changed: 18 additions & 38 deletions
Original file line numberDiff line numberDiff line change
@@ -1,43 +1,30 @@
1-
"""Own an existing B2 graph through the service lifecycle.
1+
"""Compose PostgreSQL reporting from registered adapters and trusted providers.
22
3-
This explicit bridge still requires production providers, fixed source profiles
4-
and the typed account task. It does not implement the separate adapter-first
5-
factory, durable acquisition-envelope, or database-time fencing contracts.
3+
See docs/reliable-reporting-service.md for ReportingProductionOptions inputs.
4+
The service owns the graph and source storage; the caller owns the pool and
5+
provider clients. Existing hand-built graphs can still use from_production().
66
"""
77

88
from collections.abc import Mapping
99
from typing import Any
1010

11-
from adcp.reporting.ledger import ReportingProducer
12-
from adcp.reporting.materializer import ReportingMaterializerService
13-
from adcp.reporting.production import (
14-
ReportingProductionConfigurationTask,
15-
ReportingProductionHandler,
16-
ReportingProductionOffering,
17-
ReportingProductionSourceRegistry,
18-
ReportingProductionSupport,
19-
)
20-
from adcp.reporting.projection import InMemoryReportingStatusProjection, PgReportingStatusProjection
21-
from adcp.reporting.receipts import ReceiptAccountResolver
22-
from adcp.reporting.service import ReliableReportingService, ReportingContextResolver
11+
from adcp.reporting import ReliableReportingService, ReportingAdapter, ReportingProductionOptions
12+
from adcp.reporting.service import ReportingContextResolver
2313
from adcp.server import ADCPHandler
2414

2515

2616
def compose_service(
2717
*,
28-
materializer: ReportingMaterializerService,
29-
projection: InMemoryReportingStatusProjection | PgReportingStatusProjection,
30-
offerings: tuple[ReportingProductionOffering, ...],
31-
producers: Mapping[str, ReportingProducer],
18+
pool: Any,
19+
adapters: Mapping[str, ReportingAdapter],
20+
production: ReportingProductionOptions,
3221
account_context: ReportingContextResolver,
33-
configuration_task: ReportingProductionConfigurationTask,
34-
resolve_account: ReceiptAccountResolver,
3522
application: ADCPHandler[Any],
36-
) -> tuple[ReliableReportingService, ReportingProductionHandler]:
23+
) -> tuple[ReliableReportingService, ADCPHandler[Any]]:
3724
"""Register before mounting; start/close own workers, while pools stay borrowed.
3825
39-
Each stable name identifies one fixed execution profile. Account resolution
40-
must return that profile's currency, scope, metric set and source offerings.
26+
Each stable name identifies a fixed ReportingServiceOffering.profile. Account
27+
resolution must return its currency, scope, metric set and source offerings.
4128
Recovery uses persisted generation facts and the provider's live account
4229
binding; it does not call account_context or enumerate accounts.
4330
@@ -46,17 +33,10 @@ def compose_service(
4633
settles; configure owned pools with ReportingServiceResource when needed.
4734
Reporting admission comes from the typed sync_accounts task, never configure.
4835
"""
49-
registry = ReportingProductionSourceRegistry(account_context=account_context)
50-
for name, producer in producers.items():
51-
registry.register(name, producer)
52-
support = ReportingProductionSupport(
53-
materializer,
54-
projection,
55-
offerings=offerings,
56-
configuration_task=configuration_task,
57-
resolve_account=resolve_account,
58-
source_registry=registry,
36+
service = ReliableReportingService.postgres(
37+
pool=pool, account_context=account_context, production=production
5938
)
60-
service = ReliableReportingService.from_production(support)
61-
service.install(application)
62-
return service, support.handler
39+
for name, adapter in adapters.items():
40+
service.sources.register(name, adapter)
41+
handler = service.install(application)
42+
return service, handler

‎src/adcp/reporting/__init__.py‎

Lines changed: 12 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -78,6 +78,7 @@
7878
from adcp.reporting import materializer as materializer
7979
from adcp.reporting import revision_selection as revision_selection
8080
from adcp.reporting import service as service
81+
from adcp.reporting import service_production as service_production
8182
from adcp.reporting import source as source
8283
from adcp.reporting import testing as testing
8384
from adcp.reporting.service import (
@@ -107,6 +108,12 @@
107108
from adcp.reporting.service_lifecycle import (
108109
ReportingServiceResource as ReportingServiceResource,
109110
)
111+
from adcp.reporting.service_production import (
112+
ReportingProductionOptions as ReportingProductionOptions,
113+
)
114+
from adcp.reporting.service_production import (
115+
ReportingServiceOffering as ReportingServiceOffering,
116+
)
110117
from adcp.reporting.testing import (
111118
DeterministicReportingClock as DeterministicReportingClock,
112119
)
@@ -149,6 +156,7 @@
149156
"materializer",
150157
"revision_selection",
151158
"service",
159+
"service_production",
152160
"source",
153161
"testing",
154162
}
@@ -167,6 +175,8 @@
167175
"ReportingServiceResource": ("service_lifecycle", "ReportingServiceResource"),
168176
"ReportingAccountContext": ("service", "ReportingAccountContext"),
169177
"ReportingAdapter": ("service", "ReportingAdapter"),
178+
"ReportingProductionOptions": ("service_production", "ReportingProductionOptions"),
179+
"ReportingServiceOffering": ("service_production", "ReportingServiceOffering"),
170180
"DeterministicReportingClock": ("testing", "DeterministicReportingClock"),
171181
"FaultInjectingReportingSealStore": ("testing", "FaultInjectingReportingSealStore"),
172182
"FaultInjectingReportingStagingStore": ("testing", "FaultInjectingReportingStagingStore"),
@@ -219,6 +229,8 @@ def __dir__() -> list[str]:
219229
"ReportingServiceResource",
220230
"ReportingAccountContext",
221231
"ReportingAdapter",
232+
"ReportingProductionOptions",
233+
"ReportingServiceOffering",
222234
"DeterministicReportingClock",
223235
"FaultInjectingReportingSealStore",
224236
"FaultInjectingReportingStagingStore",

0 commit comments

Comments
 (0)