From b0f02f52c95519de17bfdfd3202e28b7483e1c1b Mon Sep 17 00:00:00 2001 From: Cursor Agent Date: Sun, 16 Aug 2026 16:10:39 +0000 Subject: [PATCH 1/6] feat(response): rebuild ledgers from stored event order A two-item path can now reload the same answers after restart without re-checking live session activity. Gapped sequence, reused client or server identity, and blank session references fail closed. Co-authored-by: Seongho Bae --- src/response.rs | 102 +++++++++++++++++++++++++++ tests/response_ledger.rs | 144 ++++++++++++++++++++++++++++++++++++++- 2 files changed, 245 insertions(+), 1 deletion(-) diff --git a/src/response.rs b/src/response.rs index b5ff00cf..d41bd4ad 100644 --- a/src/response.rs +++ b/src/response.rs @@ -65,6 +65,49 @@ impl ResponseEvent { pub const fn sequence(&self) -> usize { self.sequence } + + /// Rebuild one accepted event from durable store columns. + /// + /// Use this after a process restart. The caller must still assemble events + /// into [`ResponseLedger::from_persisted`] so sequence `1..n` is checked. + /// + /// # Errors + /// + /// Returns [`WriteError::InvalidReference`] for blank or numeric-like + /// identities, [`WriteError::EmptyReference`] for a blank digest, + /// [`WriteError::InvalidPayloadDigest`] for a noncanonical digest, or + /// [`WriteError::InvalidStoredSequence`] when `sequence` is zero. + pub fn from_persisted( + server_event_ref: impl AsRef, + client_event_ref: impl AsRef, + item_version_ref: impl AsRef, + payload_digest: impl AsRef, + sequence: usize, + ) -> Result { + let server_event_ref = + normalized_reference(server_event_ref.as_ref()).ok_or(WriteError::InvalidReference)?; + let client_event_ref = + normalized_reference(client_event_ref.as_ref()).ok_or(WriteError::InvalidReference)?; + let item_version_ref = + normalized_reference(item_version_ref.as_ref()).ok_or(WriteError::InvalidReference)?; + let payload_digest = payload_digest.as_ref(); + if payload_digest.trim().is_empty() { + return Err(WriteError::EmptyReference); + } + if !is_canonical_sha256(payload_digest) { + return Err(WriteError::InvalidPayloadDigest); + } + if sequence == 0 { + return Err(WriteError::InvalidStoredSequence); + } + Ok(Self { + server_event_ref: server_event_ref.to_owned(), + client_event_ref: client_event_ref.to_owned(), + item_version_ref: item_version_ref.to_owned(), + payload_digest: payload_digest.to_owned(), + sequence, + }) + } } /// Immutable response snapshot frozen when collection completes. @@ -140,6 +183,8 @@ pub enum WriteError { ServerReferenceConflict, /// A response snapshot was requested before the session reached completion. SnapshotRequiresCompleted(SessionState), + /// Stored events were missing, gapped, or rewound relative to server sequence `1..n`. + InvalidStoredSequence, } impl Display for WriteError { @@ -164,6 +209,9 @@ impl Display for WriteError { formatter, "response snapshot requires Completed session state, found {state:?}" ), + Self::InvalidStoredSequence => formatter.write_str( + "stored response events must keep server sequence 1..n without gaps", + ), } } } @@ -212,6 +260,60 @@ impl ResponseLedger { self.events.is_empty() } + /// Return the opaque assessment-session reference bound to this ledger. + #[must_use] + pub fn session_ref(&self) -> &str { + &self.session_ref + } + + /// Return accepted response events in server-authoritative order. + #[must_use] + pub fn events(&self) -> &[ResponseEvent] { + &self.events + } + + /// Rebuild a ledger from durable events after process restart. + /// + /// Events must already be valid identities with server sequence `1..n` and + /// no reused client or server references. This does not re-check whether + /// the live session is still Active; stored answers stay stored. + /// + /// # Errors + /// + /// Returns [`WriteError::InvalidReference`] for a blank or numeric-like + /// session, [`WriteError::InvalidStoredSequence`] when sequences are not + /// exactly `1..n` in order, [`WriteError::IdempotencyConflict`] when a + /// client reference repeats, or [`WriteError::ServerReferenceConflict`] + /// when a server event reference repeats. + pub fn from_persisted( + session_ref: impl AsRef, + events: Vec, + ) -> Result { + let session_ref = + normalized_reference(session_ref.as_ref()).ok_or(WriteError::InvalidReference)?; + for (index, event) in events.iter().enumerate() { + if event.sequence != index + 1 { + return Err(WriteError::InvalidStoredSequence); + } + if events[..index] + .iter() + .any(|prior| prior.client_event_ref == event.client_event_ref) + { + return Err(WriteError::IdempotencyConflict); + } + if events[..index] + .iter() + .any(|prior| prior.server_event_ref == event.server_event_ref) + { + return Err(WriteError::ServerReferenceConflict); + } + } + Ok(Self { + session_ref: session_ref.to_owned(), + events, + }) + } + /// Record one response event or replay an identical prior event. /// /// Exact replay of an already accepted `client_event_ref` remains idempotent diff --git a/tests/response_ledger.rs b/tests/response_ledger.rs index 1aec771c..a36467da 100644 --- a/tests/response_ledger.rs +++ b/tests/response_ledger.rs @@ -1,6 +1,8 @@ //! Integration tests for response-event idempotency and immutable snapshots. -use psychometrics_commons_runtime::response::{ResponseLedger, ResponseWrite, WriteError}; +use psychometrics_commons_runtime::response::{ + ResponseEvent, ResponseLedger, ResponseWrite, WriteError, +}; use psychometrics_commons_runtime::session::SessionState; const DIGEST_A: &str = "sha256:aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa"; @@ -335,4 +337,144 @@ fn write_errors_have_stable_human_readable_context() { WriteError::SnapshotRequiresCompleted(SessionState::Active).to_string(), "response snapshot requires Completed session state, found Active" ); + assert_eq!( + WriteError::InvalidStoredSequence.to_string(), + "stored response events must keep server sequence 1..n without gaps" + ); +} + +#[test] +fn two_item_korean_path_reloads_the_same_answers_after_restart() { + let mut live = ResponseLedger::new("session_big_five_ko").unwrap(); + live.record( + SessionState::Active, + write( + "response_event_openness", + "client_openness", + "item_version_openness_ko", + DIGEST_A, + ), + ) + .unwrap(); + live.record( + SessionState::Active, + write( + "response_event_conscientiousness", + "client_conscientiousness", + "item_version_conscientiousness_ko", + DIGEST_B, + ), + ) + .unwrap(); + + assert_eq!(live.session_ref(), "session_big_five_ko"); + assert_eq!(live.events().len(), 2); + assert_eq!(live.events()[0].sequence(), 1); + assert_eq!(live.events()[1].item_version_ref(), "item_version_conscientiousness_ko"); + + let reloaded = ResponseLedger::from_persisted(live.session_ref(), live.events().to_vec()).unwrap(); + assert_eq!(reloaded, live); + assert_eq!( + reloaded.freeze_as(SessionState::Completed, "response_snapshot_big_five_ko") + .unwrap() + .event_refs(), + ["response_event_openness", "response_event_conscientiousness"] + ); +} + +#[test] +fn persisted_events_reject_rewound_or_gapped_sequence() { + let first = ResponseEvent::from_persisted( + "response_event_openness", + "client_openness", + "item_version_openness_ko", + DIGEST_A, + 2, + ) + .unwrap(); + assert_eq!( + ResponseLedger::from_persisted("session_big_five_ko", vec![first]).unwrap_err(), + WriteError::InvalidStoredSequence + ); + assert_eq!( + ResponseEvent::from_persisted( + "response_event_openness", + "client_openness", + "item_version_openness_ko", + DIGEST_A, + 0, + ) + .unwrap_err(), + WriteError::InvalidStoredSequence + ); +} + +#[test] +fn persisted_events_reject_reused_identities_and_blank_session() { + let first = ResponseEvent::from_persisted( + "response_event_openness", + "client_openness", + "item_version_openness_ko", + DIGEST_A, + 1, + ) + .unwrap(); + let duplicate_client = ResponseEvent::from_persisted( + "response_event_other", + "client_openness", + "item_version_conscientiousness_ko", + DIGEST_B, + 2, + ) + .unwrap(); + let duplicate_server = ResponseEvent::from_persisted( + "response_event_openness", + "client_conscientiousness", + "item_version_conscientiousness_ko", + DIGEST_B, + 2, + ) + .unwrap(); + + assert_eq!( + ResponseLedger::from_persisted("session_big_five_ko", vec![first.clone(), duplicate_client]) + .unwrap_err(), + WriteError::IdempotencyConflict + ); + assert_eq!( + ResponseLedger::from_persisted("session_big_five_ko", vec![first, duplicate_server]) + .unwrap_err(), + WriteError::ServerReferenceConflict + ); + assert_eq!( + ResponseLedger::from_persisted("12", Vec::new()).unwrap_err(), + WriteError::InvalidReference + ); + assert_eq!( + ResponseEvent::from_persisted("12", "client_openness", "item_version_o", DIGEST_A, 1) + .unwrap_err(), + WriteError::InvalidReference + ); + assert_eq!( + ResponseEvent::from_persisted( + "response_event_openness", + "client_openness", + "item_version_openness_ko", + " ", + 1, + ) + .unwrap_err(), + WriteError::EmptyReference + ); + assert_eq!( + ResponseEvent::from_persisted( + "response_event_openness", + "client_openness", + "item_version_openness_ko", + "sha256:not-canonical", + 1, + ) + .unwrap_err(), + WriteError::InvalidPayloadDigest + ); } From f6c5f6bc88ac74cbf605e266619d046cfff1a817 Mon Sep 17 00:00:00 2001 From: Cursor Agent Date: Sun, 16 Aug 2026 16:12:47 +0000 Subject: [PATCH 2/6] feat(response): persist in-progress events across restart Store accepted response_event rows with observed and received time so a two-item path reloads the same answers after process restart. Exact replay is idempotent; client, server, sequence, or session rebinding fails closed under READ COMMITTED. Co-authored-by: Seongho Bae --- CHANGELOG.md | 1 + docs/TRACEABILITY.md | 6 +- ...-persistence-and-transaction-boundaries.md | 10 +- docs/architecture/ERD.md | 1 + migrations/0020_response_event.sql | 53 +++ src/lib.rs | 1 + src/postgres_response_event.rs | 347 ++++++++++++++++++ .../postgres_response_event_error_contract.rs | 61 +++ tests/postgres_response_event_persistence.rs | 334 +++++++++++++++++ 9 files changed, 808 insertions(+), 6 deletions(-) create mode 100644 migrations/0020_response_event.sql create mode 100644 src/postgres_response_event.rs create mode 100644 tests/postgres_response_event_error_contract.rs create mode 100644 tests/postgres_response_event_persistence.rs diff --git a/CHANGELOG.md b/CHANGELOG.md index b2f801dd..a06af678 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -5,6 +5,7 @@ All notable product and architecture changes are recorded here. Releases use imm ## Unreleased ### Added +- PostgreSQL persistence for in-progress `response_event` rows so a two-item path reloads the same answers after restart. Exact replay is idempotent; client, server, sequence, or session rebinding fails closed. Store observed and received times with the event; completed snapshots stay on the existing snapshot adapter. - Scoring-job cancel and lease-expiry fallback classification lock the current row until the caller transaction ends, so concurrent workers cannot rewrite terminal or unleased evidence. - PostgreSQL operational-store readiness probe classifies the supported major version and write-readiness, and fails closed when a caller-declared required relation is missing. - PostgreSQL scoring-job cancellation: queued, leased, or retry-scheduled work becomes cancelled without transferring a fence, exact replay is idempotent, and completed or quarantined evidence cannot be rewritten. diff --git a/docs/TRACEABILITY.md b/docs/TRACEABILITY.md index 72bc73c2..ceef7951 100644 --- a/docs/TRACEABILITY.md +++ b/docs/TRACEABILITY.md @@ -23,7 +23,7 @@ An active PR, architecture document, conversation decision, or scheduler plan is | Anonymous core assessment | PRD §3.1, §9.1 | TRD §5, §10; UML anonymous sequence | ADR-0002, ADR-0003, ADR-0005 | Session lifecycle primitives implemented, including creation bound to one published locale-specific release; anonymous credential/HTTP flow is Target | | Pause/resume | PRD §3.1, §9.1 | TRD §5 | ADR-0005 | **Implemented** in `src/session.rs` with fail-closed transitions | | Sequence-aware item delivery evidence | PRD §3.1, §9 | TRD §5–7 | ADR-0005, ADR-0010 | **Implemented** domain primitive in `src/item_delivery.rs`; persistence/API delivery orchestration is Target | -| Idempotent response events | PRD §9.2 | TRD §6 | ADR-0005, ADR-0010 | **Implemented** in `src/response.rs` with canonical SHA-256 payload-digest identity; persistence adapter is Target | +| Idempotent response events | PRD §9.2 | TRD §6 | ADR-0005, ADR-0010 | **Implemented** in `src/response.rs` with canonical SHA-256 payload-digest identity; **Active PR** on this branch persists and reloads `response_event` under `READ COMMITTED`; HTTP response transport remains Target | | Immutable response snapshot before scoring | PRD §9.3 | TRD §5–8 | ADR-0005, ADR-0010 | **Implemented** domain semantics in `src/response.rs` | | Version-pinned scoring | PRD §9.4, §10 | TRD §8 | ADR-0004, ADR-0010 | **Implemented** reusable product-side scoring dispatch contract in `src/scoring.rs` with canonical SHA-256 engine-artifact digest provenance plus `migrations/0011_scoring_request.sql` / `src/postgres_scoring_request.rs` request-identity persistence; live fast-mlsirm integration is Target | | Bounded asynchronous scoring retry/quarantine with stale-worker fencing | PRD §9.4, §10 | TRD §8; ADR-0015 transaction boundary | ADR-0004, ADR-0010, ADR-0015 | **Implemented** product lifecycle plus PostgreSQL enqueue, claim, retry, completion, expiry recovery, and cancellation without transferring a fence; live fast-mlsirm execution remains Target | @@ -55,7 +55,7 @@ An active PR, architecture document, conversation decision, or scheduler plan is | Server-authoritative session state | TRD §5 | `src/session.rs` + session contract tests, including published-release/locale binding at creation | persistence/API concurrency test | | Only Active accepts responses | TRD §5–6 | `SessionState::accepts_responses` + response tests | transport-level rejection test | | Item delivery sequence is positive and evidence-safe | TRD §5–7 | `src/item_delivery.rs` + item-delivery domain tests | durable uniqueness/order/API integration | -| Conflicting idempotency replay fails closed | TRD §6 | `src/response.rs` | DB uniqueness/concurrency test | +| Conflicting idempotency replay fails closed | TRD §6 | `src/response.rs`; **Active PR** `migrations/0020_response_event.sql` unique `(session_ref, client_event_ref)` and `(session_ref, server_sequence)` | HTTP transport rejection test | | Snapshot requires Completed state | TRD §5–6 | `src/response.rs` | transaction atomicity test with persistence | | Scoring uses durable snapshot identity | TRD §8 | `src/scoring.rs` requires a canonical SHA-256 engine-artifact digest | live adapter + retry/outbox integration | | Stale scoring worker cannot complete a newer attempt | TRD §8; ADR-0015 | `src/scoring_job.rs` uses monotonically increasing fencing tokens and rejects stale/expired completion or failure evidence; `src/postgres_scoring_job.rs` persists enqueue, claim, retry, terminal outcomes, expired-lease recovery, and cancellation without transferring a fence | live adapter evidence | @@ -132,7 +132,7 @@ Still-Target logical modules/adapters include remaining product aggregate persis ### Active implementation work that is not protected-main truth -**Active PR** #76 data-rights processing-start persistence is not protected-main truth until an unchanged reviewed/check-clean head is integrated. Identity-verified requests persist an immutable operation identity and processing-start time under `FOR UPDATE` so later lifecycle composition cannot race the classified row. Dependent-system execution remains outside this slice. +**Active PR** response-event persist/load on this branch is not protected-main truth until an unchanged reviewed/check-clean head is integrated. In-progress answers persist in `response_event` with observed/received time, exact replay, and fail-closed rebinding so a two-item path survives restart. Completed snapshot persist remains the existing protected-main adapter. HTTP response transport remains outside this slice. Protected-main `#76` data-rights processing-start persistence is already integrated and is not restated here as active work. ## 5. ADR traceability by concern diff --git a/docs/adr/0015-persistence-and-transaction-boundaries.md b/docs/adr/0015-persistence-and-transaction-boundaries.md index c6643000..332b9ea0 100644 --- a/docs/adr/0015-persistence-and-transaction-boundaries.md +++ b/docs/adr/0015-persistence-and-transaction-boundaries.md @@ -6,9 +6,9 @@ - Scope: Psychometrics Commons-owned durable state, local transactions, migration boundaries, outbox/inbox integration - Supersedes: none - Superseded by: none -- Current/as-built status: protected main contains in-memory/domain lifecycle primitives only; active PR #24 carries the first PostgreSQL integration-evidence migration/adapter but is not protected-main truth until merged +- Current/as-built status: protected main persists integration, scoring-job, data-rights, item-delivery, consent, instrument-release, result-snapshot, response-snapshot, scoring-request, and inbox-consumption slices; in-progress `response_event` persist/load is Active PR work on this branch and is not protected-main truth until merged - Target status: upstream PostgreSQL 18.x operational persistence with real-database concurrency/crash/recovery evidence and transactional outbox/inbox semantics -- Migration status: active PR #24 introduces only the bounded integration-evidence slice; the remaining product schema still must be established from the logical ERD and this ADR without synthetic provenance backfills +- Migration status: `migrations/0020_response_event.sql` adds the in-progress event ledger with observed/received timestamps; remaining session HTTP and response HTTP families still must be established from the logical ERD without synthetic provenance backfills ## Context @@ -76,7 +76,7 @@ A physical migration may split an entity across tables or co-locate value object ### Response recording -One transaction validates the current session state, reserves server sequence, applies the idempotency/uniqueness contract, and stores the accepted response event. Two concurrent requests cannot both create the same logical `client_event_ref`. +One transaction validates the current session state, reserves server sequence, applies the idempotency/uniqueness contract, and stores the accepted response event with distinct observed and received instants (ISO 8601-1). Two concurrent requests cannot both create the same logical `client_event_ref`. The first physical classifier uses `READ COMMITTED` so a concurrent unique-key winner is visible to the exact-replay inspection (Berenson et al., 1995; PostgreSQL Global Development Group, 2026). ### Session completion @@ -303,6 +303,10 @@ The physical database technology or decomposition may change if scale, residency ## References +Berenson, H., Bernstein, P., Gray, J., Melton, J., O'Neil, E., & O'Neil, P. (1995). A critique of ANSI SQL isolation levels. *ACM SIGMOD Record, 24*(2), 1–10. https://doi.org/10.1145/568271.223785 + +International Organization for Standardization. (2019). *Date and time — Representations for information interchange — Part 1: Basic rules* (ISO 8601-1:2019). + PostgreSQL Global Development Group. (2026). *PostgreSQL 18 documentation*. PostgreSQL Global Development Group. (2026). *PostgreSQL versioning policy*. diff --git a/docs/architecture/ERD.md b/docs/architecture/ERD.md index 8f21954a..0215d8be 100644 --- a/docs/architecture/ERD.md +++ b/docs/architecture/ERD.md @@ -425,6 +425,7 @@ The target ERD deliberately includes several logical entities that are not yet p - `instrument_release` is the locale-specific publication identity already owned by `src/instrument.rs`. Physical `migrations/0006_instrument_release.sql` persists that one-row aggregate (immutable manifest columns plus `publication_state`); HTTP publication transport remains Target. - `data_rights_request` and `data_rights_propagation_state` are the first durable export/deletion slice. Physical `migrations/0003_data_rights_propagation.sql` stores requested-state identity plus one local outbox event per dependent system; verification, processing, completion, and dependent-system execution remain Target. - `item_delivery_event` reflects the already-merged `src/item_delivery.rs` domain primitive; durable persistence/API orchestration is still Target. +- `response_event` is the in-progress answer ledger. Physical `migrations/0020_response_event.sql` on this Active PR stores opaque event identity, session binding, client idempotency, item version, payload digest, server sequence, and distinct observed/received timestamps; HTTP response transport remains Target. Completed `response_snapshot` persist is already on protected main. - `consent_ledger` and `consent_event` persist the already-merged `src/consent.rs` append-only ledger. Physical persistence is carried by Active PR #49 (`migrations/0005_consent_lifecycle.sql`); HTTP consent transport and derived snapshot tables remain Target. - `participant_identity_link` is the persistence target accepted by ADR-0020. The current `src/participant.rs` `keyverse_subject_ref` field is an application-domain first-link projection, not the future mutable persistence source of truth. - `longitudinal_enrollment`, `longitudinal_observation_record`, and `temporal_analysis_submission` make the ADR-0008 Commons-owned Gyeot/TEPP orchestration boundary explicit. No TEPP analytical kernel is duplicated here. diff --git a/migrations/0020_response_event.sql b/migrations/0020_response_event.sql new file mode 100644 index 00000000..60418aa9 --- /dev/null +++ b/migrations/0020_response_event.sql @@ -0,0 +1,53 @@ +-- Opaque response-event references remain data, not SQL syntax. The persistence +-- adapter binds them as query parameters. These checks enforce the canonical +-- identity boundary (nonblank, nonnumeric-like, no outer whitespace). +CREATE TABLE IF NOT EXISTS response_event ( + response_event_ref TEXT CONSTRAINT response_event_response_event_ref_not_null NOT NULL + CONSTRAINT response_event_response_event_ref_format_check CHECK ( + response_event_ref = btrim(response_event_ref) + AND response_event_ref <> '' + AND NOT ( + response_event_ref ~ '[[:digit:]]' + AND response_event_ref ~ '^[[:digit:]+,.eE-]+$' + ) + ), + session_ref TEXT CONSTRAINT response_event_session_ref_not_null NOT NULL + CONSTRAINT response_event_session_ref_format_check CHECK ( + session_ref = btrim(session_ref) + AND session_ref <> '' + AND NOT ( + session_ref ~ '[[:digit:]]' + AND session_ref ~ '^[[:digit:]+,.eE-]+$' + ) + ), + client_event_ref TEXT CONSTRAINT response_event_client_event_ref_not_null NOT NULL + CONSTRAINT response_event_client_event_ref_format_check CHECK ( + client_event_ref = btrim(client_event_ref) + AND client_event_ref <> '' + AND NOT ( + client_event_ref ~ '[[:digit:]]' + AND client_event_ref ~ '^[[:digit:]+,.eE-]+$' + ) + ), + item_version_ref TEXT CONSTRAINT response_event_item_version_ref_not_null NOT NULL + CONSTRAINT response_event_item_version_ref_format_check CHECK ( + item_version_ref = btrim(item_version_ref) + AND item_version_ref <> '' + AND NOT ( + item_version_ref ~ '[[:digit:]]' + AND item_version_ref ~ '^[[:digit:]+,.eE-]+$' + ) + ), + payload_digest TEXT CONSTRAINT response_event_payload_digest_not_null NOT NULL + CONSTRAINT response_event_payload_digest_format_check CHECK ( + payload_digest ~ '^sha256:[0-9a-f]{64}$' + ), + server_sequence BIGINT CONSTRAINT response_event_server_sequence_not_null NOT NULL + CONSTRAINT response_event_server_sequence_positive_check CHECK (server_sequence > 0), + observed_at TIMESTAMPTZ CONSTRAINT response_event_observed_at_not_null NOT NULL, + received_at TIMESTAMPTZ CONSTRAINT response_event_received_at_not_null NOT NULL, + CONSTRAINT response_event_pkey PRIMARY KEY (response_event_ref), + CONSTRAINT response_event_session_client_unique UNIQUE (session_ref, client_event_ref), + CONSTRAINT response_event_session_sequence_unique UNIQUE (session_ref, server_sequence), + CONSTRAINT response_event_observed_not_after_received_check CHECK (observed_at <= received_at) +); diff --git a/src/lib.rs b/src/lib.rs index 8b586a68..d9fe1d10 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -27,6 +27,7 @@ pub mod postgres_inbox_consumption; pub mod postgres_instrument_release; pub mod postgres_integration; pub mod postgres_item_delivery; +pub mod postgres_response_event; pub mod postgres_response_snapshot; pub mod postgres_result_snapshot; pub mod postgres_scoring_job; diff --git a/src/postgres_response_event.rs b/src/postgres_response_event.rs new file mode 100644 index 00000000..0df1084e --- /dev/null +++ b/src/postgres_response_event.rs @@ -0,0 +1,347 @@ +//! `PostgreSQL` 18 persistence for in-progress response events. +//! +//! Completed snapshots remain in [`crate::postgres_response_snapshot`]. This +//! adapter stores the accepted event ledger so a two-item path can continue +//! after process restart. It does not store response bodies and does not score. +//! The caller owns the connection, credentials, and transaction boundary. +//! Replay requires `READ COMMITTED` so a concurrent insert that wins a unique-key +//! race is visible to the exact-replay classifier. + +use crate::reference::normalized_reference; +use crate::response::{ResponseEvent, ResponseLedger}; +use postgres::Transaction; +use std::error::Error; +use std::fmt::{Display, Formatter}; +use std::time::{Duration, SystemTime, UNIX_EPOCH}; + +const RESPONSE_EVENT_MIGRATION: &str = include_str!("../migrations/0020_response_event.sql"); + +/// Outcome of persisting one response-event ledger. +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +#[non_exhaustive] +pub enum ResponseEventPersistenceDisposition { + /// At least one new event row was inserted. + Inserted, + /// The same immutable event evidence already existed. + Duplicate, +} + +/// Fail-closed error for durable response-event persistence. +#[derive(Debug)] +#[non_exhaustive] +pub enum ResponseEventPersistenceError { + /// A session or event identity was blank, numeric-like, or unbound. + InvalidReference, + /// Event identity was replayed with different immutable evidence. + ConflictingReplay, + /// A server sequence was reused by another event identity. + SequenceConflict, + /// A sequence cannot be represented by the bounded database column. + InvalidSequence, + /// Observed or received time was zero, inverted, or out of range. + InvalidTimestamp, + /// Response-event persistence requires `PostgreSQL` `READ COMMITTED` isolation. + UnsupportedIsolationLevel, + /// Stored rows could not be rebuilt into a domain ledger. + InvalidStoredIdentity, + /// `PostgreSQL` rejected or could not execute the persistence operation. + Database(postgres::Error), +} + +impl Display for ResponseEventPersistenceError { + fn fmt(&self, formatter: &mut Formatter<'_>) -> std::fmt::Result { + formatter.write_str(match self { + Self::InvalidReference => { + "response event persistence references must be opaque durable values" + } + Self::ConflictingReplay => { + "response event identity was replayed with conflicting evidence" + } + Self::SequenceConflict => { + "response event sequence was reused by a different event identity" + } + Self::InvalidSequence => "response event sequence exceeds the PostgreSQL bigint range", + Self::InvalidTimestamp => { + "response event observed time must be positive and not after received time" + } + Self::UnsupportedIsolationLevel => { + "response event persistence requires read committed isolation" + } + Self::InvalidStoredIdentity => { + "stored response events could not be rebuilt into a ledger" + } + Self::Database(_) => "PostgreSQL response-event persistence failed", + }) + } +} + +impl Error for ResponseEventPersistenceError { + fn source(&self) -> Option<&(dyn Error + 'static)> { + match self { + Self::Database(error) => Some(error), + _ => None, + } + } +} + +impl From for ResponseEventPersistenceError { + fn from(error: postgres::Error) -> Self { + Self::Database(error) + } +} + +/// Apply the idempotent response-event migration to a `PostgreSQL` connection. +/// +/// # Errors +/// +/// Returns the `PostgreSQL` error if the migration cannot be applied. +pub fn apply_response_event_migration( + client: &mut impl postgres::GenericClient, +) -> Result<(), postgres::Error> { + client.batch_execute(RESPONSE_EVENT_MIGRATION) +} + +/// Persist every accepted event in one response ledger. +/// +/// `event_times` is `(observed_at_unix_ms, received_at_unix_ms)` aligned with +/// [`ResponseLedger::events`]. Exact replay is idempotent. Rebinding client, +/// server, item, digest, sequence, or session evidence fails closed. +/// +/// # Errors +/// +/// Returns [`ResponseEventPersistenceError`] for invalid identity, inverted +/// time, unsupported isolation, conflicting replay, sequence reuse, or a +/// database failure. +pub fn persist_response_ledger( + transaction: &mut Transaction<'_>, + ledger: &ResponseLedger, + event_times: &[(u64, u64)], +) -> Result { + require_read_committed(transaction)?; + if event_times.len() != ledger.len() { + return Err(ResponseEventPersistenceError::InvalidTimestamp); + } + let session_ref = required_reference(ledger.session_ref())?; + let mut inserted_any = false; + for (event, (observed_at_unix_ms, received_at_unix_ms)) in + ledger.events().iter().zip(event_times.iter().copied()) + { + if persist_one_event( + transaction, + session_ref, + event, + observed_at_unix_ms, + received_at_unix_ms, + )? { + inserted_any = true; + } + } + Ok(if inserted_any { + ResponseEventPersistenceDisposition::Inserted + } else { + ResponseEventPersistenceDisposition::Duplicate + }) +} + +/// Load the accepted event ledger for one session after process restart. +/// +/// Missing sessions return [`None`]. Stored order is `server_sequence`. +/// +/// # Errors +/// +/// Returns [`ResponseEventPersistenceError`] for unsupported isolation, a +/// malformed session reference, invalid stored identity, or a database failure. +pub fn load_response_ledger( + transaction: &mut Transaction<'_>, + session_ref: &str, +) -> Result, ResponseEventPersistenceError> { + require_read_committed(transaction)?; + let session_ref = required_reference(session_ref)?; + let rows = transaction.query( + "SELECT response_event_ref, client_event_ref, item_version_ref, \ + payload_digest, server_sequence \ + FROM response_event \ + WHERE session_ref = $1 \ + ORDER BY server_sequence", + &[&session_ref], + )?; + if rows.is_empty() { + return Ok(None); + } + let mut events = Vec::with_capacity(rows.len()); + for row in rows { + let sequence = postgres_loaded_sequence(row.get(4))?; + let event = ResponseEvent::from_persisted( + row.get::<_, String>(0).as_str(), + row.get::<_, String>(1).as_str(), + row.get::<_, String>(2).as_str(), + row.get::<_, String>(3).as_str(), + sequence, + ) + .map_err(|_| ResponseEventPersistenceError::InvalidStoredIdentity)?; + events.push(event); + } + ResponseLedger::from_persisted(session_ref, events) + .map(Some) + .map_err(|_| ResponseEventPersistenceError::InvalidStoredIdentity) +} + +fn persist_one_event( + transaction: &mut Transaction<'_>, + session_ref: &str, + event: &ResponseEvent, + observed_at_unix_ms: u64, + received_at_unix_ms: u64, +) -> Result { + let response_event_ref = required_reference(event.server_event_ref())?; + let client_event_ref = required_reference(event.client_event_ref())?; + let item_version_ref = required_reference(event.item_version_ref())?; + let server_sequence = postgres_sequence(event.sequence())?; + let observed_at = postgres_timestamptz(observed_at_unix_ms)?; + let received_at = postgres_timestamptz(received_at_unix_ms)?; + if observed_at_unix_ms > received_at_unix_ms { + return Err(ResponseEventPersistenceError::InvalidTimestamp); + } + let row = match transaction.query_one( + "WITH inserted AS (\ + INSERT INTO response_event (\ + response_event_ref, session_ref, client_event_ref, item_version_ref, \ + payload_digest, server_sequence, observed_at, received_at\ + ) VALUES ($1, $2, $3, $4, $5, $6, $7, $8) \ + ON CONFLICT (response_event_ref) DO NOTHING \ + RETURNING session_ref, client_event_ref, item_version_ref, payload_digest, \ + server_sequence, TRUE AS inserted\ + ) \ + SELECT session_ref, client_event_ref, item_version_ref, payload_digest, \ + server_sequence, inserted \ + FROM inserted \ + UNION ALL \ + SELECT session_ref, client_event_ref, item_version_ref, payload_digest, \ + server_sequence, FALSE AS inserted \ + FROM response_event WHERE response_event_ref = $1 \ + LIMIT 1", + &[ + &response_event_ref, + &session_ref, + &client_event_ref, + &item_version_ref, + &event.payload_digest(), + &server_sequence, + &observed_at, + &received_at, + ], + ) { + Ok(row) => row, + Err(error) => return Err(classify_unique_violation(error)), + }; + let inserted: bool = row.get(5); + if inserted { + Ok(true) + } else { + classify_existing_event(&row, session_ref, event, server_sequence) + } +} + +fn classify_existing_event( + row: &postgres::Row, + session_ref: &str, + event: &ResponseEvent, + server_sequence: i64, +) -> Result { + let stored_session: String = row.get(0); + let stored_client: String = row.get(1); + let stored_item: String = row.get(2); + let stored_digest: String = row.get(3); + let stored_sequence: i64 = row.get(4); + if stored_session == session_ref + && stored_client == event.client_event_ref() + && stored_item == event.item_version_ref() + && stored_digest == event.payload_digest() + && stored_sequence == server_sequence + { + Ok(false) + } else { + Err(ResponseEventPersistenceError::ConflictingReplay) + } +} + +fn classify_unique_violation(error: postgres::Error) -> ResponseEventPersistenceError { + match error + .as_db_error() + .and_then(postgres::error::DbError::constraint) + { + Some("response_event_session_client_unique") => { + ResponseEventPersistenceError::ConflictingReplay + } + Some("response_event_session_sequence_unique") => { + ResponseEventPersistenceError::SequenceConflict + } + _ => ResponseEventPersistenceError::Database(error), + } +} + +fn required_reference(reference: &str) -> Result<&str, ResponseEventPersistenceError> { + normalized_reference(reference).ok_or(ResponseEventPersistenceError::InvalidReference) +} + +fn postgres_sequence(value: usize) -> Result { + i64::try_from(value).map_err(|_| ResponseEventPersistenceError::InvalidSequence) +} + +fn postgres_loaded_sequence(value: i64) -> Result { + usize::try_from(value).map_err(|_| ResponseEventPersistenceError::InvalidStoredIdentity) +} + +fn postgres_timestamptz(unix_ms: u64) -> Result { + if unix_ms == 0 { + return Err(ResponseEventPersistenceError::InvalidTimestamp); + } + UNIX_EPOCH + .checked_add(Duration::from_millis(unix_ms)) + .ok_or(ResponseEventPersistenceError::InvalidTimestamp) +} + +fn require_read_committed( + transaction: &mut Transaction<'_>, +) -> Result<(), ResponseEventPersistenceError> { + let row = transaction.query_one("SHOW transaction_isolation", &[])?; + let isolation: String = row.get(0); + if isolation == "read committed" { + Ok(()) + } else { + Err(ResponseEventPersistenceError::UnsupportedIsolationLevel) + } +} + +#[cfg(test)] +mod reference_guard_tests { + use super::{ + postgres_sequence, postgres_timestamptz, required_reference, ResponseEventPersistenceError, + }; + + #[test] + fn blank_numeric_zero_time_and_overflow_fail_closed() { + assert!(matches!( + required_reference(" "), + Err(ResponseEventPersistenceError::InvalidReference) + )); + assert!(matches!( + required_reference("12"), + Err(ResponseEventPersistenceError::InvalidReference) + )); + assert_eq!( + required_reference("session_big_five_ko").unwrap(), + "session_big_five_ko" + ); + assert_eq!(postgres_sequence(1).unwrap(), 1); + assert!(matches!( + postgres_sequence(usize::MAX), + Err(ResponseEventPersistenceError::InvalidSequence) + )); + assert!(matches!( + postgres_timestamptz(0), + Err(ResponseEventPersistenceError::InvalidTimestamp) + )); + assert!(postgres_timestamptz(1_700_000_000_000).is_ok()); + } +} diff --git a/tests/postgres_response_event_error_contract.rs b/tests/postgres_response_event_error_contract.rs new file mode 100644 index 00000000..a921cb74 --- /dev/null +++ b/tests/postgres_response_event_error_contract.rs @@ -0,0 +1,61 @@ +//! Stable operator-facing error contracts for response-event persistence. + +use postgres::{Client, NoTls}; +use psychometrics_commons_runtime::postgres_response_event::ResponseEventPersistenceError; + +fn test_client() -> Client { + let connection = std::env::var("TEST_DATABASE_URL") + .expect("TEST_DATABASE_URL must identify the isolated CI PostgreSQL database"); + Client::connect(&connection, NoTls).expect("isolated CI PostgreSQL database must be reachable") +} + +#[test] +fn persistence_errors_expose_stable_messages_and_database_sources() { + for (error, expected_message) in [ + ( + ResponseEventPersistenceError::InvalidReference, + "response event persistence references must be opaque durable values", + ), + ( + ResponseEventPersistenceError::ConflictingReplay, + "response event identity was replayed with conflicting evidence", + ), + ( + ResponseEventPersistenceError::SequenceConflict, + "response event sequence was reused by a different event identity", + ), + ( + ResponseEventPersistenceError::InvalidSequence, + "response event sequence exceeds the PostgreSQL bigint range", + ), + ( + ResponseEventPersistenceError::InvalidTimestamp, + "response event observed time must be positive and not after received time", + ), + ( + ResponseEventPersistenceError::UnsupportedIsolationLevel, + "response event persistence requires read committed isolation", + ), + ( + ResponseEventPersistenceError::InvalidStoredIdentity, + "stored response events could not be rebuilt into a ledger", + ), + ] { + assert_eq!(error.to_string(), expected_message); + assert!(std::error::Error::source(&error).is_none()); + } + + let mut client = test_client(); + let database_error = client + .query_one( + "SELECT * FROM response_event_error_contract_missing_relation", + &[], + ) + .unwrap_err(); + let error = ResponseEventPersistenceError::from(database_error); + assert_eq!( + error.to_string(), + "PostgreSQL response-event persistence failed" + ); + assert!(std::error::Error::source(&error).is_some()); +} diff --git a/tests/postgres_response_event_persistence.rs b/tests/postgres_response_event_persistence.rs new file mode 100644 index 00000000..f880da93 --- /dev/null +++ b/tests/postgres_response_event_persistence.rs @@ -0,0 +1,334 @@ +//! Real `PostgreSQL` contract for in-progress response-event persistence. + +use postgres::{Client, IsolationLevel, NoTls}; +use psychometrics_commons_runtime::postgres_response_event::{ + apply_response_event_migration, load_response_ledger, persist_response_ledger, + ResponseEventPersistenceDisposition, ResponseEventPersistenceError, +}; +use psychometrics_commons_runtime::response::{ResponseEvent, ResponseLedger, ResponseWrite}; +use psychometrics_commons_runtime::session::SessionState; +use std::sync::{Mutex, MutexGuard}; + +const DIGEST_A: &str = "sha256:aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa"; +const DIGEST_B: &str = "sha256:bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb"; +const OBSERVED_OPENNESS_MS: u64 = 1_700_000_000_000; +const RECEIVED_OPENNESS_MS: u64 = 1_700_000_000_250; +const OBSERVED_CONSCIENTIOUSNESS_MS: u64 = 1_700_000_030_000; +const RECEIVED_CONSCIENTIOUSNESS_MS: u64 = 1_700_000_030_400; + +static RESPONSE_EVENT_TEST_LOCK: Mutex<()> = Mutex::new(()); + +fn response_event_test_guard() -> MutexGuard<'static, ()> { + RESPONSE_EVENT_TEST_LOCK + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) +} + +fn test_client() -> Client { + let connection = std::env::var("TEST_DATABASE_URL") + .expect("TEST_DATABASE_URL must identify the isolated CI PostgreSQL database"); + let mut client = Client::connect(&connection, NoTls) + .expect("isolated CI PostgreSQL database must be reachable"); + client + .batch_execute( + "CREATE SCHEMA IF NOT EXISTS response_event_persistence_test;\ + SET search_path TO response_event_persistence_test;", + ) + .unwrap(); + client +} + +fn reset_response_event_table(client: &mut Client) { + client + .batch_execute("DROP TABLE IF EXISTS response_event_persistence_test.response_event;") + .unwrap(); +} + +fn write<'a>( + server_event_ref: &'a str, + client_event_ref: &'a str, + item_version_ref: &'a str, + payload_digest: &'a str, +) -> ResponseWrite<'a> { + ResponseWrite { + server_event_ref, + client_event_ref, + item_version_ref, + payload_digest, + } +} + +fn two_item_korean_ledger() -> ResponseLedger { + let mut ledger = ResponseLedger::new("session_big_five_ko").unwrap(); + ledger + .record( + SessionState::Active, + write( + "response_event_openness", + "client_openness", + "item_version_openness_ko", + DIGEST_A, + ), + ) + .unwrap(); + ledger + .record( + SessionState::Active, + write( + "response_event_conscientiousness", + "client_conscientiousness", + "item_version_conscientiousness_ko", + DIGEST_B, + ), + ) + .unwrap(); + ledger +} + +fn event_times() -> [(u64, u64); 2] { + [ + (OBSERVED_OPENNESS_MS, RECEIVED_OPENNESS_MS), + (OBSERVED_CONSCIENTIOUSNESS_MS, RECEIVED_CONSCIENTIOUSNESS_MS), + ] +} + +fn persist_ok( + client: &mut Client, + ledger: &ResponseLedger, + event_times: &[(u64, u64)], +) -> ResponseEventPersistenceDisposition { + let mut transaction = client.transaction().unwrap(); + let disposition = persist_response_ledger(&mut transaction, ledger, event_times).unwrap(); + transaction.commit().unwrap(); + disposition +} + +fn persist_err( + client: &mut Client, + ledger: &ResponseLedger, + event_times: &[(u64, u64)], +) -> ResponseEventPersistenceError { + let mut transaction = client.transaction().unwrap(); + let error = persist_response_ledger(&mut transaction, ledger, event_times).unwrap_err(); + transaction.rollback().unwrap(); + error +} + +fn load_ok(client: &mut Client, session_ref: &str) -> Option { + let mut transaction = client.transaction().unwrap(); + let ledger = load_response_ledger(&mut transaction, session_ref).unwrap(); + transaction.commit().unwrap(); + ledger +} + +#[test] +fn two_item_korean_path_survives_restart_with_exact_replay() { + let _guard = response_event_test_guard(); + let mut client = test_client(); + reset_response_event_table(&mut client); + apply_response_event_migration(&mut client).unwrap(); + + let live = two_item_korean_ledger(); + assert_eq!( + persist_ok(&mut client, &live, &event_times()), + ResponseEventPersistenceDisposition::Inserted + ); + assert_eq!( + persist_ok(&mut client, &live, &event_times()), + ResponseEventPersistenceDisposition::Duplicate + ); + + let reloaded = load_ok(&mut client, "session_big_five_ko").unwrap(); + assert_eq!(reloaded, live); + assert_eq!( + reloaded.events()[1].item_version_ref(), + "item_version_conscientiousness_ko" + ); + assert!(load_ok(&mut client, "session_missing_ko").is_none()); +} + +#[test] +fn client_rebinding_and_sequence_reuse_fail_closed() { + let _guard = response_event_test_guard(); + let mut client = test_client(); + reset_response_event_table(&mut client); + apply_response_event_migration(&mut client).unwrap(); + persist_ok(&mut client, &two_item_korean_ledger(), &event_times()); + + let mut rebound_client = ResponseLedger::new("session_big_five_ko").unwrap(); + rebound_client + .record( + SessionState::Active, + write( + "response_event_openness", + "client_openness", + "item_version_openness_ko", + DIGEST_B, + ), + ) + .unwrap(); + assert!(matches!( + persist_err( + &mut client, + &rebound_client, + &[(OBSERVED_OPENNESS_MS, RECEIVED_OPENNESS_MS)] + ), + ResponseEventPersistenceError::ConflictingReplay + )); + + let mut reused_sequence = ResponseLedger::new("session_big_five_ko").unwrap(); + reused_sequence + .record( + SessionState::Active, + write( + "response_event_extraversion", + "client_extraversion", + "item_version_extraversion_ko", + DIGEST_A, + ), + ) + .unwrap(); + assert!(matches!( + persist_err( + &mut client, + &reused_sequence, + &[(OBSERVED_OPENNESS_MS, RECEIVED_OPENNESS_MS)] + ), + ResponseEventPersistenceError::SequenceConflict + )); +} + +#[test] +fn inverted_time_blank_session_and_repeatable_read_fail_closed() { + let _guard = response_event_test_guard(); + let mut client = test_client(); + reset_response_event_table(&mut client); + apply_response_event_migration(&mut client).unwrap(); + let live = two_item_korean_ledger(); + + assert!(matches!( + persist_err(&mut client, &live, &[(1, 2)]), + ResponseEventPersistenceError::InvalidTimestamp + )); + assert!(matches!( + persist_err( + &mut client, + &live, + &[ + (RECEIVED_OPENNESS_MS, OBSERVED_OPENNESS_MS), + (OBSERVED_CONSCIENTIOUSNESS_MS, RECEIVED_CONSCIENTIOUSNESS_MS) + ] + ), + ResponseEventPersistenceError::InvalidTimestamp + )); + assert!(matches!( + persist_err( + &mut client, + &live, + &[ + (0, RECEIVED_OPENNESS_MS), + (OBSERVED_CONSCIENTIOUSNESS_MS, RECEIVED_CONSCIENTIOUSNESS_MS) + ] + ), + ResponseEventPersistenceError::InvalidTimestamp + )); + + let mut transaction = client + .build_transaction() + .isolation_level(IsolationLevel::RepeatableRead) + .start() + .unwrap(); + assert!(matches!( + persist_response_ledger(&mut transaction, &live, &event_times()), + Err(ResponseEventPersistenceError::UnsupportedIsolationLevel) + )); + transaction.rollback().unwrap(); + + let mut load_transaction = client.transaction().unwrap(); + assert!(matches!( + load_response_ledger(&mut load_transaction, "12"), + Err(ResponseEventPersistenceError::InvalidReference) + )); + load_transaction.rollback().unwrap(); +} + +#[test] +fn empty_ledger_persist_is_duplicate_and_gapped_store_fails_closed() { + let _guard = response_event_test_guard(); + let mut client = test_client(); + reset_response_event_table(&mut client); + apply_response_event_migration(&mut client).unwrap(); + + let empty = ResponseLedger::new("session_empty_ko").unwrap(); + assert_eq!( + persist_ok(&mut client, &empty, &[]), + ResponseEventPersistenceDisposition::Duplicate + ); + assert!(load_ok(&mut client, "session_empty_ko").is_none()); + + persist_ok(&mut client, &two_item_korean_ledger(), &event_times()); + client + .execute( + "DELETE FROM response_event WHERE response_event_ref = 'response_event_openness'", + &[], + ) + .unwrap(); + let mut transaction = client.transaction().unwrap(); + assert!(matches!( + load_response_ledger(&mut transaction, "session_big_five_ko"), + Err(ResponseEventPersistenceError::InvalidStoredIdentity) + )); + transaction.rollback().unwrap(); +} + +#[test] +fn server_event_rebinding_to_another_session_fails_closed() { + let _guard = response_event_test_guard(); + let mut client = test_client(); + reset_response_event_table(&mut client); + apply_response_event_migration(&mut client).unwrap(); + persist_ok(&mut client, &two_item_korean_ledger(), &event_times()); + + let mut other_session = ResponseLedger::new("session_big_five_en").unwrap(); + other_session + .record( + SessionState::Active, + write( + "response_event_openness", + "client_openness_en", + "item_version_openness_en", + DIGEST_A, + ), + ) + .unwrap(); + assert!(matches!( + persist_err( + &mut client, + &other_session, + &[(OBSERVED_OPENNESS_MS, RECEIVED_OPENNESS_MS)] + ), + ResponseEventPersistenceError::ConflictingReplay + )); + + let rebound_server = ResponseLedger::from_persisted( + "session_big_five_ko", + vec![ResponseEvent::from_persisted( + "response_event_other_server", + "client_openness", + "item_version_openness_ko", + DIGEST_A, + 1, + ) + .unwrap()], + ) + .unwrap(); + assert!(matches!( + persist_err( + &mut client, + &rebound_server, + &[(OBSERVED_OPENNESS_MS, RECEIVED_OPENNESS_MS)] + ), + ResponseEventPersistenceError::ConflictingReplay + | ResponseEventPersistenceError::SequenceConflict + )); +} From 50125d21fd11971ca8080e130d63679b16c2ebc7 Mon Sep 17 00:00:00 2001 From: Cursor Agent Date: Sun, 16 Aug 2026 16:12:53 +0000 Subject: [PATCH 3/6] style(response): apply rustfmt to ledger reconstruction Co-authored-by: Seongho Bae --- src/response.rs | 5 ++--- tests/response_ledger.rs | 23 +++++++++++++++++------ 2 files changed, 19 insertions(+), 9 deletions(-) diff --git a/src/response.rs b/src/response.rs index d41bd4ad..0289ca4e 100644 --- a/src/response.rs +++ b/src/response.rs @@ -209,9 +209,8 @@ impl Display for WriteError { formatter, "response snapshot requires Completed session state, found {state:?}" ), - Self::InvalidStoredSequence => formatter.write_str( - "stored response events must keep server sequence 1..n without gaps", - ), + Self::InvalidStoredSequence => formatter + .write_str("stored response events must keep server sequence 1..n without gaps"), } } } diff --git a/tests/response_ledger.rs b/tests/response_ledger.rs index a36467da..6c6e2366 100644 --- a/tests/response_ledger.rs +++ b/tests/response_ledger.rs @@ -370,15 +370,23 @@ fn two_item_korean_path_reloads_the_same_answers_after_restart() { assert_eq!(live.session_ref(), "session_big_five_ko"); assert_eq!(live.events().len(), 2); assert_eq!(live.events()[0].sequence(), 1); - assert_eq!(live.events()[1].item_version_ref(), "item_version_conscientiousness_ko"); + assert_eq!( + live.events()[1].item_version_ref(), + "item_version_conscientiousness_ko" + ); - let reloaded = ResponseLedger::from_persisted(live.session_ref(), live.events().to_vec()).unwrap(); + let reloaded = + ResponseLedger::from_persisted(live.session_ref(), live.events().to_vec()).unwrap(); assert_eq!(reloaded, live); assert_eq!( - reloaded.freeze_as(SessionState::Completed, "response_snapshot_big_five_ko") + reloaded + .freeze_as(SessionState::Completed, "response_snapshot_big_five_ko") .unwrap() .event_refs(), - ["response_event_openness", "response_event_conscientiousness"] + [ + "response_event_openness", + "response_event_conscientiousness" + ] ); } @@ -437,8 +445,11 @@ fn persisted_events_reject_reused_identities_and_blank_session() { .unwrap(); assert_eq!( - ResponseLedger::from_persisted("session_big_five_ko", vec![first.clone(), duplicate_client]) - .unwrap_err(), + ResponseLedger::from_persisted( + "session_big_five_ko", + vec![first.clone(), duplicate_client] + ) + .unwrap_err(), WriteError::IdempotencyConflict ); assert_eq!( From 71c751497f47cfc50eb9384268b5b2fa3104a601 Mon Sep 17 00:00:00 2001 From: Cursor Agent Date: Sun, 16 Aug 2026 16:13:18 +0000 Subject: [PATCH 4/6] docs(response): record event persist as Active PR #174 Co-authored-by: Seongho Bae --- docs/TRACEABILITY.md | 6 +++--- docs/adr/0015-persistence-and-transaction-boundaries.md | 2 +- docs/architecture/ERD.md | 2 +- 3 files changed, 5 insertions(+), 5 deletions(-) diff --git a/docs/TRACEABILITY.md b/docs/TRACEABILITY.md index ceef7951..bd685543 100644 --- a/docs/TRACEABILITY.md +++ b/docs/TRACEABILITY.md @@ -23,7 +23,7 @@ An active PR, architecture document, conversation decision, or scheduler plan is | Anonymous core assessment | PRD §3.1, §9.1 | TRD §5, §10; UML anonymous sequence | ADR-0002, ADR-0003, ADR-0005 | Session lifecycle primitives implemented, including creation bound to one published locale-specific release; anonymous credential/HTTP flow is Target | | Pause/resume | PRD §3.1, §9.1 | TRD §5 | ADR-0005 | **Implemented** in `src/session.rs` with fail-closed transitions | | Sequence-aware item delivery evidence | PRD §3.1, §9 | TRD §5–7 | ADR-0005, ADR-0010 | **Implemented** domain primitive in `src/item_delivery.rs`; persistence/API delivery orchestration is Target | -| Idempotent response events | PRD §9.2 | TRD §6 | ADR-0005, ADR-0010 | **Implemented** in `src/response.rs` with canonical SHA-256 payload-digest identity; **Active PR** on this branch persists and reloads `response_event` under `READ COMMITTED`; HTTP response transport remains Target | +| Idempotent response events | PRD §9.2 | TRD §6 | ADR-0005, ADR-0010 | **Implemented** in `src/response.rs` with canonical SHA-256 payload-digest identity; **Active PR** #174 persists and reloads `response_event` under `READ COMMITTED`; HTTP response transport remains Target | | Immutable response snapshot before scoring | PRD §9.3 | TRD §5–8 | ADR-0005, ADR-0010 | **Implemented** domain semantics in `src/response.rs` | | Version-pinned scoring | PRD §9.4, §10 | TRD §8 | ADR-0004, ADR-0010 | **Implemented** reusable product-side scoring dispatch contract in `src/scoring.rs` with canonical SHA-256 engine-artifact digest provenance plus `migrations/0011_scoring_request.sql` / `src/postgres_scoring_request.rs` request-identity persistence; live fast-mlsirm integration is Target | | Bounded asynchronous scoring retry/quarantine with stale-worker fencing | PRD §9.4, §10 | TRD §8; ADR-0015 transaction boundary | ADR-0004, ADR-0010, ADR-0015 | **Implemented** product lifecycle plus PostgreSQL enqueue, claim, retry, completion, expiry recovery, and cancellation without transferring a fence; live fast-mlsirm execution remains Target | @@ -55,7 +55,7 @@ An active PR, architecture document, conversation decision, or scheduler plan is | Server-authoritative session state | TRD §5 | `src/session.rs` + session contract tests, including published-release/locale binding at creation | persistence/API concurrency test | | Only Active accepts responses | TRD §5–6 | `SessionState::accepts_responses` + response tests | transport-level rejection test | | Item delivery sequence is positive and evidence-safe | TRD §5–7 | `src/item_delivery.rs` + item-delivery domain tests | durable uniqueness/order/API integration | -| Conflicting idempotency replay fails closed | TRD §6 | `src/response.rs`; **Active PR** `migrations/0020_response_event.sql` unique `(session_ref, client_event_ref)` and `(session_ref, server_sequence)` | HTTP transport rejection test | +| Conflicting idempotency replay fails closed | TRD §6 | `src/response.rs`; **Active PR** #174 `migrations/0020_response_event.sql` unique `(session_ref, client_event_ref)` and `(session_ref, server_sequence)` | HTTP transport rejection test | | Snapshot requires Completed state | TRD §5–6 | `src/response.rs` | transaction atomicity test with persistence | | Scoring uses durable snapshot identity | TRD §8 | `src/scoring.rs` requires a canonical SHA-256 engine-artifact digest | live adapter + retry/outbox integration | | Stale scoring worker cannot complete a newer attempt | TRD §8; ADR-0015 | `src/scoring_job.rs` uses monotonically increasing fencing tokens and rejects stale/expired completion or failure evidence; `src/postgres_scoring_job.rs` persists enqueue, claim, retry, terminal outcomes, expired-lease recovery, and cancellation without transferring a fence | live adapter evidence | @@ -132,7 +132,7 @@ Still-Target logical modules/adapters include remaining product aggregate persis ### Active implementation work that is not protected-main truth -**Active PR** response-event persist/load on this branch is not protected-main truth until an unchanged reviewed/check-clean head is integrated. In-progress answers persist in `response_event` with observed/received time, exact replay, and fail-closed rebinding so a two-item path survives restart. Completed snapshot persist remains the existing protected-main adapter. HTTP response transport remains outside this slice. Protected-main `#76` data-rights processing-start persistence is already integrated and is not restated here as active work. +**Active PR** #174 response-event persist/load is not protected-main truth until an unchanged reviewed/check-clean head is integrated. In-progress answers persist in `response_event` with observed/received time, exact replay, and fail-closed rebinding so a two-item path survives restart. Completed snapshot persist remains the existing protected-main adapter. HTTP response transport remains outside this slice. Protected-main `#76` data-rights processing-start persistence is already integrated and is not restated here as active work. Prefer this head over another overlapping `response_event` persist PR. ## 5. ADR traceability by concern diff --git a/docs/adr/0015-persistence-and-transaction-boundaries.md b/docs/adr/0015-persistence-and-transaction-boundaries.md index 332b9ea0..ae80b63b 100644 --- a/docs/adr/0015-persistence-and-transaction-boundaries.md +++ b/docs/adr/0015-persistence-and-transaction-boundaries.md @@ -6,7 +6,7 @@ - Scope: Psychometrics Commons-owned durable state, local transactions, migration boundaries, outbox/inbox integration - Supersedes: none - Superseded by: none -- Current/as-built status: protected main persists integration, scoring-job, data-rights, item-delivery, consent, instrument-release, result-snapshot, response-snapshot, scoring-request, and inbox-consumption slices; in-progress `response_event` persist/load is Active PR work on this branch and is not protected-main truth until merged +- Current/as-built status: protected main persists integration, scoring-job, data-rights, item-delivery, consent, instrument-release, result-snapshot, response-snapshot, scoring-request, and inbox-consumption slices; in-progress `response_event` persist/load is Active PR #174 and is not protected-main truth until merged - Target status: upstream PostgreSQL 18.x operational persistence with real-database concurrency/crash/recovery evidence and transactional outbox/inbox semantics - Migration status: `migrations/0020_response_event.sql` adds the in-progress event ledger with observed/received timestamps; remaining session HTTP and response HTTP families still must be established from the logical ERD without synthetic provenance backfills diff --git a/docs/architecture/ERD.md b/docs/architecture/ERD.md index 0215d8be..989ea734 100644 --- a/docs/architecture/ERD.md +++ b/docs/architecture/ERD.md @@ -425,7 +425,7 @@ The target ERD deliberately includes several logical entities that are not yet p - `instrument_release` is the locale-specific publication identity already owned by `src/instrument.rs`. Physical `migrations/0006_instrument_release.sql` persists that one-row aggregate (immutable manifest columns plus `publication_state`); HTTP publication transport remains Target. - `data_rights_request` and `data_rights_propagation_state` are the first durable export/deletion slice. Physical `migrations/0003_data_rights_propagation.sql` stores requested-state identity plus one local outbox event per dependent system; verification, processing, completion, and dependent-system execution remain Target. - `item_delivery_event` reflects the already-merged `src/item_delivery.rs` domain primitive; durable persistence/API orchestration is still Target. -- `response_event` is the in-progress answer ledger. Physical `migrations/0020_response_event.sql` on this Active PR stores opaque event identity, session binding, client idempotency, item version, payload digest, server sequence, and distinct observed/received timestamps; HTTP response transport remains Target. Completed `response_snapshot` persist is already on protected main. +- `response_event` is the in-progress answer ledger. Physical `migrations/0020_response_event.sql` on Active PR #174 stores opaque event identity, session binding, client idempotency, item version, payload digest, server sequence, and distinct observed/received timestamps; HTTP response transport remains Target. Completed `response_snapshot` persist is already on protected main. - `consent_ledger` and `consent_event` persist the already-merged `src/consent.rs` append-only ledger. Physical persistence is carried by Active PR #49 (`migrations/0005_consent_lifecycle.sql`); HTTP consent transport and derived snapshot tables remain Target. - `participant_identity_link` is the persistence target accepted by ADR-0020. The current `src/participant.rs` `keyverse_subject_ref` field is an application-domain first-link projection, not the future mutable persistence source of truth. - `longitudinal_enrollment`, `longitudinal_observation_record`, and `temporal_analysis_submission` make the ADR-0008 Commons-owned Gyeot/TEPP orchestration boundary explicit. No TEPP analytical kernel is duplicated here. From 5a45869e19a190d1c11dbb2c1ce3d51227be13ea Mon Sep 17 00:00:00 2001 From: Cursor Agent Date: Sun, 16 Aug 2026 16:29:35 +0000 Subject: [PATCH 5/6] fix(response): prove restart times, recovery, and item-three prefix A Korean path can now reload stored observed/received time, continue with item 3 after restart, and keep those rows through recovery COPY. Misaligned event times fail as arity, not a clock error. Co-authored-by: Seongho Bae --- CHANGELOG.md | 2 +- docs/TRACEABILITY.md | 2 +- docs/architecture/AS_BUILT_SCHEMA.md | 4 + docs/architecture/UML.md | 2 + docs/doctoring/standards-and-evidence.md | 10 ++- src/postgres_response_event.rs | 72 ++++++++++++++++-- tests/postgres_recovery_invariants.rs | 29 ++++++++ .../postgres_response_event_error_contract.rs | 4 + tests/postgres_response_event_persistence.rs | 73 ++++++++++++++++++- 9 files changed, 187 insertions(+), 11 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index a06af678..16e5b088 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -5,7 +5,7 @@ All notable product and architecture changes are recorded here. Releases use imm ## Unreleased ### Added -- PostgreSQL persistence for in-progress `response_event` rows so a two-item path reloads the same answers after restart. Exact replay is idempotent; client, server, sequence, or session rebinding fails closed. Store observed and received times with the event; completed snapshots stay on the existing snapshot adapter. +- PostgreSQL persistence for in-progress `response_event` rows so a two-item path reloads the same answers after restart. Exact replay is idempotent; client, server, sequence, or session rebinding fails closed. Store observed and received times with the event, expose those times on reload, restore them through recovery COPY, and continue the scoring prefix after item 3. Completed snapshots stay on the existing snapshot adapter. - Scoring-job cancel and lease-expiry fallback classification lock the current row until the caller transaction ends, so concurrent workers cannot rewrite terminal or unleased evidence. - PostgreSQL operational-store readiness probe classifies the supported major version and write-readiness, and fails closed when a caller-declared required relation is missing. - PostgreSQL scoring-job cancellation: queued, leased, or retry-scheduled work becomes cancelled without transferring a fence, exact replay is idempotent, and completed or quarantined evidence cannot be rewritten. diff --git a/docs/TRACEABILITY.md b/docs/TRACEABILITY.md index bd685543..4a95b349 100644 --- a/docs/TRACEABILITY.md +++ b/docs/TRACEABILITY.md @@ -132,7 +132,7 @@ Still-Target logical modules/adapters include remaining product aggregate persis ### Active implementation work that is not protected-main truth -**Active PR** #174 response-event persist/load is not protected-main truth until an unchanged reviewed/check-clean head is integrated. In-progress answers persist in `response_event` with observed/received time, exact replay, and fail-closed rebinding so a two-item path survives restart. Completed snapshot persist remains the existing protected-main adapter. HTTP response transport remains outside this slice. Protected-main `#76` data-rights processing-start persistence is already integrated and is not restated here as active work. Prefer this head over another overlapping `response_event` persist PR. +**Active PR** #174 response-event persist/load is not protected-main truth until an unchanged reviewed/check-clean head is integrated. In-progress answers persist in `response_event` with observed/received time, exact replay, and fail-closed rebinding so a two-item path survives restart. Recovery COPY restores those rows. AS_BUILT names the physical slice. Reload exposes stored times as a persist-side projection and continues the scoring prefix after item 3. Completed snapshot persist remains the existing protected-main adapter. HTTP response transport remains outside this slice. Protected-main `#76` data-rights processing-start persistence is already integrated and is not restated here as active work. Prefer this head over another overlapping `response_event` persist PR. ## 5. ADR traceability by concern diff --git a/docs/architecture/AS_BUILT_SCHEMA.md b/docs/architecture/AS_BUILT_SCHEMA.md index 8c2c3510..5fac103f 100644 --- a/docs/architecture/AS_BUILT_SCHEMA.md +++ b/docs/architecture/AS_BUILT_SCHEMA.md @@ -25,6 +25,10 @@ The protected-main integration identity is source- and tenant-scoped. A physical PR #58 (`feat/inbox-consumption-persistence-20260814`) `migrations/0012_integration_consumption.sql` and `src/postgres_inbox_consumption.rs` adapter persist one consumption work item for an existing `integration_inbox` receipt. The slice is **Active PR**, not protected-main truth. It stores pending/processing/completed/quarantined evidence, a monotonically increasing fencing token, a time-bounded processing claim, a durable `side_effect_ref`, and optional completion or quarantine evidence. Receipt-only inbox rows remain uncompleted. A processing claim cannot be stolen by another worker. Expire-and-reclaim returns an expired claim to pending without transferring the crashed worker's fence. +## Active PR response-event physical schema + +Active PR #174 `migrations/0020_response_event.sql` and `src/postgres_response_event.rs` persist the in-progress answer ledger so a two-item path can continue after process restart. The slice is **Active PR**, not protected-main truth. It stores opaque `response_event_ref` identity, session binding, client idempotency identity, item version, canonical SHA-256 payload digest, positive `server_sequence`, and distinct `observed_at` / `received_at` timestamps. Exact replay is idempotent and keeps the original times. Client, server, sequence, or session rebinding fails closed. Reload reconstructs `ResponseLedger` in `server_sequence` order under `READ COMMITTED` and exposes stored times as a persist-side projection. HTTP response transport remains outside this slice. + ## Protected-main scoring-job physical schema `migrations/0002_scoring_job_state.sql` maps a bounded physical subset of the logical `scoring_job` aggregate into `scoring_job_state`, owned by `src/postgres_scoring_job.rs`. This is an **Implemented subset** on the named protected-main baseline. diff --git a/docs/architecture/UML.md b/docs/architecture/UML.md index 13988939..166ec8a9 100644 --- a/docs/architecture/UML.md +++ b/docs/architecture/UML.md @@ -66,6 +66,8 @@ classDiagram +item_version_ref +payload_digest +server_sequence + +observed_at + +received_at } class ResponseSnapshot { +response_snapshot_ref diff --git a/docs/doctoring/standards-and-evidence.md b/docs/doctoring/standards-and-evidence.md index 7879474c..c9dd293c 100644 --- a/docs/doctoring/standards-and-evidence.md +++ b/docs/doctoring/standards-and-evidence.md @@ -92,7 +92,11 @@ Product consequences: - clock skew, impossible ordering, and unknown precision are typed validation outcomes rather than reasons to rewrite source history; - analysis-set digests bind the exact observations and time semantics consumed by - temporal, multilevel, cross-classified, or multiple-membership analysis. + temporal, multilevel, cross-classified, or multiple-membership analysis; +- in-progress `response_event` persist uses PostgreSQL `READ COMMITTED` so a + concurrent unique-key winner is visible to exact-replay classification, and + stores observed time separately from platform receipt time (Berenson et al., + 1995; PostgreSQL Global Development Group, 2026). ## Evidence maintenance rules @@ -107,6 +111,8 @@ Product consequences: American Educational Research Association, American Psychological Association, & National Council on Measurement in Education. (2014). *Standards for educational and psychological testing*. American Educational Research Association. https://www.testingstandards.net/ +Berenson, H., Bernstein, P., Gray, J., Melton, J., O'Neil, E., & O'Neil, P. (1995). A critique of ANSI SQL isolation levels. *ACM SIGMOD Record, 24*(2), 1–10. https://doi.org/10.1145/568271.223785 + International Organization for Standardization. (2022). *ISO/IEC 27001:2022 Information security, cybersecurity and privacy protection—Information security management systems—Requirements* (3rd ed.). https://www.iso.org/standard/27001 International Organization for Standardization. (2023a). *ISO/IEC 23894:2023 Information technology—Artificial intelligence—Guidance on risk management*. https://www.iso.org/standard/77304.html @@ -119,6 +125,8 @@ International Organization for Standardization. (2025). *ISO/IEC 42005:2025 Info International Organization for Standardization. (2019). *ISO 8601-1:2019 Date and time—Representations for information interchange—Part 1: Basic rules* (with Amendment 1:2022). https://www.iso.org/standard/70907.html +PostgreSQL Global Development Group. (2026). *PostgreSQL 18 documentation: Transaction isolation*. https://www.postgresql.org/docs/18/transaction-iso.html + Temoshok, D., Proud-Madruga, D., Choong, Y.-Y., Galluzzo, R., Gupta, S., LaSalle, C., Lefkovitz, N., & Regenscheid, A. (2025). *Digital identity guidelines* (NIST Special Publication 800-63-4). National Institute of Standards and Technology. https://doi.org/10.6028/NIST.SP.800-63-4 World Wide Web Consortium. (2024). *Web Content Accessibility Guidelines (WCAG) 2.2* (W3C Recommendation, 12 December 2024). https://www.w3.org/TR/WCAG22/ diff --git a/src/postgres_response_event.rs b/src/postgres_response_event.rs index 0df1084e..c1f4a89a 100644 --- a/src/postgres_response_event.rs +++ b/src/postgres_response_event.rs @@ -40,6 +40,8 @@ pub enum ResponseEventPersistenceError { InvalidSequence, /// Observed or received time was zero, inverted, or out of range. InvalidTimestamp, + /// `event_times` was not aligned one-to-one with the ledger events. + InvalidEventTimeArity, /// Response-event persistence requires `PostgreSQL` `READ COMMITTED` isolation. UnsupportedIsolationLevel, /// Stored rows could not be rebuilt into a domain ledger. @@ -64,6 +66,9 @@ impl Display for ResponseEventPersistenceError { Self::InvalidTimestamp => { "response event observed time must be positive and not after received time" } + Self::InvalidEventTimeArity => { + "response event times must align one-to-one with the persisted ledger" + } Self::UnsupportedIsolationLevel => { "response event persistence requires read committed isolation" } @@ -110,8 +115,8 @@ pub fn apply_response_event_migration( /// # Errors /// /// Returns [`ResponseEventPersistenceError`] for invalid identity, inverted -/// time, unsupported isolation, conflicting replay, sequence reuse, or a -/// database failure. +/// time, misaligned event times, unsupported isolation, conflicting replay, +/// sequence reuse, or a database failure. pub fn persist_response_ledger( transaction: &mut Transaction<'_>, ledger: &ResponseLedger, @@ -119,7 +124,7 @@ pub fn persist_response_ledger( ) -> Result { require_read_committed(transaction)?; if event_times.len() != ledger.len() { - return Err(ResponseEventPersistenceError::InvalidTimestamp); + return Err(ResponseEventPersistenceError::InvalidEventTimeArity); } let session_ref = required_reference(ledger.session_ref())?; let mut inserted_any = false; @@ -186,6 +191,44 @@ pub fn load_response_ledger( .map_err(|_| ResponseEventPersistenceError::InvalidStoredIdentity) } +/// Load stored observed and received unix-ms pairs in `server_sequence` order. +/// +/// Missing sessions return [`None`]. Pairs align with [`load_response_ledger`] +/// so HTTP and audit can show first-write provenance without putting clocks on +/// the domain event. Exact PK replay keeps the original times. +/// +/// # Errors +/// +/// Returns [`ResponseEventPersistenceError`] for unsupported isolation, a +/// malformed session reference, invalid stored time, or a database failure. +pub fn load_response_event_times( + transaction: &mut Transaction<'_>, + session_ref: &str, +) -> Result>, ResponseEventPersistenceError> { + require_read_committed(transaction)?; + let session_ref = required_reference(session_ref)?; + let rows = transaction.query( + "SELECT observed_at, received_at \ + FROM response_event \ + WHERE session_ref = $1 \ + ORDER BY server_sequence", + &[&session_ref], + )?; + if rows.is_empty() { + return Ok(None); + } + let mut times = Vec::with_capacity(rows.len()); + for row in rows { + let observed_at: SystemTime = row.get(0); + let received_at: SystemTime = row.get(1); + times.push(( + unix_ms_from_system_time(observed_at)?, + unix_ms_from_system_time(received_at)?, + )); + } + Ok(Some(times)) +} + fn persist_one_event( transaction: &mut Transaction<'_>, session_ref: &str, @@ -301,6 +344,14 @@ fn postgres_timestamptz(unix_ms: u64) -> Result Result { + let duration = value + .duration_since(UNIX_EPOCH) + .map_err(|_| ResponseEventPersistenceError::InvalidStoredIdentity)?; + u64::try_from(duration.as_millis()) + .map_err(|_| ResponseEventPersistenceError::InvalidStoredIdentity) +} + fn require_read_committed( transaction: &mut Transaction<'_>, ) -> Result<(), ResponseEventPersistenceError> { @@ -316,8 +367,10 @@ fn require_read_committed( #[cfg(test)] mod reference_guard_tests { use super::{ - postgres_sequence, postgres_timestamptz, required_reference, ResponseEventPersistenceError, + postgres_sequence, postgres_timestamptz, required_reference, unix_ms_from_system_time, + ResponseEventPersistenceError, }; + use std::time::{Duration, UNIX_EPOCH}; #[test] fn blank_numeric_zero_time_and_overflow_fail_closed() { @@ -342,6 +395,15 @@ mod reference_guard_tests { postgres_timestamptz(0), Err(ResponseEventPersistenceError::InvalidTimestamp) )); - assert!(postgres_timestamptz(1_700_000_000_000).is_ok()); + let stored = postgres_timestamptz(1_700_000_000_000).unwrap(); + assert_eq!(unix_ms_from_system_time(stored).unwrap(), 1_700_000_000_000); + assert!(matches!( + unix_ms_from_system_time(UNIX_EPOCH - Duration::from_secs(1)), + Err(ResponseEventPersistenceError::InvalidStoredIdentity) + )); + assert_eq!( + ResponseEventPersistenceError::InvalidEventTimeArity.to_string(), + "response event times must align one-to-one with the persisted ledger" + ); } } diff --git a/tests/postgres_recovery_invariants.rs b/tests/postgres_recovery_invariants.rs index e8af1d61..dd41e5a4 100644 --- a/tests/postgres_recovery_invariants.rs +++ b/tests/postgres_recovery_invariants.rs @@ -98,6 +98,16 @@ fn seed_recovery_critical_state(client: &mut Client) { ) VALUES ( 'snapshot_recovery_alpha', 1, 'response_recovery_alpha', 'item_version_recovery_alpha', '{DIGEST_A}' + ); + INSERT INTO {SOURCE_SCHEMA}.response_event ( + response_event_ref, session_ref, client_event_ref, item_version_ref, + payload_digest, server_sequence, observed_at, received_at + ) VALUES ( + 'response_event_recovery_alpha', 'session_recovery_alpha', + 'client_event_recovery_alpha', 'item_version_recovery_alpha', + '{DIGEST_A}', 1, + TIMESTAMPTZ '2023-11-14 22:13:20+00', + TIMESTAMPTZ '2023-11-14 22:13:20.250+00' );" )) .expect("recovery fixture should satisfy all protected-main persistence constraints"); @@ -183,6 +193,24 @@ fn assert_restored_evidence(client: &mut Client) { "response_recovery_alpha" ); assert_eq!(restored_snapshot.get::<_, String>(3), DIGEST_A); + + let restored_event = client + .query_one( + &format!( + "SELECT session_ref, client_event_ref, payload_digest, server_sequence + FROM {RESTORED_SCHEMA}.response_event + WHERE response_event_ref = 'response_event_recovery_alpha'" + ), + &[], + ) + .expect("in-progress response events should survive restore"); + assert_eq!(restored_event.get::<_, String>(0), "session_recovery_alpha"); + assert_eq!( + restored_event.get::<_, String>(1), + "client_event_recovery_alpha" + ); + assert_eq!(restored_event.get::<_, String>(2), DIGEST_A); + assert_eq!(restored_event.get::<_, i64>(3), 1); } fn assert_restored_tenant_scoped_deduplication(client: &mut Client) { @@ -246,6 +274,7 @@ fn clean_restore_preserves_provenance_deduplication_and_fencing_state() { "integration_consumption", "response_snapshot", "response_snapshot_entry", + "response_event", ]; let backups: Vec<(&str, Vec)> = tables .iter() diff --git a/tests/postgres_response_event_error_contract.rs b/tests/postgres_response_event_error_contract.rs index a921cb74..09e470af 100644 --- a/tests/postgres_response_event_error_contract.rs +++ b/tests/postgres_response_event_error_contract.rs @@ -32,6 +32,10 @@ fn persistence_errors_expose_stable_messages_and_database_sources() { ResponseEventPersistenceError::InvalidTimestamp, "response event observed time must be positive and not after received time", ), + ( + ResponseEventPersistenceError::InvalidEventTimeArity, + "response event times must align one-to-one with the persisted ledger", + ), ( ResponseEventPersistenceError::UnsupportedIsolationLevel, "response event persistence requires read committed isolation", diff --git a/tests/postgres_response_event_persistence.rs b/tests/postgres_response_event_persistence.rs index f880da93..5363533b 100644 --- a/tests/postgres_response_event_persistence.rs +++ b/tests/postgres_response_event_persistence.rs @@ -2,8 +2,8 @@ use postgres::{Client, IsolationLevel, NoTls}; use psychometrics_commons_runtime::postgres_response_event::{ - apply_response_event_migration, load_response_ledger, persist_response_ledger, - ResponseEventPersistenceDisposition, ResponseEventPersistenceError, + apply_response_event_migration, load_response_event_times, load_response_ledger, + persist_response_ledger, ResponseEventPersistenceDisposition, ResponseEventPersistenceError, }; use psychometrics_commons_runtime::response::{ResponseEvent, ResponseLedger, ResponseWrite}; use psychometrics_commons_runtime::session::SessionState; @@ -11,10 +11,13 @@ use std::sync::{Mutex, MutexGuard}; const DIGEST_A: &str = "sha256:aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa"; const DIGEST_B: &str = "sha256:bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb"; +const DIGEST_C: &str = "sha256:cccccccccccccccccccccccccccccccccccccccccccccccccccccccccccccccc"; const OBSERVED_OPENNESS_MS: u64 = 1_700_000_000_000; const RECEIVED_OPENNESS_MS: u64 = 1_700_000_000_250; const OBSERVED_CONSCIENTIOUSNESS_MS: u64 = 1_700_000_030_000; const RECEIVED_CONSCIENTIOUSNESS_MS: u64 = 1_700_000_030_400; +const OBSERVED_EXTRAVERSION_MS: u64 = 1_700_000_060_000; +const RECEIVED_EXTRAVERSION_MS: u64 = 1_700_000_060_350; static RESPONSE_EVENT_TEST_LOCK: Mutex<()> = Mutex::new(()); @@ -145,6 +148,70 @@ fn two_item_korean_path_survives_restart_with_exact_replay() { "item_version_conscientiousness_ko" ); assert!(load_ok(&mut client, "session_missing_ko").is_none()); + let mut times_transaction = client.transaction().unwrap(); + assert_eq!( + load_response_event_times(&mut times_transaction, "session_big_five_ko").unwrap(), + Some(event_times().to_vec()) + ); + assert!( + load_response_event_times(&mut times_transaction, "session_missing_ko") + .unwrap() + .is_none() + ); + times_transaction.commit().unwrap(); +} + +#[test] +fn reloaded_korean_path_records_item_three_and_keeps_the_scoring_prefix() { + let _guard = response_event_test_guard(); + let mut client = test_client(); + reset_response_event_table(&mut client); + apply_response_event_migration(&mut client).unwrap(); + persist_ok(&mut client, &two_item_korean_ledger(), &event_times()); + + let mut continued = load_ok(&mut client, "session_big_five_ko").unwrap(); + continued + .record( + SessionState::Active, + write( + "response_event_extraversion", + "client_extraversion", + "item_version_extraversion_ko", + DIGEST_C, + ), + ) + .unwrap(); + assert_eq!( + persist_ok( + &mut client, + &continued, + &[ + (OBSERVED_OPENNESS_MS, RECEIVED_OPENNESS_MS), + (OBSERVED_CONSCIENTIOUSNESS_MS, RECEIVED_CONSCIENTIOUSNESS_MS), + (OBSERVED_EXTRAVERSION_MS, RECEIVED_EXTRAVERSION_MS), + ] + ), + ResponseEventPersistenceDisposition::Inserted + ); + + let reloaded = load_ok(&mut client, "session_big_five_ko").unwrap(); + assert_eq!(reloaded, continued); + assert_eq!(reloaded.events()[2].sequence(), 3); + assert_eq!( + reloaded.events()[2].item_version_ref(), + "item_version_extraversion_ko" + ); + assert_eq!( + reloaded + .freeze_as(SessionState::Completed, "response_snapshot_big_five_ko") + .unwrap() + .event_refs(), + [ + "response_event_openness", + "response_event_conscientiousness", + "response_event_extraversion" + ] + ); } #[test] @@ -208,7 +275,7 @@ fn inverted_time_blank_session_and_repeatable_read_fail_closed() { assert!(matches!( persist_err(&mut client, &live, &[(1, 2)]), - ResponseEventPersistenceError::InvalidTimestamp + ResponseEventPersistenceError::InvalidEventTimeArity )); assert!(matches!( persist_err( From 1c707dde52c53c78fe27b7661fd9290262065d16 Mon Sep 17 00:00:00 2001 From: Cursor Agent Date: Sun, 16 Aug 2026 16:30:08 +0000 Subject: [PATCH 6/6] docs(response): name persist landing as Active PR #201 Prefer this head over #174, #182, and #53 so later agents do not open another overlapping response_event persist slice. Co-authored-by: Seongho Bae --- docs/TRACEABILITY.md | 6 +++--- docs/adr/0015-persistence-and-transaction-boundaries.md | 2 +- docs/architecture/AS_BUILT_SCHEMA.md | 2 +- docs/architecture/ERD.md | 2 +- 4 files changed, 6 insertions(+), 6 deletions(-) diff --git a/docs/TRACEABILITY.md b/docs/TRACEABILITY.md index 4a95b349..ab5cfc18 100644 --- a/docs/TRACEABILITY.md +++ b/docs/TRACEABILITY.md @@ -23,7 +23,7 @@ An active PR, architecture document, conversation decision, or scheduler plan is | Anonymous core assessment | PRD §3.1, §9.1 | TRD §5, §10; UML anonymous sequence | ADR-0002, ADR-0003, ADR-0005 | Session lifecycle primitives implemented, including creation bound to one published locale-specific release; anonymous credential/HTTP flow is Target | | Pause/resume | PRD §3.1, §9.1 | TRD §5 | ADR-0005 | **Implemented** in `src/session.rs` with fail-closed transitions | | Sequence-aware item delivery evidence | PRD §3.1, §9 | TRD §5–7 | ADR-0005, ADR-0010 | **Implemented** domain primitive in `src/item_delivery.rs`; persistence/API delivery orchestration is Target | -| Idempotent response events | PRD §9.2 | TRD §6 | ADR-0005, ADR-0010 | **Implemented** in `src/response.rs` with canonical SHA-256 payload-digest identity; **Active PR** #174 persists and reloads `response_event` under `READ COMMITTED`; HTTP response transport remains Target | +| Idempotent response events | PRD §9.2 | TRD §6 | ADR-0005, ADR-0010 | **Implemented** in `src/response.rs` with canonical SHA-256 payload-digest identity; **Active PR** #201 persists and reloads `response_event` under `READ COMMITTED`; HTTP response transport remains Target | | Immutable response snapshot before scoring | PRD §9.3 | TRD §5–8 | ADR-0005, ADR-0010 | **Implemented** domain semantics in `src/response.rs` | | Version-pinned scoring | PRD §9.4, §10 | TRD §8 | ADR-0004, ADR-0010 | **Implemented** reusable product-side scoring dispatch contract in `src/scoring.rs` with canonical SHA-256 engine-artifact digest provenance plus `migrations/0011_scoring_request.sql` / `src/postgres_scoring_request.rs` request-identity persistence; live fast-mlsirm integration is Target | | Bounded asynchronous scoring retry/quarantine with stale-worker fencing | PRD §9.4, §10 | TRD §8; ADR-0015 transaction boundary | ADR-0004, ADR-0010, ADR-0015 | **Implemented** product lifecycle plus PostgreSQL enqueue, claim, retry, completion, expiry recovery, and cancellation without transferring a fence; live fast-mlsirm execution remains Target | @@ -55,7 +55,7 @@ An active PR, architecture document, conversation decision, or scheduler plan is | Server-authoritative session state | TRD §5 | `src/session.rs` + session contract tests, including published-release/locale binding at creation | persistence/API concurrency test | | Only Active accepts responses | TRD §5–6 | `SessionState::accepts_responses` + response tests | transport-level rejection test | | Item delivery sequence is positive and evidence-safe | TRD §5–7 | `src/item_delivery.rs` + item-delivery domain tests | durable uniqueness/order/API integration | -| Conflicting idempotency replay fails closed | TRD §6 | `src/response.rs`; **Active PR** #174 `migrations/0020_response_event.sql` unique `(session_ref, client_event_ref)` and `(session_ref, server_sequence)` | HTTP transport rejection test | +| Conflicting idempotency replay fails closed | TRD §6 | `src/response.rs`; **Active PR** #201 `migrations/0020_response_event.sql` unique `(session_ref, client_event_ref)` and `(session_ref, server_sequence)` | HTTP transport rejection test | | Snapshot requires Completed state | TRD §5–6 | `src/response.rs` | transaction atomicity test with persistence | | Scoring uses durable snapshot identity | TRD §8 | `src/scoring.rs` requires a canonical SHA-256 engine-artifact digest | live adapter + retry/outbox integration | | Stale scoring worker cannot complete a newer attempt | TRD §8; ADR-0015 | `src/scoring_job.rs` uses monotonically increasing fencing tokens and rejects stale/expired completion or failure evidence; `src/postgres_scoring_job.rs` persists enqueue, claim, retry, terminal outcomes, expired-lease recovery, and cancellation without transferring a fence | live adapter evidence | @@ -132,7 +132,7 @@ Still-Target logical modules/adapters include remaining product aggregate persis ### Active implementation work that is not protected-main truth -**Active PR** #174 response-event persist/load is not protected-main truth until an unchanged reviewed/check-clean head is integrated. In-progress answers persist in `response_event` with observed/received time, exact replay, and fail-closed rebinding so a two-item path survives restart. Recovery COPY restores those rows. AS_BUILT names the physical slice. Reload exposes stored times as a persist-side projection and continues the scoring prefix after item 3. Completed snapshot persist remains the existing protected-main adapter. HTTP response transport remains outside this slice. Protected-main `#76` data-rights processing-start persistence is already integrated and is not restated here as active work. Prefer this head over another overlapping `response_event` persist PR. +**Active PR** #201 response-event persist/load is not protected-main truth until an unchanged reviewed/check-clean head is integrated. In-progress answers persist in `response_event` with observed/received time, exact replay, and fail-closed rebinding so a two-item path survives restart. Recovery COPY restores those rows. AS_BUILT names the physical slice. Reload exposes stored times as a persist-side projection and continues the scoring prefix after item 3. Completed snapshot persist remains the existing protected-main adapter. HTTP response transport remains outside this slice. Protected-main `#76` data-rights processing-start persistence is already integrated and is not restated here as active work. Prefer this head over #174, #182, and #53. ## 5. ADR traceability by concern diff --git a/docs/adr/0015-persistence-and-transaction-boundaries.md b/docs/adr/0015-persistence-and-transaction-boundaries.md index ae80b63b..be62e072 100644 --- a/docs/adr/0015-persistence-and-transaction-boundaries.md +++ b/docs/adr/0015-persistence-and-transaction-boundaries.md @@ -6,7 +6,7 @@ - Scope: Psychometrics Commons-owned durable state, local transactions, migration boundaries, outbox/inbox integration - Supersedes: none - Superseded by: none -- Current/as-built status: protected main persists integration, scoring-job, data-rights, item-delivery, consent, instrument-release, result-snapshot, response-snapshot, scoring-request, and inbox-consumption slices; in-progress `response_event` persist/load is Active PR #174 and is not protected-main truth until merged +- Current/as-built status: protected main persists integration, scoring-job, data-rights, item-delivery, consent, instrument-release, result-snapshot, response-snapshot, scoring-request, and inbox-consumption slices; in-progress `response_event` persist/load is Active PR #201 and is not protected-main truth until merged - Target status: upstream PostgreSQL 18.x operational persistence with real-database concurrency/crash/recovery evidence and transactional outbox/inbox semantics - Migration status: `migrations/0020_response_event.sql` adds the in-progress event ledger with observed/received timestamps; remaining session HTTP and response HTTP families still must be established from the logical ERD without synthetic provenance backfills diff --git a/docs/architecture/AS_BUILT_SCHEMA.md b/docs/architecture/AS_BUILT_SCHEMA.md index 5fac103f..f36f5950 100644 --- a/docs/architecture/AS_BUILT_SCHEMA.md +++ b/docs/architecture/AS_BUILT_SCHEMA.md @@ -27,7 +27,7 @@ PR #58 (`feat/inbox-consumption-persistence-20260814`) `migrations/0012_integrat ## Active PR response-event physical schema -Active PR #174 `migrations/0020_response_event.sql` and `src/postgres_response_event.rs` persist the in-progress answer ledger so a two-item path can continue after process restart. The slice is **Active PR**, not protected-main truth. It stores opaque `response_event_ref` identity, session binding, client idempotency identity, item version, canonical SHA-256 payload digest, positive `server_sequence`, and distinct `observed_at` / `received_at` timestamps. Exact replay is idempotent and keeps the original times. Client, server, sequence, or session rebinding fails closed. Reload reconstructs `ResponseLedger` in `server_sequence` order under `READ COMMITTED` and exposes stored times as a persist-side projection. HTTP response transport remains outside this slice. +Active PR #201 `migrations/0020_response_event.sql` and `src/postgres_response_event.rs` persist the in-progress answer ledger so a two-item path can continue after process restart. The slice is **Active PR**, not protected-main truth. It stores opaque `response_event_ref` identity, session binding, client idempotency identity, item version, canonical SHA-256 payload digest, positive `server_sequence`, and distinct `observed_at` / `received_at` timestamps. Exact replay is idempotent and keeps the original times. Client, server, sequence, or session rebinding fails closed. Reload reconstructs `ResponseLedger` in `server_sequence` order under `READ COMMITTED` and exposes stored times as a persist-side projection. HTTP response transport remains outside this slice. ## Protected-main scoring-job physical schema diff --git a/docs/architecture/ERD.md b/docs/architecture/ERD.md index 989ea734..00ceec73 100644 --- a/docs/architecture/ERD.md +++ b/docs/architecture/ERD.md @@ -425,7 +425,7 @@ The target ERD deliberately includes several logical entities that are not yet p - `instrument_release` is the locale-specific publication identity already owned by `src/instrument.rs`. Physical `migrations/0006_instrument_release.sql` persists that one-row aggregate (immutable manifest columns plus `publication_state`); HTTP publication transport remains Target. - `data_rights_request` and `data_rights_propagation_state` are the first durable export/deletion slice. Physical `migrations/0003_data_rights_propagation.sql` stores requested-state identity plus one local outbox event per dependent system; verification, processing, completion, and dependent-system execution remain Target. - `item_delivery_event` reflects the already-merged `src/item_delivery.rs` domain primitive; durable persistence/API orchestration is still Target. -- `response_event` is the in-progress answer ledger. Physical `migrations/0020_response_event.sql` on Active PR #174 stores opaque event identity, session binding, client idempotency, item version, payload digest, server sequence, and distinct observed/received timestamps; HTTP response transport remains Target. Completed `response_snapshot` persist is already on protected main. +- `response_event` is the in-progress answer ledger. Physical `migrations/0020_response_event.sql` on Active PR #201 stores opaque event identity, session binding, client idempotency, item version, payload digest, server sequence, and distinct observed/received timestamps; HTTP response transport remains Target. Completed `response_snapshot` persist is already on protected main. - `consent_ledger` and `consent_event` persist the already-merged `src/consent.rs` append-only ledger. Physical persistence is carried by Active PR #49 (`migrations/0005_consent_lifecycle.sql`); HTTP consent transport and derived snapshot tables remain Target. - `participant_identity_link` is the persistence target accepted by ADR-0020. The current `src/participant.rs` `keyverse_subject_ref` field is an application-domain first-link projection, not the future mutable persistence source of truth. - `longitudinal_enrollment`, `longitudinal_observation_record`, and `temporal_analysis_submission` make the ADR-0008 Commons-owned Gyeot/TEPP orchestration boundary explicit. No TEPP analytical kernel is duplicated here.