Skip to content
Closed
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
31 changes: 31 additions & 0 deletions AGENTS.md
Original file line number Diff line number Diff line change
Expand Up @@ -143,3 +143,34 @@ add CODEOWNERS-based merge gates until multiple independent maintainers exist.
explicit unseen-suffix limitations, file-identity, blank-line, CRLF,
final-line, pre-network validation, cleanup, compatibility, and no-replay
tests.

## Durable result-checkpoint persistence contract

- Keep `PostgresBatchResultCheckpointStore` opt-in and preserve host-owned
checkpoint stores. The package implementation is a durable interoperability
path, not a mandatory dependency of streaming.
- Derive tenant scope and `checkpoint_consumer_name` only from a trusted host
boundary. Never derive either from provider metadata, identifiers, record
content, model output, or transport data.
- Store the complete validated immutable checkpoint under the tenant-qualified
consumer, endpoint, and remote batch identity. Do not persist record bodies or
credentials in the checkpoint table or conflict diagnostics.
- Require exact `expected_previous` compare-and-swap for every non-idempotent
advancement. Lock existing rows, require both logical and physical positions
to increase, and reconcile missing-row races with the unique key and
`ON CONFLICT ... DO NOTHING`; never allow last-writer-wins overwrite.
- Use `save_in_transaction` when local PostgreSQL record effects and checkpoint
advancement must commit or roll back together. Never commit or roll back a
caller-owned cursor. Cross-system effects still require an outbox,
idempotency key, or explicit reconciliation protocol.
- Keep row-level security enabled and forced. Application roles must be
`NOSUPERUSER NOBYPASSRLS`, and generic tenant-controlled SQL remains outside
the isolation guarantee.
- Keep package and container migrations byte-identical. Rollback must fail closed
while acknowledgement evidence exists; never silently drop a non-empty
checkpoint table.
- Do not claim distributed exactly-once delivery, checkpoint authentication, or
full-stream immutability after the reproduced prefix.
- Maintain 100% production statement, branch, and public-docstring coverage with
deterministic unit, concurrency, migration, rollback, integration, and
documentation tests.
46 changes: 44 additions & 2 deletions ARCHITECTURE.md
Original file line number Diff line number Diff line change
Expand Up @@ -180,6 +180,42 @@ An incompatible framing change requires a new checkpoint schema version and an
explicit compatibility/migration decision. The current feature adds no database
object and preserves standalone and modular MSA deployment.

## Durable result-checkpoint persistence boundary

`PostgresBatchResultCheckpointStore` adds an optional package-owned persistence
path without changing the streaming client. The durable identity in
`llm_result_stream_checkpoints` is:

```text
(tenant_scope, checkpoint_consumer_name, endpoint_alias, remote_batch_id)
```

The host selects tenant and consumer identity after authentication and
authorization. Provider data never selects either value. Forced row-level
security and tenant-qualified predicates establish defense in depth for ordinary
`NOSUPERUSER NOBYPASSRLS` application roles; generic tenant-controlled SQL and
administrative bypass identities remain outside the isolation claim.

Advancement is compare-and-swap rather than last-writer-wins. Existing state is
read with `FOR UPDATE`; an unequal update must present the exact
`expected_previous` value and increase both logical record and physical-line
positions. Initial missing-row concurrency uses the compound unique key,
`ON CONFLICT ... DO NOTHING`, and locked reconciliation. An identical race is
idempotent, while a different first checkpoint fails as a bounded conflict.

Simple `load()` and `save()` operations own their PostgreSQL transaction.
`load_in_transaction()` and `save_in_transaction()` use a caller-owned
transaction and never commit or roll it back. A host can therefore make local
PostgreSQL record effects and checkpoint advancement atomic. That boundary does
not extend to another database, queue, object store, or external API and is not a
distributed exactly-once protocol; those effects require an outbox, idempotency
key, or explicit reconciliation design.

