From 0d3366de78c6b1373588d682560d3d50a6ee8a9a Mon Sep 17 00:00:00 2001 From: Seongho Bae Date: Sun, 9 Aug 2026 09:40:55 +0900 Subject: [PATCH] feat(audit): replay bounded stable export pages --- AGENTS.md | 30 ++ ARCHITECTURE.md | 30 ++ CHANGELOG.md | 12 + CLAUDE.md | 28 ++ ...nded-checkpoint-audit-export-pagination.md | 125 +++++++ docs/checkpoint-audit.md | 56 ++- .../checkpoint-audit-export-pagination.md | 154 ++++++++ pg_llm_batch/__init__.py | 9 +- pg_llm_batch/checkpoint_audit.py | 161 ++++++++- .../test_checkpoint_audit_bigint_contract.py | 40 +++ tests/test_checkpoint_audit_integration.py | 47 ++- tests/test_checkpoint_audit_pagination.py | 331 ++++++++++++++++++ 12 files changed, 1012 insertions(+), 11 deletions(-) create mode 100644 docs/adr/0011-bounded-checkpoint-audit-export-pagination.md create mode 100644 docs/doctoring/checkpoint-audit-export-pagination.md create mode 100644 tests/test_checkpoint_audit_bigint_contract.py create mode 100644 tests/test_checkpoint_audit_pagination.py diff --git a/AGENTS.md b/AGENTS.md index 4cd97f2c9..350860330 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -270,3 +270,33 @@ add CODEOWNERS-based merge gates until multiple independent maintainers exist. - Update README, architecture, ADR, operator documentation, doctoring, and CHANGELOG whenever migration ordering, locking, evidence, failure, or compatibility semantics change. + +## Checkpoint audit export pagination contract + +- Keep pagination opt-in and preserve `list_audit_events()` behavior. Do not + replace the existing bounded convenience read with an unbounded iterator. +- Use `checkpoint_audit_event_id` as a strict positive signed-PostgreSQL-`BIGINT` + keyset cursor. Do not use `OFFSET` for multi-page audit traversal. +- Query at most the validated public limit plus one lookahead row and expose at + most the requested 1..1,000 events. A driver returning more than the SQL bound + fails closed. +- Continuations must use `checkpoint_audit_event_id < before_audit_event_id` in + exact newest-first order. When more rows exist, the next cursor is exactly the + final returned event identity. +- Revalidate every database row through `CheckpointAuditEvent`, compare its + tenant/consumer/endpoint/batch key to the trusted request, and require strict + descending identities before exposing a page. +- Treat the cursor as navigation state only. It is not proof of completeness, + chronology, authenticity, delivery, or non-repudiation; identity allocation + can contain gaps and commit ordering can differ from allocation ordering. +- Package-owned page calls do not promise a single historic snapshot. Hosts + requiring one export snapshot must begin their own PostgreSQL `REPEATABLE READ` + or stricter transaction before the first query and repeatedly call + `list_audit_event_page_in_transaction()` on the same transaction. +- Keep export destinations, credentials, retention policy, immutable/WORM + storage, delivery receipts, cryptographic manifests, and reconciliation in the + host/operator boundary. This package primitive adds no network exporter or + write authority. +- Maintain 100% production statement, branch, and public-docstring coverage with + strict cursor, page-shape, lookahead, keyset SQL, ordering, row-key, + malformed-driver, concurrency-semantics, and transaction-ownership tests. diff --git a/ARCHITECTURE.md b/ARCHITECTURE.md index 8b007a768..75193c80c 100644 --- a/ARCHITECTURE.md +++ b/ARCHITECTURE.md @@ -426,3 +426,33 @@ concurrent lock waiting, body-free JSON, unchanged `init-db`, and 100% productio statement, branch, and public-docstring coverage. Final merge evidence must be regenerated against the integrated base; successful stacked-base runs are not reusable release evidence. + +## Bounded checkpoint-audit export boundary + +`CheckpointAuditPage`, `list_audit_event_page()`, and +`list_audit_event_page_in_transaction()` add an opt-in bounded traversal layer +without changing the existing one-page audit read or database schema. Pagination +uses `checkpoint_audit_event_id` as a strict positive signed PostgreSQL `BIGINT` +keyset cursor, never `OFFSET`. + +Each query requests at most the validated public limit plus one lookahead row and +exposes at most 1,000 events. Continuations use +`checkpoint_audit_event_id < before_audit_event_id` in newest-first order. Every +row is revalidated through `CheckpointAuditEvent`, compared with the exact trusted +tenant/consumer/endpoint/batch key, and required to remain strictly descending. +Malformed collections, impossible driver overruns, cross-key rows, duplicate or +ascending identities, and cursor-domain violations fail closed before exposure. + +Keyset traversal prevents later higher-identity inserts from shifting an older +continuation window, but package-owned calls do not provide one multi-page +historic snapshot. A host that requires snapshot-stable export must begin a +caller-owned PostgreSQL `REPEATABLE READ` or stricter transaction before the first +query and repeatedly call the in-transaction method on that same transaction. + +Audit identities are navigation keys, not cryptographic chronology or completeness +proof. Sequence gaps and allocation/commit reordering are valid. External +immutable/WORM retention, delivery receipts, cryptographic manifests, +reconciliation, legal hold, and disposal remain host/operator responsibilities. +The primitive stays independently deployable and can be embedded into CWL MSA +workflows without requiring `contextual-orchestrator`, `naruon`, or a network +export service. diff --git a/CHANGELOG.md b/CHANGELOG.md index 5825217f2..c8fe222f7 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -9,6 +9,18 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ### Added +- Opt-in bounded checkpoint-audit export pagination through immutable + `CheckpointAuditPage`, `list_audit_event_page()`, and + `list_audit_event_page_in_transaction()`. Pagination uses a strict positive + PostgreSQL `BIGINT` keyset cursor, newest-first `<` continuation, and one-row + bounded lookahead rather than `OFFSET`; returned rows are revalidated against + the exact trusted tenant/consumer/endpoint/batch key and strict descending + identity order. A live PostgreSQL regression proves that a newer committed row + between pages cannot drift into the older continuation window. Package-owned + calls do not claim one multi-page snapshot; hosts needing that guarantee own a + `REPEATABLE READ` or stricter transaction. External immutable/WORM retention, + receipts, cryptographic manifests, and reconciliation remain host controls. No + migration, version bump, or release is included. - Explicit opt-in `init-checkpoint-storage` operator for existing PostgreSQL volumes. It bounded-reads, validates, and SHA-256 identifies `0007_result_stream_checkpoints` and diff --git a/CLAUDE.md b/CLAUDE.md index ee08b9d72..72b28d882 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -246,3 +246,31 @@ strict red-green-refactor tests for bounded input, exact order, one lock, one transaction, one commit, rollback, concurrency, compatibility, CLI output, documentation, and body-free diagnostics. + +## Checkpoint audit export pagination invariants + +- Keep `list_audit_events()` source-compatible and make multi-page traversal + opt-in through `CheckpointAuditPage` and the `list_audit_event_page*` methods. +- Validate `before_audit_event_id` as `None` or a strict positive signed + PostgreSQL `BIGINT`; reject booleans, coercible strings/floats, zero, negative, + and out-of-range identities before database access. +- Use primary-key keyset pagination only: newest-first ordering and strict + `checkpoint_audit_event_id < before_audit_event_id` continuation. Never use + `OFFSET` for retained audit traversal. +- Fetch no more than the validated limit plus one lookahead row and expose no + more than 1,000 events. Treat a driver overrun as an integrity failure. +- Revalidate every database row through `CheckpointAuditEvent`, require the + exact trusted tenant/consumer/endpoint/batch key, and require strictly + descending unique identities before exposing a page. +- The continuation cursor is navigation evidence only. It is not a completeness, + chronology, delivery, authenticity, or non-repudiation attestation; identity + gaps and allocation/commit reordering remain valid PostgreSQL behavior. +- Package-owned page calls do not provide one cross-page snapshot. Hosts needing + snapshot-stable export must begin a caller-owned `REPEATABLE READ` or stricter + transaction before the first query and reuse the in-transaction method. +- Keep destination credentials, immutable/WORM storage, retention, legal hold, + export receipts, cryptographic manifests, and reconciliation outside package + write authority. +- Maintain 100% production statement, branch, and public-docstring coverage with + strict cursor, immutable-page, lookahead, keyset SQL, trusted-key, ordering, + malformed-driver, no-commit, and live concurrent-insert pagination tests. diff --git a/docs/adr/0011-bounded-checkpoint-audit-export-pagination.md b/docs/adr/0011-bounded-checkpoint-audit-export-pagination.md new file mode 100644 index 000000000..20e3e513c --- /dev/null +++ b/docs/adr/0011-bounded-checkpoint-audit-export-pagination.md @@ -0,0 +1,125 @@ +# ADR 0011: Bounded checkpoint-audit export pagination + +- **Status:** Accepted for the stacked implementation +- **Date:** 2026-08-07 +- **Decision owners:** ContextualWisdomLab +- **Depends on:** ADR 0009 and the append-only checkpoint audit trail +- **Stack order:** Follows ADR 0010, the atomic checkpoint schema operator + +## Context + +The package-owned audit store previously exposed only one newest-first bounded +read. That is appropriate for an operator console, but not sufficient for a +host that must move a longer retained audit history to separately governed +storage. Leaving multi-page traversal to every embedding product would produce +incompatible cursor semantics and would encourage `OFFSET` pagination, which can +skip or duplicate rows when the visible set changes between requests. + +The existing audit identity is a PostgreSQL `BIGINT GENERATED ALWAYS AS +IDENTITY` primary key. The compound checkpoint-key index already ends in +`checkpoint_audit_event_id DESC`, so the table can support bounded keyset +pagination without a schema migration or a second ordering column. + +## Decision + +Add an opt-in `CheckpointAuditPage` contract and two audit-store methods: + +- `list_audit_event_page()` for package-owned page reads; and +- `list_audit_event_page_in_transaction()` for hosts that own the surrounding + PostgreSQL transaction. + +Each request validates a strict 1..1,000 page size and an optional positive +signed-`BIGINT` `before_audit_event_id`. The first page orders by +`checkpoint_audit_event_id DESC`. Continuations add the strict predicate +`checkpoint_audit_event_id < before_audit_event_id` and preserve the same order. +The SQL query requests only `limit + 1` rows. At most `limit` rows are returned; +the lookahead row only proves that an older continuation exists. When a +continuation exists, `next_before_audit_event_id` is exactly the identity of the +last returned event. + +Every returned row is revalidated as a `CheckpointAuditEvent`, rechecked against +the trusted tenant/consumer/endpoint/batch key, and required to be strictly +descending. A database adapter that returns a non-sequence or more than the +bounded query size fails closed. + +## Concurrency boundary + +Keyset traversal solves offset drift, not snapshot isolation. Rows committed +after page one with identities greater than the cursor cannot move or duplicate +older rows already traversed. However, package-owned calls execute as separate +transactions and therefore do not promise one historic database snapshot. + +A host that requires one snapshot for an export pass must own a PostgreSQL +`REPEATABLE READ` or stricter transaction and call the in-transaction method on +one cursor. The package does not silently alter transaction isolation because +that would change caller-owned database semantics. + +Identity allocation order is not a cryptographic chronology and may contain +gaps. A transaction can allocate an identity before another transaction and +commit later. The cursor is therefore a database navigation key, not a statement +that all real-world events before or after a wall-clock instant have been +captured. + +## Security and acquisition boundary + +The page remains tenant-qualified and inherits forced RLS and ordinary-role +append-only controls from ADR 0009. It does not create a new database object, +credential, network destination, export worker, background process, or write +permission. It never accepts provider or model output as tenant or consumer +identity. + +The feature is intended to make bounded export to separately governed immutable +or write-once retention storage practical. It does not itself provide immutable +external storage, cryptographic non-repudiation, signed completeness evidence, +administrator-proof tamper detection, or delivery acknowledgement. Those remain +host/operator controls. A later cryptographically protected export manifest can +be layered on this stable pagination contract without changing database +navigation semantics. + +## Alternatives rejected + +### SQL OFFSET/LIMIT + +Rejected because additions or visibility changes before an offset can shift the +subsequent window, producing duplicates or omissions during a long-running +export. + +### Unbounded iterator or fetch-all export + +Rejected because acquisition-grade audit volume is not bounded by a single +operator interaction. Materializing arbitrary retained history would weaken the +package's existing memory-safety policy. + +### Implicit REPEATABLE READ in the package-owned API + +Rejected because changing isolation belongs to the owner of the transaction and +must occur before the first transaction query. The in-transaction API gives a +host an explicit way to obtain one snapshot without the package surprising +other database work. + +### Timestamp cursor + +Rejected because timestamps are not unique and database wall-clock values are +not a total order. The existing primary-key identity is already indexed and +provides an unambiguous strict continuation boundary. + +## Verification + +Permanent deterministic tests require strict cursor validation, immutable page +shape, strictly descending identities, one-row lookahead, no `OFFSET`, exact +`<` keyset continuation, trusted-key row revalidation, bounded driver output, +and no package commit during owned page reads. The live PostgreSQL regression +reads page one, commits a newer event in a separate transaction, and proves that +page two continues toward older identities without replaying page-one rows or +admitting the newer row. Integration/release gates remain required after this +stack is reconciled onto protected `main`. + +## Consequences + +Operators gain a package-owned bounded traversal primitive suitable for durable +audit export and acquisition diligence. Existing `list_audit_events()` behavior +is unchanged. No migration, release version, or publication authority is added. + +The operational tradeoff is explicit: callers choosing independent page calls +accept normal transaction-to-transaction visibility changes; callers requiring a +single snapshot must own and configure that transaction deliberately. diff --git a/docs/checkpoint-audit.md b/docs/checkpoint-audit.md index 8425b797e..af63cef91 100644 --- a/docs/checkpoint-audit.md +++ b/docs/checkpoint-audit.md @@ -89,6 +89,50 @@ action, a database-generated event identity, and an insert-time database wall-clock timestamp. It does not include provider bodies, prompts, model output, credentials, DSNs, transport headers, or exception text. +## Export longer retained history + +Use the opt-in keyset page API instead of `OFFSET` or an unbounded fetch: + +```python +page = store.list_audit_event_page( + "invoice-worker", + "batch-123", + "default", + limit=250, +) + +while True: + persist_to_governed_retention(page.events) + if page.next_before_audit_event_id is None: + break + page = store.list_audit_event_page( + "invoice-worker", + "batch-123", + "default", + before_audit_event_id=page.next_before_audit_event_id, + limit=250, + ) +``` + +`before_audit_event_id` is `None` or a strict positive signed PostgreSQL `BIGINT`. +Each SQL request reads at most `limit + 1` rows and returns at most `limit` events. +Continuation uses `checkpoint_audit_event_id < before_audit_event_id` in strict +newest-first order. Returned rows are revalidated against the exact trusted +tenant, consumer, endpoint, and batch key before exposure. + +A committed event with a larger identity between page calls cannot shift the +older continuation window. Separate package-owned calls are still separate +transactions, however, so they do not form one historic database snapshot. For +an export that must observe one PostgreSQL snapshot, begin a caller-owned +`REPEATABLE READ` or stricter transaction **before the first query** and reuse +`list_audit_event_page_in_transaction()` on that transaction's cursor. + +The cursor is navigation state, not completeness, chronology, delivery, +authenticity, or non-repudiation evidence. PostgreSQL identity sequences may have +gaps and allocation order can differ from commit order. Destination credentials, +immutable/WORM storage, retention, legal hold, delivery receipts, cryptographic +manifests, and reconciliation remain host/operator responsibilities. + ## Retention and rollback The database blocks ordinary `UPDATE`, `DELETE`, and `TRUNCATE` against the audit @@ -105,8 +149,10 @@ signed/hash-chained evidence system. ## Standards and evidence -See [ADR 0009](adr/0009-append-only-checkpoint-audit-trail.md) for the decision -boundary and -[the assurance record](doctoring/checkpoint-audit-trail.md) for threat model, -verification evidence, and APA 7 references to NIST SP 800-53 Rev. 5 AU-3, the -OWASP Logging Cheat Sheet, and PostgreSQL 18 trigger and current-time semantics. +See [ADR 0009](adr/0009-append-only-checkpoint-audit-trail.md) for the accepted-save +audit boundary, [ADR 0011](adr/0011-bounded-checkpoint-audit-export-pagination.md) +for the pagination decision, [the audit assurance record](doctoring/checkpoint-audit-trail.md) +for audit threat-model evidence, and +[the export-pagination assurance record](doctoring/checkpoint-audit-export-pagination.md) +for keyset, isolation, operator, and APA 7 evidence. The latter records NIST SP +800-53 Rev. 5 AU-9 and PostgreSQL 18 transaction-isolation/concurrency guidance. diff --git a/docs/doctoring/checkpoint-audit-export-pagination.md b/docs/doctoring/checkpoint-audit-export-pagination.md new file mode 100644 index 000000000..2e1b590aa --- /dev/null +++ b/docs/doctoring/checkpoint-audit-export-pagination.md @@ -0,0 +1,154 @@ +# Checkpoint audit export pagination assurance record + +## Scope + +This record documents the bounded keyset-pagination control used to traverse +`llm_result_checkpoint_audit_events` for export or reconciliation. It supplements +the append-only audit assurance record; it does not replace the database RLS, +mutation-rejection, rollback, or retention controls described there. + +The commercial objective is practical and deliberately narrow: a host must be +able to move more than one bounded page of accepted-save audit evidence into a +separately governed retention system without inventing product-specific +pagination or materializing the entire audit history in memory. + +## Threat and failure model + +The implementation treats tenant scope, consumer name, endpoint alias, batch +identifier, page size, cursor, database row shape, row identity, and row ordering +as validation boundaries. Provider payloads and model output never select audit +scope or cursor state. + +The bounded query uses the existing tenant-qualified checkpoint key and orders by +`checkpoint_audit_event_id DESC`. A continuation uses a strict primary-key +predicate: + +```sql +AND checkpoint_audit_event_id < %s +ORDER BY checkpoint_audit_event_id DESC +LIMIT %s +``` + +The bound passed to SQL is exactly the validated public limit plus one. The extra +row is not returned to the caller; it only determines whether an older page +exists. A database adapter returning more than that bound is treated as an +internal integrity failure rather than silently expanding memory use. + +Every materialized row is reconstructed through the strict +`CheckpointAuditEvent` validator and is then compared to the exact trusted +`(tenant_scope, checkpoint_consumer_name, endpoint_alias, remote_batch_id)` key. +Rows must be strictly descending by audit identity. A cross-key row, duplicate +identity, ascending identity, malformed collection, malformed row, or cursor +outside signed PostgreSQL `BIGINT` range fails closed. + +## Concurrency semantics + +Keyset pagination intentionally avoids `OFFSET`. If page one ends at audit event +`800`, later committed events with larger identities cannot move the continuation +window because page two is anchored by `checkpoint_audit_event_id < 800`. +This prevents the common offset-shift duplicate/skip failure mode for newly +visible higher-key rows. + +This is not equivalent to one database snapshot. PostgreSQL's default isolation +is `READ COMMITTED`, under which each statement sees rows committed before that +statement begins. A package-owned page call therefore may see database changes +that were not visible to the previous call. PostgreSQL documents that +`REPEATABLE READ` instead keeps all statements in one transaction on the snapshot +established by its first query. Hosts that require a single-snapshot export pass +must explicitly begin an appropriate caller-owned transaction before the first +page and repeatedly use `list_audit_event_page_in_transaction()` on that +transaction's cursor. + +The package does not change caller transaction isolation automatically. Doing so +after the first query is prohibited by PostgreSQL, and doing so implicitly would +alter host-owned transaction semantics. + +PostgreSQL identity values are navigation keys, not trusted wall-clock sequence +numbers. Gaps are valid, and allocation/commit ordering can differ across +concurrent transactions. A continuation cursor therefore proves only where the +next keyset query resumes. It does not prove that no later transaction can reveal +an event that was invisible to an earlier export pass. + +## Security and privacy boundary + +NIST SP 800-53 Rev. 5 AU-9 requires protection of audit information against +unauthorized access, modification, and deletion. NIST Release 5.2.0, issued on +August 27, 2025, remains the current minor release of Revision 5 and does not +remove this AU-9 protection boundary. The underlying audit table implements +tenant-qualified RLS and mutation rejection for ordinary application roles. This +pagination feature does not weaken those controls and does not add a write-capable +database permission. + +NIST also describes stronger audit-protection enhancements such as write-once +media. This package does not claim that PostgreSQL alone provides administrator- +proof immutability. The bounded page API is the transport-neutral extraction +primitive by which an embedding host can copy evidence into a separately governed +immutable or write-once system under its own retention and access-control policy. +No destination, credential, exporter, background scheduler, object-store API, or +network client is introduced here. + +Audit events continue to exclude prompts, provider bodies, model output, +credentials, DSNs, transport headers, and arbitrary exception text. Cursor values +are database identities, not secrets. Diagnostics remain fixed and body-free. + +## Operational workflow + +1. Select a trusted tenant, consumer, endpoint, and batch key after host + authentication and authorization. +2. Choose a strict page limit from 1 through 1,000. +3. For ordinary incremental export, call `list_audit_event_page()` from + `before_audit_event_id=None`, persist each page to the governed destination, + and continue with `page.next_before_audit_event_id` until it is `None`. +4. If one PostgreSQL snapshot is required, begin a caller-owned `REPEATABLE READ` + or stricter transaction before the first read and use the in-transaction page + method for the entire pass. +5. Persist destination-side receipt, manifest, or cryptographic evidence when + the retention requirement needs proof of delivery or completeness. The page + cursor itself is not such proof. +6. Start a later reconciliation pass from the newest page to capture rows that + became visible only after a previous non-snapshot export. + +## Rollback + +No database migration is introduced. Rolling back this feature removes only the +new public pagination API and documentation. The audit table and retained events +are unchanged. Existing `list_audit_events()` callers remain compatible. + +Do not delete audit rows as a feature rollback. The existing rollback migration +continues to refuse removal of non-empty audit evidence. + +## Verification contract + +Permanent tests cover: + +- strict `None` or positive signed-`BIGINT` cursor validation without coercion; +- immutable page tuples and a maximum of 1,000 returned events; +- strictly descending unique audit identities; +- a continuation cursor equal to the final returned identity; +- one-row bounded lookahead and a maximum SQL row request of 1,001; +- first-page and continuation SQL without `OFFSET`; +- strict `<` continuation semantics; +- trusted tenant/consumer/endpoint/batch revalidation on returned rows; +- fail-closed malformed collection, impossible driver overrun, and row-order + behavior; +- package-owned reads that do not create an explicit package commit; and +- a live least-privilege PostgreSQL pass proving that a newer committed row + between pages cannot drift into the older continuation window. + +Final merge evidence is valid only after the full stacked dependency chain has +integrated and all required exact-head quality, security, coverage, packaging, +review, provenance, and release-acceptance gates have been regenerated against +the protected base. + +## APA 7 references + +Joint Task Force. (2020, updated 2025). *Security and privacy controls for +information systems and organizations* (NIST Special Publication 800-53, +Revision 5, Release 5.2.0). National Institute of Standards and Technology. +https://doi.org/10.6028/NIST.SP.800-53r5 + +PostgreSQL Global Development Group. (2026). *SET TRANSACTION*. PostgreSQL 18 +documentation. https://www.postgresql.org/docs/18/sql-set-transaction.html + +PostgreSQL Global Development Group. (2026). *Concurrency control*. PostgreSQL +18 documentation. https://www.postgresql.org/docs/18/mvcc.html diff --git a/pg_llm_batch/__init__.py b/pg_llm_batch/__init__.py index 56355d60a..2f3e1b0e3 100644 --- a/pg_llm_batch/__init__.py +++ b/pg_llm_batch/__init__.py @@ -10,6 +10,7 @@ BatchResultCheckpoint -- host-persistable resume evidence PostgresBatchResultCheckpointStore -- tenant-isolated durable checkpoints AuditedPostgresBatchResultCheckpointStore -- append-only accepted-save audit + CheckpointAuditPage -- stable bounded audit export page CheckpointSchemaMigration -- bounded migration identity evidence apply_checkpoint_schema_migrations -- atomic checkpoint schema operator OpenTelemetryCheckpointStore -- confidential checkpoint observability @@ -26,9 +27,12 @@ config_credentials_provider, ) from .checkpoint_audit import ( + MAX_CHECKPOINT_AUDIT_EVENT_ID, AuditedPostgresBatchResultCheckpointStore, CheckpointAuditEvent, + CheckpointAuditPage, apply_result_checkpoint_audit_schema, + validate_checkpoint_audit_cursor, validate_checkpoint_audit_limit, ) from .checkpoint_migrations import ( @@ -83,6 +87,8 @@ "PostgresBatchResultCheckpointStore", "AuditedPostgresBatchResultCheckpointStore", "CheckpointAuditEvent", + "CheckpointAuditPage", + "MAX_CHECKPOINT_AUDIT_EVENT_ID", "CheckpointSchemaMigration", "apply_checkpoint_schema_migrations", "plan_checkpoint_schema_migrations", @@ -92,6 +98,7 @@ "apply_result_checkpoint_audit_schema", "validate_checkpoint_consumer_name", "validate_checkpoint_audit_limit", + "validate_checkpoint_audit_cursor", "DurableBatchAPIClient", "TenantDurableBatchAPIClient", "DEFAULT_TENANT_SCOPE", @@ -116,4 +123,4 @@ "BatchAccumulator", "TokenCounter", "__version__", -] \ No newline at end of file +] diff --git a/pg_llm_batch/checkpoint_audit.py b/pg_llm_batch/checkpoint_audit.py index 49d0f4302..16b2f3d0e 100644 --- a/pg_llm_batch/checkpoint_audit.py +++ b/pg_llm_batch/checkpoint_audit.py @@ -22,6 +22,7 @@ AUDIT_ACTION_CHECKPOINT_SAVE_ACCEPTED = "checkpoint_save_accepted" MAX_CHECKPOINT_AUDIT_EVENTS = 1000 DEFAULT_CHECKPOINT_AUDIT_EVENTS = 100 +MAX_CHECKPOINT_AUDIT_EVENT_ID = 9_223_372_036_854_775_807 AUDIT_MIGRATION_PATH = ( Path(__file__).with_name("migrations") / "0008_result_checkpoint_audit_events.sql" ) @@ -58,8 +59,11 @@ def __post_init__(self) -> None: isinstance(self.audit_event_id, bool) or not isinstance(self.audit_event_id, int) or self.audit_event_id <= 0 + or self.audit_event_id > MAX_CHECKPOINT_AUDIT_EVENT_ID ): - raise ValueError("audit_event_id must be a positive integer") + raise ValueError( + "audit_event_id must be a positive PostgreSQL BIGINT-compatible integer" + ) validate_tenant_scope(self.tenant_scope) validate_checkpoint_consumer_name(self.consumer_name) _validated_exact_endpoint_alias(self.endpoint_alias) @@ -81,6 +85,35 @@ def __post_init__(self) -> None: raise ValueError("recorded_at must be a timezone-aware datetime") +@dataclass(frozen=True, slots=True) +class CheckpointAuditPage: + """One immutable newest-first bounded audit page and its older-row cursor.""" + + events: tuple[CheckpointAuditEvent, ...] + next_before_audit_event_id: Optional[int] + + def __post_init__(self) -> None: + """Reject mutable, oversized, unordered, or cursor-inconsistent pages.""" + if not isinstance(self.events, tuple): + raise ValueError("events must be an immutable tuple") + if len(self.events) > MAX_CHECKPOINT_AUDIT_EVENTS: + raise ValueError( + f"events must contain at most {MAX_CHECKPOINT_AUDIT_EVENTS} records" + ) + if any(not isinstance(event, CheckpointAuditEvent) for event in self.events): + raise ValueError("events must contain only CheckpointAuditEvent values") + if any( + current.audit_event_id >= previous.audit_event_id + for previous, current in zip(self.events, self.events[1:]) + ): + raise ValueError("events must be strictly descending by audit_event_id") + cursor = validate_checkpoint_audit_cursor(self.next_before_audit_event_id) + if cursor is not None and ( + not self.events or cursor != self.events[-1].audit_event_id + ): + raise ValueError("next_before_audit_event_id must equal the final event id") + + def validate_checkpoint_audit_limit(value: Any) -> int: """Validate the bounded number of audit rows one public read may return.""" if ( @@ -97,6 +130,26 @@ def validate_checkpoint_audit_limit(value: Any) -> int: return value +def validate_checkpoint_audit_cursor(value: Any) -> Optional[int]: + """Validate an optional positive PostgreSQL BIGINT audit-event keyset cursor.""" + if value is None: + return None + if ( + isinstance(value, bool) + or not isinstance(value, int) + or value < 1 + or value > MAX_CHECKPOINT_AUDIT_EVENT_ID + ): + raise ValidationError( + field="before_audit_event_id", + value=value, + reason=( + "must be a positive integer no greater than PostgreSQL BIGINT maximum" + ), + ) + return value + + def apply_result_checkpoint_audit_schema( postgres_dsn: str, migration_path: Optional[str] = None, @@ -267,3 +320,109 @@ def list_audit_events_in_transaction( if not isinstance(rows, (tuple, list)): raise RuntimeError("checkpoint audit query returned an invalid row collection") return tuple(_audit_event_from_row(row) for row in rows) + + def list_audit_event_page( + self, + consumer_name: str, + batch_id: str, + endpoint_alias: str, + *, + before_audit_event_id: Optional[int] = None, + limit: int = DEFAULT_CHECKPOINT_AUDIT_EVENTS, + ) -> CheckpointAuditPage: + """Return one stable bounded newest-first page for durable audit export.""" + _require_psycopg() + with psycopg.connect(self.postgres_dsn) as conn: + with conn.cursor() as cur: + return self.list_audit_event_page_in_transaction( + cur, + consumer_name, + batch_id, + endpoint_alias, + before_audit_event_id=before_audit_event_id, + limit=limit, + ) + + def list_audit_event_page_in_transaction( + self, + cursor: Any, + consumer_name: str, + batch_id: str, + endpoint_alias: str, + *, + before_audit_event_id: Optional[int] = None, + limit: int = DEFAULT_CHECKPOINT_AUDIT_EVENTS, + ) -> CheckpointAuditPage: + """Read one keyset page through a caller-owned PostgreSQL transaction. + + The continuation cursor is the final returned audit-event identity. A + subsequent page uses a strict ``<`` predicate, so rows inserted later + with larger identities cannot shift already traversed older rows. Hosts + needing one database snapshot across multiple pages should call this + method repeatedly inside their own PostgreSQL REPEATABLE READ or stricter + transaction. + """ + consumer = validate_checkpoint_consumer_name(consumer_name) + remote_batch_id = _validated_batch_id(batch_id) + alias = _validated_exact_endpoint_alias(endpoint_alias) + bounded_limit = validate_checkpoint_audit_limit(limit) + before = validate_checkpoint_audit_cursor(before_audit_event_id) + fetch_size = bounded_limit + 1 + _set_transaction_tenant_scope(cursor, self.tenant_scope) + + query = ( + f"SELECT {_AUDIT_COLUMNS} " + "FROM llm_result_checkpoint_audit_events " + "WHERE tenant_scope = %s " + "AND checkpoint_consumer_name = %s " + "AND endpoint_alias = %s " + "AND remote_batch_id = %s " + ) + params: list[Any] = [ + self.tenant_scope, + consumer, + alias, + remote_batch_id, + ] + if before is not None: + query += "AND checkpoint_audit_event_id < %s " + params.append(before) + query += "ORDER BY checkpoint_audit_event_id DESC LIMIT %s" + params.append(fetch_size) + cursor.execute(query, tuple(params)) + + rows = cursor.fetchall() + if not isinstance(rows, (tuple, list)): + raise RuntimeError("checkpoint audit query returned an invalid row collection") + if len(rows) > fetch_size: + raise RuntimeError("checkpoint audit query exceeded its bounded query size") + + events = tuple(_audit_event_from_row(row) for row in rows) + expected_key = (self.tenant_scope, consumer, alias, remote_batch_id) + for event in events: + event_key = ( + event.tenant_scope, + event.consumer_name, + event.endpoint_alias, + event.batch_id, + ) + if event_key != expected_key: + raise RuntimeError("checkpoint audit query returned a row outside the requested key") + if before is not None and any( + event.audit_event_id >= before for event in events + ): + raise RuntimeError("checkpoint audit query violated the continuation cursor") + if any( + current.audit_event_id >= previous.audit_event_id + for previous, current in zip(events, events[1:]) + ): + raise RuntimeError("checkpoint audit query was not strictly descending") + + page_events = events[:bounded_limit] + next_before = ( + page_events[-1].audit_event_id if len(events) > bounded_limit else None + ) + return CheckpointAuditPage( + events=page_events, + next_before_audit_event_id=next_before, + ) diff --git a/tests/test_checkpoint_audit_bigint_contract.py b/tests/test_checkpoint_audit_bigint_contract.py new file mode 100644 index 000000000..7c85e10a3 --- /dev/null +++ b/tests/test_checkpoint_audit_bigint_contract.py @@ -0,0 +1,40 @@ +# SPDX-License-Identifier: Apache-2.0 +"""PostgreSQL BIGINT compatibility tests for public checkpoint-audit identities.""" + +from __future__ import annotations + +from datetime import datetime, timezone + +import pytest + +from pg_llm_batch.checkpoint_audit import ( + MAX_CHECKPOINT_AUDIT_EVENT_ID, + CheckpointAuditEvent, +) + + +def _event(event_id: int) -> CheckpointAuditEvent: + """Build one otherwise-valid accepted-save event with a selected identity.""" + return CheckpointAuditEvent( + audit_event_id=event_id, + tenant_scope="tenant-a", + consumer_name="worker-a", + endpoint_alias="default", + batch_id="batch-1", + action="checkpoint_save_accepted", + schema_version=1, + file_kind="result", + file_id="file-1", + file_line_number=1, + batch_line_count=1, + record_count=1, + prefix_sha256="a" * 64, + recorded_at=datetime(2026, 8, 7, tzinfo=timezone.utc), + ) + + +def test_public_audit_event_identity_matches_postgresql_bigint_domain() -> None: + """Direct event construction cannot represent an identity PostgreSQL cannot store.""" + assert _event(MAX_CHECKPOINT_AUDIT_EVENT_ID).audit_event_id == MAX_CHECKPOINT_AUDIT_EVENT_ID + with pytest.raises(ValueError, match="PostgreSQL BIGINT"): + _event(MAX_CHECKPOINT_AUDIT_EVENT_ID + 1) diff --git a/tests/test_checkpoint_audit_integration.py b/tests/test_checkpoint_audit_integration.py index 4660a945c..78d46ab33 100644 --- a/tests/test_checkpoint_audit_integration.py +++ b/tests/test_checkpoint_audit_integration.py @@ -145,9 +145,7 @@ def test_live_audit_is_tenant_isolated_append_only_and_rollback_safe() -> None: raise RuntimeError("audit identity sequence could not be resolved") audit_sequence_parts = tuple(sequence_row[0]) cursor.execute( - sql.SQL( - "GRANT USAGE, SELECT ON SEQUENCE {} TO {}" - ).format( + sql.SQL("GRANT USAGE, SELECT ON SEQUENCE {} TO {}").format( sql.Identifier(*audit_sequence_parts), sql.Identifier(application_role_name), ) @@ -182,6 +180,47 @@ def test_live_audit_is_tenant_isolated_append_only_and_rollback_safe() -> None: assert {event.tenant_scope for event in events_a} == {"tenant-a"} assert {event.tenant_scope for event in events_b} == {"tenant-b"} + # Real separate transactions exercise the operational export case: page + # one is read, a newer accepted-save event commits, and page two continues + # strictly toward older identities without OFFSET drift. + assert tenant_a.save(consumer, first_a) == first_a + assert tenant_a.save(consumer, first_a) == first_a + first_page = tenant_a.list_audit_event_page( + consumer, + batch_id, + "default", + limit=2, + ) + first_page_ids = tuple(event.audit_event_id for event in first_page.events) + assert len(first_page_ids) == 2 + assert first_page_ids[0] > first_page_ids[1] + assert first_page.next_before_audit_event_id == first_page_ids[-1] + + assert tenant_a.save(consumer, first_a) == first_a + newest_after_first_page = tenant_a.list_audit_events( + consumer, + batch_id, + "default", + limit=1, + )[0].audit_event_id + assert newest_after_first_page > first_page_ids[0] + + second_page = tenant_a.list_audit_event_page( + consumer, + batch_id, + "default", + before_audit_event_id=first_page.next_before_audit_event_id, + limit=2, + ) + second_page_ids = tuple(event.audit_event_id for event in second_page.events) + assert len(second_page_ids) == 2 + assert set(first_page_ids).isdisjoint(second_page_ids) + assert newest_after_first_page not in second_page_ids + assert all( + event_id < first_page.next_before_audit_event_id + for event_id in second_page_ids + ) + timed_batch_id = f"timed-{suffix}" timed_checkpoint = _checkpoint(timed_batch_id, "c" * 64) with psycopg.connect(application_role_dsn) as timed_connection: @@ -258,7 +297,7 @@ def test_live_audit_is_tenant_isolated_append_only_and_rollback_safe() -> None: assert rollback_error.value.sqlstate == "55000" database_admin.rollback() - assert len(tenant_a.list_audit_events(consumer, batch_id, "default")) == 2 + assert len(tenant_a.list_audit_events(consumer, batch_id, "default")) == 5 assert len(tenant_b.list_audit_events(consumer, batch_id, "default")) == 1 finally: with psycopg.connect(ADMIN_DSN, autocommit=True) as cluster_admin: diff --git a/tests/test_checkpoint_audit_pagination.py b/tests/test_checkpoint_audit_pagination.py new file mode 100644 index 000000000..e00ee5467 --- /dev/null +++ b/tests/test_checkpoint_audit_pagination.py @@ -0,0 +1,331 @@ +# SPDX-License-Identifier: Apache-2.0 +"""Behavior contracts for bounded stable checkpoint-audit export pages.""" + +from __future__ import annotations + +from datetime import datetime, timezone +from typing import Any + +import pytest + +import pg_llm_batch.checkpoint_audit as checkpoint_audit +from pg_llm_batch.checkpoint_audit import ( + MAX_CHECKPOINT_AUDIT_EVENT_ID, + MAX_CHECKPOINT_AUDIT_EVENTS, + CheckpointAuditEvent, + CheckpointAuditPage, + validate_checkpoint_audit_cursor, +) +from pg_llm_batch.exceptions import ValidationError + + +def _event(event_id: int) -> CheckpointAuditEvent: + """Build one valid audit event with a selectable monotonically increasing identity.""" + return CheckpointAuditEvent( + audit_event_id=event_id, + tenant_scope="tenant-a", + consumer_name="worker-a", + endpoint_alias="default", + batch_id="batch-1", + action="checkpoint_save_accepted", + schema_version=1, + file_kind="result", + file_id="file-1", + file_line_number=event_id, + batch_line_count=event_id, + record_count=event_id, + prefix_sha256=f"{event_id:064x}", + recorded_at=datetime(2026, 8, 7, event_id % 24, tzinfo=timezone.utc), + ) + + +def _row(event_id: int) -> tuple[Any, ...]: + """Return one database-shaped row matching :func:`_event`.""" + event = _event(event_id) + return ( + event.audit_event_id, + event.tenant_scope, + event.consumer_name, + event.endpoint_alias, + event.batch_id, + event.action, + event.schema_version, + event.file_kind, + event.file_id, + event.file_line_number, + event.batch_line_count, + event.record_count, + event.prefix_sha256, + event.recorded_at, + ) + + +class FakeCursor: + """Capture audit-page SQL and return deterministic rows.""" + + def __init__(self, rows: Any) -> None: + self.rows = rows + self.calls: list[tuple[str, tuple[Any, ...]]] = [] + + 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: + """Capture one normalized SQL statement and its parameters.""" + self.calls.append((" ".join(sql.split()), params or ())) + + def fetchall(self) -> Any: + """Return the configured database rows.""" + return self.rows + + +class FakeConnection: + """Expose one cursor without inventing a commit boundary for audit reads.""" + + def __init__(self, cursor: FakeCursor) -> None: + self.fake_cursor = cursor + self.commits = 0 + + 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: + """Return the deterministic fake cursor.""" + return self.fake_cursor + + def commit(self) -> None: + """Record an unexpected explicit package commit if production attempts one.""" + self.commits += 1 + + +class FakePsycopg: + """Return one deterministic connection and record requested DSNs.""" + + def __init__(self, connection: FakeConnection) -> None: + self.connection = connection + self.dsns: list[str] = [] + + def connect(self, dsn: str) -> FakeConnection: + """Record one DSN and return the configured connection.""" + self.dsns.append(dsn) + return self.connection + + +def test_audit_cursor_is_strict_positive_postgres_bigint_or_none() -> None: + """Export cursors reject coercion and values PostgreSQL cannot represent.""" + assert validate_checkpoint_audit_cursor(None) is None + assert validate_checkpoint_audit_cursor(1) == 1 + assert ( + validate_checkpoint_audit_cursor(MAX_CHECKPOINT_AUDIT_EVENT_ID) + == 9_223_372_036_854_775_807 + ) + for value in (True, 0, -1, 9_223_372_036_854_775_808, 1.0, "7", object()): + with pytest.raises(ValidationError): + validate_checkpoint_audit_cursor(value) + + +def test_audit_page_is_immutable_descending_and_cursor_bound() -> None: + """A public page cannot advertise an invalid continuation boundary.""" + first = _event(9) + second = _event(8) + page = CheckpointAuditPage(events=(first, second), next_before_audit_event_id=8) + assert page.events == (first, second) + assert page.next_before_audit_event_id == 8 + + invalid_pages = ( + ((second, first), None), + ((first, first), None), + ((first, second), 7), + ((), 1), + ) + for events, cursor in invalid_pages: + with pytest.raises(ValueError): + CheckpointAuditPage(events=events, next_before_audit_event_id=cursor) + + +def test_audit_page_rejects_mutable_wrong_type_and_oversized_event_collections() -> None: + """Direct page construction cannot bypass tuple, type, or public-size contracts.""" + with pytest.raises(ValueError, match="tuple"): + CheckpointAuditPage(events=[_event(1)], next_before_audit_event_id=None) # type: ignore[arg-type] + with pytest.raises(ValueError, match="CheckpointAuditEvent"): + CheckpointAuditPage(events=(object(),), next_before_audit_event_id=None) # type: ignore[arg-type] + oversized = tuple(_event(index + 1) for index in range(MAX_CHECKPOINT_AUDIT_EVENTS + 1)) + with pytest.raises(ValueError, match="at most"): + CheckpointAuditPage(events=oversized, next_before_audit_event_id=None) + + +def test_first_export_page_uses_one_bounded_lookahead_row( + monkeypatch: pytest.MonkeyPatch, +) -> None: + """A first page returns at most the requested rows and one stable next cursor.""" + cursor = FakeCursor(rows=(_row(9), _row(8), _row(7))) + scopes: list[str] = [] + monkeypatch.setattr( + checkpoint_audit, + "_set_transaction_tenant_scope", + lambda _cursor, tenant: scopes.append(tenant), + ) + store = checkpoint_audit.AuditedPostgresBatchResultCheckpointStore( + "postgresql://unit", + tenant_scope="tenant-a", + ) + + page = store.list_audit_event_page_in_transaction( + cursor, + "worker-a", + "batch-1", + "default", + limit=2, + ) + + assert tuple(event.audit_event_id for event in page.events) == (9, 8) + assert page.next_before_audit_event_id == 8 + assert scopes == ["tenant-a"] + query, params = cursor.calls[-1] + assert "checkpoint_audit_event_id < %s" not in query + assert "ORDER BY checkpoint_audit_event_id DESC LIMIT %s" in query + assert "OFFSET" not in query + assert params == ("tenant-a", "worker-a", "default", "batch-1", 3) + + +def test_next_export_page_is_keyset_bound_and_ignores_newer_concurrent_rows( + monkeypatch: pytest.MonkeyPatch, +) -> None: + """A continuation cursor advances toward older rows without OFFSET drift.""" + cursor = FakeCursor(rows=(_row(7), _row(6))) + monkeypatch.setattr(checkpoint_audit, "_set_transaction_tenant_scope", lambda *_: None) + store = checkpoint_audit.AuditedPostgresBatchResultCheckpointStore( + "postgresql://unit", + tenant_scope="tenant-a", + ) + + page = store.list_audit_event_page_in_transaction( + cursor, + "worker-a", + "batch-1", + "default", + before_audit_event_id=8, + limit=2, + ) + + assert tuple(event.audit_event_id for event in page.events) == (7, 6) + assert page.next_before_audit_event_id is None + query, params = cursor.calls[-1] + assert "checkpoint_audit_event_id < %s" in query + assert "OFFSET" not in query + assert params == ("tenant-a", "worker-a", "default", "batch-1", 8, 3) + + +def test_continuation_page_revalidates_returned_identity_against_cursor( + monkeypatch: pytest.MonkeyPatch, +) -> None: + """A faulty adapter cannot return an event at or newer than the strict cursor.""" + monkeypatch.setattr(checkpoint_audit, "_set_transaction_tenant_scope", lambda *_: None) + store = checkpoint_audit.AuditedPostgresBatchResultCheckpointStore( + "postgresql://unit", + tenant_scope="tenant-a", + ) + + with pytest.raises(RuntimeError, match="violated the continuation cursor"): + store.list_audit_event_page_in_transaction( + FakeCursor(rows=(_row(8), _row(7))), + "worker-a", + "batch-1", + "default", + before_audit_event_id=8, + limit=2, + ) + + +def test_export_page_revalidates_database_key_and_descending_order( + monkeypatch: pytest.MonkeyPatch, +) -> None: + """Corrupt or cross-key database rows cannot become trusted export evidence.""" + monkeypatch.setattr(checkpoint_audit, "_set_transaction_tenant_scope", lambda *_: None) + store = checkpoint_audit.AuditedPostgresBatchResultCheckpointStore( + "postgresql://unit", + tenant_scope="tenant-a", + ) + + wrong_key = list(_row(9)) + wrong_key[1] = "tenant-b" + with pytest.raises(RuntimeError, match="outside the requested key"): + store.list_audit_event_page_in_transaction( + FakeCursor(rows=(tuple(wrong_key),)), + "worker-a", + "batch-1", + "default", + ) + + with pytest.raises(RuntimeError, match="strictly descending"): + store.list_audit_event_page_in_transaction( + FakeCursor(rows=(_row(8), _row(9))), + "worker-a", + "batch-1", + "default", + ) + + +def test_export_page_fails_closed_on_invalid_collection_or_driver_overrun( + monkeypatch: pytest.MonkeyPatch, +) -> None: + """Unexpected driver output cannot bypass bounded export pagination.""" + monkeypatch.setattr(checkpoint_audit, "_set_transaction_tenant_scope", lambda *_: None) + store = checkpoint_audit.AuditedPostgresBatchResultCheckpointStore("postgresql://unit") + + with pytest.raises(RuntimeError, match="invalid row collection"): + store.list_audit_event_page_in_transaction( + FakeCursor(rows=object()), + "worker-a", + "batch-1", + "default", + limit=2, + ) + + with pytest.raises(RuntimeError, match="exceeded its bounded query size"): + store.list_audit_event_page_in_transaction( + FakeCursor(rows=(_row(9), _row(8), _row(7), _row(6))), + "worker-a", + "batch-1", + "default", + limit=2, + ) + + +def test_owned_export_page_uses_one_connection_without_explicit_commit( + monkeypatch: pytest.MonkeyPatch, +) -> None: + """The package-owned read delegates through one connection without calling commit().""" + cursor = FakeCursor(rows=(_row(3),)) + connection = FakeConnection(cursor) + fake_psycopg = FakePsycopg(connection) + monkeypatch.setattr(checkpoint_audit, "psycopg", fake_psycopg) + monkeypatch.setattr(checkpoint_audit, "_require_psycopg", lambda: None) + monkeypatch.setattr(checkpoint_audit, "_set_transaction_tenant_scope", lambda *_: None) + store = checkpoint_audit.AuditedPostgresBatchResultCheckpointStore( + "postgresql://unit", + tenant_scope="tenant-a", + ) + + page = store.list_audit_event_page( + "worker-a", + "batch-1", + "default", + before_audit_event_id=4, + limit=2, + ) + + assert tuple(event.audit_event_id for event in page.events) == (3,) + assert page.next_before_audit_event_id is None + assert fake_psycopg.dsns == ["postgresql://unit"] + assert connection.commits == 0