Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
23 commits
Select commit Hold shift + click to select a range
02125d5
fix(decisioning): serve account-scoped task polling and listing
bokelley Oct 2, 2026
e382615
feat(decisioning): observe committed task lifecycle transitions
bokelley Oct 2, 2026
d1c2b77
chore(decisioning): combine task registry prerequisites
bokelley Oct 2, 2026
d5319d7
fix(decisioning): preserve optional gates and canonical task failures
bokelley Oct 2, 2026
e4b9fb5
chore(decisioning): incorporate task polling review fixes
bokelley Oct 2, 2026
b5c7d42
feat(pg): lazily resolve decisioning stores on the serving loop
bokelley Oct 2, 2026
37a69cd
fix(pg): type terminal task rows without optional driver imports
bokelley Oct 2, 2026
dedb338
chore(pg): propagate optional-driver task row typing repair
bokelley Oct 2, 2026
a022fb2
test(tasks): gate polling fixtures on configured PostgreSQL
bokelley Oct 3, 2026
4b66cd5
test(tasks): configure PostgreSQL lifecycle coverage explicitly
bokelley Oct 3, 2026
eafc667
chore(tasks): propagate polling fixture CI repair
bokelley Oct 3, 2026
2698030
chore(tasks): propagate observer fixture CI repair
bokelley Oct 3, 2026
8247a3a
test(pg): separate fake lazy delegates from real database fixtures
bokelley Oct 3, 2026
e66829f
Merge branch 'main' into feat/task-lifecycle-observers
bokelley Oct 3, 2026
ad3b824
Merge branch 'main' into conductor/sdk-task-registry
bokelley Oct 3, 2026
f2d3cb7
chore(tasks): merge main while retaining PostgreSQL observer coverage
bokelley Oct 3, 2026
a0d7541
chore(tasks): merge main while retaining all PostgreSQL task coverage
bokelley Oct 3, 2026
39933f7
chore(tasks): merge landed polling with lifecycle observers
bokelley Oct 3, 2026
3af0228
chore(tasks): compose landed polling with lifecycle and lazy stores
bokelley Oct 3, 2026
ead8d38
perf(reporting): accelerate canonical surrogate detection
bokelley Oct 3, 2026
5fada42
chore(tasks): merge current main into lifecycle observers
bokelley Oct 4, 2026
05f8c41
chore(tasks): retain lazy store coverage across lifecycle main bridge
bokelley Oct 4, 2026
cbe62c9
chore(tasks): merge landed lifecycle while preserving lazy stores
bokelley Oct 4, 2026
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
1 change: 1 addition & 0 deletions .github/workflows/ci.yml
Original file line number Diff line number Diff line change
Expand Up @@ -201,6 +201,7 @@ jobs:
tests/conformance/decisioning/test_pg_idempotency_reservation.py \
tests/test_decisioning_task_polling.py \
tests/test_decisioning_task_lifecycle.py \
tests/test_decisioning_lazy_pg_stores.py \
tests/conformance/decisioning/test_pg_task_webhook_outbox.py \
tests/conformance/decisioning/test_pg_notification_outbox.py \
tests/test_notification_outbox.py \
Expand Down
71 changes: 71 additions & 0 deletions docs/lazy-pg-stores.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,71 @@
# Lazy PostgreSQL decisioning stores

`adcp.decisioning.pg` exports `LazyProposalStore`, `LazyTaskRegistry`, and
`LazyTaskWebhookOutbox`. Each accepts a synchronous or asynchronous zero-argument
factory returning its original concrete PostgreSQL store. Task factories may
also return an original `(PgTaskRegistry, PgTaskWebhookOutbox)` pair.

Factories are called on the serving event loop, on the first async operation or
explicit `await wrapper.resolve()`. Concurrent calls share one successful
resolution. Failed or canceled initialization is not cached, so later calls
retry. Factories own pool opening and cleanup on failed initialization; wrappers
never open, close, replace or reopen their caller-owned pools. Schema setup stays
explicit: do it in the factory or call the wrapper's `create_schema()`.