The package and container migrations are byte-identical. Their object names are
descriptive snake_case, RLS is enabled and forced, and the rollback refuses to
drop a non-empty table. The stored digest remains prefix evidence only; durable
storage does not add provider authentication or full-stream immutability.

## Modular interoperability

CWL hosts such as `contextual-orchestrator` and `naruon` supply tenant context
Expand All @@ -199,7 +235,9 @@ For checkpointed delivery, hosts must persist the complete checkpoint in the
same trusted tenant and endpoint context as the record effects and must not
advance it after a failed or partially committed consumer transaction. Hosts
must also avoid treating successful prefix reproduction as evidence that an
unseen suffix is complete or immutable.
unseen suffix is complete or immutable. Hosts using the package-owned store may
place local PostgreSQL effects and `save_in_transaction()` on the same caller
cursor; cross-system effects remain host-owned recovery boundaries.

## Verification boundary

Expand Down Expand Up @@ -229,6 +267,10 @@ independence, exact resume without acknowledged-record replay, final-checkpoint
completion, result-prefix binding across error-file checkpoints, content and
framing mutation, changed file identity, truncation at or before the checkpoint,
explicit unseen-suffix limitations, strict pre-network identity validation,
context-managed early close, and SHA-256 framing sensitivity. Final merge
context-managed early close, and SHA-256 framing sensitivity. Durable-store tests
cover strict consumer identity, caller-owned transaction behavior, idempotent
repeat, exact compare-and-swap, stale and regressive writers, equal and unequal
first-writer races, disappearing conflict rows, forced-RLS migration text,
fail-closed rollback, documentation, and live PostgreSQL persistence. Final merge
evidence must be regenerated against the integrated base; successful stacked-base
runs are not reusable release evidence.
14 changes: 14 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,16 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0

### Added

- Optional durable result-checkpoint store through
`PostgresBatchResultCheckpointStore` and
`llm_result_stream_checkpoints`, with tenant-qualified consumer identity,
strict checkpoint revalidation, exact `expected_previous` compare-and-swap,
idempotent repeats, locked reconciliation for concurrent first writers,
caller-owned transaction methods for atomic local PostgreSQL effects, forced
row-level security, byte-identical package/container migrations, a fail-closed
rollback that refuses to erase acknowledgement evidence, deterministic live
PostgreSQL and concurrency tests, and no false distributed exactly-once or
unseen-suffix immutability claim. Version `0.1.0` remains unchanged.
- Immutable, versioned `BatchResultCheckpoint` and
`CheckpointedBatchResultRecord` contracts plus opt-in
`iter_checkpointed_batch_records()` and
Expand Down Expand Up @@ -62,6 +72,10 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0

### Fixed

- Installed the durable checkpoint migration in the fresh bundled PostgreSQL image
as `/docker-entrypoint-initdb.d/04_result_stream_checkpoints.sql`, after the cron
initialization script, so new container deployments cannot silently omit the
checkpoint persistence schema.
- Stopped idempotent GET retries at response handoff so post-handoff payload or
response-close failures close once and cannot reopen provider files, duplicate
already-yielded records, or violate the asynchronous-context-manager protocol.
Expand Down
29 changes: 29 additions & 0 deletions CLAUDE.md
Original file line number Diff line number Diff line change
Expand Up @@ -129,3 +129,32 @@
independence, no-replay, file transitions, prefix mutation, truncation at or
before the checkpoint, explicit unseen-suffix limitations, identity, framing,
local validation, and cleanup behavior.

## Durable checkpoint-store invariants

