diff --git a/CHANGELOG.md b/CHANGELOG.md index a52850e1..18aada75 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 +- Active PR #248 adds PostgreSQL 18 durability for immutable normalized longitudinal observation evidence: exact source identity, separate validity/recorded/received/ingested clocks, explicit membership shares whose deferred transaction invariant totals 10,000 basis points, exact replay, tenant isolation, and fail-closed rebinding/immutability. This entry is active-PR evidence only and is not a protected-main or release claim. - Active PR #287 rejects padded or otherwise noncanonical published narrative, style-mapping, interpretation-unit, and approved-selection references before deterministic rendering. Numeric score authority and the protected-main narrative boundary are unchanged; this is active-PR evidence only and not a release claim. - Personal result export delivery authorization (`authorize_result_export_read`) reuses the stored-record `ReadOwnResult` check on the product-owned participant and immutable result, then requires the export's result snapshot reference and copied participant reference to match that exact snapshot. A cross-tenant caller fails closed with the ordinary result-authorization denial before export-binding details are evaluated, so a mismatched export is not an existence oracle. No new permission or persistence is introduced; authorized HTTP transport ships through merged `src/result_export_http.rs` and `openapi/result-exports.yaml`. This is the post-#231 export-delivery guard (ADR-0010 provenance; ADR-0003 tenant-bound authorization). - Personal result export copies one immutable snapshot into JSON and a human-readable report so a purchaser can archive the same Extraversion estimate, standard error, and version provenance they were shown. The owner participant reference stays in both artifacts. Abstained or failed constructs keep their disposition and do not receive an invented score. Approved limitation text is required. diff --git a/docs/TRACEABILITY.md b/docs/TRACEABILITY.md index 9d546879..1a0aff21 100644 --- a/docs/TRACEABILITY.md +++ b/docs/TRACEABILITY.md @@ -29,7 +29,7 @@ An active PR, architecture document, conversation decision, or scheduler plan is | 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 | | Immutable result provenance | PRD §3.1, §9.4 | TRD §9 | ADR-0004, ADR-0010 | **Implemented** in `src/result.rs`; result-serving transport is Target | | Personal JSON and human-readable result export | PRD §3.1, §9.4 | TRD §18 `POST /v1/results/{result_ref}/exports`; ADR-0010 export provenance | ADR-0010 | **Implemented** domain copy through merged #231, delivery guard through merged #249 (`src/result_export_authorization.rs`), and authorized HTTP transport through merged #256 (`src/result_export_http.rs`, `openapi/result-exports.yaml`) | -| Deterministic narrative fallback | PRD §3.2, §9.5 | TRD §17; Architecture narrative view | ADR-0009, ADR-0010, ADR-0018 | **Active PR #287** hardens the non-protected narrative slice so published narrative/rule references must already use canonical opaque spelling; protected-main narrative runtime remains Target | +| Deterministic narrative fallback | PRD §3.2, §9.5 | TRD §17; Architecture narrative view | ADR-0009, ADR-0010, ADR-0018 | **Implemented** through merged #287 (`src/deterministic_narrative.rs`, `src/style_mapping.rs`): published narrative/rule references must already use canonical opaque spelling before deterministic rendering; numeric score authority is unchanged | | Continuous scores remain source of truth; Personality Style is presentation | PRD §3.2 | Measurement Governance; AI Governance | ADR-0018 | Target product narrative mapping; numeric source remains External fast-mlsirm contract | | Immutable instrument release/version lifecycle | PRD §6, §9 | TRD §7; UML publication state | ADR-0005, ADR-0010 | **Implemented** in `src/instrument.rs` plus `migrations/0006_instrument_release.sql` and `src/postgres_instrument_release.rs`: immutable release manifest, exact version/digest/locale/item set, fail-closed Draft/Review/Published/Suspended/Retired lifecycle, idempotent publication events, and new-session eligibility | | Instrument publication requires intended-use scientific/right/locale evidence | PRD §6, §9, §10 | Measurement Governance; publication evidence gate | ADR-0004, ADR-0013, ADR-0019 | **Implemented** policy gate and immutable evidence provenance in `src/instrument.rs`; each real instrument still requires its own rights/locale/scientific evidence artifacts before publication | @@ -44,7 +44,7 @@ An active PR, architecture document, conversation decision, or scheduler plan is | Operation-scoped capability health | PRD §7, §13 | `docs/OPERABILITY.md` §3–4; Deployment/Operations | ADR-0011, ADR-0017 | **Implemented** domain health/readiness contract in `src/health.rs` plus `src/postgres_health.rs` PostgreSQL major/write-readiness and caller-declared relation presence; HTTP probes, measured thresholds, and deployment evidence remain Target | | Korean/English exact locale versions | PRD §3.1, §9.9 | TRD §28; instrument release + locale governance | ADR-0013, ADR-0019 | **Partially implemented**: locale is pinned/validated by `src/instrument.rs`; merged #259 ships protected-main exact `ko-KR`/`en-US` participant report labels that copy immutable scores and provenance; real form content, rights, translation, invariance, HTTP delivery, and accessible reference-client serving remain Target | | WCAG 2.2 AA supported reference client | PRD §9.10 | TRD §27; Quality Attributes | ADR-0002, ADR-0013 | Target; no reference client implementation on evaluated main | -| EMA/ESM longitudinal flow | PRD §4 | TRD §16; UML longitudinal sequence; logical ERD extension | ADR-0008 | External Gyeot/TEPP dependencies + Target Commons enrollment/orchestration adapter; `src/longitudinal_observation.rs` records validity, recorded, received, and ingested clocks with explicit membership shares. Enrollment state, PostgreSQL persistence, HTTP, Gyeot collection, and TEPP kernels remain Target | +| EMA/ESM longitudinal flow | PRD §4 | TRD §16; UML longitudinal sequence; logical ERD extension | ADR-0008 | External Gyeot/TEPP dependencies + Target Commons enrollment/orchestration adapter; `src/longitudinal_observation.rs` records validity, recorded, received, and ingested clocks with explicit membership shares. **Active PR #248 / IMPLEMENTED_ON_ACTIVE_PR** adds PostgreSQL 18 persistence for immutable normalized observation records and membership shares via `migrations/0031_longitudinal_observation.sql` and `src/postgres_longitudinal_observation.rs`; enrollment persistence, HTTP, Gyeot collection, and TEPP kernels remain Target | | Measurement Workbench | PRD §6 | C4/component view; UML publication-evidence sequence; Measurement Governance | ADR-0001, ADR-0002, ADR-0004, ADR-0019 | Target; fast-mlsirm/Inkspan/RankWeave are External dependencies | | Headless replaceable clients | PRD §7 | TRD §1, §18; C4 | ADR-0001, ADR-0002 | Architecture established; public transport is Target | | Community/Hosted/Enterprise profiles | PRD §7, §13 | TRD deployment sections; Deployment/Operations | ADR-0011, ADR-0017 | Target deployment packaging/evidence | diff --git a/docs/architecture/AS_BUILT_SCHEMA.md b/docs/architecture/AS_BUILT_SCHEMA.md index 0d515d25..a6b8290e 100644 --- a/docs/architecture/AS_BUILT_SCHEMA.md +++ b/docs/architecture/AS_BUILT_SCHEMA.md @@ -26,6 +26,17 @@ The protected-main integration identity is source- and tenant-scoped. A physical PR #218 (`migrations/0014_assessment_session.sql`, `migrations/0016_assessment_session_command.sql`, and `src/postgres_assessment_session.rs`) persist and load one assessment-session identity bound to a published locale-specific release, plus append-only command history. New sessions start only through `created_session_for_start` / `start_created_assessment_session` / `start_created_assessment_session_from_stored_release`. Durable start locks `instrument_release` with `SELECT … FOR UPDATE` so a stale in-memory Published object cannot insert after persist Suspend or Retire. First insert through `persist_assessment_session` takes the same lock, so a reconstituted Created aggregate cannot insert after that later persist. When that lock finds a missing or unpublished release, persist still classifies an exact stored Created row as duplicate so a concurrent retry after the first insert commits cannot turn a later Suspend or Retire into a false unpublished failure. Exact replay of an already stored start or Created row still returns the original session after a later persist Suspend or Retire. The slice is **Active PR**, not protected-main truth. It stores participant, release, version, digest, locale, current state, and creation time. Exact replay is idempotent. Rebinding any stored field or command evidence, or persisting a shorter command history than already stored, fails closed so a stale Activate-only worker cannot rewind Pause/Resume. Command persist locks the `assessment_session` header row with `SELECT … FOR UPDATE` before inserting or counting commands. Load restores created identity without asking whether the release still accepts new sessions, then replays commands so Activate/Pause/Resume survive restart. Isolation is the global opaque `session_ref` primary key; this slice does not add `tenant_ref` because the domain `AssessmentSession` aggregate does not carry tenant. Persist-backed `POST /v1/sessions` / `GET /v1/sessions/{session_ref}` (`src/session_http.rs`, `openapi/sessions.yaml`) sit on this start path. Command HTTP remains outside this slice. #205 is the unlocked-peek first-insert-seal predecessor; #209 is the weaker NotFound-allows-insert competitor; #198 is the exact start-replay predecessor; #180 is the stored-publication lock predecessor; #188 is the in-memory replay predecessor that still lacks the store lock; #153 is the in-memory-start predecessor; #164 is the unlocked stored-load predecessor; #146 is the header-lock predecessor; #129 is the sequential stale-prefix predecessor; #125 is the command-history predecessor that still rewinds on a stale shorter persist; #109 is the persist-and-load predecessor. +## Active PR longitudinal-observation physical schema + +PR #248 (`migrations/0031_longitudinal_observation.sql` and `src/postgres_longitudinal_observation.rs`) is **Active PR / IMPLEMENTED_ON_ACTIVE_PR**, not protected-main truth. It maps the logical `longitudinal_observation_record` semantics in `ERD.md` onto two product-owned PostgreSQL 18 relations without moving Gyeot collection or TEPP temporal/multilevel/multiple-membership analysis into this repository. + +- `longitudinal_observation` is the immutable parent record. Its opaque observation, tenant, enrollment, source-system, and source-observation references preserve the logical record identity and source provenance; validity, source-recorded, platform-received, and durable-ingestion clocks remain separate; civil-time/UTC-offset and clock-anomaly evidence are retained rather than collapsed. +- `longitudinal_membership_share` is the immutable one-to-many membership relation owned by one observation record. Each opaque membership-context reference carries an integer share, and the migration defers the aggregate invariant that one observation's shares total exactly 10,000 basis points until transaction completion. +- Exact source identity is unique on `(tenant_ref, enrollment_ref, source_system_ref, source_observation_ref)`. Exact replay is idempotent; source rebinding or stored-evidence mutation fails closed. The same source tuple may repeat under a different tenant only with a different observation-record identity. UPDATE, DELETE, and TRUNCATE guards preserve the logical ERD's immutable-evidence semantics. +- Real PostgreSQL tests exercise persistence/reload, tenant isolation, replay/conflict classification, source-identity rebinding rejection, schema integrity, membership totals, immutability, and missing-schema failure. No direct Gyeot or TEPP database access is introduced. + +This is an explicit logical-to-physical reconciliation: the physical split preserves one logical observation record with zero-or-more explicit membership-share evidence rows; it does not create a second collection or analysis kernel. Durable longitudinal enrollment, HTTP transport, live Gyeot/TEPP adapters, recovery acceptance, and research-release registration remain Target. + ## Active PR outbox delivery-lease physical schema PR #60 (`feat/outbox-delivery-lease-20260814`) extends the protected-main `integration_outbox` relation through `migrations/0013_outbox_delivery_lease.sql` and `src/postgres_integration.rs`. The extension is **Active PR**, not protected-main truth. diff --git a/docs/architecture/ERD.md b/docs/architecture/ERD.md index 7e12c367..bf825bae 100644 --- a/docs/architecture/ERD.md +++ b/docs/architecture/ERD.md @@ -508,7 +508,7 @@ A physical schema must enforce equivalents of the following constraints: | at most one current Active `participant_identity_link` per participant under the accepted single-account-link policy | unambiguous current account projection | | unique active `(tenant_ref, identity_issuer, identity_subject_ref)` unless an explicit account-merge ADR permits otherwise | prevent one external subject from silently owning multiple product participants | | unique `(request_ref, retained_scope_ref)` for retained data-rights completion evidence | prevent duplicate retained-scope evidence while preserving tenant/request binding | -| unique `(enrollment_ref, source_system_ref, source_observation_ref)` | longitudinal ingestion replay safety | +| unique `(tenant_ref, enrollment_ref, source_system_ref, source_observation_ref)` | tenant-scoped longitudinal ingestion replay safety | | unique analysis-submission idempotency identity per `(enrollment_ref, analysis_spec_ref, observation_set_digest)` | repeatable TEPP dispatch | | unique `(tenant_ref, source, outbox_event_ref)` or an equivalently stronger globally unique event identity with tenant binding | durable outbound event identity | | unique `(tenant_ref, consumer_name, source, source_event_ref)` | tenant-bound inbox deduplication | diff --git a/migrations/0031_longitudinal_observation.sql b/migrations/0031_longitudinal_observation.sql new file mode 100644 index 00000000..6e3860e6 --- /dev/null +++ b/migrations/0031_longitudinal_observation.sql @@ -0,0 +1,270 @@ +-- Durable, immutable normalized longitudinal observation evidence. + +-- Keep the database boundary aligned with Rust 1.97's Unicode 17 char::is_numeric contract. +-- PostgreSQL POSIX [[:digit:]] covers decimal digits but does not prove parity for Unicode letter +-- numbers and other numbers such as Roman numerals, superscripts, and vulgar fractions. The +-- generated int4multirange below is the exact Unicode 17 numeric code-point set used by rustc 1.97. +-- Mixed opaque identifiers remain valid because a value is numeric-like only when it contains at +-- least one numeric code point and every character is numeric or an allowed numeric spelling token. +-- The fixed pg_catalog search path prevents caller-controlled schemas from changing function +-- resolution inside this CHECK helper. +CREATE OR REPLACE FUNCTION longitudinal_reference_is_valid(reference_value text) +RETURNS boolean +LANGUAGE sql +IMMUTABLE +PARALLEL SAFE +SET search_path = pg_catalog +AS $longitudinal_reference$ + WITH reference_character AS ( + SELECT substr(reference_value, character_index, 1) AS character_text + FROM generate_series(1, character_length(reference_value)) AS character_index + ), + reference_classification AS ( + SELECT + character_text, + ascii(character_text) <@ '{[48,58),[178,180),[185,186),[188,191),[1632,1642),[1776,1786),[1984,1994),[2406,2416),[2534,2544),[2548,2554),[2662,2672),[2790,2800),[2918,2928),[2930,2936),[3046,3059),[3174,3184),[3192,3199),[3302,3312),[3416,3423),[3430,3449),[3558,3568),[3664,3674),[3792,3802),[3872,3892),[4160,4170),[4240,4250),[4969,4989),[5870,5873),[6112,6122),[6128,6138),[6160,6170),[6470,6480),[6608,6619),[6784,6794),[6800,6810),[6992,7002),[7088,7098),[7232,7242),[7248,7258),[8304,8305),[8308,8314),[8320,8330),[8528,8579),[8581,8586),[9312,9372),[9450,9472),[10102,10132),[11517,11518),[12295,12296),[12321,12330),[12344,12347),[12690,12694),[12832,12842),[12872,12880),[12881,12896),[12928,12938),[12977,12992),[42528,42538),[42726,42736),[43056,43062),[43216,43226),[43264,43274),[43472,43482),[43504,43514),[43600,43610),[44016,44026),[65296,65306),[65799,65844),[65856,65913),[65930,65932),[66273,66300),[66336,66340),[66369,66370),[66378,66379),[66513,66518),[66720,66730),[67672,67680),[67705,67712),[67751,67760),[67835,67840),[67862,67868),[68028,68030),[68032,68048),[68050,68096),[68160,68169),[68221,68223),[68253,68256),[68331,68336),[68440,68448),[68472,68480),[68521,68528),[68858,68864),[68912,68922),[68928,68938),[69216,69247),[69405,69415),[69457,69461),[69573,69580),[69714,69744),[69872,69882),[69942,69952),[70096,70106),[70113,70133),[70384,70394),[70736,70746),[70864,70874),[71248,71258),[71360,71370),[71376,71396),[71472,71484),[71904,71923),[72016,72026),[72688,72698),[72784,72813),[73040,73050),[73120,73130),[73184,73194),[73552,73562),[73664,73685),[74752,74863),[90416,90426),[92768,92778),[92864,92874),[93008,93018),[93019,93026),[93552,93562),[93824,93847),[94196,94199),[118000,118010),[119488,119508),[119520,119540),[119648,119673),[120782,120832),[123200,123210),[123632,123642),[124144,124154),[124401,124411),[125127,125136),[125264,125274),[126065,126124),[126125,126128),[126129,126133),[126209,126254),[126255,126270),[127232,127245),[130032,130042)}'::int4multirange + AS is_numeric + FROM reference_character + ) + SELECT + reference_value IS NOT NULL + AND reference_value <> '' + AND left(reference_value, 1) !~ '[[:space:]]' + AND right(reference_value, 1) !~ '[[:space:]]' + AND reference_value !~ '[[:cntrl:]]' + AND NOT COALESCE( + bool_or(is_numeric) + AND bool_and( + is_numeric + OR character_text = ANY ( + ARRAY[ + '+', + '-', + '.', + ',', + 'e', + 'E', + U&'\066B', + U&'\066C', + U&'\FF0E', + U&'\FF0C' + ] + ) + ), + FALSE + ) + FROM reference_classification; +$longitudinal_reference$; + +CREATE TABLE IF NOT EXISTS longitudinal_observation ( + observation_record_ref text PRIMARY KEY, + tenant_ref text NOT NULL, + enrollment_ref text NOT NULL, + source_system_ref text NOT NULL, + source_observation_ref text NOT NULL, + construct_ref text NOT NULL, + measure_ref text NOT NULL, + validity_start_at_unix_ms bigint NOT NULL, + validity_end_at_unix_ms bigint NOT NULL, + recorded_at_unix_ms bigint NOT NULL, + received_at_unix_ms bigint NOT NULL, + ingested_at_unix_ms bigint NOT NULL, + timezone_name text NOT NULL, + utc_offset_minutes smallint NOT NULL, + clock_anomaly_code text, + CONSTRAINT longitudinal_observation_source_identity_unique + UNIQUE (tenant_ref, enrollment_ref, source_system_ref, source_observation_ref), + CONSTRAINT longitudinal_observation_reference_check CHECK ( + longitudinal_reference_is_valid(observation_record_ref) + AND longitudinal_reference_is_valid(tenant_ref) + AND longitudinal_reference_is_valid(enrollment_ref) + AND longitudinal_reference_is_valid(source_system_ref) + AND longitudinal_reference_is_valid(source_observation_ref) + AND longitudinal_reference_is_valid(construct_ref) + AND longitudinal_reference_is_valid(measure_ref) + ), + CONSTRAINT longitudinal_observation_time_check CHECK ( + validity_start_at_unix_ms > 0 AND validity_end_at_unix_ms >= validity_start_at_unix_ms + AND recorded_at_unix_ms > 0 AND received_at_unix_ms > 0 AND ingested_at_unix_ms >= received_at_unix_ms + ), + CONSTRAINT longitudinal_observation_timezone_check CHECK ( + timezone_name = btrim(timezone_name) AND timezone_name <> '' AND timezone_name !~ '^[0-9]+$' + AND utc_offset_minutes BETWEEN -720 AND 840 + ), + CONSTRAINT longitudinal_observation_anomaly_check CHECK ( + (clock_anomaly_code IS NULL AND recorded_at_unix_ms <= received_at_unix_ms) + OR ( + clock_anomaly_code = 'recorded_after_received' + AND recorded_at_unix_ms > received_at_unix_ms + ) + ) +); + +CREATE TABLE IF NOT EXISTS longitudinal_membership_share ( + observation_record_ref text NOT NULL + REFERENCES longitudinal_observation(observation_record_ref) ON DELETE RESTRICT, + membership_sequence bigint NOT NULL, + membership_context_ref text NOT NULL, + weight_parts_per_10_000 integer NOT NULL, + PRIMARY KEY (observation_record_ref, membership_sequence), + CONSTRAINT longitudinal_membership_context_unique + UNIQUE (observation_record_ref, membership_context_ref), + CONSTRAINT longitudinal_membership_sequence_check CHECK (membership_sequence > 0), + CONSTRAINT longitudinal_membership_reference_check CHECK ( + longitudinal_reference_is_valid(membership_context_ref) + ), + CONSTRAINT longitudinal_membership_weight_check CHECK ( + weight_parts_per_10_000 BETWEEN 1 AND 10000 + ) +); + +-- Reapply evolving CHECK definitions so an idempotent rerun upgrades a schema created by an +-- earlier iteration of this not-yet-shipped migration rather than trusting constraint names alone. +ALTER TABLE longitudinal_observation + DROP CONSTRAINT IF EXISTS longitudinal_observation_reference_check; +ALTER TABLE longitudinal_observation + ADD CONSTRAINT longitudinal_observation_reference_check CHECK ( + longitudinal_reference_is_valid(observation_record_ref) + AND longitudinal_reference_is_valid(tenant_ref) + AND longitudinal_reference_is_valid(enrollment_ref) + AND longitudinal_reference_is_valid(source_system_ref) + AND longitudinal_reference_is_valid(source_observation_ref) + AND longitudinal_reference_is_valid(construct_ref) + AND longitudinal_reference_is_valid(measure_ref) + ); + +ALTER TABLE longitudinal_observation + DROP CONSTRAINT IF EXISTS longitudinal_observation_anomaly_check; +ALTER TABLE longitudinal_observation + ADD CONSTRAINT longitudinal_observation_anomaly_check CHECK ( + (clock_anomaly_code IS NULL AND recorded_at_unix_ms <= received_at_unix_ms) + OR ( + clock_anomaly_code = 'recorded_after_received' + AND recorded_at_unix_ms > received_at_unix_ms + ) + ) NOT VALID; + +ALTER TABLE longitudinal_membership_share + DROP CONSTRAINT IF EXISTS longitudinal_membership_reference_check; +ALTER TABLE longitudinal_membership_share + ADD CONSTRAINT longitudinal_membership_reference_check CHECK ( + longitudinal_reference_is_valid(membership_context_ref) + ); + +-- Existing rows from a partial pre-merge rollout must already satisfy the clock/code relation. +-- Fail the migration with an operator-readable error before validating the strengthened CHECK. +DO $longitudinal_anomaly_preflight$ +BEGIN + IF EXISTS ( + SELECT 1 + FROM longitudinal_observation + WHERE NOT ( + (clock_anomaly_code IS NULL AND recorded_at_unix_ms <= received_at_unix_ms) + OR ( + clock_anomaly_code = 'recorded_after_received' + AND recorded_at_unix_ms > received_at_unix_ms + ) + ) + ) THEN + RAISE EXCEPTION 'stored longitudinal observation clock anomaly evidence is inconsistent' + USING ERRCODE = '23514', + CONSTRAINT = 'longitudinal_observation_anomaly_check'; + END IF; +END; +$longitudinal_anomaly_preflight$; + +ALTER TABLE longitudinal_observation + VALIDATE CONSTRAINT longitudinal_observation_anomaly_check; + +-- Legitimate rows are append-only. Keep the INSERT guard as a named, early diagnostic in addition +-- to the CHECK constraint; the CHECK remains authoritative even if immutability is disabled for an +-- operator repair or integrity exercise. +CREATE OR REPLACE FUNCTION validate_longitudinal_observation_anomaly_insert() +RETURNS trigger +LANGUAGE plpgsql +AS $$ +BEGIN + IF ( + (NEW.clock_anomaly_code IS NULL AND NEW.recorded_at_unix_ms <= NEW.received_at_unix_ms) + OR ( + NEW.clock_anomaly_code = 'recorded_after_received' + AND NEW.recorded_at_unix_ms > NEW.received_at_unix_ms + ) + ) THEN + RETURN NEW; + END IF; + + RAISE EXCEPTION 'longitudinal observation clock anomaly code does not match clock order' + USING ERRCODE = '23514', + CONSTRAINT = 'longitudinal_observation_anomaly_check'; +END; +$$; + +DROP TRIGGER IF EXISTS longitudinal_observation_anomaly_insert ON longitudinal_observation; +CREATE TRIGGER longitudinal_observation_anomaly_insert +BEFORE INSERT ON longitudinal_observation +FOR EACH ROW EXECUTE FUNCTION validate_longitudinal_observation_anomaly_insert(); + +CREATE OR REPLACE FUNCTION reject_longitudinal_mutation() +RETURNS trigger +LANGUAGE plpgsql +AS $$ +BEGIN + RAISE EXCEPTION 'longitudinal observation evidence is immutable' + USING ERRCODE = '55000'; +END; +$$; + +DROP TRIGGER IF EXISTS longitudinal_observation_immutable_update ON longitudinal_observation; +CREATE TRIGGER longitudinal_observation_immutable_update +BEFORE UPDATE OR DELETE ON longitudinal_observation +FOR EACH ROW EXECUTE FUNCTION reject_longitudinal_mutation(); + +DROP TRIGGER IF EXISTS longitudinal_observation_immutable_truncate ON longitudinal_observation; +CREATE TRIGGER longitudinal_observation_immutable_truncate +BEFORE TRUNCATE ON longitudinal_observation +FOR EACH STATEMENT EXECUTE FUNCTION reject_longitudinal_mutation(); + +DROP TRIGGER IF EXISTS longitudinal_membership_immutable_update ON longitudinal_membership_share; +CREATE TRIGGER longitudinal_membership_immutable_update +BEFORE UPDATE OR DELETE ON longitudinal_membership_share +FOR EACH ROW EXECUTE FUNCTION reject_longitudinal_mutation(); + +DROP TRIGGER IF EXISTS longitudinal_membership_immutable_truncate ON longitudinal_membership_share; +CREATE TRIGGER longitudinal_membership_immutable_truncate +BEFORE TRUNCATE ON longitudinal_membership_share +FOR EACH STATEMENT EXECUTE FUNCTION reject_longitudinal_mutation(); + +CREATE OR REPLACE FUNCTION enforce_longitudinal_membership_total() +RETURNS trigger +LANGUAGE plpgsql +AS $$ +DECLARE + total_weight bigint; +BEGIN + SELECT COALESCE(SUM(weight_parts_per_10_000), 0) + INTO total_weight + FROM longitudinal_membership_share + WHERE observation_record_ref = NEW.observation_record_ref; + IF total_weight <> 10000 THEN + RAISE EXCEPTION 'longitudinal membership shares must sum to 10000' + USING ERRCODE = '23514'; + END IF; + RETURN NULL; +END; +$$; + +-- A membership-only trigger cannot observe an observation whose membership vector is empty, +-- because no child INSERT exists to fire it. Defer the same invariant from the parent INSERT so +-- the complete vector can be inserted in the transaction while a header-only commit still fails. +DROP TRIGGER IF EXISTS longitudinal_observation_membership_total_check ON longitudinal_observation; +CREATE CONSTRAINT TRIGGER longitudinal_observation_membership_total_check +AFTER INSERT ON longitudinal_observation +DEFERRABLE INITIALLY DEFERRED +FOR EACH ROW EXECUTE FUNCTION enforce_longitudinal_membership_total(); + +DROP TRIGGER IF EXISTS longitudinal_membership_total_check ON longitudinal_membership_share; +CREATE CONSTRAINT TRIGGER longitudinal_membership_total_check +AFTER INSERT ON longitudinal_membership_share +DEFERRABLE INITIALLY DEFERRED +FOR EACH ROW EXECUTE FUNCTION enforce_longitudinal_membership_total(); \ No newline at end of file diff --git a/src/lib.rs b/src/lib.rs index 85957052..4ab4983f 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -37,6 +37,7 @@ pub mod postgres_inbox_consumption; pub mod postgres_instrument_release; pub mod postgres_integration; pub mod postgres_item_delivery; +pub mod postgres_longitudinal_observation; pub mod postgres_response_snapshot; pub mod postgres_result_snapshot; pub mod postgres_scoring_job; diff --git a/src/postgres_longitudinal_observation.rs b/src/postgres_longitudinal_observation.rs new file mode 100644 index 00000000..a236c5cb --- /dev/null +++ b/src/postgres_longitudinal_observation.rs @@ -0,0 +1,471 @@ +//! `PostgreSQL` persistence for immutable normalized longitudinal observations. +//! +//! Commons persists tenant-bound ingestion evidence only. Gyeot is the collection +//! service that captures longitudinal observations, while TEPP is the analysis +//! service responsible for temporal and multiple-membership models. This module +//! stores the normalized evidence without taking over either service's responsibility. + +use crate::longitudinal_observation::{ + ClockAnomaly, LongitudinalObservationInput, LongitudinalObservationRecord, + LongitudinalObservationSet, MembershipShareInput, ObservationTimeInput, +}; +use crate::reference::normalized_reference; +use postgres::Transaction; +use std::error::Error; +use std::fmt::{Display, Formatter}; + +const LONGITUDINAL_OBSERVATION_MIGRATION: &str = + include_str!("../migrations/0031_longitudinal_observation.sql"); + +/// Outcome of persisting one immutable normalized observation. +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +#[non_exhaustive] +pub enum LongitudinalObservationPersistenceDisposition { + /// A new observation and all membership shares were inserted. + Inserted, + /// The exact immutable observation was already present. + Duplicate, +} + +/// Fail-closed error for longitudinal observation persistence. +#[derive(Debug)] +#[non_exhaustive] +pub enum LongitudinalObservationPersistenceError { + /// A tenant or observation-record reference was blank, numeric-like, padded, or noncanonical. + InvalidReference, + /// A Rust clock or sequence cannot be represented by `PostgreSQL` `bigint`. + InvalidNumericRange, + /// An observation or tenant-scoped source identity was replayed with different evidence. + ConflictingReplay, + /// Persisted rows cannot be reconstructed into one valid immutable domain record. + CorruptHistory, + /// Persistence requires `PostgreSQL` `READ COMMITTED` isolation. + UnsupportedIsolationLevel, + /// `PostgreSQL` rejected or could not execute the operation. + Database(postgres::Error), +} + +impl Display for LongitudinalObservationPersistenceError { + fn fmt(&self, formatter: &mut Formatter<'_>) -> std::fmt::Result { + formatter.write_str(match self { + Self::InvalidReference => { + "longitudinal observation tenant and record references must use their exact opaque form" + } + Self::InvalidNumericRange => { + "longitudinal observation clocks or membership sequence exceed the PostgreSQL bigint range" + } + Self::ConflictingReplay => { + "longitudinal observation identity was replayed with conflicting immutable evidence" + } + Self::CorruptHistory => { + "stored longitudinal observation evidence cannot be reconstructed safely" + } + Self::UnsupportedIsolationLevel => { + "longitudinal observation persistence requires read committed isolation" + } + Self::Database(_) => "PostgreSQL longitudinal observation persistence failed", + }) + } +} + +impl Error for LongitudinalObservationPersistenceError { + fn source(&self) -> Option<&(dyn Error + 'static)> { + match self { + Self::Database(error) => Some(error), + _ => None, + } + } +} + +impl From for LongitudinalObservationPersistenceError { + fn from(error: postgres::Error) -> Self { + Self::Database(error) + } +} + +/// Apply the idempotent longitudinal-observation migration. +/// +/// The migration uses unqualified database object names. Callers must set the +/// intended `PostgreSQL` `search_path` before applying it and use the same schema +/// context for related persistence queries. +/// +/// # Errors +/// +/// Returns the `PostgreSQL` error when the schema cannot be installed. +pub fn apply_longitudinal_observation_migration( + client: &mut impl postgres::GenericClient, +) -> Result<(), postgres::Error> { + client.batch_execute(LONGITUDINAL_OBSERVATION_MIGRATION) +} + +/// Persist one tenant-bound normalized observation and its complete membership vector. +/// +/// Exact replay is idempotent. Reusing either the global Commons observation +/// identity or the tenant-scoped `(enrollment, source system, source observation)` +/// identity with different evidence fails closed. The same source tuple may exist +/// in a different tenant only under a different observation-record identity. +/// The caller owns the transaction and must use `READ COMMITTED` so a concurrent +/// winner is visible to replay classification. +/// +/// # Errors +/// +/// Returns [`LongitudinalObservationPersistenceError`] for a noncanonical tenant, +/// unsupported isolation, unrepresentable numeric values, conflicting replay, or +/// a database failure. +pub fn persist_longitudinal_observation( + transaction: &mut Transaction<'_>, + tenant_ref: &str, + record: &LongitudinalObservationRecord, +) -> Result +{ + let tenant_ref = required_reference(tenant_ref)?; + require_read_committed(transaction)?; + let validity_start = postgres_u64(record.validity_start_at_unix_ms())?; + let validity_end = postgres_u64(record.validity_end_at_unix_ms())?; + let recorded_at = postgres_u64(record.recorded_at_unix_ms())?; + let received_at = postgres_u64(record.received_at_unix_ms())?; + let ingested_at = postgres_u64(record.ingested_at_unix_ms())?; + let anomaly_code = clock_anomaly_code(record.clock_anomaly()); + + let inserted = transaction.execute( + "INSERT INTO longitudinal_observation (\ + observation_record_ref, tenant_ref, enrollment_ref, source_system_ref, \ + source_observation_ref, construct_ref, measure_ref, validity_start_at_unix_ms, \ + validity_end_at_unix_ms, recorded_at_unix_ms, received_at_unix_ms, \ + ingested_at_unix_ms, timezone_name, utc_offset_minutes, clock_anomaly_code\ + ) VALUES ($1,$2,$3,$4,$5,$6,$7,$8,$9,$10,$11,$12,$13,$14,$15) \ + ON CONFLICT DO NOTHING", + &[ + &record.observation_record_ref(), + &tenant_ref, + &record.enrollment_ref(), + &record.source_system_ref(), + &record.source_observation_ref(), + &record.construct_ref(), + &record.measure_ref(), + &validity_start, + &validity_end, + &recorded_at, + &received_at, + &ingested_at, + &record.timezone_name(), + &record.utc_offset_minutes(), + &anomaly_code, + ], + )?; + if inserted == 0 { + return classify_existing(transaction, tenant_ref, record); + } + + for (index, share) in record.membership_shares().iter().enumerate() { + let sequence = postgres_usize(index + 1)?; + let weight = i32::from(share.weight_parts_per_10_000()); + transaction.execute( + "INSERT INTO longitudinal_membership_share (\ + observation_record_ref, membership_sequence, membership_context_ref, \ + weight_parts_per_10_000\ + ) VALUES ($1,$2,$3,$4)", + &[ + &record.observation_record_ref(), + &sequence, + &share.membership_context_ref(), + &weight, + ], + )?; + } + Ok(LongitudinalObservationPersistenceDisposition::Inserted) +} + +/// Load one immutable observation after process restart within a tenant boundary. +/// +/// The loader rebuilds the public domain record through the same validation path +/// used for fresh ingestion. Tenant mismatch and a missing record return `None`; +/// incomplete, non-contiguous, numerically invalid, or internally inconsistent +/// stored evidence fails closed as [`LongitudinalObservationPersistenceError::CorruptHistory`]. +/// +/// # Errors +/// +/// Returns [`LongitudinalObservationPersistenceError`] for noncanonical references, +/// corrupt stored evidence, or a database failure. +pub fn load_longitudinal_observation( + transaction: &mut Transaction<'_>, + tenant_ref: &str, + observation_record_ref: &str, +) -> Result, LongitudinalObservationPersistenceError> { + let tenant_ref = required_reference(tenant_ref)?; + let observation_record_ref = required_reference(observation_record_ref)?; + let row = transaction.query_opt( + "SELECT enrollment_ref, source_system_ref, source_observation_ref, construct_ref, \ + measure_ref, validity_start_at_unix_ms, validity_end_at_unix_ms, \ + recorded_at_unix_ms, received_at_unix_ms, ingested_at_unix_ms, \ + timezone_name, utc_offset_minutes, clock_anomaly_code \ + FROM longitudinal_observation \ + WHERE tenant_ref = $1 AND observation_record_ref = $2", + &[&tenant_ref, &observation_record_ref], + )?; + let Some(row) = row else { + return Ok(None); + }; + + let membership_rows = transaction.query( + "SELECT membership_sequence, membership_context_ref, weight_parts_per_10_000 \ + FROM longitudinal_membership_share \ + WHERE observation_record_ref = $1 ORDER BY membership_sequence", + &[&observation_record_ref], + )?; + let mut stored_memberships = Vec::with_capacity(membership_rows.len()); + for (index, membership_row) in membership_rows.iter().enumerate() { + require_membership_sequence(membership_row.get(0), postgres_usize(index + 1)?)?; + let membership_context_ref: String = membership_row.get(1); + let weight_parts_per_10_000 = database_u16(membership_row.get(2))?; + stored_memberships.push((membership_context_ref, weight_parts_per_10_000)); + } + let membership_inputs = stored_memberships + .iter() + .map( + |(membership_context_ref, weight_parts_per_10_000)| MembershipShareInput { + membership_context_ref, + weight_parts_per_10_000: *weight_parts_per_10_000, + }, + ) + .collect::>(); + + let enrollment_ref: String = row.get(0); + let source_system_ref: String = row.get(1); + let source_observation_ref: String = row.get(2); + let construct_ref: String = row.get(3); + let measure_ref: String = row.get(4); + let validity_start_at_unix_ms = database_u64(row.get(5))?; + let validity_end_at_unix_ms = database_u64(row.get(6))?; + let recorded_at_unix_ms = database_u64(row.get(7))?; + let received_at_unix_ms = database_u64(row.get(8))?; + let ingested_at_unix_ms = database_u64(row.get(9))?; + let timezone_name: String = row.get(10); + let utc_offset_minutes: i16 = row.get(11); + let stored_anomaly_code: Option = row.get(12); + + let record = LongitudinalObservationSet::new() + .ingest(LongitudinalObservationInput { + observation_record_ref, + enrollment_ref: &enrollment_ref, + source_system_ref: &source_system_ref, + source_observation_ref: &source_observation_ref, + construct_ref: &construct_ref, + measure_ref: &measure_ref, + membership_shares: &membership_inputs, + time: ObservationTimeInput { + validity_start_at_unix_ms, + validity_end_at_unix_ms, + recorded_at_unix_ms, + received_at_unix_ms, + ingested_at_unix_ms, + timezone_name: &timezone_name, + utc_offset_minutes, + }, + }) + .map_err(|_| LongitudinalObservationPersistenceError::CorruptHistory)?; + require_clock_anomaly_code(stored_anomaly_code.as_deref(), record.clock_anomaly())?; + Ok(Some(record)) +} + +fn classify_existing( + transaction: &mut Transaction<'_>, + tenant_ref: &str, + record: &LongitudinalObservationRecord, +) -> Result +{ + let rows = transaction.query( + "SELECT observation_record_ref, tenant_ref, enrollment_ref, source_system_ref, \ + source_observation_ref, construct_ref, measure_ref, validity_start_at_unix_ms, \ + validity_end_at_unix_ms, recorded_at_unix_ms, received_at_unix_ms, \ + ingested_at_unix_ms, timezone_name, utc_offset_minutes, clock_anomaly_code \ + FROM longitudinal_observation \ + WHERE observation_record_ref = $1 \ + OR (tenant_ref = $2 AND enrollment_ref = $3 AND source_system_ref = $4 \ + AND source_observation_ref = $5)", + &[ + &record.observation_record_ref(), + &tenant_ref, + &record.enrollment_ref(), + &record.source_system_ref(), + &record.source_observation_ref(), + ], + )?; + if rows.len() != 1 { + return Err(LongitudinalObservationPersistenceError::ConflictingReplay); + } + let row = &rows[0]; + let anomaly_code = clock_anomaly_code(record.clock_anomaly()).map(str::to_owned); + let exact_header = row.get::<_, String>(0) == record.observation_record_ref() + && row.get::<_, String>(1) == tenant_ref + && row.get::<_, String>(2) == record.enrollment_ref() + && row.get::<_, String>(3) == record.source_system_ref() + && row.get::<_, String>(4) == record.source_observation_ref() + && row.get::<_, String>(5) == record.construct_ref() + && row.get::<_, String>(6) == record.measure_ref() + && row.get::<_, i64>(7) == postgres_u64(record.validity_start_at_unix_ms())? + && row.get::<_, i64>(8) == postgres_u64(record.validity_end_at_unix_ms())? + && row.get::<_, i64>(9) == postgres_u64(record.recorded_at_unix_ms())? + && row.get::<_, i64>(10) == postgres_u64(record.received_at_unix_ms())? + && row.get::<_, i64>(11) == postgres_u64(record.ingested_at_unix_ms())? + && row.get::<_, String>(12) == record.timezone_name() + && row.get::<_, i16>(13) == record.utc_offset_minutes() + && row.get::<_, Option>(14) == anomaly_code; + if !exact_header { + return Err(LongitudinalObservationPersistenceError::ConflictingReplay); + } + + let stored = transaction.query( + "SELECT membership_sequence, membership_context_ref, weight_parts_per_10_000 \ + FROM longitudinal_membership_share \ + WHERE observation_record_ref = $1 ORDER BY membership_sequence", + &[&record.observation_record_ref()], + )?; + if stored.len() != record.membership_shares().len() { + return Err(LongitudinalObservationPersistenceError::ConflictingReplay); + } + for (index, (row, share)) in stored.iter().zip(record.membership_shares()).enumerate() { + if row.get::<_, i64>(0) != postgres_usize(index + 1)? + || row.get::<_, String>(1) != share.membership_context_ref() + || row.get::<_, i32>(2) != i32::from(share.weight_parts_per_10_000()) + { + return Err(LongitudinalObservationPersistenceError::ConflictingReplay); + } + } + Ok(LongitudinalObservationPersistenceDisposition::Duplicate) +} + +fn required_reference(reference: &str) -> Result<&str, LongitudinalObservationPersistenceError> { + match normalized_reference(reference) { + Some(normalized) if normalized == reference => Ok(reference), + _ => Err(LongitudinalObservationPersistenceError::InvalidReference), + } +} + +fn clock_anomaly_code(anomaly: Option) -> Option<&'static str> { + anomaly.map(|value| match value { + ClockAnomaly::RecordedAfterReceived => "recorded_after_received", + }) +} + +fn require_clock_anomaly_code( + stored_code: Option<&str>, + anomaly: Option, +) -> Result<(), LongitudinalObservationPersistenceError> { + if stored_code == clock_anomaly_code(anomaly) { + Ok(()) + } else { + Err(LongitudinalObservationPersistenceError::CorruptHistory) + } +} + +fn postgres_u64(value: u64) -> Result { + i64::try_from(value).map_err(|_| LongitudinalObservationPersistenceError::InvalidNumericRange) +} + +fn postgres_usize(value: usize) -> Result { + i64::try_from(value).map_err(|_| LongitudinalObservationPersistenceError::InvalidNumericRange) +} + +fn database_u64(value: i64) -> Result { + u64::try_from(value).map_err(|_| LongitudinalObservationPersistenceError::CorruptHistory) +} + +fn database_u16(value: i32) -> Result { + u16::try_from(value).map_err(|_| LongitudinalObservationPersistenceError::CorruptHistory) +} + +fn require_membership_sequence( + stored_sequence: i64, + expected_sequence: i64, +) -> Result<(), LongitudinalObservationPersistenceError> { + if stored_sequence == expected_sequence { + Ok(()) + } else { + Err(LongitudinalObservationPersistenceError::CorruptHistory) + } +} + +fn require_read_committed( + transaction: &mut Transaction<'_>, +) -> Result<(), LongitudinalObservationPersistenceError> { + let row = transaction.query_one("SHOW transaction_isolation", &[])?; + let isolation: String = row.get(0); + if isolation == "read committed" { + Ok(()) + } else { + Err(LongitudinalObservationPersistenceError::UnsupportedIsolationLevel) + } +} + +#[cfg(test)] +mod numeric_guard_tests { + use super::{ + database_u16, database_u64, postgres_u64, postgres_usize, require_clock_anomaly_code, + require_membership_sequence, required_reference, LongitudinalObservationPersistenceError, + }; + use crate::longitudinal_observation::ClockAnomaly; + use std::error::Error; + + #[test] + fn reference_numeric_conversion_and_error_sources_fail_closed() { + assert_eq!( + required_reference("tenant_clinic_seoul").unwrap(), + "tenant_clinic_seoul" + ); + assert!(matches!( + required_reference(" tenant_clinic_seoul "), + Err(LongitudinalObservationPersistenceError::InvalidReference) + )); + assert!(matches!( + required_reference("12"), + Err(LongitudinalObservationPersistenceError::InvalidReference) + )); + assert_eq!(postgres_u64(7).unwrap(), 7); + assert!(matches!( + postgres_u64(u64::MAX), + Err(LongitudinalObservationPersistenceError::InvalidNumericRange) + )); + assert_eq!(postgres_usize(1).unwrap(), 1); + #[cfg(target_pointer_width = "64")] + assert!(matches!( + postgres_usize(usize::MAX), + Err(LongitudinalObservationPersistenceError::InvalidNumericRange) + )); + #[cfg(not(target_pointer_width = "64"))] + assert_eq!( + postgres_usize(usize::MAX).unwrap(), + i64::try_from(usize::MAX).unwrap() + ); + assert_eq!(database_u64(7).unwrap(), 7); + assert!(matches!( + database_u64(-1), + Err(LongitudinalObservationPersistenceError::CorruptHistory) + )); + assert_eq!(database_u16(10_000).unwrap(), 10_000); + assert!(matches!( + database_u16(-1), + Err(LongitudinalObservationPersistenceError::CorruptHistory) + )); + assert!(require_membership_sequence(1, 1).is_ok()); + assert!(matches!( + require_membership_sequence(2, 1), + Err(LongitudinalObservationPersistenceError::CorruptHistory) + )); + assert!(require_clock_anomaly_code(None, None).is_ok()); + assert!(require_clock_anomaly_code( + Some("recorded_after_received"), + Some(ClockAnomaly::RecordedAfterReceived) + ) + .is_ok()); + assert!(matches!( + require_clock_anomaly_code(None, Some(ClockAnomaly::RecordedAfterReceived)), + Err(LongitudinalObservationPersistenceError::CorruptHistory) + )); + let error = LongitudinalObservationPersistenceError::ConflictingReplay; + assert!(Error::source(&error).is_none()); + assert!(error.to_string().contains("conflicting")); + let corrupt = LongitudinalObservationPersistenceError::CorruptHistory; + assert!(corrupt.to_string().contains("reconstructed")); + } +} diff --git a/tests/postgres_longitudinal_observation_persistence.rs b/tests/postgres_longitudinal_observation_persistence.rs new file mode 100644 index 00000000..ff0860e4 --- /dev/null +++ b/tests/postgres_longitudinal_observation_persistence.rs @@ -0,0 +1,472 @@ +//! Real `PostgreSQL` contract for durable longitudinal observation evidence. + +use postgres::{Client, IsolationLevel, NoTls}; +use psychometrics_commons_runtime::longitudinal_observation::{ + LongitudinalObservationInput, LongitudinalObservationRecord, LongitudinalObservationSet, + MembershipShareInput, ObservationTimeInput, +}; +use psychometrics_commons_runtime::postgres_longitudinal_observation::{ + apply_longitudinal_observation_migration, load_longitudinal_observation, + persist_longitudinal_observation, LongitudinalObservationPersistenceDisposition, + LongitudinalObservationPersistenceError, +}; +use std::sync::{Mutex, MutexGuard}; + +static TEST_LOCK: Mutex<()> = Mutex::new(()); + +fn guard() -> MutexGuard<'static, ()> { + TEST_LOCK + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) +} + +fn client() -> Client { + let url = std::env::var("TEST_DATABASE_URL").expect("TEST_DATABASE_URL is required"); + let mut client = Client::connect(&url, NoTls).expect("CI PostgreSQL must be reachable"); + client + .batch_execute( + "CREATE SCHEMA IF NOT EXISTS longitudinal_observation_persistence_test; \ + SET search_path TO longitudinal_observation_persistence_test;", + ) + .unwrap(); + client +} + +fn reset(client: &mut Client) { + client + .batch_execute( + "DROP TABLE IF EXISTS longitudinal_observation_persistence_test.longitudinal_membership_share CASCADE; \ + DROP TABLE IF EXISTS longitudinal_observation_persistence_test.longitudinal_observation CASCADE; \ + DROP FUNCTION IF EXISTS longitudinal_observation_persistence_test.reject_longitudinal_mutation() CASCADE; \ + DROP FUNCTION IF EXISTS longitudinal_observation_persistence_test.enforce_longitudinal_membership_total() CASCADE;", + ) + .unwrap(); +} + +fn observation(record_ref: &str, construct_ref: &str) -> LongitudinalObservationRecord { + let memberships = [ + MembershipShareInput { + membership_context_ref: "clinic_ward_seoul_01", + weight_parts_per_10_000: 6_000, + }, + MembershipShareInput { + membership_context_ref: "night_shift_team_alpha", + weight_parts_per_10_000: 4_000, + }, + ]; + LongitudinalObservationSet::new() + .ingest(LongitudinalObservationInput { + observation_record_ref: record_ref, + enrollment_ref: "longitudinal_enrollment_ko_001", + source_system_ref: "gyeot_mobile_collection", + source_observation_ref: "gyeot_observation_20260818_001", + construct_ref, + measure_ref: "measure_ipip_extraversion_ko_v1", + membership_shares: &memberships, + time: ObservationTimeInput { + validity_start_at_unix_ms: 1_776_661_900_000, + validity_end_at_unix_ms: 1_776_662_200_000, + recorded_at_unix_ms: 1_776_662_200_000, + received_at_unix_ms: 1_776_662_260_000, + ingested_at_unix_ms: 1_776_662_270_000, + timezone_name: "Asia/Seoul", + utc_offset_minutes: 540, + }, + }) + .unwrap() +} + +fn persist( + client: &mut Client, + tenant_ref: &str, + record: &LongitudinalObservationRecord, +) -> Result +{ + let mut tx = client.transaction().unwrap(); + let result = persist_longitudinal_observation(&mut tx, tenant_ref, record); + match result { + Ok(disposition) => { + tx.commit().unwrap(); + Ok(disposition) + } + Err(error) => { + tx.rollback().unwrap(); + Err(error) + } + } +} + +fn load( + client: &mut Client, + tenant_ref: &str, + observation_record_ref: &str, +) -> Result, LongitudinalObservationPersistenceError> { + let mut tx = client.transaction().unwrap(); + let result = load_longitudinal_observation(&mut tx, tenant_ref, observation_record_ref); + match result { + Ok(record) => { + tx.commit().unwrap(); + Ok(record) + } + Err(error) => { + tx.rollback().unwrap(); + Err(error) + } + } +} + +#[test] +fn seoul_multiple_membership_observation_is_tenant_bound_durable_and_idempotent() { + let _guard = guard(); + let mut client = client(); + reset(&mut client); + apply_longitudinal_observation_migration(&mut client).unwrap(); + let record = observation( + "longitudinal_observation_record_001", + "construct_extraversion", + ); + assert_eq!( + persist(&mut client, "tenant_clinic_seoul", &record).unwrap(), + LongitudinalObservationPersistenceDisposition::Inserted + ); + assert_eq!( + persist(&mut client, "tenant_clinic_seoul", &record).unwrap(), + LongitudinalObservationPersistenceDisposition::Duplicate + ); + let row = client + .query_one( + "SELECT tenant_ref, timezone_name, utc_offset_minutes \ + FROM longitudinal_observation WHERE observation_record_ref = $1", + &[&record.observation_record_ref()], + ) + .unwrap(); + assert_eq!(row.get::<_, String>(0), "tenant_clinic_seoul"); + assert_eq!(row.get::<_, String>(1), "Asia/Seoul"); + assert_eq!(row.get::<_, i16>(2), 540); + let memberships = client + .query( + "SELECT membership_context_ref, weight_parts_per_10_000 \ + FROM longitudinal_membership_share WHERE observation_record_ref = $1 \ + ORDER BY membership_sequence", + &[&record.observation_record_ref()], + ) + .unwrap(); + assert_eq!(memberships.len(), 2); + assert_eq!(memberships[0].get::<_, String>(0), "clinic_ward_seoul_01"); + assert_eq!(memberships[0].get::<_, i32>(1), 6_000); + assert_eq!(memberships[1].get::<_, String>(0), "night_shift_team_alpha"); + assert_eq!(memberships[1].get::<_, i32>(1), 4_000); +} + +#[test] +fn restart_recovery_is_exact_and_tenant_scoped() { + let _guard = guard(); + let mut client = client(); + reset(&mut client); + apply_longitudinal_observation_migration(&mut client).unwrap(); + let record = observation( + "longitudinal_observation_record_recovery", + "construct_extraversion", + ); + persist(&mut client, "tenant_clinic_seoul", &record).unwrap(); + + assert_eq!( + load( + &mut client, + "tenant_clinic_seoul", + "longitudinal_observation_record_recovery" + ) + .unwrap(), + Some(record) + ); + assert_eq!( + load( + &mut client, + "tenant_clinic_busan", + "longitudinal_observation_record_recovery" + ) + .unwrap(), + None + ); + assert_eq!( + load( + &mut client, + "tenant_clinic_seoul", + "longitudinal_observation_record_missing" + ) + .unwrap(), + None + ); + assert!(matches!( + load( + &mut client, + " tenant_clinic_seoul ", + "longitudinal_observation_record_recovery" + ), + Err(LongitudinalObservationPersistenceError::InvalidReference) + )); +} + +#[test] +fn restart_recovery_rejects_incomplete_persisted_history() { + let _guard = guard(); + let mut client = client(); + reset(&mut client); + apply_longitudinal_observation_migration(&mut client).unwrap(); + + // Simulate externally restored legacy/corrupt evidence without weakening the production + // migration: normal writes keep this constraint trigger enabled and cannot commit this state. + client + .batch_execute( + "ALTER TABLE longitudinal_observation \ + DISABLE TRIGGER longitudinal_observation_membership_total_check;", + ) + .unwrap(); + client + .execute( + "INSERT INTO longitudinal_observation (\ + observation_record_ref, tenant_ref, enrollment_ref, source_system_ref, \ + source_observation_ref, construct_ref, measure_ref, validity_start_at_unix_ms, \ + validity_end_at_unix_ms, recorded_at_unix_ms, received_at_unix_ms, \ + ingested_at_unix_ms, timezone_name, utc_offset_minutes, clock_anomaly_code\ + ) VALUES ($1,$2,$3,$4,$5,$6,$7,$8,$9,$10,$11,$12,$13,$14,$15)", + &[ + &"longitudinal_observation_record_corrupt", + &"tenant_clinic_seoul", + &"longitudinal_enrollment_ko_corrupt", + &"gyeot_mobile_collection", + &"gyeot_observation_corrupt", + &"construct_extraversion", + &"measure_ipip_extraversion_ko_v1", + &1_776_661_900_000_i64, + &1_776_662_200_000_i64, + &1_776_662_200_000_i64, + &1_776_662_260_000_i64, + &1_776_662_270_000_i64, + &"Asia/Seoul", + &540_i16, + &Option::<&str>::None, + ], + ) + .unwrap(); + client + .batch_execute( + "ALTER TABLE longitudinal_observation \ + ENABLE TRIGGER longitudinal_observation_membership_total_check;", + ) + .unwrap(); + + assert!(matches!( + load( + &mut client, + "tenant_clinic_seoul", + "longitudinal_observation_record_corrupt" + ), + Err(LongitudinalObservationPersistenceError::CorruptHistory) + )); +} + +#[test] +fn tenant_and_source_rebinding_fail_closed_without_cross_tenant_aliasing() { + let _guard = guard(); + let mut client = client(); + reset(&mut client); + apply_longitudinal_observation_migration(&mut client).unwrap(); + let first = observation( + "longitudinal_observation_record_002", + "construct_extraversion", + ); + persist(&mut client, "tenant_clinic_seoul", &first).unwrap(); + + assert!(matches!( + persist(&mut client, " tenant_clinic_seoul ", &first), + Err(LongitudinalObservationPersistenceError::InvalidReference) + )); + assert!(matches!( + persist(&mut client, "tenant_clinic_busan", &first), + Err(LongitudinalObservationPersistenceError::ConflictingReplay) + )); + + let rebound = observation( + "longitudinal_observation_record_rebound", + "construct_agreeableness", + ); + assert!(matches!( + persist(&mut client, "tenant_clinic_seoul", &rebound), + Err(LongitudinalObservationPersistenceError::ConflictingReplay) + )); + assert_eq!( + persist(&mut client, "tenant_clinic_busan", &rebound).unwrap(), + LongitudinalObservationPersistenceDisposition::Inserted + ); + + let update_error = client + .execute( + "UPDATE longitudinal_observation SET construct_ref = 'construct_rebound' \ + WHERE observation_record_ref = $1", + &[&first.observation_record_ref()], + ) + .unwrap_err(); + assert_eq!( + update_error + .as_db_error() + .map(postgres::error::DbError::code), + Some(&postgres::error::SqlState::OBJECT_NOT_IN_PREREQUISITE_STATE) + ); + let delete_error = client + .execute( + "DELETE FROM longitudinal_membership_share WHERE observation_record_ref = $1", + &[&first.observation_record_ref()], + ) + .unwrap_err(); + assert_eq!( + delete_error + .as_db_error() + .map(postgres::error::DbError::code), + Some(&postgres::error::SqlState::OBJECT_NOT_IN_PREREQUISITE_STATE) + ); +} + +#[test] +fn persistence_requires_read_committed_and_live_schema() { + let _guard = guard(); + let mut client = client(); + reset(&mut client); + apply_longitudinal_observation_migration(&mut client).unwrap(); + let record = observation( + "longitudinal_observation_record_003", + "construct_extraversion", + ); + let mut tx = client + .build_transaction() + .isolation_level(IsolationLevel::Serializable) + .start() + .unwrap(); + assert!(matches!( + persist_longitudinal_observation(&mut tx, "tenant_clinic_seoul", &record), + Err(LongitudinalObservationPersistenceError::UnsupportedIsolationLevel) + )); + tx.rollback().unwrap(); + reset(&mut client); + assert!(matches!( + persist(&mut client, "tenant_clinic_seoul", &record), + Err(LongitudinalObservationPersistenceError::Database(_)) + )); +} + +#[test] +fn persist_fails_closed_when_membership_relation_is_missing() { + let _guard = guard(); + let mut client = client(); + reset(&mut client); + apply_longitudinal_observation_migration(&mut client).unwrap(); + client + .batch_execute("DROP TABLE longitudinal_membership_share CASCADE") + .unwrap(); + let record = observation( + "longitudinal_observation_record_missing_membership_relation", + "construct_extraversion", + ); + assert!(matches!( + persist(&mut client, "tenant_clinic_seoul", &record), + Err(LongitudinalObservationPersistenceError::Database(_)) + )); +} + +#[test] +fn load_fails_closed_when_observation_or_membership_relations_are_missing() { + let _guard = guard(); + let mut client = client(); + reset(&mut client); + let record = observation( + "longitudinal_observation_record_missing_load_relation", + "construct_extraversion", + ); + assert!(matches!( + load( + &mut client, + "tenant_clinic_seoul", + record.observation_record_ref() + ), + Err(LongitudinalObservationPersistenceError::Database(_)) + )); + + apply_longitudinal_observation_migration(&mut client).unwrap(); + persist(&mut client, "tenant_clinic_seoul", &record).unwrap(); + client + .batch_execute("DROP TABLE longitudinal_membership_share CASCADE") + .unwrap(); + assert!(matches!( + load( + &mut client, + "tenant_clinic_seoul", + record.observation_record_ref() + ), + Err(LongitudinalObservationPersistenceError::Database(_)) + )); +} + +#[test] +fn replay_header_select_failure_is_a_database_failure() { + let _guard = guard(); + let mut client = client(); + reset(&mut client); + apply_longitudinal_observation_migration(&mut client).unwrap(); + let record = observation( + "longitudinal_observation_record_hidden_header_select", + "construct_extraversion", + ); + persist(&mut client, "tenant_clinic_seoul", &record).unwrap(); + + let sink = "longitudinal_observation_header_select_sink"; + client + .batch_execute(&format!( + "DROP SCHEMA IF EXISTS {sink} CASCADE; \ + CREATE SCHEMA {sink}; \ + CREATE OR REPLACE FUNCTION longitudinal_observation_redirect_after_insert() \ + RETURNS trigger LANGUAGE plpgsql AS $$ \ + BEGIN \ + PERFORM set_config('search_path', '{sink}', false); \ + RETURN NULL; \ + END $$; \ + CREATE TRIGGER longitudinal_observation_redirect_after_insert \ + AFTER INSERT ON longitudinal_observation \ + FOR EACH STATEMENT EXECUTE FUNCTION longitudinal_observation_redirect_after_insert();" + )) + .unwrap(); + + let error = persist(&mut client, "tenant_clinic_seoul", &record) + .expect_err("header replay must fail closed when classify-select cannot see stored rows"); + client + .batch_execute(&format!( + "DROP TRIGGER IF EXISTS longitudinal_observation_redirect_after_insert \ + ON longitudinal_observation; \ + DROP FUNCTION IF EXISTS longitudinal_observation_redirect_after_insert(); \ + DROP SCHEMA IF EXISTS {sink} CASCADE;" + )) + .unwrap(); + assert!(matches!( + error, + LongitudinalObservationPersistenceError::Database(_) + )); +} + +#[test] +fn replay_membership_select_failure_is_a_database_failure() { + let _guard = guard(); + let mut client = client(); + reset(&mut client); + apply_longitudinal_observation_migration(&mut client).unwrap(); + let record = observation( + "longitudinal_observation_record_hidden_membership_select", + "construct_extraversion", + ); + persist(&mut client, "tenant_clinic_seoul", &record).unwrap(); + client + .batch_execute("DROP TABLE longitudinal_membership_share CASCADE") + .unwrap(); + assert!(matches!( + persist(&mut client, "tenant_clinic_seoul", &record), + Err(LongitudinalObservationPersistenceError::Database(_)) + )); +} diff --git a/tests/postgres_longitudinal_observation_replay_coverage.rs b/tests/postgres_longitudinal_observation_replay_coverage.rs new file mode 100644 index 00000000..ec1ae6a3 --- /dev/null +++ b/tests/postgres_longitudinal_observation_replay_coverage.rs @@ -0,0 +1,404 @@ +//! Exhaustive replay-boundary coverage for durable longitudinal observations. +//! +//! These tests exercise immutable-evidence mismatches through the public +//! persistence API instead of weakening the production immutability contract. + +use postgres::{Client, NoTls}; +use psychometrics_commons_runtime::longitudinal_observation::{ + LongitudinalObservationInput, LongitudinalObservationRecord, LongitudinalObservationSet, + MembershipShareInput, ObservationTimeInput, +}; +use psychometrics_commons_runtime::postgres_longitudinal_observation::{ + apply_longitudinal_observation_migration, load_longitudinal_observation, + persist_longitudinal_observation, LongitudinalObservationPersistenceDisposition, + LongitudinalObservationPersistenceError, +}; +use std::error::Error; +use std::sync::{Mutex, MutexGuard}; + +static TEST_LOCK: Mutex<()> = Mutex::new(()); + +const BASE_RECORD_REF: &str = "longitudinal_observation_record_coverage"; +const BASE_ENROLLMENT_REF: &str = "longitudinal_enrollment_coverage"; +const BASE_SOURCE_SYSTEM_REF: &str = "gyeot_mobile_collection"; +const BASE_SOURCE_OBSERVATION_REF: &str = "gyeot_observation_coverage"; +const BASE_CONSTRUCT_REF: &str = "construct_extraversion"; +const BASE_MEASURE_REF: &str = "measure_ipip_extraversion_ko_v1"; +const BASE_VALIDITY_START: u64 = 1_776_661_900_000; +const BASE_VALIDITY_END: u64 = 1_776_662_200_000; +const BASE_RECORDED_AT: u64 = 1_776_662_200_000; +const BASE_RECEIVED_AT: u64 = 1_776_662_260_000; +const BASE_INGESTED_AT: u64 = 1_776_662_270_000; +const BASE_TIMEZONE: &str = "Asia/Seoul"; +const BASE_OFFSET: i16 = 540; + +#[derive(Clone, Copy)] +struct ObservationSpec<'a> { + observation_record_ref: &'a str, + enrollment_ref: &'a str, + source_system_ref: &'a str, + source_observation_ref: &'a str, + construct_ref: &'a str, + measure_ref: &'a str, + memberships: &'a [MembershipShareInput<'a>], + validity_start_at_unix_ms: u64, + validity_end_at_unix_ms: u64, + recorded_at_unix_ms: u64, + received_at_unix_ms: u64, + ingested_at_unix_ms: u64, + timezone_name: &'a str, + utc_offset_minutes: i16, +} + +impl<'a> ObservationSpec<'a> { + fn base(memberships: &'a [MembershipShareInput<'a>]) -> Self { + Self { + observation_record_ref: BASE_RECORD_REF, + enrollment_ref: BASE_ENROLLMENT_REF, + source_system_ref: BASE_SOURCE_SYSTEM_REF, + source_observation_ref: BASE_SOURCE_OBSERVATION_REF, + construct_ref: BASE_CONSTRUCT_REF, + measure_ref: BASE_MEASURE_REF, + memberships, + validity_start_at_unix_ms: BASE_VALIDITY_START, + validity_end_at_unix_ms: BASE_VALIDITY_END, + recorded_at_unix_ms: BASE_RECORDED_AT, + received_at_unix_ms: BASE_RECEIVED_AT, + ingested_at_unix_ms: BASE_INGESTED_AT, + timezone_name: BASE_TIMEZONE, + utc_offset_minutes: BASE_OFFSET, + } + } + + fn build(self) -> LongitudinalObservationRecord { + LongitudinalObservationSet::new() + .ingest(LongitudinalObservationInput { + observation_record_ref: self.observation_record_ref, + enrollment_ref: self.enrollment_ref, + source_system_ref: self.source_system_ref, + source_observation_ref: self.source_observation_ref, + construct_ref: self.construct_ref, + measure_ref: self.measure_ref, + membership_shares: self.memberships, + time: ObservationTimeInput { + validity_start_at_unix_ms: self.validity_start_at_unix_ms, + validity_end_at_unix_ms: self.validity_end_at_unix_ms, + recorded_at_unix_ms: self.recorded_at_unix_ms, + received_at_unix_ms: self.received_at_unix_ms, + ingested_at_unix_ms: self.ingested_at_unix_ms, + timezone_name: self.timezone_name, + utc_offset_minutes: self.utc_offset_minutes, + }, + }) + .unwrap() + } +} + +fn guard() -> MutexGuard<'static, ()> { + TEST_LOCK + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) +} + +fn fresh_client() -> Client { + let url = std::env::var("TEST_DATABASE_URL").expect("TEST_DATABASE_URL is required"); + let mut client = Client::connect(&url, NoTls).expect("CI PostgreSQL must be reachable"); + client + .batch_execute( + "DROP SCHEMA IF EXISTS longitudinal_observation_replay_coverage_test CASCADE; \ + CREATE SCHEMA longitudinal_observation_replay_coverage_test; \ + SET search_path TO longitudinal_observation_replay_coverage_test;", + ) + .unwrap(); + apply_longitudinal_observation_migration(&mut client).unwrap(); + client +} + +fn base_memberships() -> [MembershipShareInput<'static>; 2] { + [ + MembershipShareInput { + membership_context_ref: "clinic_ward_seoul_01", + weight_parts_per_10_000: 6_000, + }, + MembershipShareInput { + membership_context_ref: "night_shift_team_alpha", + weight_parts_per_10_000: 4_000, + }, + ] +} + +fn base_record() -> LongitudinalObservationRecord { + let memberships = base_memberships(); + ObservationSpec::base(&memberships).build() +} + +fn persist( + client: &mut Client, + tenant_ref: &str, + record: &LongitudinalObservationRecord, +) -> Result +{ + let mut transaction = client.transaction().unwrap(); + let result = persist_longitudinal_observation(&mut transaction, tenant_ref, record); + match result { + Ok(disposition) => { + transaction.commit().unwrap(); + Ok(disposition) + } + Err(error) => { + transaction.rollback().unwrap(); + Err(error) + } + } +} + +fn assert_conflict( + client: &mut Client, + tenant_ref: &str, + candidate: &LongitudinalObservationRecord, +) { + assert!(matches!( + persist(client, tenant_ref, candidate), + Err(LongitudinalObservationPersistenceError::ConflictingReplay) + )); +} + +#[test] +fn every_immutable_header_dimension_rejects_rebinding() { + let _guard = guard(); + let mut client = fresh_client(); + let base = base_record(); + assert_eq!( + persist(&mut client, "tenant_clinic_seoul", &base).unwrap(), + LongitudinalObservationPersistenceDisposition::Inserted + ); + assert_eq!( + persist(&mut client, "tenant_clinic_seoul", &base).unwrap(), + LongitudinalObservationPersistenceDisposition::Duplicate + ); + + let memberships = base_memberships(); + let spec = ObservationSpec::base(&memberships); + let variants = [ + ObservationSpec { + observation_record_ref: "longitudinal_observation_record_other", + ..spec + } + .build(), + ObservationSpec { + enrollment_ref: "longitudinal_enrollment_other", + ..spec + } + .build(), + ObservationSpec { + source_system_ref: "gyeot_collection_other", + ..spec + } + .build(), + ObservationSpec { + source_observation_ref: "gyeot_observation_other", + ..spec + } + .build(), + ObservationSpec { + construct_ref: "construct_agreeableness", + ..spec + } + .build(), + ObservationSpec { + measure_ref: "measure_ipip_extraversion_ko_v2", + ..spec + } + .build(), + ObservationSpec { + validity_start_at_unix_ms: BASE_VALIDITY_START + 1, + ..spec + } + .build(), + ObservationSpec { + validity_end_at_unix_ms: BASE_VALIDITY_END + 1, + ..spec + } + .build(), + ObservationSpec { + recorded_at_unix_ms: BASE_RECORDED_AT + 1, + ..spec + } + .build(), + ObservationSpec { + received_at_unix_ms: BASE_RECEIVED_AT + 1, + ..spec + } + .build(), + ObservationSpec { + ingested_at_unix_ms: BASE_INGESTED_AT + 1, + ..spec + } + .build(), + ObservationSpec { + timezone_name: "Asia/Tokyo", + ..spec + } + .build(), + ObservationSpec { + utc_offset_minutes: 480, + ..spec + } + .build(), + ObservationSpec { + recorded_at_unix_ms: BASE_RECEIVED_AT + 1, + ..spec + } + .build(), + ]; + for candidate in &variants { + assert_conflict(&mut client, "tenant_clinic_seoul", candidate); + } + assert_conflict(&mut client, "tenant_clinic_busan", &base); +} + +#[test] +fn every_membership_dimension_and_source_alias_rejects_rebinding() { + let _guard = guard(); + let mut client = fresh_client(); + let base = base_record(); + assert_eq!( + persist(&mut client, "tenant_clinic_seoul", &base).unwrap(), + LongitudinalObservationPersistenceDisposition::Inserted + ); + + let one_membership = [MembershipShareInput { + membership_context_ref: "clinic_ward_seoul_01", + weight_parts_per_10_000: 10_000, + }]; + let one_membership_record = ObservationSpec::base(&one_membership).build(); + assert_conflict(&mut client, "tenant_clinic_seoul", &one_membership_record); + + let context_mismatch_memberships = [ + MembershipShareInput { + membership_context_ref: "clinic_ward_seoul_02", + weight_parts_per_10_000: 6_000, + }, + MembershipShareInput { + membership_context_ref: "night_shift_team_alpha", + weight_parts_per_10_000: 4_000, + }, + ]; + let context_mismatch = ObservationSpec::base(&context_mismatch_memberships).build(); + assert_conflict(&mut client, "tenant_clinic_seoul", &context_mismatch); + + let weight_mismatch_memberships = [ + MembershipShareInput { + membership_context_ref: "clinic_ward_seoul_01", + weight_parts_per_10_000: 5_000, + }, + MembershipShareInput { + membership_context_ref: "night_shift_team_alpha", + weight_parts_per_10_000: 5_000, + }, + ]; + let weight_mismatch = ObservationSpec::base(&weight_mismatch_memberships).build(); + assert_conflict(&mut client, "tenant_clinic_seoul", &weight_mismatch); + + let memberships = base_memberships(); + let busan_record = ObservationSpec { + observation_record_ref: "longitudinal_observation_record_busan_source_alias", + ..ObservationSpec::base(&memberships) + } + .build(); + assert_eq!( + persist(&mut client, "tenant_clinic_busan", &busan_record).unwrap(), + LongitudinalObservationPersistenceDisposition::Inserted + ); + assert_conflict(&mut client, "tenant_clinic_busan", &base); +} + +#[test] +fn corrupted_sequence_and_anomaly_evidence_fail_closed_after_restart() { + let _guard = guard(); + let mut client = fresh_client(); + let base = base_record(); + persist(&mut client, "tenant_clinic_seoul", &base).unwrap(); + + client + .batch_execute( + "ALTER TABLE longitudinal_membership_share \ + DISABLE TRIGGER longitudinal_membership_immutable_update; \ + UPDATE longitudinal_membership_share SET membership_sequence = 3 \ + WHERE observation_record_ref = 'longitudinal_observation_record_coverage' \ + AND membership_sequence = 1; \ + ALTER TABLE longitudinal_membership_share \ + ENABLE TRIGGER longitudinal_membership_immutable_update;", + ) + .unwrap(); + assert_conflict(&mut client, "tenant_clinic_seoul", &base); + + let mut transaction = client.transaction().unwrap(); + assert!(matches!( + load_longitudinal_observation(&mut transaction, "tenant_clinic_seoul", BASE_RECORD_REF), + Err(LongitudinalObservationPersistenceError::CorruptHistory) + )); + transaction.rollback().unwrap(); + + let mut client = fresh_client(); + let base = base_record(); + persist(&mut client, "tenant_clinic_seoul", &base).unwrap(); + client + .batch_execute( + "ALTER TABLE longitudinal_observation \ + DISABLE TRIGGER longitudinal_observation_immutable_update; \ + ALTER TABLE longitudinal_observation \ + DROP CONSTRAINT longitudinal_observation_anomaly_check; \ + UPDATE longitudinal_observation \ + SET clock_anomaly_code = 'recorded_after_received' \ + WHERE observation_record_ref = 'longitudinal_observation_record_coverage'; \ + ALTER TABLE longitudinal_observation \ + ENABLE TRIGGER longitudinal_observation_immutable_update;", + ) + .unwrap(); + let mut transaction = client.transaction().unwrap(); + assert!(matches!( + load_longitudinal_observation(&mut transaction, "tenant_clinic_seoul", BASE_RECORD_REF), + Err(LongitudinalObservationPersistenceError::CorruptHistory) + )); + transaction.rollback().unwrap(); + + let mut transaction = client.transaction().unwrap(); + assert!(matches!( + load_longitudinal_observation(&mut transaction, "tenant_clinic_seoul", " 123 "), + Err(LongitudinalObservationPersistenceError::InvalidReference) + )); + transaction.rollback().unwrap(); +} + +#[test] +fn all_public_error_variants_have_stable_messages_and_database_sources() { + for error in [ + LongitudinalObservationPersistenceError::InvalidReference, + LongitudinalObservationPersistenceError::InvalidNumericRange, + LongitudinalObservationPersistenceError::ConflictingReplay, + LongitudinalObservationPersistenceError::CorruptHistory, + LongitudinalObservationPersistenceError::UnsupportedIsolationLevel, + ] { + assert!(!error.to_string().is_empty()); + assert!(Error::source(&error).is_none()); + } + + let _guard = guard(); + let mut client = fresh_client(); + let base = base_record(); + client + .batch_execute("DROP SCHEMA longitudinal_observation_replay_coverage_test CASCADE") + .unwrap(); + let error = persist(&mut client, "tenant_clinic_seoul", &base) + .expect_err("missing persistence schema must surface the PostgreSQL source"); + assert!(matches!( + error, + LongitudinalObservationPersistenceError::Database(_) + )); + assert_eq!( + error.to_string(), + "PostgreSQL longitudinal observation persistence failed" + ); + assert!(Error::source(&error).is_some()); +} diff --git a/tests/postgres_longitudinal_observation_schema_integrity.rs b/tests/postgres_longitudinal_observation_schema_integrity.rs new file mode 100644 index 00000000..616d14cc --- /dev/null +++ b/tests/postgres_longitudinal_observation_schema_integrity.rs @@ -0,0 +1,277 @@ +//! Real `PostgreSQL` integrity contracts for longitudinal observation evidence. + +use postgres::{Client, NoTls}; +use psychometrics_commons_runtime::postgres_longitudinal_observation::apply_longitudinal_observation_migration; +use std::sync::{Mutex, MutexGuard}; + +static TEST_LOCK: Mutex<()> = Mutex::new(()); + +fn guard() -> MutexGuard<'static, ()> { + TEST_LOCK + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) +} + +fn client(schema_name: &str) -> Client { + assert!( + schema_name + .chars() + .all(|character| character.is_ascii_lowercase() || character == '_'), + "schema names must be two-word snake_case identifiers" + ); + let url = std::env::var("TEST_DATABASE_URL").expect("TEST_DATABASE_URL is required"); + let mut client = Client::connect(&url, NoTls).expect("CI PostgreSQL must be reachable"); + client + .batch_execute(&format!( + "DROP SCHEMA IF EXISTS {schema_name} CASCADE; \ + CREATE SCHEMA {schema_name}; \ + SET search_path TO {schema_name};" + )) + .unwrap(); + client +} + +fn insert_observation( + client: &mut Client, + tenant_ref: &str, + recorded_at_unix_ms: i64, + received_at_unix_ms: i64, + clock_anomaly_code: Option<&str>, +) -> Result { + client.execute( + "INSERT INTO longitudinal_observation (\ + observation_record_ref, tenant_ref, enrollment_ref, source_system_ref, \ + source_observation_ref, construct_ref, measure_ref, validity_start_at_unix_ms, \ + validity_end_at_unix_ms, recorded_at_unix_ms, received_at_unix_ms, \ + ingested_at_unix_ms, timezone_name, utc_offset_minutes, clock_anomaly_code\ + ) VALUES ($1,$2,$3,$4,$5,$6,$7,$8,$9,$10,$11,$12,$13,$14,$15)", + &[ + &"longitudinal_observation_schema_integrity", + &tenant_ref, + &"longitudinal_enrollment_schema_integrity", + &"gyeot_mobile_collection", + &"gyeot_observation_schema_integrity", + &"construct_extraversion", + &"measure_ipip_extraversion_ko_v1", + &1_776_661_900_000_i64, + &1_776_662_200_000_i64, + &recorded_at_unix_ms, + &received_at_unix_ms, + &1_776_662_270_000_i64, + &"Asia/Seoul", + &540_i16, + &clock_anomaly_code, + ], + ) +} + +fn assert_check_constraint(error: &postgres::Error, expected_constraint: &str) { + let database_error = error + .as_db_error() + .expect("schema integrity rejection must be a PostgreSQL database error"); + assert_eq!( + database_error.code(), + &postgres::error::SqlState::CHECK_VIOLATION + ); + assert_eq!(database_error.constraint(), Some(expected_constraint)); +} + +#[test] +fn observation_header_cannot_commit_without_complete_membership_vector() { + let _guard = guard(); + let mut client = client("longitudinal_observation_schema_integrity_membership_test"); + apply_longitudinal_observation_migration(&mut client).unwrap(); + + let mut transaction = client.transaction().unwrap(); + transaction + .execute( + "INSERT INTO longitudinal_observation (\ + observation_record_ref, tenant_ref, enrollment_ref, source_system_ref, \ + source_observation_ref, construct_ref, measure_ref, validity_start_at_unix_ms, \ + validity_end_at_unix_ms, recorded_at_unix_ms, received_at_unix_ms, \ + ingested_at_unix_ms, timezone_name, utc_offset_minutes, clock_anomaly_code\ + ) VALUES ($1,$2,$3,$4,$5,$6,$7,$8,$9,$10,$11,$12,$13,$14,$15)", + &[ + &"longitudinal_observation_missing_membership", + &"tenant_clinic_seoul", + &"longitudinal_enrollment_missing_membership", + &"gyeot_mobile_collection", + &"gyeot_observation_missing_membership", + &"construct_extraversion", + &"measure_ipip_extraversion_ko_v1", + &1_776_661_900_000_i64, + &1_776_662_200_000_i64, + &1_776_662_200_000_i64, + &1_776_662_260_000_i64, + &1_776_662_270_000_i64, + &"Asia/Seoul", + &540_i16, + &Option::<&str>::None, + ], + ) + .unwrap(); + + let error = transaction + .commit() + .expect_err("an observation with no membership rows must fail at commit"); + assert_eq!( + error.as_db_error().map(postgres::error::DbError::code), + Some(&postgres::error::SqlState::CHECK_VIOLATION) + ); +} + +#[test] +fn digit_bearing_opaque_references_remain_valid_at_the_database_boundary() { + let _guard = guard(); + let mut client = client("longitudinal_observation_schema_integrity_opaque_reference_test"); + apply_longitudinal_observation_migration(&mut client).unwrap(); + + for reference in [ + "longitudinal_observation_record_001", + "gyeot_observation_20260818_001", + "clinic_ward_seoul_01", + ] { + let is_valid: bool = client + .query_one("SELECT longitudinal_reference_is_valid($1)", &[&reference]) + .unwrap() + .get(0); + assert!( + is_valid, + "opaque references containing digits must not be misclassified as numeric-like: {reference:?}" + ); + } +} + +#[test] +fn numeric_like_references_are_rejected_by_the_database_boundary() { + let _guard = guard(); + let mut client = client("longitudinal_observation_schema_integrity_reference_test"); + apply_longitudinal_observation_migration(&mut client).unwrap(); + + let error = insert_observation( + &mut client, + "1.5", + 1_776_662_200_000, + 1_776_662_260_000, + None, + ) + .expect_err("numeric-like opaque references must fail before evidence is stored"); + assert_check_constraint(&error, "longitudinal_observation_reference_check"); +} + +#[test] +fn unicode_numeric_separator_aliases_are_rejected_by_the_database_boundary() { + let _guard = guard(); + let mut client = client("longitudinal_observation_schema_integrity_unicode_reference_test"); + apply_longitudinal_observation_migration(&mut client).unwrap(); + + for reference in ["1.5", "1,000", "1٫5", "1٬000"] { + let is_valid: bool = client + .query_one("SELECT longitudinal_reference_is_valid($1)", &[&reference]) + .unwrap() + .get(0); + assert!( + !is_valid, + "Unicode numeric separator alias must not be accepted: {reference:?}" + ); + } +} + +#[test] +fn unicode_numeric_categories_match_the_rust_reference_boundary() { + let _guard = guard(); + let mut client = client("longitudinal_observation_schema_integrity_numeric_category_test"); + apply_longitudinal_observation_migration(&mut client).unwrap(); + + // Rust char::is_numeric covers Unicode decimal digits, letter numbers, and other numbers. + // PostgreSQL's POSIX [[:digit:]] class alone does not prove that same boundary. + for reference in ["½", "²", "Ⅳ", "+Ⅳ", "½e²"] { + let is_valid: bool = client + .query_one("SELECT longitudinal_reference_is_valid($1)", &[&reference]) + .unwrap() + .get(0); + assert!( + !is_valid, + "Unicode numeric reference accepted by PostgreSQL but rejected by Rust: {reference:?}" + ); + } +} + +#[test] +fn embedded_control_characters_are_rejected_by_the_database_boundary() { + let _guard = guard(); + let mut client = client("longitudinal_observation_schema_integrity_control_reference_test"); + apply_longitudinal_observation_migration(&mut client).unwrap(); + + for reference in ["tenant_\u{0001}_clinic", "tenant_\u{001f}_clinic"] { + let is_valid: bool = client + .query_one("SELECT longitudinal_reference_is_valid($1)", &[&reference]) + .unwrap() + .get(0); + assert!( + !is_valid, + "control-character reference alias must not be accepted: {reference:?}" + ); + } +} + +#[test] +fn anomaly_relation_remains_a_check_constraint_when_immutability_is_disabled() { + let _guard = guard(); + let mut client = client("longitudinal_observation_schema_integrity_anomaly_update_test"); + apply_longitudinal_observation_migration(&mut client).unwrap(); + + client + .batch_execute( + "BEGIN; \ + INSERT INTO longitudinal_observation (\ + observation_record_ref, tenant_ref, enrollment_ref, source_system_ref, \ + source_observation_ref, construct_ref, measure_ref, validity_start_at_unix_ms, \ + validity_end_at_unix_ms, recorded_at_unix_ms, received_at_unix_ms, \ + ingested_at_unix_ms, timezone_name, utc_offset_minutes, clock_anomaly_code\ + ) VALUES (\ + 'longitudinal_observation_anomaly_update', 'tenant_clinic_seoul', \ + 'longitudinal_enrollment_anomaly_update', 'gyeot_mobile_collection', \ + 'gyeot_observation_anomaly_update', 'construct_extraversion', \ + 'measure_ipip_extraversion_ko_v1', 1776661900000, 1776662200000, \ + 1776662200000, 1776662260000, 1776662270000, 'Asia/Seoul', 540, NULL\ + ); \ + INSERT INTO longitudinal_membership_share (\ + observation_record_ref, membership_sequence, membership_context_ref, \ + weight_parts_per_10_000\ + ) VALUES (\ + 'longitudinal_observation_anomaly_update', 1, 'clinic_ward_seoul_01', 10000\ + ); \ + COMMIT; \ + ALTER TABLE longitudinal_observation \ + DISABLE TRIGGER longitudinal_observation_immutable_update;", + ) + .unwrap(); + + let error = client + .execute( + "UPDATE longitudinal_observation \ + SET clock_anomaly_code = 'recorded_after_received' \ + WHERE observation_record_ref = 'longitudinal_observation_anomaly_update'", + &[], + ) + .expect_err("the CHECK constraint must reject a code without a clock inversion"); + assert_check_constraint(&error, "longitudinal_observation_anomaly_check"); +} + +#[test] +fn anomaly_code_must_match_the_observed_clock_order() { + let _guard = guard(); + let mut client = client("longitudinal_observation_schema_integrity_anomaly_test"); + apply_longitudinal_observation_migration(&mut client).unwrap(); + + let error = insert_observation( + &mut client, + "tenant_clinic_seoul", + 1_776_662_200_000, + 1_776_662_260_000, + Some("recorded_after_received"), + ) + .expect_err("an anomaly code without the corresponding clock inversion must fail"); + assert_check_constraint(&error, "longitudinal_observation_anomaly_check"); +} diff --git a/tests/session_http_framing.rs b/tests/session_http_framing.rs index 59b12038..cf490036 100644 --- a/tests/session_http_framing.rs +++ b/tests/session_http_framing.rs @@ -57,10 +57,11 @@ fn framing_error(request: &[u8]) -> std::io::ErrorKind { let payload = request.to_vec(); let client = std::thread::spawn(move || { let mut stream = TcpStream::connect(address).unwrap(); - stream.write_all(&payload).unwrap(); - // The server is expected to close malformed requests immediately; the - // client-side half-close can therefore race with the peer close. - stream.shutdown(Shutdown::Write).ok(); + // The listener may reject and close as soon as the invalid Content-Length is + // parsed. A later write or shutdown then races with that close under coverage + // instrumentation; the contract under test is the accept-side InvalidData. + let _ = stream.write_all(&payload); + let _ = stream.shutdown(Shutdown::Write); let mut response = String::new(); let _ = stream.read_to_string(&mut response); });