```python
from psycopg_pool import AsyncConnectionPool
from adcp.decisioning.pg import (
LazyTaskRegistry, LazyTaskWebhookOutbox, PgTaskRegistry, PgTaskWebhookOutbox,
)

pool = AsyncConnectionPool("postgresql://localhost/seller", open=False)

async def task_stores() -> tuple[PgTaskRegistry, PgTaskWebhookOutbox]:
await pool.open()
outbox = PgTaskWebhookOutbox(
pool=pool, sender=signed_sender, encryption_key=encryption_key,
delivery_retry_horizon_seconds=86400,
)
registry = PgTaskRegistry(pool=pool, task_webhook_outbox=outbox)
await registry.create_schema()
await outbox.create_schema()
return registry, outbox

registry = LazyTaskRegistry(task_stores)
worker_outbox = LazyTaskWebhookOutbox.from_registry(registry)
# await worker_outbox.run_worker() and registry.issue(...) share the same pair.
# Cancel/join workers before closing pool during application shutdown.
```

Construct the original registry and outbox together, passing the same pool to
both. Their existing constructor checks (including signing-scope wiring) remain
authoritative, and a returned pair is checked again for identity. The worker
facade from `from_registry()` resolves the same original outbox, rather than
running an independent pool factory. Do not pass an unresolved worker facade to
the concrete `PgTaskRegistry` constructor.

The wrappers declare the concrete stores' durability before resolution and
reject a factory returning a lossy store. `LazyTaskRegistry` retains the concrete
SDK registry's optional listing protocol; no arbitrary delegate attributes or
optional APIs are forwarded. Its additive observer methods can be registered or
removed before and after resolution. They receive the first submitted event and
retain the concrete registry's commit/no-op/failure semantics.

`resolved` inspects the cached store without running the factory. A task registry
exposes its original `task_webhook_outbox` and atomic-outbox marker after
resolution. Servers advertising SDK task-webhook signing must await resolution
before synchronous server construction: the boot validator needs the actual
sender, retry horizon and shared-pool proof. Resolved exact SDK wrappers are
accepted by that proof; unaudited subclasses remain rejected. Polling-only
servers can resolve on first use. Outbox synchronous registration/crypto helpers
and retry-horizon inspection also require explicit resolution; they raise a
clear error when called too early.

After success, a wrapper is bound to the event loop that resolved it. A different
loop raises an error, even if the previous loop has closed; create a fresh wrapper
and pool for the new loop. A failed initialization with no cached store may retry
on a new loop after the old loop closes. Closing the borrowed pool causes normal
pool errors on later operations, rather than rerunning the factory.

No database migration or change to the existing `LazyBackend` is required. The
new wrappers reuse the concrete stores' current SQL and schema assets.
14 changes: 14 additions & 0 deletions src/adcp/decisioning/pg/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -35,6 +35,14 @@
PG_AVAILABLE,
PgBuyerAgentRegistry,
)
from adcp.decisioning.pg.lazy import (
LazyProposalStore,
LazyProposalStoreFactory,
LazyTaskRegistry,
LazyTaskRegistryFactory,
LazyTaskWebhookOutbox,
LazyTaskWebhookOutboxFactory,
)
from adcp.decisioning.pg.proposal_store import PgProposalStore
from adcp.decisioning.pg.task_registry import (
PgTaskRegistry,
Expand All @@ -46,6 +54,12 @@
__all__ = [
"DEFAULT_TABLE_NAME",
"PG_AVAILABLE",
"LazyProposalStore",
"LazyProposalStoreFactory",
"LazyTaskRegistry",
"LazyTaskRegistryFactory",
"LazyTaskWebhookOutbox",
"LazyTaskWebhookOutboxFactory",
"PgBuyerAgentRegistry",
"PgProposalStore",
"PgTaskRegistry",
Expand Down
Loading
Loading