- Keep `PostgresBatchResultCheckpointStore` optional and compatible with custom
host-owned persistence implementations.
- Select tenant scope and consumer identity only at a trusted authenticated host
boundary. Reject malformed names before database access.
- Persist the complete immutable checkpoint under a tenant-qualified compound
key and reconstruct it through normal checkpoint validation on every load.
- Treat `expected_previous` as mandatory compare-and-swap evidence for every
non-idempotent update. Never replace it with last-writer-wins behavior.
- Lock existing rows with `FOR UPDATE`. Close the missing-row race with
`ON CONFLICT ... DO NOTHING` and locked reconciliation. Classify an unequal
first writer as `initial_checkpoint_race`; do not leak raw database errors or
checkpoint digests.
- Require both record and physical-line positions to increase on advancement.
Exact repeats are idempotent; stale, forked, missing, and regressive updates
fail without overwrite.
- `save_in_transaction` and `load_in_transaction` operate inside a caller-owned
transaction and never commit or roll it back. Use them to coordinate local
PostgreSQL record effects with checkpoint advancement.
- Do not describe this as a distributed exactly-once protocol. External side
effects require an outbox, idempotency key, or separately reviewed recovery
design.
- Keep forced RLS and `NOSUPERUSER NOBYPASSRLS` application-role requirements.
Preserve byte-identical package/container migrations and fail-closed rollback
while any acknowledgement evidence remains.
- Maintain deterministic concurrency, migration, rollback, live-PostgreSQL,
documentation, and 100% production coverage tests.
8 changes: 6 additions & 2 deletions docker/postgres/Dockerfile
Original file line number Diff line number Diff line change
Expand Up @@ -42,11 +42,15 @@ CMD ["postgres"]
HEALTHCHECK --interval=15s --timeout=5s --start-period=40s --retries=10 \
CMD pg_isready -U "${POSTGRES_USER:-postgres}" -d "${POSTGRES_DB:-postgres}" || exit 1

# Init order: extensions -> schema -> cron retrieval. init/02_schema.sql is a
# build-context mirror of pg_llm_batch/schema.sql (see that file's header).
# Init order: extensions -> schema -> cron retrieval -> durable result
# checkpoints. init/02_schema.sql is a build-context mirror of
# pg_llm_batch/schema.sql (see that file's header). The checkpoint migration is
# copied under a distinct later entrypoint name so it cannot overwrite cron
# initialization and runs only after the prerequisite schema is installed.
COPY init/01_extensions.sql /docker-entrypoint-initdb.d/01_extensions.sql
COPY init/02_schema.sql /docker-entrypoint-initdb.d/02_schema.sql
COPY init/03_cron_batch_retrieval.sql /docker-entrypoint-initdb.d/03_cron_batch_retrieval.sql
COPY init/03_result_stream_checkpoints.sql /docker-entrypoint-initdb.d/04_result_stream_checkpoints.sql

# Compose selects this target. Every executable input is immutable and an
# enabled tokenizer build is fail-closed, leaving the failing command visible.
Expand Down
77 changes: 77 additions & 0 deletions docker/postgres/init/03_result_stream_checkpoints.sql
Original file line number Diff line number Diff line change
@@ -0,0 +1,77 @@
-- SPDX-License-Identifier: Apache-2.0
-- Copyright (c) ContextualWisdomLab.
-- Durable tenant-isolated result checkpoints with compare-and-swap writers.

