From 6b2e19e3503ef3671aa2d2f974df7561ba1e1ca0 Mon Sep 17 00:00:00 2001 From: Seongho Bae Date: Sun, 9 Aug 2026 08:48:34 +0900 Subject: [PATCH] chore(stack): replay durable checkpoints onto current base --- AGENTS.md | 31 ++ ARCHITECTURE.md | 46 +- CHANGELOG.md | 14 + CLAUDE.md | 29 + docker/postgres/Dockerfile | 8 +- .../init/03_result_stream_checkpoints.sql | 77 +++ .../0007-durable-result-checkpoint-store.md | 128 +++++ .../durable-result-checkpoint-store.md | 111 ++++ .../postgresql-checkpoint-counter-bounds.md | 39 ++ docs/result-streaming.md | 103 +++- pg_llm_batch/__init__.py | 11 + pg_llm_batch/checkpoint_store.py | 403 ++++++++++++++ .../0007_result_stream_checkpoints.sql | 77 +++ .../0007_result_stream_checkpoints.sql | 23 + tests/test_checkpoint_store.py | 498 ++++++++++++++++++ ...checkpoint_store_container_installation.py | 59 +++ tests/test_checkpoint_store_documentation.py | 72 +++ tests/test_checkpoint_store_integration.py | 74 +++ .../test_checkpoint_store_postgres_bounds.py | 123 +++++ tests/test_checkpoint_store_schema.py | 47 ++ 20 files changed, 1959 insertions(+), 14 deletions(-) create mode 100644 docker/postgres/init/03_result_stream_checkpoints.sql create mode 100644 docs/adr/0007-durable-result-checkpoint-store.md create mode 100644 docs/doctoring/durable-result-checkpoint-store.md create mode 100644 docs/doctoring/postgresql-checkpoint-counter-bounds.md create mode 100644 pg_llm_batch/checkpoint_store.py create mode 100644 pg_llm_batch/migrations/0007_result_stream_checkpoints.sql create mode 100644 pg_llm_batch/migrations/rollback/0007_result_stream_checkpoints.sql create mode 100644 tests/test_checkpoint_store.py create mode 100644 tests/test_checkpoint_store_container_installation.py create mode 100644 tests/test_checkpoint_store_documentation.py create mode 100644 tests/test_checkpoint_store_integration.py create mode 100644 tests/test_checkpoint_store_postgres_bounds.py create mode 100644 tests/test_checkpoint_store_schema.py diff --git a/AGENTS.md b/AGENTS.md index 4312a2109..93cb54e1f 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -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. diff --git a/ARCHITECTURE.md b/ARCHITECTURE.md index 3372b3d7b..f265ff06e 100644 --- a/ARCHITECTURE.md +++ b/ARCHITECTURE.md @@ -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 @@ -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 @@ -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. diff --git a/CHANGELOG.md b/CHANGELOG.md index d550f0f78..6a83f1f40 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -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 @@ -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. diff --git a/CLAUDE.md b/CLAUDE.md index 4f8d8843e..6e66882af 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -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. diff --git a/docker/postgres/Dockerfile b/docker/postgres/Dockerfile index a1a7ae969..8f77cef53 100644 --- a/docker/postgres/Dockerfile +++ b/docker/postgres/Dockerfile @@ -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. diff --git a/docker/postgres/init/03_result_stream_checkpoints.sql b/docker/postgres/init/03_result_stream_checkpoints.sql new file mode 100644 index 000000000..f136259c6 --- /dev/null +++ b/docker/postgres/init/03_result_stream_checkpoints.sql @@ -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 $$; diff --git a/docs/adr/0007-durable-result-checkpoint-store.md b/docs/adr/0007-durable-result-checkpoint-store.md new file mode 100644 index 000000000..1348a5210 --- /dev/null +++ b/docs/adr/0007-durable-result-checkpoint-store.md @@ -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. diff --git a/docs/doctoring/durable-result-checkpoint-store.md b/docs/doctoring/durable-result-checkpoint-store.md new file mode 100644 index 000000000..07e482d54 --- /dev/null +++ b/docs/doctoring/durable-result-checkpoint-store.md @@ -0,0 +1,111 @@ +# Durable result-checkpoint store assurance record + +## Scope and claim boundary + +This record covers tenant-qualified PostgreSQL persistence for immutable +`BatchResultCheckpoint` values. The implementation provides deterministic +compare-and-swap, idempotent repeat, caller-owned local transaction support, +forced row-level security, bounded conflict diagnostics, and a fail-closed +rollback migration. + +It does **not** claim provider authentication, checkpoint signature or MAC, +full-stream immutability after the reproduced prefix, distributed exactly-once +delivery, or isolation for superusers, `BYPASSRLS` roles, generic tenant-controlled +SQL, or an incorrectly authorized tenant-to-scope mapping. + +## Control rationale + +### Concurrent writer control + +PostgreSQL 18 documents that row-level `FOR UPDATE` locks block competing writers +and lockers until the current transaction ends. Existing-row advancement uses +that lock and exact `expected_previous` equality. Because a missing row cannot be +row-locked, initial creation additionally uses the compound unique key and +`ON CONFLICT ... DO NOTHING`, followed by locked reconciliation. This separates +identical idempotent races from conflicting first acknowledgements. + +### Tenant isolation + +PostgreSQL row security applies a default-deny posture when row-level security is +enabled and no applicable policy permits a row. The migration enables and forces +RLS, while package operations bind the trusted tenant with transaction-local +`set_config` and repeat the tenant key in each predicate. Production application +roles must be `NOSUPERUSER NOBYPASSRLS`. + +NIST SP 800-53 Rev. 5 control families relevant to this slice include AC (Access +Control), AU (Audit and Accountability), CP (Contingency Planning), SC (System and +Communications Protection), and SI (System and Information Integrity). This +mapping is design evidence, not a certification or assertion that the package +alone satisfies an organization's full control implementation. + +### Recovery and rollback + +A checkpoint is stored only as the complete validated immutable value. A caller +may use `save_in_transaction()` so local PostgreSQL record effects and checkpoint +advancement commit or roll back together. The standalone `save()` method owns its +transaction for simpler deployments. External side effects still require a +transactional outbox, stable idempotency key, or explicit reconciliation. + +A non-empty rollback guard cannot query through forced RLS with no tenant setting: +that context sees no rows and could falsely authorize destruction. The rollback +therefore executes as one atomic `DO` block, temporarily applies +`NO FORCE ROW LEVEL SECURITY`, and performs an owner-visible table-wide emptiness +check. If any acknowledgement exists, SQLSTATE 55000 aborts the transaction, so +the owner-enforcement relaxation is rolled back with the failed drop attempt. A +role lacking table-owner authority fails before it can relax RLS or drop the +table. Operators must export, reconcile, or explicitly remove checkpoint evidence +before schema rollback. + +## Threat and failure matrix + +| Threat or failure | Deterministic control | Residual boundary | +|---|---|---| +| Stale writer overwrites newer acknowledgement | `FOR UPDATE` plus exact `expected_previous` comparison | Administrative direct writes remain outside package guarantees | +| Two writers create the first checkpoint | Compound unique key, `ON CONFLICT`, locked reconciliation | Database outage still aborts the operation | +| Duplicate retry of the same acknowledgement | Exact checkpoint equality is idempotent | Duplicate external side effects require host idempotency | +| Regressive logical or physical position | Both `record_count` and `batch_line_count` must increase | A malicious database administrator can alter rows | +| Cross-tenant lookup or write | Forced RLS, transaction-local scope, tenant-qualified predicates | Superuser, `BYPASSRLS`, arbitrary SQL, and bad authorization mapping are excluded | +| Malformed database row | Reconstruct and revalidate `BatchResultCheckpoint`; invalid shape fails closed | Recovery requires operator repair | +| Forced RLS hides rows during rollback | Atomic owner-visible `NO FORCE ROW LEVEL SECURITY` check; non-empty evidence raises SQLSTATE 55000 and restores RLS by rollback | An authorized owner can still explicitly delete evidence before rerunning rollback | +| Provider suffix changes after checkpoint | Explicitly not attested by the prefix digest | Requires provider validator, authenticated digest, or full-stream manifest | +| Side effect and checkpoint split across systems | No false exactly-once claim | Requires outbox/idempotency/reconciliation at the host boundary | + +## Verification evidence + +Deterministic unit tests cover strict consumer and tenant validation, compound-key +SQL parameters, malformed rows, package-owned and caller-owned transaction paths, +idempotent repeats, stale and forked writers, logical and physical regressions, +identical and conflicting initial races, disappearing conflict rows, schema +installation, and 100% production statement and branch coverage. + +Static migration tests require byte-identical package and container SQL, forced +RLS, tenant policy text, bounded digest and position constraints, descriptive +snake_case object names, no `BYPASSRLS`, and an owner-visible destructive rollback +guard ordered before both the evidence check and table drop. A live PostgreSQL +integration test exercises idempotent creation, exact advancement, load, +stale-writer rejection, and cleanup when `PG_LLM_BATCH_TEST_DSN` is set. + +The bundled PostgreSQL image installs that same reviewed checkpoint migration as +`/docker-entrypoint-initdb.d/04_result_stream_checkpoints.sql`, after the cron +initialization script and without reusing another init destination. A permanent +container-installation regression requires both byte identity and this exact +ordered Dockerfile copy, so a fresh bundled PostgreSQL image cannot silently omit +the durable checkpoint schema. + +Final merge evidence must be regenerated on the integrated exact head and base. +A successful stacked-branch run is development evidence only and cannot authorize +release, provenance, or reuse of an older artifact. + +## APA 7th references + +National Institute of Standards and Technology. (2020). *Security and privacy +controls for information systems and organizations* (NIST Special Publication +800-53 Rev. 5). https://doi.org/10.6028/NIST.SP.800-53r5 + +PostgreSQL Global Development Group. (n.d.). *PostgreSQL 18 documentation: +Explicit locking*. Retrieved August 6, 2026, from +https://www.postgresql.org/docs/18/explicit-locking.html + +PostgreSQL Global Development Group. (n.d.). *PostgreSQL 18 documentation: Row +security policies*. Retrieved August 6, 2026, from +https://www.postgresql.org/docs/18/ddl-rowsecurity.html diff --git a/docs/doctoring/postgresql-checkpoint-counter-bounds.md b/docs/doctoring/postgresql-checkpoint-counter-bounds.md new file mode 100644 index 000000000..3c517a8ca --- /dev/null +++ b/docs/doctoring/postgresql-checkpoint-counter-bounds.md @@ -0,0 +1,39 @@ +# PostgreSQL checkpoint counter bounds + +## Decision + +`PostgresBatchResultCheckpointStore` validates every persisted checkpoint counter +before tenant binding or any other SQL statement. `file_line_number`, +`batch_line_count`, and `record_count` must not exceed `9,223,372,036,854,775,807`, +the maximum signed eight-byte integer accepted by PostgreSQL `BIGINT`. + +The general in-memory `BatchResultCheckpoint` remains storage-independent. The +PostgreSQL adapter owns this narrower persistence boundary because other host +stores may support wider integers. Both the candidate checkpoint and +`expected_previous` compare-and-swap evidence are checked before database access. +The exact maximum remains valid; the first larger value fails as a structured +`ValidationError` naming the offending nested field. This prevents a driver or +server numeric-overflow exception from aborting a caller-owned transaction after +business work has already begun. + +## Verification + +Deterministic tests cover overflow through physical-line and record-count paths, +compare-and-swap evidence, the exact legal maximum, and the requirement that no +transaction-local tenant-setting statement or checkpoint query executes for +invalid storage values. Production statement, branch, and public-docstring +coverage remain at 100%. + +## Operational consequence + +A host approaching the signed `BIGINT` ceiling must rotate to a new logical batch +or adopt a separately versioned schema using an explicitly reviewed wider numeric +representation. Silently coercing, saturating, wrapping, or changing the existing +migration column types is prohibited because it would invalidate checkpoint +identity, ordering, and rollback evidence. + +## Reference + +PostgreSQL Global Development Group. (n.d.). *8.1. Numeric types*. In +*PostgreSQL 18 documentation*. Retrieved August 7, 2026, from +https://www.postgresql.org/docs/18/datatype-numeric.html diff --git a/docs/result-streaming.md b/docs/result-streaming.md index 2f6a59041..34e5682b2 100644 --- a/docs/result-streaming.md +++ b/docs/result-streaming.md @@ -99,6 +99,87 @@ mismatch as an operator reconciliation event rather than silently advancing it. See [ADR 0006](adr/0006-resumable-result-checkpoints.md) and the [assurance record](doctoring/resumable-result-checkpoints.md). +## Package-owned durable checkpoint storage + +Apply the dedicated migration after the base package schema: + +```python +from pg_llm_batch import ( + PostgresBatchResultCheckpointStore, + apply_result_checkpoint_schema, +) + +apply_result_checkpoint_schema(dsn) +checkpoint_store = PostgresBatchResultCheckpointStore( + dsn, + tenant_scope="tenant-a", +) +``` + +For simple standalone processing, `load()` and `save()` own their PostgreSQL +transactions: + +```python +resume_after = checkpoint_store.load("invoice-worker", "batch-123", "default") + +async with client.open_checkpointed_batch_records( + "batch-123", + "default", + resume_after=resume_after, +) as records: + async for item in records: + apply_idempotent_record(item.record) + resume_after = checkpoint_store.save( + "invoice-worker", + item.checkpoint, + expected_previous=resume_after, + ) +``` + +When record effects are stored in the same PostgreSQL database, use a +caller-owned transaction so the effect and acknowledgement cannot split: + +```python +import psycopg + +with psycopg.connect(dsn) as connection: + with connection.cursor() as cursor: + apply_record_with_cursor(cursor, item.record) + resume_after = checkpoint_store.save_in_transaction( + cursor, + "invoice-worker", + item.checkpoint, + expected_previous=resume_after, + ) + connection.commit() +``` + +`save_in_transaction()` never commits or rolls back the caller's cursor. An exact +repeat is idempotent. Every different durable row requires the exact +`expected_previous` value and strictly increasing record and physical-line +positions. A stale, forked, regressive, missing, or conflicting first writer +raises `CheckpointConflictError` without overwrite. + +Tenant scope and consumer identity must come from the host's authenticated and +authorized control plane. Production application roles must be +`NOSUPERUSER NOBYPASSRLS`, must not expose arbitrary tenant-controlled SQL, and +must have only the table and schema privileges required by the deployment. +Forced row-level security is defense in depth, not a credential or substitute for +authorization. + +This is not a distributed exactly-once protocol. A queue, another database, +object store, webhook, or provider-side effect cannot share the local PostgreSQL +transaction and still requires a stable idempotency key, transactional outbox, or +operator reconciliation. Durable storage also does not authenticate the +checkpoint or prove full-stream immutability after the reproduced prefix. + +The rollback file +`pg_llm_batch/migrations/rollback/0007_result_stream_checkpoints.sql` refuses to +drop a non-empty table. Export or reconcile acknowledgement evidence before an +operator deliberately removes it. See +[ADR 0007](adr/0007-durable-result-checkpoint-store.md) and the +[durable-store assurance record](doctoring/durable-result-checkpoint-store.md). + ## Resource and trust boundaries - The inherited total decoded-byte limit is enforced independently for each @@ -148,16 +229,18 @@ verification. The streaming client subclasses `BatchAPIClient`, so credentials, gateway URL validation, timeouts, pre-handoff retry policy, and session lifecycle remain -identical. The checkpoint feature does not require PostgreSQL schema changes or -another CWL service. Embedding hosts may store each record and checkpoint in -their own durable queue, tenant-qualified database, transactional outbox, or -bounded transformation pipeline. - -The existing OpenTelemetry subclass does not automatically wrap this opt-in -iterator. Hosts that need per-record or resume-reconciliation telemetry should -instrument the consumer boundary with low-cardinality attributes and must not -attach provider identifiers, checkpoint digests, prompts, response bodies, or -model output. +identical. The package-owned checkpoint store is optional; custom host stores +remain supported, and the streaming client itself does not open the checkpoint +table. + +Embedding hosts may store each record and checkpoint in their own durable queue, +tenant-qualified database, transactional outbox, or bounded transformation +pipeline. The existing OpenTelemetry subclass does not automatically wrap this +opt-in iterator or checkpoint store. Hosts that need per-record, +resume-reconciliation, or checkpoint-conflict telemetry should instrument the +consumer boundary with low-cardinality attributes and must not attach provider +identifiers, checkpoint digests, prompts, response bodies, model output, or raw +database exception text. ## References diff --git a/pg_llm_batch/__init__.py b/pg_llm_batch/__init__.py index bf9743816..74385e196 100644 --- a/pg_llm_batch/__init__.py +++ b/pg_llm_batch/__init__.py @@ -8,6 +8,7 @@ BatchAPIClient -- submit, poll, and retrieve StreamingBatchAPIClient -- bounded incremental result records BatchResultCheckpoint -- host-persistable resume evidence + PostgresBatchResultCheckpointStore -- tenant-isolated durable checkpoints DurableBatchAPIClient -- standalone durable lifecycle state TenantDurableBatchAPIClient -- tenant-isolated lifecycle state PostgresConfigStore, SecretStore -- database configuration and secrets @@ -20,6 +21,12 @@ GatewayCredentials, config_credentials_provider, ) +from .checkpoint_store import ( + CheckpointConflictError, + PostgresBatchResultCheckpointStore, + apply_result_checkpoint_schema, + validate_checkpoint_consumer_name, +) from .config import PostgresConfigStore, SecretStore, get_config_store from .db import ( DEFAULT_TENANT_SCOPE, @@ -57,6 +64,10 @@ "BatchResultRecord", "BatchResultCheckpoint", "CheckpointedBatchResultRecord", + "PostgresBatchResultCheckpointStore", + "CheckpointConflictError", + "apply_result_checkpoint_schema", + "validate_checkpoint_consumer_name", "DurableBatchAPIClient", "TenantDurableBatchAPIClient", "DEFAULT_TENANT_SCOPE", diff --git a/pg_llm_batch/checkpoint_store.py b/pg_llm_batch/checkpoint_store.py new file mode 100644 index 000000000..8d9da18ee --- /dev/null +++ b/pg_llm_batch/checkpoint_store.py @@ -0,0 +1,403 @@ +# SPDX-License-Identifier: Apache-2.0 +# Copyright (c) ContextualWisdomLab. +"""Tenant-isolated PostgreSQL persistence for resumable result checkpoints.""" + +from __future__ import annotations + +import re +from pathlib import Path +from typing import Any, Optional + +from .db import ( + DEFAULT_TENANT_SCOPE, + _require_psycopg, + _set_transaction_tenant_scope, + psycopg, + validate_endpoint_alias, + validate_remote_resource_id, + validate_tenant_scope, +) +from .exceptions import PgLlmBatchError, ValidationError +from .result_streaming import BatchResultCheckpoint + +MIGRATION_PATH = ( + Path(__file__).with_name("migrations") / "0007_result_stream_checkpoints.sql" +) +MAX_CHECKPOINT_CONSUMER_CHARACTERS = 128 +POSTGRES_BIGINT_MAX = (1 << 63) - 1 +CHECKPOINT_CONSUMER_PATTERN = re.compile( + rf"[A-Za-z0-9][A-Za-z0-9._:-]{{0,{MAX_CHECKPOINT_CONSUMER_CHARACTERS - 1}}}\Z" +) +_CHECKPOINT_COLUMNS = ( + "schema_version, remote_batch_id, endpoint_alias, file_kind, file_id, " + "file_line_number, batch_line_count, record_count, prefix_sha256" +) +_POSTGRES_BIGINT_CHECKPOINT_FIELDS = ( + "file_line_number", + "batch_line_count", + "record_count", +) + + +class CheckpointConflictError(PgLlmBatchError): + """Raised when a durable checkpoint compare-and-swap cannot proceed safely.""" + + def __init__(self, consumer_name: str, batch_id: str, reason: str) -> None: + """Describe one bounded durable checkpoint concurrency conflict.""" + super().__init__( + message="Result checkpoint update conflicted with durable state", + error_code="CHECKPOINT_CONFLICT", + details={ + "consumer_name": consumer_name, + "batch_id": batch_id, + "reason": reason, + }, + ) + self.consumer_name = consumer_name + self.batch_id = batch_id + self.reason = reason + + +def validate_checkpoint_consumer_name(value: Any) -> str: + """Validate one host-selected checkpoint consumer name without coercion.""" + if ( + not isinstance(value, str) + or CHECKPOINT_CONSUMER_PATTERN.fullmatch(value) is None + ): + raise ValidationError( + field="consumer_name", + value=value, + reason=( + "must be 1-128 ASCII characters beginning with an alphanumeric " + "character and containing only letters, digits, dot, underscore, " + "colon, or hyphen" + ), + ) + return value + + +def _validated_checkpoint(value: Any, field: str) -> BatchResultCheckpoint: + """Require one immutable checkpoint whose counters fit PostgreSQL storage.""" + if not isinstance(value, BatchResultCheckpoint): + raise ValidationError( + field=field, + value=value, + reason="must be a BatchResultCheckpoint", + ) + for checkpoint_field in _POSTGRES_BIGINT_CHECKPOINT_FIELDS: + count = getattr(value, checkpoint_field) + if count > POSTGRES_BIGINT_MAX: + raise ValidationError( + field=f"{field}.{checkpoint_field}", + value=count, + reason=( + "must be no greater than PostgreSQL BIGINT maximum " + f"{POSTGRES_BIGINT_MAX}" + ), + ) + return value + + +def _validated_exact_endpoint_alias(value: Any) -> str: + """Require one endpoint alias that is already in canonical form.""" + try: + normalized = validate_endpoint_alias(value) + except ValidationError as exc: + raise ValidationError( + field="endpoint_alias", + value=value, + reason="must be a supported endpoint alias", + ) from exc + if normalized != value: + raise ValidationError( + field="endpoint_alias", + value=value, + reason="must already be normalized without surrounding whitespace", + ) + return normalized + + +def _validated_batch_id(value: Any) -> str: + """Require one supported provider batch identifier.""" + try: + return validate_remote_resource_id(value, "batch_id") + except ValidationError as exc: + raise ValidationError( + field="batch_id", + value=value, + reason="must be a supported provider identifier", + ) from exc + + +def _checkpoint_from_row(row: Any) -> BatchResultCheckpoint: + """Revalidate one database row as an immutable checkpoint.""" + if not isinstance(row, (tuple, list)) or len(row) != 9: + raise RuntimeError("result checkpoint row has an invalid shape") + return BatchResultCheckpoint( + schema_version=row[0], + batch_id=row[1], + endpoint_alias=row[2], + file_kind=row[3], + file_id=row[4], + file_line_number=row[5], + batch_line_count=row[6], + record_count=row[7], + prefix_sha256=row[8], + ) + + +def _checkpoint_values(checkpoint: BatchResultCheckpoint) -> tuple[Any, ...]: + """Return one stable SQL value tuple for a validated checkpoint.""" + return ( + checkpoint.schema_version, + checkpoint.file_kind, + checkpoint.file_id, + checkpoint.file_line_number, + checkpoint.batch_line_count, + checkpoint.record_count, + checkpoint.prefix_sha256, + ) + + +def apply_result_checkpoint_schema( + postgres_dsn: str, + migration_path: Optional[str] = None, +) -> None: + """Apply the idempotent durable result-checkpoint migration.""" + _require_psycopg() + path = Path(migration_path) if migration_path else MIGRATION_PATH + sql = path.read_text(encoding="utf-8") + with psycopg.connect(postgres_dsn) as conn: + with conn.cursor() as cur: + cur.execute(sql) + conn.commit() + + +class PostgresBatchResultCheckpointStore: + """Persist tenant-qualified streaming checkpoints with compare-and-swap safety.""" + + def __init__( + self, + postgres_dsn: str, + *, + tenant_scope: str = DEFAULT_TENANT_SCOPE, + ) -> None: + """Bind one database and trusted local tenant scope to the store.""" + self.postgres_dsn = postgres_dsn + try: + self.tenant_scope = validate_tenant_scope(tenant_scope) + except ValidationError as exc: + raise ValidationError( + field="tenant_scope", + value=tenant_scope, + reason="must be a supported trusted tenant scope", + ) from exc + + def load( + self, + consumer_name: str, + batch_id: str, + endpoint_alias: str, + ) -> Optional[BatchResultCheckpoint]: + """Load the current checkpoint in one package-owned transaction.""" + _require_psycopg() + with psycopg.connect(self.postgres_dsn) as conn: + with conn.cursor() as cur: + return self.load_in_transaction( + cur, + consumer_name, + batch_id, + endpoint_alias, + ) + + def load_in_transaction( + self, + cursor: Any, + consumer_name: str, + batch_id: str, + endpoint_alias: str, + ) -> Optional[BatchResultCheckpoint]: + """Load a checkpoint through a caller-owned database transaction. + + The caller owns commit and rollback. This method binds the store's trusted + tenant scope with transaction-local PostgreSQL configuration before the + tenant-qualified read. + """ + consumer = validate_checkpoint_consumer_name(consumer_name) + remote_batch_id = _validated_batch_id(batch_id) + alias = _validated_exact_endpoint_alias(endpoint_alias) + _set_transaction_tenant_scope(cursor, self.tenant_scope) + cursor.execute( + f"SELECT {_CHECKPOINT_COLUMNS} " + "FROM llm_result_stream_checkpoints " + "WHERE tenant_scope = %s " + "AND checkpoint_consumer_name = %s " + "AND endpoint_alias = %s " + "AND remote_batch_id = %s", + (self.tenant_scope, consumer, alias, remote_batch_id), + ) + row = cursor.fetchone() + return None if row is None else _checkpoint_from_row(row) + + def save( + self, + consumer_name: str, + checkpoint: BatchResultCheckpoint, + *, + expected_previous: Optional[BatchResultCheckpoint] = None, + ) -> BatchResultCheckpoint: + """Create or advance a checkpoint in one package-owned transaction.""" + _require_psycopg() + with psycopg.connect(self.postgres_dsn) as conn: + with conn.cursor() as cur: + saved = self.save_in_transaction( + cur, + consumer_name, + checkpoint, + expected_previous=expected_previous, + ) + conn.commit() + return saved + + def save_in_transaction( + self, + cursor: Any, + consumer_name: str, + checkpoint: BatchResultCheckpoint, + *, + expected_previous: Optional[BatchResultCheckpoint] = None, + ) -> BatchResultCheckpoint: + """Compare and swap a checkpoint in a caller-owned transaction. + + An identical repeat is idempotent. A different existing row requires the + caller's exact previously loaded checkpoint and strictly increasing + record and physical-line counts. Missing, stale, forked, or regressive + updates fail without overwriting durable evidence. The caller owns commit + and rollback, enabling local business effects and checkpoint advancement + to share one PostgreSQL transaction. + """ + consumer = validate_checkpoint_consumer_name(consumer_name) + current_candidate = _validated_checkpoint(checkpoint, "checkpoint") + previous_candidate = ( + None + if expected_previous is None + else _validated_checkpoint(expected_previous, "expected_previous") + ) + if previous_candidate is not None and ( + previous_candidate.batch_id != current_candidate.batch_id + or previous_candidate.endpoint_alias != current_candidate.endpoint_alias + ): + raise ValidationError( + field="expected_previous", + value=expected_previous, + reason="must identify the same batch and endpoint as checkpoint", + ) + + _set_transaction_tenant_scope(cursor, self.tenant_scope) + cursor.execute( + f"SELECT {_CHECKPOINT_COLUMNS} " + "FROM llm_result_stream_checkpoints " + "WHERE tenant_scope = %s " + "AND checkpoint_consumer_name = %s " + "AND endpoint_alias = %s " + "AND remote_batch_id = %s FOR UPDATE", + ( + self.tenant_scope, + consumer, + current_candidate.endpoint_alias, + current_candidate.batch_id, + ), + ) + row = cursor.fetchone() + if row is None: + if previous_candidate is not None: + raise CheckpointConflictError( + consumer, + current_candidate.batch_id, + "expected_previous_missing", + ) + cursor.execute( + "INSERT INTO llm_result_stream_checkpoints (" + "tenant_scope, checkpoint_consumer_name, endpoint_alias, " + "remote_batch_id, schema_version, file_kind, file_id, " + "file_line_number, batch_line_count, record_count, " + "prefix_sha256) VALUES (" + "%s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s) " + "ON CONFLICT (tenant_scope, checkpoint_consumer_name, " + "endpoint_alias, remote_batch_id) DO NOTHING " + "RETURNING remote_batch_id", + ( + self.tenant_scope, + consumer, + current_candidate.endpoint_alias, + current_candidate.batch_id, + *_checkpoint_values(current_candidate), + ), + ) + if cursor.fetchone() is None: + cursor.execute( + f"SELECT {_CHECKPOINT_COLUMNS} " + "FROM llm_result_stream_checkpoints " + "WHERE tenant_scope = %s " + "AND checkpoint_consumer_name = %s " + "AND endpoint_alias = %s " + "AND remote_batch_id = %s FOR UPDATE", + ( + self.tenant_scope, + consumer, + current_candidate.endpoint_alias, + current_candidate.batch_id, + ), + ) + concurrent_row = cursor.fetchone() + if concurrent_row is None: + raise RuntimeError( + "result checkpoint insert conflict row disappeared" + ) + concurrent = _checkpoint_from_row(concurrent_row) + if concurrent == current_candidate: + return concurrent + raise CheckpointConflictError( + consumer, + current_candidate.batch_id, + "initial_checkpoint_race", + ) + return current_candidate + + durable = _checkpoint_from_row(row) + if durable == current_candidate: + return durable + if previous_candidate is None or durable != previous_candidate: + raise CheckpointConflictError( + consumer, + current_candidate.batch_id, + "expected_previous_stale", + ) + if ( + current_candidate.record_count <= durable.record_count + or current_candidate.batch_line_count <= durable.batch_line_count + ): + raise CheckpointConflictError( + consumer, + current_candidate.batch_id, + "checkpoint_regression", + ) + cursor.execute( + "UPDATE llm_result_stream_checkpoints SET " + "schema_version = %s, file_kind = %s, file_id = %s, " + "file_line_number = %s, batch_line_count = %s, " + "record_count = %s, prefix_sha256 = %s, " + "updated_at = NOW() " + "WHERE tenant_scope = %s " + "AND checkpoint_consumer_name = %s " + "AND endpoint_alias = %s " + "AND remote_batch_id = %s", + ( + *_checkpoint_values(current_candidate), + self.tenant_scope, + consumer, + current_candidate.endpoint_alias, + current_candidate.batch_id, + ), + ) + return current_candidate diff --git a/pg_llm_batch/migrations/0007_result_stream_checkpoints.sql b/pg_llm_batch/migrations/0007_result_stream_checkpoints.sql new file mode 100644 index 000000000..f136259c6 --- /dev/null +++ b/pg_llm_batch/migrations/0007_result_stream_checkpoints.sql @@ -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 $$; diff --git a/pg_llm_batch/migrations/rollback/0007_result_stream_checkpoints.sql b/pg_llm_batch/migrations/rollback/0007_result_stream_checkpoints.sql new file mode 100644 index 000000000..3678aff8f --- /dev/null +++ b/pg_llm_batch/migrations/rollback/0007_result_stream_checkpoints.sql @@ -0,0 +1,23 @@ +-- SPDX-License-Identifier: Apache-2.0 +-- Refuse destructive rollback while durable acknowledgement evidence exists. + +DO $$ +BEGIN + IF to_regclass('llm_result_stream_checkpoints') IS NOT NULL THEN + -- FORCE RLS would hide every row when no tenant setting is bound. A role + -- capable of dropping this table must first become subject to the normal + -- owner-bypass rule so the emptiness check observes every tenant. The DO + -- block is one transaction: a raised exception rolls this relaxation back. + ALTER TABLE llm_result_stream_checkpoints NO FORCE ROW LEVEL SECURITY; + + IF EXISTS ( + SELECT 1 FROM llm_result_stream_checkpoints LIMIT 1 + ) THEN + RAISE EXCEPTION + 'Refusing to drop non-empty llm_result_stream_checkpoints; export or reconcile checkpoints first' + USING ERRCODE = '55000'; + END IF; + END IF; + + DROP TABLE IF EXISTS llm_result_stream_checkpoints; +END $$; diff --git a/tests/test_checkpoint_store.py b/tests/test_checkpoint_store.py new file mode 100644 index 000000000..92c87be42 --- /dev/null +++ b/tests/test_checkpoint_store.py @@ -0,0 +1,498 @@ +# SPDX-License-Identifier: Apache-2.0 +"""Tests for durable tenant-isolated result-checkpoint persistence.""" + +from __future__ import annotations + +from dataclasses import replace +from typing import Any + +import pytest + +import pg_llm_batch.checkpoint_store as checkpoint_store +from pg_llm_batch import ( + CheckpointConflictError, + PostgresBatchResultCheckpointStore, + apply_result_checkpoint_schema, + validate_checkpoint_consumer_name, +) +from pg_llm_batch.exceptions import ValidationError +from pg_llm_batch.result_streaming import BatchResultCheckpoint + + +def checkpoint( + *, + record_count: int = 1, + batch_line_count: int = 2, + digest: str = "a" * 64, +) -> BatchResultCheckpoint: + """Build one valid immutable checkpoint for persistence tests.""" + return BatchResultCheckpoint( + schema_version=1, + batch_id="batch-1", + endpoint_alias="default", + file_kind="result", + file_id="file-1", + file_line_number=batch_line_count, + batch_line_count=batch_line_count, + record_count=record_count, + prefix_sha256=digest, + ) + + +class FakeCursor: + """Execute the checkpoint store's bounded SQL contract in memory.""" + + def __init__(self, database: "FakeDatabase") -> None: + self.database = database + self.result: Any = None + + def __enter__(self) -> "FakeCursor": + """Enter the fake cursor context.""" + return self + + def __exit__(self, *_exc: Any) -> None: + """Exit the fake cursor context.""" + return None + + def execute(self, sql: str, params: tuple[Any, ...] | None = None) -> None: + """Execute one expected tenant-qualified checkpoint statement.""" + normalized = " ".join(sql.split()) + parameters = params or () + self.database.calls.append((normalized, parameters)) + if normalized.startswith("SELECT set_config"): + self.result = (parameters[0],) + return + if normalized.startswith("SELECT schema_version"): + self.result = self.database.rows.get(parameters) + return + if normalized.startswith("INSERT INTO"): + ( + tenant, + consumer, + alias, + batch_id, + schema, + kind, + file_id, + file_line, + batch_lines, + records, + digest, + ) = parameters + key = (tenant, consumer, alias, batch_id) + if self.database.insert_conflict_without_row: + self.database.insert_conflict_without_row = False + self.result = None + return + if self.database.insert_race_row is not None: + self.database.rows[key] = self.database.insert_race_row + self.database.insert_race_row = None + self.result = None + return + if key in self.database.rows: + self.result = None + return + self.database.rows[key] = ( + schema, + batch_id, + alias, + kind, + file_id, + file_line, + batch_lines, + records, + digest, + ) + self.result = (batch_id,) + return + if normalized.startswith("UPDATE"): + ( + schema, + kind, + file_id, + file_line, + batch_lines, + records, + digest, + tenant, + consumer, + alias, + batch_id, + ) = parameters + self.database.rows[(tenant, consumer, alias, batch_id)] = ( + schema, + batch_id, + alias, + kind, + file_id, + file_line, + batch_lines, + records, + digest, + ) + self.result = None + return + if not parameters: + self.result = None + return + raise AssertionError(normalized) + + def fetchone(self) -> Any: + """Return the result from the preceding fake statement.""" + return self.result + + +class FakeConnection: + """Provide cursor and commit accounting for one fake database.""" + + def __init__(self, database: "FakeDatabase") -> None: + self.database = database + + def __enter__(self) -> "FakeConnection": + """Enter the fake connection context.""" + return self + + def __exit__(self, *_exc: Any) -> None: + """Exit the fake connection context.""" + return None + + def cursor(self) -> FakeCursor: + """Create one fake cursor.""" + return FakeCursor(self.database) + + def commit(self) -> None: + """Record one explicit commit.""" + self.database.commits += 1 + + +class FakePsycopg: + """Connect checkpoint store calls to one in-memory fake database.""" + + def __init__(self, database: "FakeDatabase") -> None: + self.database = database + + def connect(self, dsn: str) -> FakeConnection: + """Connect to one deterministic fake database.""" + self.database.dsns.append(dsn) + return FakeConnection(self.database) + + +class FakeDatabase: + """Hold deterministic rows, statements, and race simulation state.""" + + def __init__(self) -> None: + self.rows: dict[tuple[Any, ...], tuple[Any, ...]] = {} + self.calls: list[tuple[str, tuple[Any, ...]]] = [] + self.dsns: list[str] = [] + self.commits = 0 + self.insert_race_row: tuple[Any, ...] | None = None + self.insert_conflict_without_row = False + + +@pytest.fixture +def database(monkeypatch: pytest.MonkeyPatch) -> FakeDatabase: + """Install one deterministic psycopg replacement for each test.""" + fake_database = FakeDatabase() + monkeypatch.setattr(checkpoint_store, "psycopg", FakePsycopg(fake_database)) + monkeypatch.setattr(checkpoint_store, "_require_psycopg", lambda: None) + return fake_database + + +def test_consumer_name_validation_is_strict_and_noncoercive() -> None: + """Consumer names accept only the bounded durable-key grammar.""" + assert validate_checkpoint_consumer_name("worker.primary") == "worker.primary" + for value in (None, 1, "", " leading", "bad/name", "a" * 129): + with pytest.raises(ValidationError): + validate_checkpoint_consumer_name(value) + + +def test_store_validates_tenant_before_database_access() -> None: + """An invalid trusted tenant never reaches PostgreSQL.""" + with pytest.raises(ValidationError): + PostgresBatchResultCheckpointStore("postgresql://unit", tenant_scope=" bad") + + +def test_load_returns_none_and_uses_tenant_qualified_key( + database: FakeDatabase, +) -> None: + """Package-owned loads bind tenant context and every compound-key field.""" + store = PostgresBatchResultCheckpointStore( + "postgresql://unit", + tenant_scope="tenant-a", + ) + assert store.load("worker-a", "batch-1", "default") is None + assert database.dsns == ["postgresql://unit"] + assert database.calls[0][1] == ("tenant-a",) + assert database.calls[1][1] == ( + "tenant-a", + "worker-a", + "default", + "batch-1", + ) + + +def test_load_in_transaction_uses_caller_cursor_without_commit( + database: FakeDatabase, +) -> None: + """Caller-owned reads remain inside the caller's transaction boundary.""" + database.rows[("tenant-a", "worker-a", "default", "batch-1")] = ( + 1, + "batch-1", + "default", + "result", + "file-1", + 2, + 2, + 1, + "a" * 64, + ) + store = PostgresBatchResultCheckpointStore( + "postgresql://unit", + tenant_scope="tenant-a", + ) + cursor = FakeCursor(database) + + assert ( + store.load_in_transaction(cursor, "worker-a", "batch-1", "default") + == checkpoint() + ) + assert database.dsns == [] + assert database.commits == 0 + + +def test_load_revalidates_database_rows(database: FakeDatabase) -> None: + """Malformed durable rows fail closed before becoming public checkpoints.""" + database.rows[("tenant-a", "worker-a", "default", "batch-1")] = ( + 1, + "batch-1", + "default", + "result", + "file-1", + 2, + 2, + 1, + "a" * 64, + ) + store = PostgresBatchResultCheckpointStore( + "postgresql://unit", + tenant_scope="tenant-a", + ) + assert store.load("worker-a", "batch-1", "default") == checkpoint() + database.rows[("tenant-a", "worker-a", "default", "batch-1")] = (1, 2) + with pytest.raises(RuntimeError, match="invalid shape"): + store.load("worker-a", "batch-1", "default") + + +def test_load_rejects_noncanonical_identifiers_before_database( + database: FakeDatabase, +) -> None: + """Invalid batch and endpoint identities fail before SQL execution.""" + store = PostgresBatchResultCheckpointStore("postgresql://unit") + with pytest.raises(ValidationError): + store.load("worker-a", "bad/batch", "default") + with pytest.raises(ValidationError): + store.load("worker-a", "batch-1", "") + with pytest.raises(ValidationError): + store.load("worker-a", "batch-1", " default ") + assert database.calls == [] + + +def test_save_creates_and_idempotently_repeats_checkpoint( + database: FakeDatabase, +) -> None: + """Identical package-owned saves write once and remain idempotent.""" + store = PostgresBatchResultCheckpointStore( + "postgresql://unit", + tenant_scope="tenant-a", + ) + first = checkpoint() + assert store.save("worker-a", first) == first + assert store.save("worker-a", first) == first + assert database.commits == 2 + statements = [call[0] for call in database.calls] + assert sum(item.startswith("INSERT") for item in statements) == 1 + assert sum(item.startswith("UPDATE") for item in statements) == 0 + + +def test_save_in_transaction_does_not_commit_caller_work( + database: FakeDatabase, +) -> None: + """Caller-owned saves never commit unrelated local business effects.""" + store = PostgresBatchResultCheckpointStore("postgresql://unit") + cursor = FakeCursor(database) + + assert store.save_in_transaction(cursor, "worker-a", checkpoint()) == checkpoint() + assert database.dsns == [] + assert database.commits == 0 + + +def test_save_handles_same_checkpoint_initial_insert_race( + database: FakeDatabase, +) -> None: + """A concurrent identical first writer remains an idempotent success.""" + first = checkpoint() + database.insert_race_row = ( + first.schema_version, + first.batch_id, + first.endpoint_alias, + first.file_kind, + first.file_id, + first.file_line_number, + first.batch_line_count, + first.record_count, + first.prefix_sha256, + ) + store = PostgresBatchResultCheckpointStore("postgresql://unit") + + assert store.save("worker-a", first) == first + statements = [call[0] for call in database.calls] + assert any("ON CONFLICT" in statement for statement in statements) + assert ( + sum(statement.startswith("SELECT schema_version") for statement in statements) + == 2 + ) + + +def test_save_rejects_disappearing_insert_conflict_row( + database: FakeDatabase, +) -> None: + """A conflict without a visible durable row fails closed.""" + database.insert_conflict_without_row = True + store = PostgresBatchResultCheckpointStore("postgresql://unit") + + with pytest.raises(RuntimeError, match="conflict row disappeared"): + store.save("worker-a", checkpoint()) + + assert database.rows == {} + + +def test_save_fails_closed_on_different_initial_insert_race( + database: FakeDatabase, +) -> None: + """A concurrent different first writer becomes one bounded conflict.""" + database.insert_race_row = ( + 1, + "batch-1", + "default", + "result", + "file-1", + 3, + 3, + 2, + "b" * 64, + ) + store = PostgresBatchResultCheckpointStore("postgresql://unit") + + with pytest.raises(CheckpointConflictError) as raised: + store.save("worker-a", checkpoint()) + + assert raised.value.reason == "initial_checkpoint_race" + assert database.commits == 0 + + +def test_save_requires_missing_expected_checkpoint_to_remain_missing( + database: FakeDatabase, +) -> None: + """A claimed previous checkpoint cannot create a missing durable row.""" + store = PostgresBatchResultCheckpointStore("postgresql://unit") + with pytest.raises(CheckpointConflictError) as raised: + store.save( + "worker-a", + checkpoint(record_count=2, batch_line_count=3), + expected_previous=checkpoint(), + ) + assert raised.value.reason == "expected_previous_missing" + assert raised.value.error_code == "CHECKPOINT_CONFLICT" + assert database.rows == {} + + +def test_save_rejects_stale_or_forked_expected_checkpoint( + database: FakeDatabase, +) -> None: + """Writers cannot overwrite state they did not observe exactly.""" + store = PostgresBatchResultCheckpointStore("postgresql://unit") + first = checkpoint() + store.save("worker-a", first) + advanced = checkpoint(record_count=2, batch_line_count=3, digest="b" * 64) + with pytest.raises(CheckpointConflictError) as missing: + store.save("worker-a", advanced) + assert missing.value.reason == "expected_previous_stale" + fork = replace(first, prefix_sha256="c" * 64) + with pytest.raises(CheckpointConflictError) as stale: + store.save("worker-a", advanced, expected_previous=fork) + assert stale.value.reason == "expected_previous_stale" + stored = database.rows[("standalone", "worker-a", "default", "batch-1")] + assert stored[8] == "a" * 64 + + +def test_save_rejects_regressive_counts(database: FakeDatabase) -> None: + """Checkpoint advancement requires both logical and physical progress.""" + store = PostgresBatchResultCheckpointStore("postgresql://unit") + first = checkpoint(record_count=2, batch_line_count=3) + store.save("worker-a", first) + for candidate in ( + checkpoint(record_count=2, batch_line_count=4, digest="b" * 64), + checkpoint(record_count=3, batch_line_count=3, digest="c" * 64), + ): + with pytest.raises(CheckpointConflictError) as raised: + store.save("worker-a", candidate, expected_previous=first) + assert raised.value.reason == "checkpoint_regression" + + +def test_save_advances_exact_expected_checkpoint(database: FakeDatabase) -> None: + """An exact compare-and-swap advances the tenant-qualified durable row.""" + store = PostgresBatchResultCheckpointStore( + "postgresql://unit", + tenant_scope="tenant-a", + ) + first = checkpoint() + second = checkpoint(record_count=2, batch_line_count=4, digest="b" * 64) + store.save("worker-a", first) + assert store.save("worker-a", second, expected_previous=first) == second + assert database.commits == 2 + assert store.load("worker-a", "batch-1", "default") == second + update = next(call for call in database.calls if call[0].startswith("UPDATE")) + assert update[1][-4:] == ("tenant-a", "worker-a", "default", "batch-1") + + +def test_save_rejects_wrong_types_and_mismatched_expected_identity( + database: FakeDatabase, +) -> None: + """Type and identity mismatches fail before durable state access.""" + store = PostgresBatchResultCheckpointStore("postgresql://unit") + with pytest.raises(ValidationError): + store.save("worker-a", object()) + with pytest.raises(ValidationError): + store.save("worker-a", checkpoint(), expected_previous=object()) + with pytest.raises(ValidationError): + store.save( + "worker-a", + checkpoint(), + expected_previous=replace(checkpoint(), batch_id="batch-2"), + ) + assert database.calls == [] + + +def test_apply_schema_uses_explicit_migration( + database: FakeDatabase, + tmp_path: Any, +) -> None: + """Operators may apply one explicitly selected migration file.""" + migration = tmp_path / "checkpoint.sql" + migration.write_text("SELECT 1;", encoding="utf-8") + apply_result_checkpoint_schema("postgresql://unit", str(migration)) + assert database.calls[-1] == ("SELECT 1;", ()) + assert database.commits == 1 + + +def test_apply_schema_uses_packaged_default( + database: FakeDatabase, + monkeypatch: pytest.MonkeyPatch, + tmp_path: Any, +) -> None: + """The default installer reads the package-owned migration path.""" + migration = tmp_path / "default.sql" + migration.write_text("SELECT 2;", encoding="utf-8") + monkeypatch.setattr(checkpoint_store, "MIGRATION_PATH", migration) + apply_result_checkpoint_schema("postgresql://unit") + assert database.calls[-1] == ("SELECT 2;", ()) diff --git a/tests/test_checkpoint_store_container_installation.py b/tests/test_checkpoint_store_container_installation.py new file mode 100644 index 000000000..4f8cd3490 --- /dev/null +++ b/tests/test_checkpoint_store_container_installation.py @@ -0,0 +1,59 @@ +# SPDX-License-Identifier: Apache-2.0 +"""Regression tests for installing the durable checkpoint migration in the image.""" + +from __future__ import annotations + +from pathlib import Path + + +REPOSITORY_ROOT = Path(__file__).resolve().parents[1] +DOCKERFILE_PATH = REPOSITORY_ROOT / "docker" / "postgres" / "Dockerfile" +PACKAGE_MIGRATION_PATH = ( + REPOSITORY_ROOT + / "pg_llm_batch" + / "migrations" + / "0007_result_stream_checkpoints.sql" +) +CONTAINER_MIGRATION_PATH = ( + REPOSITORY_ROOT + / "docker" + / "postgres" + / "init" + / "03_result_stream_checkpoints.sql" +) +CONTAINER_CHECKPOINT_COPY = ( + "COPY init/03_result_stream_checkpoints.sql " + "/docker-entrypoint-initdb.d/04_result_stream_checkpoints.sql" +) +CONTAINER_CRON_COPY = ( + "COPY init/03_cron_batch_retrieval.sql " + "/docker-entrypoint-initdb.d/03_cron_batch_retrieval.sql" +) + + +def test_checkpoint_migration_is_byte_identical_and_installed_in_image() -> None: + """The bundled PostgreSQL image must execute the reviewed checkpoint migration.""" + package_sql = PACKAGE_MIGRATION_PATH.read_bytes() + container_sql = CONTAINER_MIGRATION_PATH.read_bytes() + dockerfile = DOCKERFILE_PATH.read_text(encoding="utf-8") + + assert container_sql == package_sql + assert dockerfile.count(CONTAINER_CHECKPOINT_COPY) == 1 + assert dockerfile.index(CONTAINER_CRON_COPY) < dockerfile.index( + CONTAINER_CHECKPOINT_COPY + ) + + +def test_checkpoint_migration_uses_a_unique_init_destination() -> None: + """Checkpoint initialization must not overwrite another entrypoint script.""" + dockerfile = DOCKERFILE_PATH.read_text(encoding="utf-8") + init_copy_lines = [ + line.strip() + for line in dockerfile.splitlines() + if line.startswith("COPY init/") + and "/docker-entrypoint-initdb.d/" in line + ] + destinations = [line.split()[-1] for line in init_copy_lines] + + assert CONTAINER_CHECKPOINT_COPY in init_copy_lines + assert len(destinations) == len(set(destinations)) diff --git a/tests/test_checkpoint_store_documentation.py b/tests/test_checkpoint_store_documentation.py new file mode 100644 index 000000000..3bf3d71c7 --- /dev/null +++ b/tests/test_checkpoint_store_documentation.py @@ -0,0 +1,72 @@ +# SPDX-License-Identifier: Apache-2.0 +"""Authoritative documentation contracts for durable result checkpoints.""" + +from pathlib import Path + + +def _text(path: str) -> str: + """Read one authoritative UTF-8 project document with normalized spacing.""" + return " ".join(Path(path).read_text(encoding="utf-8").split()) + + +def test_authoritative_documents_define_durable_checkpoint_contract() -> None: + """Contributor and operator contracts agree on persistence boundaries.""" + contracts = { + "AGENTS.md": ( + "PostgresBatchResultCheckpointStore", + "save_in_transaction", + "NOBYPASSRLS", + ), + "CLAUDE.md": ( + "PostgresBatchResultCheckpointStore", + "expected_previous", + "initial_checkpoint_race", + ), + "ARCHITECTURE.md": ( + "llm_result_stream_checkpoints", + "compare-and-swap", + "caller-owned transaction", + ), + "CHANGELOG.md": ( + "durable result-checkpoint store", + "fail-closed rollback", + "fresh bundled PostgreSQL image", + "04_result_stream_checkpoints.sql", + ), + "docs/result-streaming.md": ( + "apply_result_checkpoint_schema", + "save_in_transaction", + "not a distributed exactly-once protocol", + ), + "docs/adr/0007-durable-result-checkpoint-store.md": ( + "Status: Accepted", + "FOR UPDATE", + "ON CONFLICT", + "full-stream immutability", + ), + "docs/doctoring/durable-result-checkpoint-store.md": ( + "PostgreSQL 18", + "NIST SP 800-53 Rev. 5", + "Retrieved August 6, 2026", + "04_result_stream_checkpoints.sql", + "after the cron initialization script", + ), + } + for path, required_phrases in contracts.items(): + content = _text(path) + for phrase in required_phrases: + assert phrase in content, f"{path} is missing {phrase!r}" + + +def test_doctoring_records_primary_sources_in_apa_7_style() -> None: + """The assurance record cites current primary sources with stable details.""" + doctoring = _text("docs/doctoring/durable-result-checkpoint-store.md") + required_references = ( + "National Institute of Standards and Technology. (2020).", + "https://doi.org/10.6028/NIST.SP.800-53r5", + "PostgreSQL Global Development Group. (n.d.).", + "Explicit locking", + "Row security policies", + ) + for reference in required_references: + assert reference in doctoring diff --git a/tests/test_checkpoint_store_integration.py b/tests/test_checkpoint_store_integration.py new file mode 100644 index 000000000..df41f28f9 --- /dev/null +++ b/tests/test_checkpoint_store_integration.py @@ -0,0 +1,74 @@ +# SPDX-License-Identifier: Apache-2.0 +"""Live PostgreSQL integration tests for durable result checkpoints.""" + +from __future__ import annotations + +import os +import uuid + +import pytest + +from pg_llm_batch import ( + BatchResultCheckpoint, + CheckpointConflictError, + PostgresBatchResultCheckpointStore, + apply_result_checkpoint_schema, +) + +pytestmark = pytest.mark.integration + +DSN = os.environ.get("PG_LLM_BATCH_TEST_DSN") +skip_no_db = pytest.mark.skipif( + not DSN, + reason="PG_LLM_BATCH_TEST_DSN not set; skipping live-DB integration", +) + + +def _checkpoint(batch_id: str, *, record_count: int, digest: str) -> BatchResultCheckpoint: + """Build one live-database checkpoint with monotonic positions.""" + return BatchResultCheckpoint( + schema_version=1, + batch_id=batch_id, + endpoint_alias="default", + file_kind="result", + file_id="file-live", + file_line_number=record_count, + batch_line_count=record_count, + record_count=record_count, + prefix_sha256=digest, + ) + + +@skip_no_db +def test_live_checkpoint_compare_and_swap_is_idempotent_and_fail_closed() -> None: + """PostgreSQL preserves exact CAS state and rejects a stale writer.""" + import psycopg + + apply_result_checkpoint_schema(DSN) + suffix = uuid.uuid4().hex + consumer = f"integration-{suffix}" + batch_id = f"batch-{suffix}" + store = PostgresBatchResultCheckpointStore(DSN, tenant_scope="standalone") + first = _checkpoint(batch_id, record_count=1, digest="a" * 64) + second = _checkpoint(batch_id, record_count=2, digest="b" * 64) + + try: + assert store.save(consumer, first) == first + assert store.save(consumer, first) == first + assert store.save(consumer, second, expected_previous=first) == second + assert store.load(consumer, batch_id, "default") == second + with pytest.raises(CheckpointConflictError): + store.save(consumer, first, expected_previous=first) + finally: + with psycopg.connect(DSN) as connection: + with connection.cursor() as cursor: + cursor.execute( + "SELECT set_config('pg_llm_batch.tenant_scope', %s, true)", + ("standalone",), + ) + cursor.execute( + "DELETE FROM llm_result_stream_checkpoints " + "WHERE checkpoint_consumer_name = %s", + (consumer,), + ) + connection.commit() diff --git a/tests/test_checkpoint_store_postgres_bounds.py b/tests/test_checkpoint_store_postgres_bounds.py new file mode 100644 index 000000000..e9678a7fe --- /dev/null +++ b/tests/test_checkpoint_store_postgres_bounds.py @@ -0,0 +1,123 @@ +# SPDX-License-Identifier: Apache-2.0 +"""PostgreSQL storage-boundary tests for durable result checkpoints.""" + +from __future__ import annotations + +from dataclasses import replace +from typing import Any + +import pytest + +from pg_llm_batch import PostgresBatchResultCheckpointStore +from pg_llm_batch.exceptions import ValidationError +from pg_llm_batch.result_streaming import BatchResultCheckpoint + +POSTGRES_BIGINT_MAX = (1 << 63) - 1 + + +class RefusingCursor: + """Fail whenever invalid checkpoint input reaches database execution.""" + + def execute(self, _sql: str, _params: tuple[Any, ...] | None = None) -> None: + """Prove storage-bound validation completed before any SQL statement.""" + raise AssertionError("database access occurred before validation") + + def fetchone(self) -> Any: + """Prevent accidental result access in a pre-database validation test.""" + raise AssertionError("database result read occurred before validation") + + +def checkpoint() -> BatchResultCheckpoint: + """Build one valid checkpoint within PostgreSQL signed BIGINT limits.""" + return BatchResultCheckpoint( + schema_version=1, + batch_id="batch-1", + endpoint_alias="default", + file_kind="result", + file_id="file-1", + file_line_number=1, + batch_line_count=1, + record_count=1, + prefix_sha256="a" * 64, + ) + + +@pytest.mark.parametrize( + ("changes", "expected_field"), + ( + ( + { + "file_line_number": POSTGRES_BIGINT_MAX + 1, + "batch_line_count": POSTGRES_BIGINT_MAX + 1, + }, + "checkpoint.file_line_number", + ), + ( + {"batch_line_count": POSTGRES_BIGINT_MAX + 1}, + "checkpoint.batch_line_count", + ), + ( + { + "record_count": POSTGRES_BIGINT_MAX + 1, + "batch_line_count": POSTGRES_BIGINT_MAX + 1, + }, + "checkpoint.batch_line_count", + ), + ), +) +def test_save_rejects_checkpoint_counts_above_postgres_bigint_before_sql( + changes: dict[str, int], + expected_field: str, +) -> None: + """Oversized durable counts fail deterministically before tenant SQL binding.""" + candidate = replace(checkpoint(), **changes) + store = PostgresBatchResultCheckpointStore("postgresql://unit") + + with pytest.raises(ValidationError) as raised: + store.save_in_transaction(RefusingCursor(), "worker-a", candidate) + + assert raised.value.field == expected_field + assert raised.value.reason == ( + f"must be no greater than PostgreSQL BIGINT maximum {POSTGRES_BIGINT_MAX}" + ) + + +def test_save_rejects_oversized_expected_previous_before_sql() -> None: + """Compare-and-swap evidence must also fit PostgreSQL before database access.""" + previous = replace( + checkpoint(), + file_line_number=POSTGRES_BIGINT_MAX + 1, + batch_line_count=POSTGRES_BIGINT_MAX + 1, + ) + candidate = replace( + checkpoint(), + file_line_number=POSTGRES_BIGINT_MAX, + batch_line_count=POSTGRES_BIGINT_MAX, + record_count=2, + prefix_sha256="b" * 64, + ) + store = PostgresBatchResultCheckpointStore("postgresql://unit") + + with pytest.raises(ValidationError) as raised: + store.save_in_transaction( + RefusingCursor(), + "worker-a", + candidate, + expected_previous=previous, + ) + + assert raised.value.field == "expected_previous.file_line_number" + + +def test_save_accepts_postgres_bigint_maximum_before_sql() -> None: + """The exact signed BIGINT maximum remains a supported durable value.""" + candidate = replace( + checkpoint(), + file_line_number=POSTGRES_BIGINT_MAX, + batch_line_count=POSTGRES_BIGINT_MAX, + record_count=POSTGRES_BIGINT_MAX, + ) + store = PostgresBatchResultCheckpointStore("postgresql://unit") + + with pytest.raises(AssertionError, match="database access occurred"): + store.save_in_transaction(RefusingCursor(), "worker-a", candidate) diff --git a/tests/test_checkpoint_store_schema.py b/tests/test_checkpoint_store_schema.py new file mode 100644 index 000000000..7f4a0cbe6 --- /dev/null +++ b/tests/test_checkpoint_store_schema.py @@ -0,0 +1,47 @@ +# SPDX-License-Identifier: Apache-2.0 +"""Static migration and rollback contracts for durable result checkpoints.""" + +from pathlib import Path + +PACKAGE_SQL = Path("pg_llm_batch/migrations/0007_result_stream_checkpoints.sql") +DOCKER_SQL = Path("docker/postgres/init/03_result_stream_checkpoints.sql") +ROLLBACK_SQL = Path("pg_llm_batch/migrations/rollback/0007_result_stream_checkpoints.sql") + + +def test_packaged_and_container_checkpoint_migrations_are_identical() -> None: + """Package and container installs execute the same migration bytes.""" + assert PACKAGE_SQL.is_file() + assert PACKAGE_SQL.read_bytes() == DOCKER_SQL.read_bytes() + + +def test_checkpoint_schema_is_tenant_isolated_and_fail_closed() -> None: + """The forward migration enforces tenant isolation and bounded fields.""" + sql = PACKAGE_SQL.read_text(encoding="utf-8") + required = ( + "CREATE TABLE IF NOT EXISTS llm_result_stream_checkpoints", + "checkpoint_consumer_name TEXT NOT NULL", + "uq_llm_result_stream_checkpoints_tenant_consumer_batch", + "UNIQUE (\n tenant_scope,\n checkpoint_consumer_name,\n endpoint_alias,\n remote_batch_id", + "ENABLE ROW LEVEL SECURITY", + "FORCE ROW LEVEL SECURITY", + "plc_llm_result_stream_checkpoints_tenant_scope", + "current_setting('pg_llm_batch.tenant_scope', true)", + "prefix_sha256 ~ '^[0-9a-f]{64}$'", + "record_count <= batch_line_count", + "idx_llm_result_stream_checkpoints_tenant_updated", + ) + for contract in required: + assert contract in sql + assert "BYPASSRLS" not in sql + + +def test_rollback_refuses_to_destroy_acknowledgement_evidence() -> None: + """Rollback exposes owner-visible rows before its destructive emptiness check.""" + sql = ROLLBACK_SQL.read_text(encoding="utf-8") + no_force = "ALTER TABLE llm_result_stream_checkpoints NO FORCE ROW LEVEL SECURITY" + assert no_force in sql + assert "EXISTS (" in sql + assert "Refusing to drop non-empty llm_result_stream_checkpoints" in sql + assert "ERRCODE = '55000'" in sql + assert sql.index(no_force) < sql.index("EXISTS (") + assert sql.index("RAISE EXCEPTION") < sql.index("DROP TABLE")