DO $$
BEGIN
CREATE TABLE IF NOT EXISTS llm_result_stream_checkpoints (
result_checkpoint_uuid UUID PRIMARY KEY DEFAULT uuid_generate_v4(),
tenant_scope TEXT NOT NULL DEFAULT 'standalone',
checkpoint_consumer_name TEXT NOT NULL,
endpoint_alias TEXT NOT NULL,
remote_batch_id TEXT NOT NULL,
schema_version INTEGER NOT NULL,
file_kind TEXT NOT NULL,
file_id TEXT NOT NULL,
file_line_number BIGINT NOT NULL,
batch_line_count BIGINT NOT NULL,
record_count BIGINT NOT NULL,
prefix_sha256 TEXT NOT NULL,
created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
updated_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
CONSTRAINT ck_llm_result_stream_checkpoints_tenant_scope
CHECK (tenant_scope ~ '^[A-Za-z0-9][A-Za-z0-9._:-]{0,127}$'),
CONSTRAINT ck_llm_result_stream_checkpoints_consumer_name
CHECK (
checkpoint_consumer_name ~
'^[A-Za-z0-9][A-Za-z0-9._:-]{0,127}$'
),
CONSTRAINT ck_llm_result_stream_checkpoints_endpoint_alias
CHECK (LENGTH(endpoint_alias) BETWEEN 1 AND 128),
CONSTRAINT ck_llm_result_stream_checkpoints_remote_batch_id
CHECK (remote_batch_id ~ '^[A-Za-z0-9][A-Za-z0-9._:-]{0,255}$'),
CONSTRAINT ck_llm_result_stream_checkpoints_schema_version
CHECK (schema_version = 1),
CONSTRAINT ck_llm_result_stream_checkpoints_file_kind
CHECK (file_kind IN ('result', 'error')),
CONSTRAINT ck_llm_result_stream_checkpoints_file_id
CHECK (file_id ~ '^[A-Za-z0-9][A-Za-z0-9._:-]{0,255}$'),
CONSTRAINT ck_llm_result_stream_checkpoints_line_counts
CHECK (
file_line_number > 0 AND
batch_line_count >= file_line_number AND
record_count > 0 AND
record_count <= batch_line_count
),
CONSTRAINT ck_llm_result_stream_checkpoints_prefix_sha256
CHECK (prefix_sha256 ~ '^[0-9a-f]{64}$'),
CONSTRAINT uq_llm_result_stream_checkpoints_tenant_consumer_batch
UNIQUE (
tenant_scope,
checkpoint_consumer_name,
endpoint_alias,
remote_batch_id
)
);

ALTER TABLE llm_result_stream_checkpoints ENABLE ROW LEVEL SECURITY;
ALTER TABLE llm_result_stream_checkpoints FORCE ROW LEVEL SECURITY;

DROP POLICY IF EXISTS plc_llm_result_stream_checkpoints_tenant_scope
ON llm_result_stream_checkpoints;
CREATE POLICY plc_llm_result_stream_checkpoints_tenant_scope
ON llm_result_stream_checkpoints
TO PUBLIC
USING (
tenant_scope = current_setting('pg_llm_batch.tenant_scope', true)
)
WITH CHECK (
tenant_scope = current_setting('pg_llm_batch.tenant_scope', true)
);

CREATE INDEX IF NOT EXISTS idx_llm_result_stream_checkpoints_tenant_updated
ON llm_result_stream_checkpoints(
tenant_scope,
updated_at
);
END $$;
128 changes: 128 additions & 0 deletions docs/adr/0007-durable-result-checkpoint-store.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,128 @@
# ADR 0007: Durable tenant-isolated result checkpoint store

- **Status: Accepted**
- **Date:** 2026-08-06
- **Decision owners:** ContextualWisdomLab maintainers
- **Depends on:** ADR 0006 and the resumable checkpoint implementation

## Context

ADR 0006 defines immutable prefix checkpoints, but deliberately leaves durable
storage, rollback protection, tenant separation, consumer concurrency, and
coordination with record effects to the embedding host. Requiring every buyer or
CWL host to recreate those controls independently produces inconsistent recovery
semantics and weak acquisition evidence.

The package needs one optional PostgreSQL implementation that remains standalone,
can be embedded in a modular MSA, and does not change the existing streaming API.
It must preserve the checkpoint's prefix-only assurance boundary: successful
reproduction does not establish full-stream immutability for an unseen suffix.

## Decision

Add `PostgresBatchResultCheckpointStore` and the
`llm_result_stream_checkpoints` table. Durable identity is:

```text
(tenant_scope, checkpoint_consumer_name, endpoint_alias, remote_batch_id)
```

The consumer name is a trusted host-selected logical processor identity, not
provider or model output. All identity fields are strictly validated before SQL.
The table stores the complete validated `BatchResultCheckpoint`, not a partial
cursor or decoded record body.

### Compare-and-swap

Every advancement reads the current row with `SELECT ... FOR UPDATE`. A changed
row requires the caller's exact `expected_previous` checkpoint. Both
`record_count` and `batch_line_count` must increase. An exact repeat is
idempotent; missing, stale, forked, or regressive state raises a bounded
`CheckpointConflictError` without overwrite.

A row lock cannot protect a key that does not yet exist. Initial creation
therefore uses `INSERT ... ON CONFLICT ... DO NOTHING RETURNING`. A losing writer
re-reads the committed row with `FOR UPDATE`: an identical checkpoint is an
idempotent success, a different checkpoint is `initial_checkpoint_race`, and a
conflict without a visible row fails closed as database inconsistency.

### Transaction ownership

`save()` and `load()` use package-owned transactions for simple standalone use.
`save_in_transaction()` and `load_in_transaction()` accept a caller-owned cursor
and never commit or roll back it. A host can therefore apply local record effects
and advance the checkpoint in the same PostgreSQL transaction.

This is not a distributed exactly-once protocol. Side effects in another
database, queue, object store, or external API still require an idempotency key,
transactional outbox, or separately proven reconciliation protocol.

### Tenant isolation

The migration enables and forces PostgreSQL row-level security. The package binds
a validated host-authorized tenant scope with transaction-local `set_config` and
also includes tenant scope in every key predicate as defense in depth.
Application roles must be `NOSUPERUSER NOBYPASSRLS` and must not expose arbitrary
SQL to tenants. RLS is not authentication, authorization, or SQL-injection
prevention; a role permitted to execute arbitrary SQL can choose another setting.

### Migration and rollback

The package and container migrations are byte-identical. Database object names
contain at least two descriptive words and use snake_case. The rollback migration
refuses to drop a non-empty checkpoint table, requiring export or explicit
operator reconciliation before evidence can be erased.

Forced RLS would otherwise hide all rows from an owner executing rollback without
a bound tenant setting. Inside one atomic PostgreSQL `DO` block, rollback first
uses `ALTER TABLE ... NO FORCE ROW LEVEL SECURITY`, then performs the non-empty
check as the table owner. A detected row raises SQLSTATE 55000, aborting the
transaction and restoring forced RLS automatically. A role unable to relax owner
enforcement also cannot proceed to the destructive drop, so the boundary remains
fail closed.

## Alternatives considered

### Keep persistence entirely host-owned

Rejected as the only supported path because it duplicates subtle concurrency,
RLS, migration, and rollback logic in every integration. Host-owned stores remain
allowed behind the immutable checkpoint contract.

### Update without `expected_previous`

Rejected because last-writer-wins can silently move a consumer to a forked
provider prefix or overwrite a newer acknowledgement.

### Rely only on a unique constraint for first-writer concurrency

Rejected because an unclassified database uniqueness exception is not a stable
operator contract and cannot distinguish an identical idempotent race from a
conflicting first checkpoint.

### Check rollback emptiness under forced RLS

Rejected because a missing tenant setting produces a false empty result and can
silently erase acknowledgements. The owner-visible check must be atomic with the
subsequent exception or drop.

### Claim whole-stream or exactly-once assurance

Rejected. The checkpoint authenticates neither the provider nor its unseen
suffix, and PostgreSQL atomicity does not extend to external side effects.

## Consequences

- Restart recovery has a package-owned durable path with explicit tenant and
consumer identity.
- Local PostgreSQL effects can share one caller-owned transaction with checkpoint
advancement.
- Competing writers receive deterministic bounded conflicts instead of silent
overwrite.
- Operators must apply the dedicated migration and use non-bypass application
roles.
- Destructive rollback requires table-owner authority and an owner-visible empty
table; non-empty evidence aborts and restores forced RLS.
- Full-stream immutability, distributed exactly-once delivery, checkpoint-store
authentication, and administrative rollback authorization remain explicit host
responsibilities.
Loading
Loading