diff --git a/ARCHITECTURE.md b/ARCHITECTURE.md index bba2f658..ce9f4899 100644 --- a/ARCHITECTURE.md +++ b/ARCHITECTURE.md @@ -84,6 +84,7 @@ research_contribution tenant_authorization integration_outbox integration_inbox +health_probes ``` Splitting a module into a separate service later must not change its domain semantics or bypass existing versioned contracts. diff --git a/CHANGELOG.md b/CHANGELOG.md index 33407ed7..7a5a3497 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -5,6 +5,11 @@ All notable product and architecture changes are recorded here. Releases use imm ## Unreleased ### Added +- Operator health probes can start from `HEALTH_LISTEN_ADDR` or platform `PORT`, optionally observe `DATABASE_URL` only for GET `/ready`, and keep GET `/live` free of store I/O when the store is down. Blank, padded, or unknown listen/store/backlog values fail closed. Caller-measured `HEALTH_BACKLOG_HEALTH` is required before readiness can be true; the process does not invent a backlog threshold. +- Operator health probes can observe a live PostgreSQL operational snapshot and answer GET `/live` and GET `/ready` without exposing driver errors. GET `/live` does not perform store I/O. Bare GET `/ready` on the PostgreSQL adapter requires `postgres_operational_store`. The bound listener applies a 2-second I/O timeout and rejects oversized requests without echoing them. Measured backlog thresholds remain caller-supplied. +- PostgreSQL operational health snapshot composes runtime and relation probes with caller-supplied backlog into one fail-closed `RuntimeHealthSnapshot` without exposing driver errors. +- A bound TCP listener can serve operator GET `/live` and GET `/ready` probes in a blocking accept loop until accept fails, or one HTTP/1.1 request per accepted connection. Accept retries Interrupted, ConnectionAborted, and ConnectionReset. A dropped probe connection does not stop later probes on either the in-memory or PostgreSQL-backed serve loop. It does not add public product routes, TLS, keep-alive, or measured SLO values. +- Operator HTTP probes translate the domain health snapshot into GET `/live` and GET `/ready` responses with an as-built OpenAPI 3.2.0 contract. Liveness stays independent of operation readiness; named required capabilities, stalled backlog, unknown integrity, and unknown capabilities fail closed. Unsupported methods and paths return RFC 9457 problem details with explicit `urn:psychometrics-commons:problem:` types, not `about:blank`, and without raw store errors. - 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/README.md b/README.md index 46b10021..24ab7be5 100644 --- a/README.md +++ b/README.md @@ -52,6 +52,10 @@ It does **not** duplicate psychometric numerical kernels, identity credentials, - [CLAUDE.md](CLAUDE.md) — concise coding-agent entry point into the same normative contracts. - [Changelog](CHANGELOG.md) — unreleased and released product/architecture changes. +## Operator health probes + +To run the first process surface, set `HEALTH_LISTEN_ADDR` (for example `127.0.0.1:8080`) or platform `PORT`, then call `psychometrics_commons_runtime::health_process::run_health_process`. Point liveness at GET `/live` and readiness at GET `/ready`. Set `DATABASE_URL` only when the process should observe PostgreSQL for readiness, and set `HEALTH_BACKLOG_HEALTH=within_bounds` only after backlog is actually measured. A down store must not take `/live` with it. + ## Architecture authority and implementation status An accepted ADR must define concrete ownership, interfaces, invariants, failure behavior, security/privacy/tenancy boundaries, migration and rollback, validation evidence, alternatives, and reversal conditions. Material implementation that contradicts an accepted ADR requires an explicit superseding decision rather than silent architectural drift. diff --git a/docs/DOCUMENTATION_ASSESSMENT.md b/docs/DOCUMENTATION_ASSESSMENT.md index ccd3c7c1..64dbb4e9 100644 --- a/docs/DOCUMENTATION_ASSESSMENT.md +++ b/docs/DOCUMENTATION_ASSESSMENT.md @@ -31,7 +31,7 @@ This reconciliation closes those architecture-definition gaps without promoting | Quality / risk / compliance readiness | **Sufficient as assurance baseline** | Evidence scenarios and risk state are explicit; readiness is not certification. | | Traceability | **Repaired in this reconciliation** | Baseline now names exact protected-main `748876…`, marks newly merged domain modules Implemented/Partial, and isolates PR #24 as Active PR. Must be updated after every material merge. | | Roadmap / agent guidance / changelog | **Sufficient for continued delivery** | Must remain code-current; documentation completion is not a terminal condition for the execution loop. | -| Machine-readable OpenAPI / AsyncAPI | **Not yet applicable as as-built evidence** | Add and validate with the first implemented HTTP/event transport. Do not publish aspirational operations as deployed. | +| Machine-readable OpenAPI / AsyncAPI | **Active PR for operator health probes only** | `openapi/health-probes.yaml` lists GET `/live` and GET `/ready`. Public/admin product routes and AsyncAPI remain unimplemented and must not be listed as deployed. | | Physical schema / as-built topology | **Partial / implementation-gated** | Logical ERD is authoritative target semantics. Actual migrations/topology/rollback/restore evidence must be compared to it as those artifacts land. | | Instrument-release evidence bundles | **Target** | Every publishable consumer instrument needs immutable rights, locale/translation, scoring/calibration/norm, DIF/invariance/linking where claimed, scoreability, intended-use and narrative-rule evidence. | diff --git a/docs/OPERABILITY.md b/docs/OPERABILITY.md index b729f0c5..2ac87128 100644 --- a/docs/OPERABILITY.md +++ b/docs/OPERABILITY.md @@ -44,6 +44,8 @@ The implementation must distinguish at least: Readiness must not fail solely because an optional capability is unavailable if the selected operation can safely proceed without it. Conversely, a process can be live while not ready to accept new state-changing requests. +Operator HTTP probes, when implemented, are GET `/live` and GET `/ready`. `/live` answers process liveness only and must not perform store I/O; a hung or failed PostgreSQL connection must not restart a still-live process. `/ready` answers operation-scoped readiness and may name required capabilities as repeated `capability` query parameters. When the PostgreSQL adapter answers a bare GET `/ready` (no `capability=`), it requires `postgres_operational_store` so a read-only or unsupported store cannot advertise readiness to a load balancer. These probes do not publish measured SLO values. A bound TCP listener, when present, serves those same operations in a blocking accept loop until accept fails, or one request per accepted connection. Interrupted, aborted, or reset accepts retry, including `ConnectionReset` before `accept` returns, matching TCP reset processing in RFC 9293. A dropped probe connection does not stop later probes. It applies a bounded read/write timeout, rejects oversized requests without echoing them, and is not a measured availability claim. PostgreSQL observation happens after accept and only for GET `/ready`. Probe failure is unknown/unready and must not expose driver errors. Unsupported methods and paths use explicit `urn:psychometrics-commons:problem:` types rather than `about:blank`. Operators start the probe process with `run_health_process` after setting `HEALTH_LISTEN_ADDR` or platform `PORT`. Optional `DATABASE_URL` is observed only for GET `/ready`. Optional `HEALTH_BACKLOG_HEALTH` must be `within_bounds`, `stalled`, or `unknown`; missing backlog stays unknown and not ready. Point liveness at GET `/live` and readiness at GET `/ready`. Do not treat a single `accept_one_*` call as a running probe server. + ## 4. Capability degradation matrix | Dependency/capability failure | Required product behavior | @@ -229,6 +231,12 @@ Never collapse these maturity levels. SOC 2/CSAP readiness work may map evidence ## 15. References +Eddy, W. (Ed.). (2022). *Transmission Control Protocol (TCP)* (RFC 9293). Internet Engineering Task Force. https://doi.org/10.17487/RFC9293 + +Fielding, R., Nottingham, M., & Reschke, J. (Eds.). (2022). *HTTP Semantics* (RFC 9110). Internet Engineering Task Force. https://doi.org/10.17487/RFC9110 + International Organization for Standardization & International Electrotechnical Commission. (2023). *ISO/IEC 25010:2023 Systems and software engineering—Systems and software Quality Requirements and Evaluation (SQuaRE)—Product quality model*. +Kubernetes Authors. (2024). *Configure liveness, readiness and startup probes*. Kubernetes Documentation. https://kubernetes.io/docs/tasks/configure-pod-container/configure-liveness-readiness-startup-probes/ + National Institute of Standards and Technology. (2022). *Secure Software Development Framework (SSDF) Version 1.1: Recommendations for mitigating the risk of software vulnerabilities* (NIST SP 800-218). https://doi.org/10.6028/NIST.SP.800-218 diff --git a/docs/QUALITY_ATTRIBUTES.md b/docs/QUALITY_ATTRIBUTES.md index 3176e540..8ed67510 100644 --- a/docs/QUALITY_ATTRIBUTES.md +++ b/docs/QUALITY_ATTRIBUTES.md @@ -68,6 +68,12 @@ This document converts broad quality goals into **stimulus → environment → r - **Response:** completed response snapshot remains durable; scoring job waits/retries; no invented score. - **Evidence:** worker/job state and recovery test. +### QA-AVL-04 — Operational store down during health probes + +- **Stimulus:** `DATABASE_URL` is configured and the operational store refuses connections while the health-probe process is running. +- **Response:** GET `/live` remains HTTP 200 without store I/O or driver text; GET `/ready` returns HTTP 503 without echoing the URL or driver error. +- **Evidence:** `tests/health_process_contract.rs` unreachable-store listener contract. + ## 4. Security ### QA-SEC-01 — Cross-tenant object reference diff --git a/docs/TRACEABILITY.md b/docs/TRACEABILITY.md index 8e9a9c2c..8b86d1c0 100644 --- a/docs/TRACEABILITY.md +++ b/docs/TRACEABILITY.md @@ -40,7 +40,7 @@ An active PR, architecture document, conversation decision, or scheduler plan is | Research identity separation | PRD §5, §11 | TRD §14; ERD restricted linkage | ADR-0003, ADR-0006, ADR-0007, ADR-0020 | Partially implemented via research-contribution identity separation; restricted linkage persistence is Target | | Research release manifests | PRD §5 | TRD §15 | ADR-0007, ADR-0010 | Target; semantic-data-portal is External dependency | | Durable outbox/inbox delivery semantics | PRD §7, §9 | TRD §19–20 | ADR-0014, ADR-0015 | **Partially implemented**: domain contracts in `src/integration.rs`; PostgreSQL 18 outbox/inbox identity, delivery-attempt persistence, and inbox consumption distinct from receipt; live side-effect execution remains Target | -| 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 | +| 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; **Active PR** #132 (successor to #122) binds GET `/live` and GET `/ready` from `HEALTH_LISTEN_ADDR` or platform `PORT`, keeps serving after a dropped probe, answers `/live` without store I/O even when `DATABASE_URL` is down, and observes PostgreSQL only for `/ready` (bare `/ready` requires `postgres_operational_store`) without exposing driver errors; measured thresholds and deployment-profile 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`; actual English/Korean form content, rights, translation, invariance and serving are 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/normalized-ingestion/orchestration adapter | @@ -77,8 +77,8 @@ An active PR, architecture document, conversation decision, or scheduler plan is | No default tenant for writes | TRD §11; Security/Data | authorization-domain primitive exists; persistence remains Target | persistence/API tenant negative tests | | Tenant-bound transactional outbox/inbox | TRD §19–20; ADR-0014/0015 | `src/integration.rs` domain envelope/inbox/retry contracts plus PostgreSQL tenant/source-scoped integration evidence, delivery-attempt persistence, and inbox consumption | durable side-effect processing completion, poison-message/crash recovery, broader aggregate transaction integration | | Inbox receipt is not side-effect completion | ADR-0014/0015; UML integration sequence | `src/integration.rs` states/retry semantics; PostgreSQL inbox consumption persists pending/processing/completed and expire-and-reclaim | live adapter crash/retry tests | -| Liveness is distinct from operation readiness | Operability §3–4; ADR-0017 | **Implemented** in `src/health.rs` and `src/postgres_health.rs`: liveness is modeled independently from operation-scoped readiness and PostgreSQL write-readiness | live transport probes, metrics, and deployment-profile acceptance | -| Optional capability outage does not fail unrelated work | Operability §3–4; ADR-0011/0017 | **Implemented** in `src/health.rs` and `src/postgres_health.rs`: readiness evaluates only capabilities required by the selected operation and maps PostgreSQL evidence onto that contract | degraded-mode transport/integration tests | +| Liveness is distinct from operation readiness | Operability §3–4; ADR-0017 | **Implemented** in `src/health.rs` and `src/postgres_health.rs`; **Active PR** #132 (successor to #122) exposes GET `/live` independently from GET `/ready` from `HEALTH_LISTEN_ADDR` or `PORT`, keeps serving after a dropped probe, and does not observe PostgreSQL for `/live` even when `DATABASE_URL` is configured | metrics and deployment-profile acceptance | +| Optional capability outage does not fail unrelated work | Operability §3–4; ADR-0011/0017 | **Implemented** in `src/health.rs` and `src/postgres_health.rs`; **Active PR** #132 (successor to #122) keeps `/ready?capability=` fail-closed for unknown or unsafe required capabilities and for an unreachable `DATABASE_URL` | degraded-mode transport/integration tests | | Unknown/stalled backlog or unknown/incompatible integrity blocks new state-changing work | Operability §3, §6, §8 | **Implemented** domain contract in `src/health.rs`; `src/postgres_health.rs` fails closed on unsupported/read-only PostgreSQL or a missing required relation | persistence/job backlog metrics, stronger schema probes, alerting, and failure-injection evidence | | No operational IDs in public research release | TRD §14–15; Research Governance | architecture policy | release fixture/static/runtime leakage tests | | AI optional; deterministic core remains | PRD §9.5; TRD §17; AI Governance | architecture policy | narrative fallback end-to-end test | @@ -97,6 +97,8 @@ src/lib.rs ├── consent.rs # purpose-specific consent + research contribution lifecycle ├── data_rights.rs # export/deletion lifecycle and retention evidence ├── health.rs # operation-scoped liveness/readiness and capability-state contract +├── health_http.rs # Active PR #132 (successor to #122) operator GET /live and GET /ready probes plus a bounded-timeout TCP listener that keeps serving after a dropped probe (not protected-main truth) +├── health_process.rs # Active PR #132 (successor to #122) binds HEALTH_LISTEN_ADDR or PORT and observes DATABASE_URL only for /ready (not protected-main truth) ├── instrument.rs # immutable release manifest + scientific publication-evidence gate ├── integration.rs # outbox/inbox/retry/quarantine domain contracts ├── item_delivery.rs # sequence-aware delivery evidence without confidential response data @@ -104,7 +106,8 @@ src/lib.rs ├── participant.rs # stable participant identity + issuer-scoped optional Keyverse account link ├── postgres_consent.rs # PostgreSQL purpose-specific consent ledger persistence ├── postgres_data_rights.rs # PostgreSQL data-rights request and local propagation persistence -├── postgres_health.rs # PostgreSQL major/write-readiness and relation-integrity probe +├── postgres_health.rs # PostgreSQL major/write-readiness, relation-integrity, and Active PR #132 (successor to #122) operational snapshot +├── postgres_health_http.rs # Active PR #132 (successor to #122) observes PostgreSQL only for GET /ready; /live stays store-I/O free; serve loop keeps accepting after a dropped probe (not protected-main truth) ├── postgres_inbox_consumption.rs # PostgreSQL inbox consumption distinct from receipt ├── postgres_instrument_release.rs # PostgreSQL locale-specific instrument-release persistence ├── postgres_integration.rs # PostgreSQL integration evidence/delivery-attempt persistence adapter @@ -128,10 +131,12 @@ migrations/ └── 0012_integration_consumption.sql ``` -Still-Target logical modules/adapters include remaining product aggregate persistence/repositories, public/admin HTTP and event transports, live fast-mlsirm/Keyverse/Gyeot/TEPP/semantic-data-portal adapters, research-release staging, deterministic narrative mapping, longitudinal normalized ingestion, participant identity-link history persistence, runtime health transports/metrics, and Measurement Workbench orchestration. +Still-Target logical modules/adapters include remaining product aggregate persistence/repositories, public/admin product HTTP and event transports, live fast-mlsirm/Keyverse/Gyeot/TEPP/semantic-data-portal adapters, research-release staging, deterministic narrative mapping, longitudinal normalized ingestion, participant identity-link history persistence, TLS/keep-alive/metrics for the health listener, and Measurement Workbench orchestration. ### Active implementation work that is not protected-main truth +**Active PR** #132 (successor to #122) PostgreSQL-backed operator health HTTP is not protected-main truth until an unchanged reviewed/check-clean head is integrated. Operators start the process with `run_health_process` after setting `HEALTH_LISTEN_ADDR` or `PORT`. GET `/live` answers process liveness without store I/O even when `DATABASE_URL` is down. GET `/ready` observes a live operational snapshot after accept; a bare `/ready` requires `postgres_operational_store`. Caller-measured `HEALTH_BACKLOG_HEALTH` is required before readiness can be true. Driver errors are not exposed. Public/admin product routes, measured backlog thresholds, TLS, and deployment-profile evidence remain outside this slice. + **Active PR** #60 exclusive outbox delivery-lease persistence is not protected-main truth until an unchanged reviewed/check-clean head is integrated. Pending outbox rows accept one fenced worker lease, recover expiry from the database clock without transferring the fence, reject a future caller timestamp that would steal a still-live lease, and reject stale or zero-window claims. Live side-effect execution remains outside this slice. ## 5. ADR traceability by concern @@ -190,7 +195,7 @@ Whenever a durable conversation decision changes one of those boundaries, the ap The prose API/event families in TRD are architecture requirements, not evidence of an implemented transport. -When the first HTTP API is implemented, the same PR or a prerequisite PR must add and validate an OpenAPI 3.2.x document whose operations and problem responses match the actual implementation. HTTP errors use RFC 9457 problem details unless a documented domain representation is more appropriate. +The first implemented HTTP surface is the operator health-probe pair. Active PR #132 (successor to #122) adds `openapi/health-probes.yaml` (OpenAPI 3.2.0) listing only GET `/live` and GET `/ready`, plus a process entrypoint that binds those probes from `HEALTH_LISTEN_ADDR` or `PORT`. Probe success/unready bodies use the documented health-snapshot representation; malformed or unsupported requests use RFC 9457 `application/problem+json`. Public/admin product routes remain unimplemented and are not listed. When durable message transport is implemented, the same PR or a prerequisite PR must add and validate an AsyncAPI 3.1.x document for actually produced/consumed event channels and message schemas. It must encode/reference ADR-0014 canonical UTF-8 payload hashing, SHA-256 payload digest semantics, tenant/resource binding, deduplication identity, pending/processing/completed consumption, replay retention, and quarantine behavior. @@ -220,6 +225,12 @@ CI should validate linked documentation paths and status/name consistency now an ## 10. References +Eddy, W. (Ed.). (2022). *Transmission Control Protocol (TCP)* (RFC 9293). Internet Engineering Task Force. https://doi.org/10.17487/RFC9293 + +Fielding, R., Nottingham, M., & Reschke, J. (Eds.). (2022). *HTTP Semantics* (RFC 9110). Internet Engineering Task Force. https://doi.org/10.17487/RFC9110 + +Kubernetes Authors. (2024). *Configure liveness, readiness and startup probes*. Kubernetes Documentation. https://kubernetes.io/docs/tasks/configure-pod-container/configure-liveness-readiness-startup-probes/ + Nottingham, M., Wilde, E., & Dalal, S. (2023). *Problem Details for HTTP APIs* (RFC 9457). Internet Engineering Task Force. https://doi.org/10.17487/RFC9457 OpenAPI Initiative. (2025). *OpenAPI Specification, Version 3.2.0*. diff --git a/docs/adr/0014-api-and-event-contract-representation.md b/docs/adr/0014-api-and-event-contract-representation.md index b8f3bbeb..44f74255 100644 --- a/docs/adr/0014-api-and-event-contract-representation.md +++ b/docs/adr/0014-api-and-event-contract-representation.md @@ -6,7 +6,7 @@ - Scope: Psychometrics Commons public/admin HTTP APIs, product-owned durable domain events, errors, schema/version negotiation - Supersedes: none - Superseded by: none -- Current/as-built status: public/admin HTTP transport and durable external event transport are not yet implemented on protected main; current Rust domain contracts are transport-neutral +- Current/as-built status: public/admin product HTTP transport and durable external event transport are not yet implemented on protected main; operator GET `/live` and GET `/ready` probes exist only on Active PR #132 (successor to #122) with `openapi/health-probes.yaml`; that same PR binds a process from `HEALTH_LISTEN_ADDR` or `PORT`, keeps serving after a dropped probe, answers `/live` without store I/O even when `DATABASE_URL` is down, and observes `observe_postgres_operational_snapshot` only for `/ready` without exposing driver errors - Target status: every implemented HTTP/event surface has an exact versioned machine-readable as-built contract and deterministic integrity/idempotency semantics - Migration status: no deployed HTTP/event transport requires migration yet; the first implementation must introduce the contract in the same or prerequisite PR @@ -254,6 +254,12 @@ A future major transport change may supersede this ADR if OpenAPI/AsyncAPI no lo Bray, T. (Ed.). (2017). *The JavaScript Object Notation (JSON) Data Interchange Format* (RFC 8259). Internet Engineering Task Force. https://doi.org/10.17487/RFC8259 +Eddy, W. (Ed.). (2022). *Transmission Control Protocol (TCP)* (RFC 9293). Internet Engineering Task Force. https://doi.org/10.17487/RFC9293 + +Fielding, R., Nottingham, M., & Reschke, J. (Eds.). (2022). *HTTP Semantics* (RFC 9110). Internet Engineering Task Force. https://doi.org/10.17487/RFC9110 + +Kubernetes Authors. (2024). *Configure liveness, readiness and startup probes*. Kubernetes Documentation. https://kubernetes.io/docs/tasks/configure-pod-container/configure-liveness-readiness-startup-probes/ + Rundgren, A., Jordan, B., & Erdtman, S. (2020). *JSON Canonicalization Scheme (JCS)* (RFC 8785). Internet Engineering Task Force. https://doi.org/10.17487/RFC8785 AsyncAPI Initiative. (2026). *AsyncAPI Specification, Version 3.1.0*. diff --git a/docs/architecture/DEPLOYMENT_AND_OPERATIONS.md b/docs/architecture/DEPLOYMENT_AND_OPERATIONS.md index c839c6cc..8cbb489d 100644 --- a/docs/architecture/DEPLOYMENT_AND_OPERATIONS.md +++ b/docs/architecture/DEPLOYMENT_AND_OPERATIONS.md @@ -117,6 +117,8 @@ backup/restore process A single process implementation must not blur transactional semantics merely because components are co-located. +The first runnable operator surface is the health-probe process. Set `HEALTH_LISTEN_ADDR` or platform `PORT` (binds `0.0.0.0:$PORT`), then call `run_health_process`. Point liveness at GET `/live` and readiness at GET `/ready`. Optional `DATABASE_URL` is observed only for readiness. Do not treat probe success as a measured SLO. + ## 4. Health and readiness ### Liveness diff --git a/docs/doctoring/standards-and-evidence.md b/docs/doctoring/standards-and-evidence.md index 7879474c..4922413e 100644 --- a/docs/doctoring/standards-and-evidence.md +++ b/docs/doctoring/standards-and-evidence.md @@ -1,7 +1,7 @@ # Standards and Evidence Baseline - Status: Living doctoring record -- Last reviewed: 2026-08-11 +- Last reviewed: 2026-08-16 - Scope: Psychometrics Commons product, hosted runtime, reference clients, optional AI, identity integration, and assessment governance This record identifies authoritative standards and primary guidance that materially constrain product design. It is not a certification claim. Each implementation PR that relies on one of these sources must translate the source into a concrete requirement, test, control, or ADR rather than citing it decoratively. @@ -30,6 +30,16 @@ Product consequences: - timing accommodations are part of instrument-version evidence when timing can affect the response process; - automated accessibility checks are supplemented by manual and assistive-technology acceptance testing. +## Operator HTTP + +RFC 9110 defines HTTP request-target and origin-server authority. The first process surface binds an explicit listen address from `HEALTH_LISTEN_ADDR` or platform `PORT` and answers only GET `/live` and GET `/ready`. It is not a public assessment origin and must not invent routes or echo request/store text. + +Product consequences: + +- operators set an unpadded listen address or port before the process starts; +- liveness remains independent of store I/O; +- unsupported methods and paths use RFC 9457 problem types, not `about:blank`. + ## Digital identity and federation NIST SP 800-63 Revision 4 is the current NIST Digital Identity Guidelines suite. The final Revision 4 was published in July 2025 and supersedes Revision 3. The suite covers identity proofing, authentication, authenticator management, federation, assertions, security/privacy, and customer-experience considerations. @@ -107,6 +117,10 @@ 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/ +Eddy, W. (Ed.). (2022). *Transmission Control Protocol (TCP)* (RFC 9293). Internet Engineering Task Force. https://doi.org/10.17487/RFC9293 + +Fielding, R., Nottingham, M., & Reschke, J. (Eds.). (2022). *HTTP Semantics* (RFC 9110). Internet Engineering Task Force. https://doi.org/10.17487/RFC9110 + 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 diff --git a/openapi/health-probes.yaml b/openapi/health-probes.yaml new file mode 100644 index 00000000..0ed9d2cb --- /dev/null +++ b/openapi/health-probes.yaml @@ -0,0 +1,130 @@ +openapi: 3.2.0 +info: + title: Psychometrics Commons operator health probes + version: 0.1.0 + description: > + As-built operator liveness and readiness probes. These operations expose a + RuntimeHealthSnapshot. They are not public product APIs and do not claim + measured availability SLOs. +paths: + /live: + get: + operationId: getLive + summary: Process liveness + description: > + Returns 200 when the process is live and 503 when it is not. Optional + capability outages do not change liveness. + responses: + "200": + description: Process is live. + content: + application/json: + schema: + $ref: "#/components/schemas/HealthSnapshot" + "400": + $ref: "#/components/responses/Problem" + "405": + $ref: "#/components/responses/Problem" + "503": + description: Process is not live. + content: + application/json: + schema: + $ref: "#/components/schemas/HealthSnapshot" + /ready: + get: + operationId: getReady + summary: Operation-scoped readiness + description: > + Returns 200 only when the process is live, backlog and integrity accept + new work, and every named capability can accept new work. Unknown + required capabilities fail closed. + parameters: + - name: capability + in: query + required: false + explode: true + schema: + type: array + items: + type: string + description: Opaque required capability references for this operation. + responses: + "200": + description: New work for the named capabilities is safe. + content: + application/json: + schema: + $ref: "#/components/schemas/HealthSnapshot" + "400": + $ref: "#/components/responses/Problem" + "405": + $ref: "#/components/responses/Problem" + "503": + description: New state-changing work is not safe. + content: + application/json: + schema: + $ref: "#/components/schemas/HealthSnapshot" +components: + schemas: + HealthSnapshot: + type: object + additionalProperties: false + required: + - live + - ready + - backlog_health + - data_integrity_health + - capabilities + properties: + live: + type: boolean + ready: + type: boolean + backlog_health: + type: string + enum: [within_bounds, stalled, unknown] + data_integrity_health: + type: string + enum: [verified, incompatible, unknown] + capabilities: + type: array + items: + $ref: "#/components/schemas/CapabilityHealth" + CapabilityHealth: + type: object + additionalProperties: false + required: [capability_ref, state, accepts_new_work] + properties: + capability_ref: + type: string + state: + type: string + enum: [available, degraded, unavailable, unknown] + accepts_new_work: + type: boolean + Problem: + type: object + additionalProperties: false + required: [type, title, status, detail] + properties: + type: + type: string + enum: + - urn:psychometrics-commons:problem:bad-request + - urn:psychometrics-commons:problem:not-found + - urn:psychometrics-commons:problem:method-not-allowed + title: + type: string + status: + type: integer + detail: + type: string + responses: + Problem: + description: Safe RFC 9457 problem details without raw store errors. + content: + application/problem+json: + schema: + $ref: "#/components/schemas/Problem" diff --git a/src/health_http.rs b/src/health_http.rs new file mode 100644 index 00000000..3efea37b --- /dev/null +++ b/src/health_http.rs @@ -0,0 +1,752 @@ +//! Operator HTTP probes for process liveness and operation-scoped readiness. +//! +//! These probes translate [`RuntimeHealthSnapshot`] into load-balancer-safe +//! HTTP responses. They do not invent availability SLOs, execute store I/O, or +//! expose raw database or provider errors. + +use crate::health::{ + BacklogHealth, CapabilityHealth, CapabilityState, DataIntegrityHealth, RuntimeHealthSnapshot, +}; +use std::fmt::Write; +use std::io::{self, Read}; +use std::net::{SocketAddr, TcpListener, TcpStream}; +use std::time::Duration; + +/// Process-liveness probe path. +pub const HEALTH_LIVE_PATH: &str = "/live"; +/// Operation-readiness probe path. +pub const HEALTH_READY_PATH: &str = "/ready"; +/// Bounded read/write timeout for one accepted probe connection. +pub const HEALTH_HTTP_IO_TIMEOUT: Duration = Duration::from_secs(2); +/// Maximum accepted probe request size, including headers. +pub const HEALTH_HTTP_MAX_REQUEST_BYTES: usize = 8_192; + +/// HTTP response produced by a health probe request. +#[derive(Clone, Debug, Eq, PartialEq)] +pub struct HealthHttpResponse { + status: u16, + content_type: &'static str, + body: String, +} + +impl HealthHttpResponse { + fn json(status: u16, body: String) -> Self { + Self { + status, + content_type: "application/json", + body, + } + } + + fn problem(status: u16, type_uri: &str, title: &str, detail: &str) -> Self { + Self { + status, + content_type: "application/problem+json", + body: format!( + "{{\"type\":{},\"title\":{},\"status\":{status},\"detail\":{}}}", + json_string(type_uri), + json_string(title), + json_string(detail) + ), + } + } + + /// Return the HTTP status code. + #[must_use] + pub const fn status(&self) -> u16 { + self.status + } + + /// Return the response content type. + #[must_use] + pub const fn content_type(&self) -> &'static str { + self.content_type + } + + /// Return the response body. + #[must_use] + pub fn body(&self) -> &str { + &self.body + } +} + +/// Translate one raw HTTP/1.1 request into a liveness or readiness response. +/// +/// Liveness answers whether the process is live. Readiness answers whether new +/// state-changing work is safe for the caller-named required capabilities. +/// Unknown required capabilities, stalled backlog, unknown integrity, or a +/// non-live process fail closed with HTTP 503. Unsupported methods and paths +/// return RFC 9457 problem details without echoing the raw request. +#[must_use] +pub fn handle_health_http_request( + request: &str, + snapshot: &RuntimeHealthSnapshot, +) -> HealthHttpResponse { + let Some((method, target)) = parse_request_line(request) else { + return HealthHttpResponse::problem( + 400, + "urn:psychometrics-commons:problem:bad-request", + "Bad Request", + "health probe request must include an HTTP method and target", + ); + }; + if method != "GET" { + return HealthHttpResponse::problem( + 405, + "urn:psychometrics-commons:problem:method-not-allowed", + "Method Not Allowed", + "health probes accept GET /live and GET /ready only", + ); + } + let (path, query) = split_target(target); + match path { + HEALTH_LIVE_PATH => { + let status = if snapshot.is_live() { 200 } else { 503 }; + HealthHttpResponse::json(status, snapshot_body(snapshot, snapshot.is_ready_for(&[]))) + } + HEALTH_READY_PATH => health_ready_response(snapshot, &required_capabilities(query)), + _ => HealthHttpResponse::problem( + 404, + "urn:psychometrics-commons:problem:not-found", + "Not Found", + "health probes accept GET /live and GET /ready only", + ), + } +} + +/// Bind a blocking TCP listener for operator health probes. +/// +/// The caller chooses the address. Tests and local operators typically bind +/// `127.0.0.1:0`. This function does not start accepting connections. +/// +/// # Errors +/// +/// Returns the I/O error if the operating system cannot bind the address. +pub fn bind_health_http(addr: SocketAddr) -> io::Result { + TcpListener::bind(addr) +} + +/// Accept one TCP connection and serve a single health-probe request. +/// +/// The connection is closed after the response. Keep-alive, TLS, and public +/// product routes are outside this slice. +/// +/// # Errors +/// +/// Returns the I/O error if accept, read, or write fails. +pub fn accept_one_health_http( + listener: &TcpListener, + snapshot: &RuntimeHealthSnapshot, +) -> io::Result<()> { + accept_one_health_http_with(listener, |request| { + handle_health_http_request(request, snapshot) + }) +} + +/// Accept one TCP connection and answer it with `handler`. +/// +/// The handler runs after a bounded read so store observation cannot start +/// before the connection is accepted. Incomplete or oversized requests become +/// empty request text and fail closed as HTTP 400 without echoing input. +/// +/// # Errors +/// +/// Returns the I/O error if accept, read, or write fails. +pub fn accept_one_health_http_with(listener: &TcpListener, handler: F) -> io::Result<()> +where + F: FnOnce(&str) -> HealthHttpResponse, +{ + let (mut stream, _) = listener.accept()?; + stream.set_read_timeout(Some(HEALTH_HTTP_IO_TIMEOUT))?; + stream.set_write_timeout(Some(HEALTH_HTTP_IO_TIMEOUT))?; + let request = read_http_request(&mut stream)?; + let response = handler(&request); + write_http_response(&mut stream, &response) +} + +/// Serve probe requests until `accept` fails. +/// +/// Operators run this loop so Kubernetes or a load balancer can keep asking +/// GET `/live` and GET `/ready`. Interrupted, aborted, or reset accepts retry. +/// Per-connection read or write errors do not stop later probes. Other accept +/// errors, including `WouldBlock` on a closed or non-blocking listener, stop +/// the loop so it does not spin. TLS, keep-alive, and measured SLO values +/// remain outside this slice. +/// +/// # Errors +/// +/// Returns the I/O error that stopped the loop. +pub fn serve_health_http( + listener: &TcpListener, + snapshot: &RuntimeHealthSnapshot, +) -> io::Result<()> { + serve_health_http_with(listener, |request| { + handle_health_http_request(request, snapshot) + }) +} + +/// Serve probe requests with `handler` until `accept` fails. +/// +/// # Errors +/// +/// Returns the I/O error that stopped the loop. +pub fn serve_health_http_with(listener: &TcpListener, mut handler: F) -> io::Result<()> +where + F: FnMut(&str) -> HealthHttpResponse, +{ + loop { + match listener.accept() { + Ok((mut stream, _)) => { + // Per-connection I/O never stops later probes; a dropped client + // must not take the operator listener down. + let _ = serve_accepted_health_http(&mut stream, &mut handler); + } + Err(error) => { + if !should_continue_after_serve_error(&error, ServeIoSource::Accept) { + return Err(error); + } + } + } + } +} + +fn serve_accepted_health_http(stream: &mut TcpStream, handler: &mut F) -> io::Result<()> +where + F: FnMut(&str) -> HealthHttpResponse, +{ + stream.set_read_timeout(Some(HEALTH_HTTP_IO_TIMEOUT))?; + stream.set_write_timeout(Some(HEALTH_HTTP_IO_TIMEOUT))?; + let request = read_http_request(stream)?; + let response = handler(&request); + write_http_response(stream, &response) +} + +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +enum ServeIoSource { + Accept, +} + +fn should_continue_after_serve_error(error: &io::Error, source: ServeIoSource) -> bool { + match source { + ServeIoSource::Accept => matches!( + error.kind(), + io::ErrorKind::Interrupted + | io::ErrorKind::ConnectionAborted + | io::ErrorKind::ConnectionReset + ), + } +} + +/// Return whether this request is GET `/ready` and must observe operational state. +#[must_use] +pub fn health_request_requires_readiness_snapshot(request: &str) -> bool { + let Some((method, target)) = parse_request_line(request) else { + return false; + }; + method == "GET" && split_target(target).0 == HEALTH_READY_PATH +} + +/// Return caller-named `capability` query values from one probe request. +#[must_use] +pub fn health_request_required_capabilities(request: &str) -> Vec<&str> { + parse_request_line(request) + .map(|(_, target)| required_capabilities(split_target(target).1)) + .unwrap_or_default() +} + +/// Answer readiness for caller-named required capabilities. +#[must_use] +pub fn health_ready_response( + snapshot: &RuntimeHealthSnapshot, + required_capabilities: &[&str], +) -> HealthHttpResponse { + let ready = snapshot.is_ready_for(required_capabilities); + let status = if ready { 200 } else { 503 }; + HealthHttpResponse::json(status, snapshot_body(snapshot, ready)) +} + +fn read_http_request(stream: &mut TcpStream) -> io::Result { + let mut buffer = Vec::new(); + let mut chunk = [0_u8; 512]; + loop { + let read_result = stream.read(&mut chunk); + match apply_request_read(&mut buffer, &chunk, read_result)? { + RequestReadProgress::Continue => {} + RequestReadProgress::Complete => break, + } + } + if buffer.len() > HEALTH_HTTP_MAX_REQUEST_BYTES + || !buffer.windows(4).any(|window| window == b"\r\n\r\n") + { + return Ok(String::new()); + } + Ok(String::from_utf8_lossy(&buffer).into_owned()) +} + +#[derive(Debug)] +enum RequestReadProgress { + Continue, + Complete, +} + +fn apply_request_read( + buffer: &mut Vec, + chunk: &[u8], + read_result: io::Result, +) -> io::Result { + match read_result { + Ok(0) => Ok(RequestReadProgress::Complete), + Ok(read) => { + buffer.extend_from_slice(&chunk[..read]); + if buffer.windows(4).any(|window| window == b"\r\n\r\n") + || buffer.len() > HEALTH_HTTP_MAX_REQUEST_BYTES + { + Ok(RequestReadProgress::Complete) + } else { + Ok(RequestReadProgress::Continue) + } + } + Err(error) + if matches!( + error.kind(), + io::ErrorKind::TimedOut | io::ErrorKind::WouldBlock + ) => + { + Ok(RequestReadProgress::Complete) + } + Err(error) => Err(error), + } +} + +fn write_http_response(stream: &mut TcpStream, response: &HealthHttpResponse) -> io::Result<()> { + let body = response.body().as_bytes(); + let allow = if response.status() == 405 { + "Allow: GET\r\n" + } else { + "" + }; + let header = format!( + "HTTP/1.1 {} {}\r\nContent-Type: {}\r\nContent-Length: {}\r\nCache-Control: no-store\r\n{allow}Connection: close\r\n\r\n", + response.status(), + reason_phrase(response.status()), + response.content_type(), + body.len() + ); + io::Write::write_all(stream, header.as_bytes())?; + io::Write::write_all(stream, body) +} + +const fn reason_phrase(status: u16) -> &'static str { + match status { + 200 => "OK", + 400 => "Bad Request", + 404 => "Not Found", + 405 => "Method Not Allowed", + 503 => "Service Unavailable", + _ => "Error", + } +} + +fn parse_request_line(request: &str) -> Option<(&str, &str)> { + let line = request.lines().next()?; + let mut parts = line.split_whitespace(); + let method = parts.next()?; + let target = parts.next()?; + let version = parts.next()?; + if !version.starts_with("HTTP/") || parts.next().is_some() { + return None; + } + Some((method, target)) +} + +fn split_target(target: &str) -> (&str, &str) { + target.split_once('?').unwrap_or((target, "")) +} + +fn required_capabilities(query: &str) -> Vec<&str> { + query + .split('&') + .filter_map(|pair| { + let (key, value) = pair.split_once('=')?; + (key == "capability").then_some(value) + }) + .collect() +} + +fn snapshot_body(snapshot: &RuntimeHealthSnapshot, ready: bool) -> String { + let mut capabilities = String::from("["); + for (index, capability) in snapshot.capabilities().iter().enumerate() { + if index > 0 { + capabilities.push(','); + } + capabilities.push_str(&capability_body(capability)); + } + capabilities.push(']'); + format!( + "{{\"live\":{},\"ready\":{ready},\"backlog_health\":{},\"data_integrity_health\":{},\"capabilities\":{capabilities}}}", + json_bool(snapshot.is_live()), + json_string(backlog_label(snapshot.backlog_health())), + json_string(integrity_label(snapshot.data_integrity_health())), + ) +} + +fn capability_body(capability: &CapabilityHealth) -> String { + format!( + "{{\"capability_ref\":{},\"state\":{},\"accepts_new_work\":{}}}", + json_string(capability.capability_ref()), + json_string(capability_state_label(capability.state())), + json_bool(capability.accepts_new_work()) + ) +} + +const fn backlog_label(health: BacklogHealth) -> &'static str { + match health { + BacklogHealth::WithinBounds => "within_bounds", + BacklogHealth::Stalled => "stalled", + BacklogHealth::Unknown => "unknown", + } +} + +const fn integrity_label(health: DataIntegrityHealth) -> &'static str { + match health { + DataIntegrityHealth::Verified => "verified", + DataIntegrityHealth::Incompatible => "incompatible", + DataIntegrityHealth::Unknown => "unknown", + } +} + +const fn capability_state_label(state: CapabilityState) -> &'static str { + match state { + CapabilityState::Available => "available", + CapabilityState::Degraded => "degraded", + CapabilityState::Unavailable => "unavailable", + CapabilityState::Unknown => "unknown", + } +} + +fn json_bool(value: bool) -> &'static str { + if value { + "true" + } else { + "false" + } +} + +fn json_string(value: &str) -> String { + let mut escaped = String::from("\""); + for character in value.chars() { + match character { + '"' => escaped.push_str("\\\""), + '\\' => escaped.push_str("\\\\"), + '\n' => escaped.push_str("\\n"), + '\r' => escaped.push_str("\\r"), + '\t' => escaped.push_str("\\t"), + character if character.is_control() => { + let _ = write!(escaped, "\\u{:04x}", character as u32); + } + character => escaped.push(character), + } + } + escaped.push('"'); + escaped +} + +#[cfg(test)] +mod tests { + use super::{ + apply_request_read, backlog_label, bind_health_http, capability_state_label, + handle_health_http_request, health_ready_response, health_request_required_capabilities, + health_request_requires_readiness_snapshot, integrity_label, json_string, reason_phrase, + serve_health_http_with, should_continue_after_serve_error, RequestReadProgress, + ServeIoSource, HEALTH_HTTP_MAX_REQUEST_BYTES, HEALTH_LIVE_PATH, HEALTH_READY_PATH, + }; + use crate::health::{ + BacklogHealth, CapabilityHealth, CapabilityState, DataIntegrityHealth, + RuntimeHealthSnapshot, + }; + use std::io; + use std::net::{IpAddr, Ipv4Addr, SocketAddr}; + + #[test] + fn remaining_labels_and_json_escapes_are_stable() { + assert_eq!(backlog_label(BacklogHealth::Unknown), "unknown"); + assert_eq!( + integrity_label(DataIntegrityHealth::Incompatible), + "incompatible" + ); + assert_eq!( + capability_state_label(CapabilityState::Degraded), + "degraded" + ); + assert_eq!(capability_state_label(CapabilityState::Unknown), "unknown"); + assert_eq!(integrity_label(DataIntegrityHealth::Unknown), "unknown"); + assert_eq!(json_string("a\"b\\c"), "\"a\\\"b\\\\c\""); + assert_eq!(json_string("a\n\r\t"), "\"a\\n\\r\\t\""); + assert_eq!(json_string("\u{0001}"), "\"\\u0001\""); + assert_eq!(reason_phrase(200), "OK"); + assert_eq!(reason_phrase(400), "Bad Request"); + assert_eq!(reason_phrase(404), "Not Found"); + assert_eq!(reason_phrase(405), "Method Not Allowed"); + assert_eq!(reason_phrase(503), "Service Unavailable"); + assert_eq!(reason_phrase(418), "Error"); + } + + #[test] + fn request_line_rejects_extra_tokens_and_non_http_versions() { + let snapshot = RuntimeHealthSnapshot::new( + true, + BacklogHealth::Unknown, + DataIntegrityHealth::Incompatible, + vec![ + CapabilityHealth::new("research_registration", CapabilityState::Degraded, true) + .unwrap(), + ], + ) + .unwrap(); + assert_eq!( + handle_health_http_request("GET /live HTTP/1.1 extra\r\n\r\n", &snapshot).status(), + 400 + ); + assert_eq!( + handle_health_http_request("GET /live SMTP/1.0\r\n\r\n", &snapshot).status(), + 400 + ); + let live = handle_health_http_request( + &format!("GET {HEALTH_LIVE_PATH}?capability=ignored HTTP/1.1\r\n\r\n"), + &snapshot, + ); + assert_eq!(live.status(), 200); + assert!(live.body().contains("\"backlog_health\":\"unknown\"")); + assert!(live + .body() + .contains("\"data_integrity_health\":\"incompatible\"")); + assert!(live.body().contains("\"state\":\"degraded\"")); + assert_eq!(live.content_type(), "application/json"); + + let ready_snapshot = RuntimeHealthSnapshot::new( + true, + BacklogHealth::WithinBounds, + DataIntegrityHealth::Verified, + vec![ + CapabilityHealth::new("research_registration", CapabilityState::Degraded, true) + .unwrap(), + ], + ) + .unwrap(); + let ready = handle_health_http_request( + "GET /ready?capability=research_registration HTTP/1.1\r\n\r\n", + &ready_snapshot, + ); + assert_eq!(ready.status(), 200); + assert_eq!(ready.content_type(), "application/json"); + assert!(ready.body().contains("\"ready\":true")); + + let not_allowed = handle_health_http_request("POST /live HTTP/1.1\r\n\r\n", &snapshot); + assert_eq!(not_allowed.status(), 405); + assert_eq!(not_allowed.content_type(), "application/problem+json"); + assert!(not_allowed + .body() + .contains("\"title\":\"Method Not Allowed\"")); + assert!( + not_allowed + .body() + .contains("\"type\":\"urn:psychometrics-commons:problem:method-not-allowed\""), + "health problems must use an explicit product type, not about:blank: {}", + not_allowed.body() + ); + assert!(!not_allowed.body().contains("about:blank")); + + let missing = handle_health_http_request("GET /v1/sessions HTTP/1.1\r\n\r\n", &snapshot); + assert_eq!(missing.status(), 404); + assert_eq!(missing.content_type(), "application/problem+json"); + assert!(missing.body().contains("\"title\":\"Not Found\"")); + assert!(missing + .body() + .contains("\"type\":\"urn:psychometrics-commons:problem:not-found\"")); + assert!(!missing.body().contains("about:blank")); + + let bad = handle_health_http_request("NOT-A-REQUEST", &snapshot); + assert_eq!(bad.status(), 400); + assert!(bad + .body() + .contains("\"type\":\"urn:psychometrics-commons:problem:bad-request\"")); + assert!(!bad.body().contains("about:blank")); + + let not_live = RuntimeHealthSnapshot::new( + false, + BacklogHealth::Stalled, + DataIntegrityHealth::Unknown, + vec![ + CapabilityHealth::new("research_registration", CapabilityState::Unavailable, false) + .unwrap(), + CapabilityHealth::new("authenticated_linking", CapabilityState::Available, true) + .unwrap(), + ], + ) + .unwrap(); + let dead = handle_health_http_request( + &format!("GET {HEALTH_LIVE_PATH} HTTP/1.1\r\n\r\n"), + ¬_live, + ); + assert_eq!(dead.status(), 503); + assert!(dead.body().contains("\"live\":false")); + assert!(dead.body().contains("\"ready\":false")); + let not_ready = handle_health_http_request( + &format!("GET {HEALTH_READY_PATH} HTTP/1.1\r\n\r\n"), + ¬_live, + ); + assert_eq!(not_ready.status(), 503); + assert!(not_ready.body().contains("\"ready\":false")); + } + + #[test] + fn readiness_helpers_classify_requests_and_required_capabilities() { + let snapshot = RuntimeHealthSnapshot::new( + true, + BacklogHealth::WithinBounds, + DataIntegrityHealth::Verified, + vec![ + CapabilityHealth::new("research_registration", CapabilityState::Degraded, true) + .unwrap(), + ], + ) + .unwrap(); + assert!(!health_request_requires_readiness_snapshot( + "GET /live HTTP/1.1\r\n\r\n" + )); + assert!(!health_request_requires_readiness_snapshot("NOT-A-REQUEST")); + assert!(!health_request_requires_readiness_snapshot( + "POST /ready HTTP/1.1\r\n\r\n" + )); + assert!(health_request_requires_readiness_snapshot( + "GET /ready?capability=scoring HTTP/1.1\r\n\r\n" + )); + assert!(health_request_required_capabilities("NOT-A-REQUEST").is_empty()); + assert_eq!( + health_request_required_capabilities( + "GET /ready?capability=scoring&capability=research_registration HTTP/1.1\r\n\r\n" + ), + vec!["scoring", "research_registration"] + ); + let ready = health_ready_response(&snapshot, &["research_registration"]); + assert_eq!(ready.status(), 200); + let unready = health_ready_response(&snapshot, &["unregistered_capability"]); + assert_eq!(unready.status(), 503); + assert!(health_request_requires_readiness_snapshot( + "GET /ready HTTP/1.1\r\n\r\n" + )); + assert!(health_request_required_capabilities("GET /ready HTTP/1.1\r\n\r\n").is_empty()); + } + + #[test] + fn request_read_progress_covers_eof_timeout_and_io_failure() { + let mut buffer = Vec::new(); + assert!(matches!( + apply_request_read(&mut buffer, b"", Ok(0)).unwrap(), + RequestReadProgress::Complete + )); + buffer.clear(); + assert!(matches!( + apply_request_read(&mut buffer, b"GET /l", Ok(6)).unwrap(), + RequestReadProgress::Continue + )); + assert_eq!(buffer, b"GET /l"); + buffer.clear(); + assert!(matches!( + apply_request_read(&mut buffer, b"GET /live HTTP/1.1\r\n\r\n", Ok(22)).unwrap(), + RequestReadProgress::Complete + )); + buffer.clear(); + let oversized = vec![b'A'; HEALTH_HTTP_MAX_REQUEST_BYTES + 1]; + assert!(matches!( + apply_request_read(&mut buffer, &oversized, Ok(oversized.len())).unwrap(), + RequestReadProgress::Complete + )); + buffer.clear(); + assert!(matches!( + apply_request_read( + &mut buffer, + b"", + Err(io::Error::new(io::ErrorKind::TimedOut, "timeout")) + ) + .unwrap(), + RequestReadProgress::Complete + )); + buffer.clear(); + assert!(matches!( + apply_request_read( + &mut buffer, + b"", + Err(io::Error::new(io::ErrorKind::WouldBlock, "block")) + ) + .unwrap(), + RequestReadProgress::Complete + )); + let error = apply_request_read( + &mut buffer, + b"", + Err(io::Error::new(io::ErrorKind::ConnectionReset, "reset")), + ) + .expect_err("non-timeout I/O errors must propagate"); + assert_eq!(error.kind(), io::ErrorKind::ConnectionReset); + } + + #[test] + fn serve_loop_stops_when_accept_returns_a_fatal_error() { + let listener = + bind_health_http(SocketAddr::new(IpAddr::V4(Ipv4Addr::LOCALHOST), 0)).unwrap(); + listener + .set_nonblocking(true) + .expect("the test must force accept to return WouldBlock"); + let error = serve_health_http_with(&listener, |_| { + panic!("a fatal accept must not invoke the request handler") + }) + .expect_err("WouldBlock on accept must stop the probe process"); + assert_eq!(error.kind(), io::ErrorKind::WouldBlock); + } + + #[test] + fn serve_loop_retries_interrupted_accepts_and_stops_on_other_errors() { + assert!(should_continue_after_serve_error( + &io::Error::new(io::ErrorKind::Interrupted, "signal"), + ServeIoSource::Accept, + )); + assert!(should_continue_after_serve_error( + &io::Error::new(io::ErrorKind::ConnectionAborted, "client gone"), + ServeIoSource::Accept, + )); + assert!( + should_continue_after_serve_error( + &io::Error::new(io::ErrorKind::ConnectionReset, "rst before accept"), + ServeIoSource::Accept, + ), + "accept ConnectionReset must retry so a reset handshake does not stop later probes" + ); + assert!( + !should_continue_after_serve_error( + &io::Error::new(io::ErrorKind::WouldBlock, "nonblocking"), + ServeIoSource::Accept, + ), + "accept WouldBlock must stop a closed or non-blocking listener" + ); + } + + #[test] + fn serve_loop_keeps_running_after_per_connection_io_errors() { + assert!(should_continue_after_serve_error( + &io::Error::new(io::ErrorKind::ConnectionAborted, "client gone"), + ServeIoSource::Accept, + )); + assert!(should_continue_after_serve_error( + &io::Error::new(io::ErrorKind::Interrupted, "signal"), + ServeIoSource::Accept, + )); + assert!( + !should_continue_after_serve_error( + &io::Error::new(io::ErrorKind::WouldBlock, "nonblocking"), + ServeIoSource::Accept, + ), + "accept WouldBlock must still stop a closed or non-blocking listener" + ); + } +} diff --git a/src/health_process.rs b/src/health_process.rs new file mode 100644 index 00000000..93dfc2d1 --- /dev/null +++ b/src/health_process.rs @@ -0,0 +1,430 @@ +//! Process entrypoint that binds operator health probes from environment input. +//! +//! Operators start this process so a load balancer can keep calling GET `/live` +//! and GET `/ready`. Liveness never opens a database connection. Readiness +//! observes `DATABASE_URL` only after accept, and never echoes driver errors. + +use crate::health::{ + BacklogHealth, CapabilityHealth, CapabilityState, DataIntegrityHealth, RuntimeHealthSnapshot, +}; +use crate::health_http::{ + bind_health_http, handle_health_http_request, health_ready_response, + health_request_required_capabilities, health_request_requires_readiness_snapshot, + serve_health_http_with, HealthHttpResponse, HEALTH_HTTP_IO_TIMEOUT, +}; +use crate::postgres_health::POSTGRES_OPERATIONAL_STORE_CAPABILITY_REF; +use crate::postgres_health_http::handle_postgres_health_http_request; +use postgres::{Client, NoTls}; +use std::error::Error; +use std::fmt::{Display, Formatter}; +use std::io; +use std::net::{IpAddr, Ipv4Addr, SocketAddr, TcpListener}; +use std::str::FromStr; + +/// Full listen-address environment variable. Wins over [`HEALTH_LISTEN_PORT_ENV`]. +pub const HEALTH_LISTEN_ADDR_ENV: &str = "HEALTH_LISTEN_ADDR"; +/// Platform TCP port environment variable. Binds `0.0.0.0:$PORT` when set alone. +pub const HEALTH_LISTEN_PORT_ENV: &str = "PORT"; +/// Optional operational-store URL observed only for GET `/ready`. +pub const HEALTH_DATABASE_URL_ENV: &str = "DATABASE_URL"; +/// Optional caller-measured backlog label. Missing means unknown and not ready. +pub const HEALTH_BACKLOG_HEALTH_ENV: &str = "HEALTH_BACKLOG_HEALTH"; + +/// Fail-closed configuration error for the health-probe process. +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +#[non_exhaustive] +pub enum HealthProcessConfigError { + /// Neither `HEALTH_LISTEN_ADDR` nor `PORT` was set. + MissingListenAddress, + /// `HEALTH_LISTEN_ADDR` was blank, padded, or not a socket address. + InvalidListenAddress, + /// `PORT` was blank, padded, or not a TCP port. + InvalidListenPort, + /// `DATABASE_URL` was set but was not an unpadded postgres URL or libpq string. + InvalidDatabaseUrl, + /// `HEALTH_BACKLOG_HEALTH` was set to an unknown label. + InvalidBacklogHealth, +} + +impl Display for HealthProcessConfigError { + fn fmt(&self, formatter: &mut Formatter<'_>) -> std::fmt::Result { + formatter.write_str(match self { + Self::MissingListenAddress => { + "set HEALTH_LISTEN_ADDR to a socket address, or set PORT to a TCP port" + } + Self::InvalidListenAddress => { + "HEALTH_LISTEN_ADDR must be an unpadded host:port socket address" + } + Self::InvalidListenPort => "PORT must be an unpadded TCP port from 0 to 65535", + Self::InvalidDatabaseUrl => { + "DATABASE_URL must be an unpadded postgres URL or libpq keyword/value string when set" + } + Self::InvalidBacklogHealth => { + "HEALTH_BACKLOG_HEALTH must be within_bounds, stalled, or unknown when set" + } + }) + } +} + +impl Error for HealthProcessConfigError {} + +/// Validated listen and optional store configuration for one health-probe process. +pub struct HealthProcessConfig { + listen_addr: SocketAddr, + database_url: Option, + connect_config: Option, + backlog_health: BacklogHealth, +} + +impl HealthProcessConfig { + /// Return the address the process will bind. + #[must_use] + pub const fn listen_addr(&self) -> SocketAddr { + self.listen_addr + } + + /// Return the configured store URL without printing it. + #[must_use] + pub fn database_url(&self) -> Option<&str> { + self.database_url.as_deref() + } + + /// Return the caller-supplied backlog evidence used for readiness. + #[must_use] + pub const fn backlog_health(&self) -> BacklogHealth { + self.backlog_health + } +} + +impl std::fmt::Debug for HealthProcessConfig { + fn fmt(&self, formatter: &mut Formatter<'_>) -> std::fmt::Result { + formatter + .debug_struct("HealthProcessConfig") + .field("listen_addr", &self.listen_addr) + .field("database_url_present", &self.database_url.is_some()) + .field("backlog_health", &self.backlog_health) + .finish_non_exhaustive() + } +} + +/// Parse listen, store, and backlog environment values without starting I/O. +/// +/// Unknown, blank, or whitespace-padded values fail closed so a load balancer +/// cannot point at a process that guessed its listen target or store URL. +/// +/// # Errors +/// +/// Returns [`HealthProcessConfigError`] when required listen input is missing +/// or any provided value is padded, empty, or semantically unknown. +pub fn parse_health_process_config( + getenv: F, +) -> Result +where + F: Fn(&str) -> Option, +{ + let listen_addr = parse_listen_addr( + getenv(HEALTH_LISTEN_ADDR_ENV), + getenv(HEALTH_LISTEN_PORT_ENV), + )?; + let (database_url, connect_config) = parse_database_url(getenv(HEALTH_DATABASE_URL_ENV))?; + let backlog_health = parse_backlog_health(getenv(HEALTH_BACKLOG_HEALTH_ENV))?; + Ok(HealthProcessConfig { + listen_addr, + database_url, + connect_config, + backlog_health, + }) +} + +/// Bind the configured listen address for operator probes. +/// +/// # Errors +/// +/// Returns the I/O error if the operating system cannot bind the address. +pub fn bind_health_process(config: &HealthProcessConfig) -> io::Result { + bind_health_http(config.listen_addr()) +} + +/// Serve GET `/live` and GET `/ready` until `accept` fails. +/// +/// GET `/live` uses a process-liveness snapshot and never opens a store +/// connection. GET `/ready` connects only when `DATABASE_URL` is configured. +/// Connect or probe failure becomes HTTP 503 without driver text. +/// +/// # Errors +/// +/// Returns the I/O error that stopped the accept loop. +pub fn serve_health_process( + listener: &TcpListener, + config: &HealthProcessConfig, +) -> io::Result<()> { + serve_health_http_with(listener, |request| { + answer_health_process_request(request, config) + }) +} + +/// Run the process from environment values: parse, bind, then serve. +/// +/// # Errors +/// +/// Returns a configuration error or the listen/serve I/O error. +pub fn run_health_process(getenv: F) -> Result<(), HealthProcessRunError> +where + F: Fn(&str) -> Option, +{ + let config = + parse_health_process_config(getenv).map_err(HealthProcessRunError::InvalidConfig)?; + let listener = bind_health_process(&config).map_err(HealthProcessRunError::Listen)?; + serve_bound_health_process(&listener, &config) +} + +fn serve_bound_health_process( + listener: &std::net::TcpListener, + config: &HealthProcessConfig, +) -> Result<(), HealthProcessRunError> { + serve_health_process(listener, config).map_err(HealthProcessRunError::Listen) +} + +/// Runtime failure after configuration has been parsed. +#[derive(Debug)] +pub enum HealthProcessRunError { + /// Environment values were missing or unknown. + InvalidConfig(HealthProcessConfigError), + /// Binding or serving the probe listener failed. + Listen(io::Error), +} + +impl Display for HealthProcessRunError { + fn fmt(&self, formatter: &mut Formatter<'_>) -> std::fmt::Result { + match self { + Self::InvalidConfig(error) => Display::fmt(error, formatter), + Self::Listen(error) => write!( + formatter, + "bind HEALTH_LISTEN_ADDR or PORT and keep the process running: {error}" + ), + } + } +} + +impl Error for HealthProcessRunError { + fn source(&self) -> Option<&(dyn Error + 'static)> { + match self { + Self::InvalidConfig(error) => Some(error), + Self::Listen(error) => Some(error), + } + } +} + +fn parse_listen_addr( + listen_addr: Option, + port: Option, +) -> Result { + match exact_env_value(listen_addr) { + Err(()) => return Err(HealthProcessConfigError::InvalidListenAddress), + Ok(Some(value)) => { + return value + .parse() + .map_err(|_| HealthProcessConfigError::InvalidListenAddress); + } + Ok(None) => {} + } + match exact_env_value(port) { + Err(()) => Err(HealthProcessConfigError::InvalidListenPort), + Ok(None) => Err(HealthProcessConfigError::MissingListenAddress), + Ok(Some(value)) => { + let port = value + .parse::() + .map_err(|_| HealthProcessConfigError::InvalidListenPort)?; + Ok(SocketAddr::new(IpAddr::V4(Ipv4Addr::UNSPECIFIED), port)) + } + } +} + +fn parse_database_url( + raw: Option, +) -> Result<(Option, Option), HealthProcessConfigError> { + match exact_env_value(raw) { + Err(()) => Err(HealthProcessConfigError::InvalidDatabaseUrl), + Ok(None) => Ok((None, None)), + Ok(Some(value)) => { + let config = postgres::Config::from_str(&value) + .map_err(|_| HealthProcessConfigError::InvalidDatabaseUrl)?; + Ok((Some(value), Some(config))) + } + } +} + +fn parse_backlog_health(raw: Option) -> Result { + match exact_env_value(raw) { + Err(()) => Err(HealthProcessConfigError::InvalidBacklogHealth), + Ok(None) => Ok(BacklogHealth::Unknown), + Ok(Some(value)) => match value.as_str() { + "within_bounds" => Ok(BacklogHealth::WithinBounds), + "stalled" => Ok(BacklogHealth::Stalled), + "unknown" => Ok(BacklogHealth::Unknown), + _ => Err(HealthProcessConfigError::InvalidBacklogHealth), + }, + } +} + +fn exact_env_value(raw: Option) -> Result, ()> { + match raw { + None => Ok(None), + Some(value) if value.is_empty() || value.trim() != value.as_str() => Err(()), + Some(value) => Ok(Some(value)), + } +} + +fn answer_health_process_request( + request: &str, + config: &HealthProcessConfig, +) -> HealthHttpResponse { + if !health_request_requires_readiness_snapshot(request) { + return handle_health_http_request(request, &process_liveness_snapshot()); + } + match config.connect_config.as_ref() { + None => handle_health_http_request(request, &process_liveness_snapshot()), + Some(connect_config) => match connect_operational_store(connect_config) { + Some(mut client) => handle_postgres_health_http_request( + request, + &mut client, + &[], + config.backlog_health(), + ), + None => health_ready_response( + &unavailable_store_snapshot(config.backlog_health()), + &ready_required_capabilities(request), + ), + }, + } +} + +fn connect_operational_store(connect_config: &postgres::Config) -> Option { + let mut connect_config = connect_config.clone(); + connect_config.connect_timeout(HEALTH_HTTP_IO_TIMEOUT); + connect_config.connect(NoTls).ok() +} + +fn ready_required_capabilities(request: &str) -> Vec<&str> { + let named = health_request_required_capabilities(request); + if named.is_empty() { + vec![POSTGRES_OPERATIONAL_STORE_CAPABILITY_REF] + } else { + named + } +} + +fn process_liveness_snapshot() -> RuntimeHealthSnapshot { + RuntimeHealthSnapshot::new( + true, + BacklogHealth::Unknown, + DataIntegrityHealth::Unknown, + Vec::new(), + ) + .expect("empty process-liveness snapshot is valid") +} + +fn unavailable_store_snapshot(backlog_health: BacklogHealth) -> RuntimeHealthSnapshot { + let capability = CapabilityHealth::new( + POSTGRES_OPERATIONAL_STORE_CAPABILITY_REF, + CapabilityState::Unknown, + false, + ) + .expect("repository-owned postgres capability reference must remain valid"); + RuntimeHealthSnapshot::new( + true, + backlog_health, + DataIntegrityHealth::Unknown, + vec![capability], + ) + .expect("unavailable store snapshot contains one unique capability") +} + +#[cfg(test)] +mod tests { + use super::{ + bind_health_process, exact_env_value, parse_health_process_config, + ready_required_capabilities, run_health_process, serve_bound_health_process, + HealthProcessConfigError, HealthProcessRunError, HEALTH_LISTEN_ADDR_ENV, + }; + use crate::postgres_health::POSTGRES_OPERATIONAL_STORE_CAPABILITY_REF; + use std::error::Error; + use std::io; + + #[test] + fn config_debug_redacts_the_database_url() { + let config = parse_health_process_config(|key| match key { + HEALTH_LISTEN_ADDR_ENV => Some("127.0.0.1:0".to_owned()), + "DATABASE_URL" => Some("postgres://operator:secret@db/product".to_owned()), + _ => None, + }) + .unwrap(); + let rendered = format!("{config:?}"); + assert!(rendered.contains("database_url_present: true")); + assert!(!rendered.contains("secret")); + assert!(!rendered.contains("operator")); + } + + #[test] + fn exact_env_value_treats_absence_as_optional() { + assert_eq!(exact_env_value(None), Ok(None)); + assert_eq!( + exact_env_value(Some("ready".to_owned())), + Ok(Some("ready".to_owned())) + ); + assert_eq!(exact_env_value(Some(String::new())), Err(())); + assert_eq!(exact_env_value(Some(" padded".to_owned())), Err(())); + } + + #[test] + fn ready_required_capabilities_default_to_the_operational_store() { + assert_eq!( + ready_required_capabilities("GET /ready HTTP/1.1\r\n\r\n"), + vec![POSTGRES_OPERATIONAL_STORE_CAPABILITY_REF] + ); + assert_eq!( + ready_required_capabilities("GET /ready?capability=scoring HTTP/1.1\r\n\r\n"), + vec!["scoring"] + ); + } + + #[test] + fn run_and_exit_code_fail_closed_for_missing_listen_config() { + let error = run_health_process(|_| None).expect_err("missing listen config must not run"); + assert!(matches!( + error, + HealthProcessRunError::InvalidConfig(HealthProcessConfigError::MissingListenAddress) + )); + assert!(error.to_string().contains("HEALTH_LISTEN_ADDR")); + assert!(error.source().is_some()); + } + + #[test] + fn serve_bound_process_maps_accept_failure_to_a_listen_error() { + let config = parse_health_process_config(|key| match key { + HEALTH_LISTEN_ADDR_ENV => Some("127.0.0.1:0".to_owned()), + _ => None, + }) + .unwrap(); + let listener = bind_health_process(&config).unwrap(); + listener + .set_nonblocking(true) + .expect("the test must force accept to return WouldBlock"); + let error = serve_bound_health_process(&listener, &config) + .expect_err("a non-blocking accept must become a listen error"); + assert!(matches!(error, HealthProcessRunError::Listen(_))); + assert!(error.to_string().contains("HEALTH_LISTEN_ADDR")); + assert!(error.source().is_some()); + } + + #[test] + fn run_error_display_covers_listen_failures() { + let error = HealthProcessRunError::Listen(io::Error::new( + io::ErrorKind::AddrInUse, + "address already in use", + )); + assert!(error.to_string().contains("HEALTH_LISTEN_ADDR")); + assert!(error.to_string().contains("address already in use")); + assert!(error.source().is_some()); + } +} diff --git a/src/lib.rs b/src/lib.rs index aa17f0a3..22acfed4 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -16,6 +16,8 @@ pub mod data_rights; pub mod data_rights_authorization; pub mod deterministic_narrative; pub mod health; +pub mod health_http; +pub mod health_process; pub mod instrument; pub mod integration; pub mod item_delivery; @@ -25,6 +27,7 @@ pub mod postgres_consent; pub mod postgres_data_rights; pub mod postgres_data_rights_processing; pub mod postgres_health; +pub mod postgres_health_http; pub mod postgres_inbox_consumption; pub mod postgres_instrument_release; pub mod postgres_integration; diff --git a/src/postgres_health.rs b/src/postgres_health.rs index 0abe56c9..69e79cb2 100644 --- a/src/postgres_health.rs +++ b/src/postgres_health.rs @@ -5,7 +5,10 @@ //! can safely accept product-owned state changes. It does not own credentials, //! connection pooling, migrations, backup, or recovery. -use crate::health::{CapabilityHealth, CapabilityState, DataIntegrityHealth, HealthContractError}; +use crate::health::{ + BacklogHealth, CapabilityHealth, CapabilityState, DataIntegrityHealth, HealthContractError, + RuntimeHealthSnapshot, +}; use postgres::GenericClient; /// Initial supported `PostgreSQL` server major version from ADR-0015. @@ -158,3 +161,50 @@ pub fn probe_postgres_relation_integrity( } Ok(DataIntegrityHealth::Verified) } + +/// Compose live `PostgreSQL` probes into one operation-scoped runtime snapshot. +/// +/// The process is treated as live because this observation is executing. Probe +/// failures become unknown capability or integrity evidence and fail readiness +/// closed. The caller supplies backlog health so this function does not invent +/// a measured threshold. Raw driver errors are not returned. +/// +/// # Panics +/// +/// Panics only if the crate-owned `postgres_operational_store` capability +/// reference is no longer a valid health identity. +#[must_use] +pub fn observe_postgres_operational_snapshot( + client: &mut impl GenericClient, + required_relations: &[&str], + backlog_health: BacklogHealth, +) -> RuntimeHealthSnapshot { + match probe_postgres_runtime(client) { + Ok(runtime) => { + let integrity = probe_postgres_relation_integrity(client, required_relations) + .unwrap_or(DataIntegrityHealth::Unknown); + let capability = runtime + .capability_health() + .expect("repository-owned postgres capability reference must remain valid"); + RuntimeHealthSnapshot::new(true, backlog_health, integrity, vec![capability]) + .expect("postgres snapshot contains one unique capability") + } + Err(_) => unknown_postgres_snapshot(backlog_health), + } +} + +fn unknown_postgres_snapshot(backlog_health: BacklogHealth) -> RuntimeHealthSnapshot { + let capability = CapabilityHealth::new( + POSTGRES_OPERATIONAL_STORE_CAPABILITY_REF, + CapabilityState::Unknown, + false, + ) + .expect("repository-owned postgres capability reference must remain valid"); + RuntimeHealthSnapshot::new( + true, + backlog_health, + DataIntegrityHealth::Unknown, + vec![capability], + ) + .expect("unknown postgres snapshot contains one unique capability") +} diff --git a/src/postgres_health_http.rs b/src/postgres_health_http.rs new file mode 100644 index 00000000..9f2e7420 --- /dev/null +++ b/src/postgres_health_http.rs @@ -0,0 +1,128 @@ +//! Compose live `PostgreSQL` operational snapshots into operator HTTP probes. +//! +//! GET `/live` answers process liveness without store I/O. GET `/ready` +//! observes the caller-owned connection after the request is accepted and +//! reuses the transport-neutral health HTTP translator. The adapter does not +//! invent backlog thresholds or expose raw driver errors. + +use crate::health::{BacklogHealth, DataIntegrityHealth, RuntimeHealthSnapshot}; +use crate::health_http::{ + accept_one_health_http_with, handle_health_http_request, health_ready_response, + health_request_required_capabilities, health_request_requires_readiness_snapshot, + serve_health_http_with, HealthHttpResponse, +}; +use crate::postgres_health::{ + observe_postgres_operational_snapshot, POSTGRES_OPERATIONAL_STORE_CAPABILITY_REF, +}; +use postgres::GenericClient; +use std::io; +use std::net::TcpListener; + +/// Observe `PostgreSQL` only when answering GET `/ready`. +/// +/// GET `/live` and RFC 9457 problem responses use a process-liveness snapshot +/// and never touch the caller-owned connection. Bare GET `/ready` requires +/// `postgres_operational_store` so a read-only or unsupported store cannot +/// advertise readiness to a load balancer that omits `capability=`. +#[must_use] +pub fn handle_postgres_health_http_request( + request: &str, + client: &mut impl GenericClient, + required_relations: &[&str], + backlog_health: BacklogHealth, +) -> HealthHttpResponse { + if !health_request_requires_readiness_snapshot(request) { + return handle_health_http_request(request, &process_liveness_snapshot()); + } + let snapshot = + observe_postgres_operational_snapshot(client, required_relations, backlog_health); + let required = + postgres_ready_required_capabilities(health_request_required_capabilities(request)); + health_ready_response(&snapshot, &required) +} + +fn postgres_ready_required_capabilities(named: Vec<&str>) -> Vec<&str> { + if named.is_empty() { + vec![POSTGRES_OPERATIONAL_STORE_CAPABILITY_REF] + } else { + named + } +} + +/// Accept one TCP connection, then observe `PostgreSQL` only if the request is GET `/ready`. +/// +/// # Errors +/// +/// Returns the I/O error if accept, read, or write fails. +pub fn accept_one_postgres_health_http( + listener: &TcpListener, + client: &mut impl GenericClient, + required_relations: &[&str], + backlog_health: BacklogHealth, +) -> io::Result<()> { + accept_one_health_http_with(listener, |request| { + handle_postgres_health_http_request(request, client, required_relations, backlog_health) + }) +} + +/// Serve PostgreSQL-backed probes until `accept` fails. +/// +/// GET `/live` still answers without store I/O. GET `/ready` observes the +/// caller-owned connection after each accept. A dropped probe connection does +/// not stop later probes. +/// +/// # Errors +/// +/// Returns the I/O error that stopped the loop. +pub fn serve_postgres_health_http( + listener: &TcpListener, + client: &mut impl GenericClient, + required_relations: &[&str], + backlog_health: BacklogHealth, +) -> io::Result<()> { + serve_health_http_with(listener, |request| { + handle_postgres_health_http_request(request, client, required_relations, backlog_health) + }) +} + +fn process_liveness_snapshot() -> RuntimeHealthSnapshot { + RuntimeHealthSnapshot::new( + true, + BacklogHealth::Unknown, + DataIntegrityHealth::Unknown, + Vec::new(), + ) + .expect("empty process-liveness snapshot is valid") +} + +#[cfg(test)] +mod tests { + use super::{postgres_ready_required_capabilities, process_liveness_snapshot}; + use crate::health::{BacklogHealth, DataIntegrityHealth}; + use crate::postgres_health::POSTGRES_OPERATIONAL_STORE_CAPABILITY_REF; + + #[test] + fn process_liveness_snapshot_does_not_claim_readiness() { + let snapshot = process_liveness_snapshot(); + assert!(snapshot.is_live()); + assert_eq!(snapshot.backlog_health(), BacklogHealth::Unknown); + assert_eq!( + snapshot.data_integrity_health(), + DataIntegrityHealth::Unknown + ); + assert!(snapshot.capabilities().is_empty()); + assert!(!snapshot.is_ready_for(&[])); + } + + #[test] + fn bare_ready_requires_the_postgres_operational_store() { + assert_eq!( + postgres_ready_required_capabilities(Vec::new()), + vec![POSTGRES_OPERATIONAL_STORE_CAPABILITY_REF] + ); + assert_eq!( + postgres_ready_required_capabilities(vec!["scoring"]), + vec!["scoring"] + ); + } +} diff --git a/tests/health_http_listener_contract.rs b/tests/health_http_listener_contract.rs new file mode 100644 index 00000000..35998b43 --- /dev/null +++ b/tests/health_http_listener_contract.rs @@ -0,0 +1,266 @@ +//! Bound TCP listener contract for operator health probes. +//! +//! The listener is a transport around the existing request translator. It does +//! not add public/admin product routes or invent availability SLOs. + +use psychometrics_commons_runtime::health::{ + BacklogHealth, CapabilityHealth, CapabilityState, DataIntegrityHealth, RuntimeHealthSnapshot, +}; +use psychometrics_commons_runtime::health_http::{ + accept_one_health_http, bind_health_http, serve_health_http, HEALTH_LIVE_PATH, + HEALTH_READY_PATH, +}; +use std::io::{Read, Write}; +use std::net::{IpAddr, Ipv4Addr, SocketAddr, TcpStream}; +use std::sync::mpsc; +use std::thread; +use std::time::Duration; + +fn healthy_snapshot() -> RuntimeHealthSnapshot { + RuntimeHealthSnapshot::new( + true, + BacklogHealth::WithinBounds, + DataIntegrityHealth::Verified, + vec![CapabilityHealth::new("scoring", CapabilityState::Available, true).unwrap()], + ) + .unwrap() +} + +fn exchange(addr: SocketAddr, request: &str) -> String { + let mut stream = TcpStream::connect_timeout(&addr, Duration::from_secs(2)).unwrap(); + stream + .set_read_timeout(Some(Duration::from_secs(2))) + .unwrap(); + stream.write_all(request.as_bytes()).unwrap(); + stream.shutdown(std::net::Shutdown::Write).unwrap(); + let mut body = String::new(); + stream.read_to_string(&mut body).unwrap(); + body +} + +#[test] +fn bound_listener_serves_live_and_ready_probes_over_tcp() { + let listener = bind_health_http(SocketAddr::new(IpAddr::V4(Ipv4Addr::LOCALHOST), 0)).unwrap(); + let addr = listener.local_addr().unwrap(); + let snapshot = healthy_snapshot(); + let server = thread::spawn(move || { + accept_one_health_http(&listener, &snapshot).unwrap(); + accept_one_health_http(&listener, &snapshot).unwrap(); + accept_one_health_http(&listener, &snapshot).unwrap(); + }); + + let live = exchange( + addr, + &format!("GET {HEALTH_LIVE_PATH} HTTP/1.1\r\nHost: localhost\r\n\r\n"), + ); + assert!(live.starts_with("HTTP/1.1 200 OK\r\n")); + assert!(live.contains("Content-Type: application/json\r\n")); + assert!(live.contains("Connection: close\r\n")); + assert!(live.contains("\"live\":true")); + + let ready = exchange( + addr, + &format!("GET {HEALTH_READY_PATH}?capability=scoring HTTP/1.1\r\nHost: localhost\r\n\r\n"), + ); + assert!(ready.starts_with("HTTP/1.1 200 OK\r\n")); + assert!(ready.contains("\"ready\":true")); + + let unready = exchange( + addr, + &format!( + "GET {HEALTH_READY_PATH}?capability=unregistered_capability HTTP/1.1\r\nHost: localhost\r\n\r\n" + ), + ); + assert!(unready.starts_with("HTTP/1.1 503 Service Unavailable\r\n")); + assert!(unready.contains("Content-Type: application/json\r\n")); + assert!(unready.contains("\"ready\":false")); + + server.join().unwrap(); +} + +#[test] +fn bound_listener_returns_problem_details_for_unsupported_methods() { + let listener = bind_health_http(SocketAddr::new(IpAddr::V4(Ipv4Addr::LOCALHOST), 0)).unwrap(); + let addr = listener.local_addr().unwrap(); + let snapshot = healthy_snapshot(); + let server = thread::spawn(move || { + accept_one_health_http(&listener, &snapshot).unwrap(); + }); + + let response = exchange( + addr, + &format!("POST {HEALTH_LIVE_PATH} HTTP/1.1\r\nHost: localhost\r\n\r\n"), + ); + assert!(response.starts_with("HTTP/1.1 405 Method Not Allowed\r\n")); + assert!(response.contains("Content-Type: application/problem+json\r\n")); + assert!(response.contains("Allow: GET\r\n")); + assert!(response.contains("Cache-Control: no-store\r\n")); + assert!(response.contains("\"title\":\"Method Not Allowed\"")); + assert!(!response.contains("postgres")); + + server.join().unwrap(); +} + +#[test] +fn bound_listener_fails_closed_for_unknown_paths_and_truncated_requests() { + let listener = bind_health_http(SocketAddr::new(IpAddr::V4(Ipv4Addr::LOCALHOST), 0)).unwrap(); + let addr = listener.local_addr().unwrap(); + let snapshot = healthy_snapshot(); + let server = thread::spawn(move || { + accept_one_health_http(&listener, &snapshot).unwrap(); + accept_one_health_http(&listener, &snapshot).unwrap(); + }); + + let missing = exchange(addr, "GET /v1/sessions HTTP/1.1\r\nHost: localhost\r\n\r\n"); + assert!(missing.starts_with("HTTP/1.1 404 Not Found\r\n")); + assert!(missing.contains("Content-Type: application/problem+json\r\n")); + assert!(!missing.contains("/v1/instruments")); + + let truncated = exchange(addr, "GET /live"); + assert!(truncated.starts_with("HTTP/1.1 400 Bad Request\r\n")); + assert!(truncated.contains("Content-Type: application/problem+json\r\n")); + assert!(!truncated.contains("GET /live HTTP")); + + server.join().unwrap(); +} + +#[test] +fn bound_listener_rejects_an_oversized_request_without_echoing_it() { + let listener = bind_health_http(SocketAddr::new(IpAddr::V4(Ipv4Addr::LOCALHOST), 0)).unwrap(); + let addr = listener.local_addr().unwrap(); + let snapshot = healthy_snapshot(); + let server = thread::spawn(move || { + accept_one_health_http(&listener, &snapshot).unwrap(); + }); + + let oversized = format!( + "GET /live HTTP/1.1\r\nHost: localhost\r\nX-Pad: {}\r\n\r\n", + "A".repeat(9_000) + ); + let mut stream = TcpStream::connect_timeout(&addr, Duration::from_secs(2)).unwrap(); + stream + .set_read_timeout(Some(Duration::from_secs(2))) + .unwrap(); + let _ = stream.write_all(oversized.as_bytes()); + let _ = stream.shutdown(std::net::Shutdown::Write); + let mut response = String::new(); + match stream.read_to_string(&mut response) { + Ok(_) => {} + Err(error) if error.kind() == std::io::ErrorKind::ConnectionReset => {} + Err(error) => panic!("unexpected oversized-request read error: {error}"), + } + assert!( + response.starts_with("HTTP/1.1 400 Bad Request\r\n"), + "{response}" + ); + assert!(!response.contains(&"A".repeat(32))); + + server.join().unwrap(); +} + +#[test] +fn bound_listener_fails_closed_when_the_client_never_finishes_the_request() { + let listener = bind_health_http(SocketAddr::new(IpAddr::V4(Ipv4Addr::LOCALHOST), 0)).unwrap(); + let addr = listener.local_addr().unwrap(); + let snapshot = healthy_snapshot(); + let server = thread::spawn(move || { + accept_one_health_http(&listener, &snapshot).unwrap(); + }); + + let mut stream = TcpStream::connect_timeout(&addr, Duration::from_secs(2)).unwrap(); + stream + .set_read_timeout(Some(Duration::from_secs(4))) + .unwrap(); + stream.write_all(b"GET /live HTTP/1.1\r\nHost: ").unwrap(); + let mut response = String::new(); + stream + .read_to_string(&mut response) + .expect("the probe listener must answer an incomplete request instead of hanging"); + assert!( + response.starts_with("HTTP/1.1 400 Bad Request\r\n"), + "{response}" + ); + assert!(response.contains("Cache-Control: no-store\r\n")); + assert!(!response.contains("GET /live HTTP")); + + server.join().unwrap(); +} + +#[test] +fn serve_loop_answers_successive_probes_until_accept_fails() { + let listener = bind_health_http(SocketAddr::new(IpAddr::V4(Ipv4Addr::LOCALHOST), 0)).unwrap(); + let addr = listener.local_addr().unwrap(); + let snapshot = healthy_snapshot(); + let stop = listener.try_clone().unwrap(); + let (done_tx, done_rx) = mpsc::channel(); + let server = thread::spawn(move || { + let result = serve_health_http(&listener, &snapshot); + let _ = done_tx.send(result); + }); + + let live = exchange( + addr, + &format!("GET {HEALTH_LIVE_PATH} HTTP/1.1\r\nHost: localhost\r\n\r\n"), + ); + assert!(live.starts_with("HTTP/1.1 200 OK\r\n")); + assert!(live.contains("\"live\":true")); + + let ready = exchange( + addr, + &format!("GET {HEALTH_READY_PATH}?capability=scoring HTTP/1.1\r\nHost: localhost\r\n\r\n"), + ); + assert!(ready.starts_with("HTTP/1.1 200 OK\r\n")); + assert!(ready.contains("\"ready\":true")); + + stop.set_nonblocking(true) + .expect("the test must be able to stop the shared listener"); + let _ = TcpStream::connect_timeout(&addr, Duration::from_millis(200)); + let stopped = done_rx + .recv_timeout(Duration::from_secs(3)) + .expect("serve_health_http must return after accept can no longer block"); + assert!( + stopped.is_err(), + "serve_health_http must surface the accept failure instead of hanging" + ); + server.join().unwrap(); +} + +#[test] +fn serve_loop_keeps_accepting_after_a_client_drops_the_connection() { + let listener = bind_health_http(SocketAddr::new(IpAddr::V4(Ipv4Addr::LOCALHOST), 0)).unwrap(); + let addr = listener.local_addr().unwrap(); + let snapshot = healthy_snapshot(); + let stop = listener.try_clone().unwrap(); + let (done_tx, done_rx) = mpsc::channel(); + let server = thread::spawn(move || { + let result = serve_health_http(&listener, &snapshot); + let _ = done_tx.send(result); + }); + + { + let mut dropped = TcpStream::connect_timeout(&addr, Duration::from_secs(2)).unwrap(); + dropped + .write_all( + format!("GET {HEALTH_LIVE_PATH} HTTP/1.1\r\nHost: localhost\r\n\r\n").as_bytes(), + ) + .unwrap(); + } + + let live = exchange( + addr, + &format!("GET {HEALTH_LIVE_PATH} HTTP/1.1\r\nHost: localhost\r\n\r\n"), + ); + assert!( + live.starts_with("HTTP/1.1 200 OK\r\n"), + "a dropped probe must not stop later GET /live answers: {live}" + ); + assert!(live.contains("\"live\":true")); + + stop.set_nonblocking(true) + .expect("the test must be able to stop the shared listener"); + let _ = TcpStream::connect_timeout(&addr, Duration::from_millis(200)); + let _ = done_rx + .recv_timeout(Duration::from_secs(3)) + .expect("serve_health_http must still return after accept can no longer block"); + server.join().unwrap(); +} diff --git a/tests/health_http_probe_contract.rs b/tests/health_http_probe_contract.rs new file mode 100644 index 00000000..eaf76c76 --- /dev/null +++ b/tests/health_http_probe_contract.rs @@ -0,0 +1,170 @@ +//! Contract tests for operator liveness and readiness HTTP probes. +//! +//! These probes are the first implemented HTTP surface. They expose the existing +//! domain health snapshot without inventing SLO values or leaking raw store errors. + +use psychometrics_commons_runtime::health::{ + BacklogHealth, CapabilityHealth, CapabilityState, DataIntegrityHealth, RuntimeHealthSnapshot, +}; +use psychometrics_commons_runtime::health_http::{ + handle_health_http_request, HEALTH_LIVE_PATH, HEALTH_READY_PATH, +}; +use std::fs; +use std::path::PathBuf; + +fn healthy_snapshot() -> RuntimeHealthSnapshot { + RuntimeHealthSnapshot::new( + true, + BacklogHealth::WithinBounds, + DataIntegrityHealth::Verified, + vec![ + CapabilityHealth::new("scoring", CapabilityState::Available, true).unwrap(), + CapabilityHealth::new("authenticated_linking", CapabilityState::Unavailable, false) + .unwrap(), + ], + ) + .unwrap() +} + +fn request(method: &str, target: &str) -> String { + format!("{method} {target} HTTP/1.1\r\nHost: localhost\r\n\r\n") +} + +#[test] +fn liveness_probe_is_independent_from_operation_readiness() { + let live = handle_health_http_request(&request("GET", HEALTH_LIVE_PATH), &healthy_snapshot()); + assert_eq!(live.status(), 200); + assert_eq!(live.content_type(), "application/json"); + assert!(live.body().contains("\"live\":true")); + assert!(live.body().contains("\"ready\":true")); + + let not_live = RuntimeHealthSnapshot::new( + false, + BacklogHealth::WithinBounds, + DataIntegrityHealth::Verified, + vec![], + ) + .unwrap(); + let response = handle_health_http_request(&request("GET", HEALTH_LIVE_PATH), ¬_live); + assert_eq!(response.status(), 503); + assert!(response.body().contains("\"live\":false")); + assert!(response.body().contains("\"ready\":false")); +} + +#[test] +fn readiness_probe_fails_closed_for_named_or_unknown_required_capabilities() { + let snapshot = healthy_snapshot(); + let ready = handle_health_http_request(&request("GET", HEALTH_READY_PATH), &snapshot); + assert_eq!(ready.status(), 200); + assert!(ready.body().contains("\"ready\":true")); + + let scoring = + handle_health_http_request(&request("GET", "/ready?capability=scoring"), &snapshot); + assert_eq!(scoring.status(), 200); + + let linking = handle_health_http_request( + &request("GET", "/ready?capability=authenticated_linking"), + &snapshot, + ); + assert_eq!(linking.status(), 503); + assert!(linking.body().contains("\"ready\":false")); + assert!(linking + .body() + .contains("\"capability_ref\":\"authenticated_linking\"")); + + let unknown = handle_health_http_request( + &request("GET", "/ready?capability=unregistered_capability"), + &snapshot, + ); + assert_eq!(unknown.status(), 503); +} + +#[test] +fn stalled_backlog_or_unknown_integrity_makes_readiness_unavailable() { + let stalled = RuntimeHealthSnapshot::new( + true, + BacklogHealth::Stalled, + DataIntegrityHealth::Verified, + vec![], + ) + .unwrap(); + let stalled_response = handle_health_http_request(&request("GET", HEALTH_READY_PATH), &stalled); + assert_eq!(stalled_response.status(), 503); + assert!(stalled_response + .body() + .contains("\"backlog_health\":\"stalled\"")); + + let unknown_integrity = RuntimeHealthSnapshot::new( + true, + BacklogHealth::WithinBounds, + DataIntegrityHealth::Unknown, + vec![], + ) + .unwrap(); + let integrity_response = + handle_health_http_request(&request("GET", HEALTH_READY_PATH), &unknown_integrity); + assert_eq!(integrity_response.status(), 503); + assert!(integrity_response + .body() + .contains("\"data_integrity_health\":\"unknown\"")); +} + +#[test] +fn unsupported_method_or_path_returns_safe_problem_details() { + let snapshot = healthy_snapshot(); + let not_allowed = handle_health_http_request(&request("POST", HEALTH_LIVE_PATH), &snapshot); + assert_eq!(not_allowed.status(), 405); + assert_eq!(not_allowed.content_type(), "application/problem+json"); + assert!(not_allowed + .body() + .contains("\"title\":\"Method Not Allowed\"")); + assert!(not_allowed + .body() + .contains("\"type\":\"urn:psychometrics-commons:problem:method-not-allowed\"")); + assert!(!not_allowed.body().contains("about:blank")); + assert!(!not_allowed.body().contains("postgres")); + assert!(!not_allowed.body().contains("sql")); + + let missing = handle_health_http_request(&request("GET", "/v1/sessions"), &snapshot); + assert_eq!(missing.status(), 404); + assert_eq!(missing.content_type(), "application/problem+json"); + assert!(missing.body().contains("\"title\":\"Not Found\"")); + assert!(missing + .body() + .contains("\"type\":\"urn:psychometrics-commons:problem:not-found\"")); + assert!(!missing.body().contains("/v1/instruments")); +} + +#[test] +fn malformed_request_fails_closed_without_echoing_raw_input() { + let response = handle_health_http_request("NOT-A-REQUEST", &healthy_snapshot()); + assert_eq!(response.status(), 400); + assert_eq!(response.content_type(), "application/problem+json"); + assert!(response.body().contains("\"title\":\"Bad Request\"")); + assert!(response + .body() + .contains("\"type\":\"urn:psychometrics-commons:problem:bad-request\"")); + assert!(!response.body().contains("NOT-A-REQUEST")); + assert!(!response.body().contains("about:blank")); +} + +#[test] +fn as_built_openapi_lists_only_implemented_health_probe_operations() { + let openapi_path = PathBuf::from(env!("CARGO_MANIFEST_DIR")).join("openapi/health-probes.yaml"); + let openapi = fs::read_to_string(openapi_path).expect("as-built health OpenAPI must exist"); + assert!(openapi.contains("openapi: 3.2.0")); + assert!(openapi.contains(HEALTH_LIVE_PATH)); + assert!(openapi.contains(HEALTH_READY_PATH)); + assert!(openapi.contains("urn:psychometrics-commons:problem:bad-request")); + assert!(!openapi.contains("about:blank")); + assert!(!openapi.contains("/v1/sessions")); + assert!(!openapi.contains("/v1/instruments")); + let live_section = openapi + .split("/ready:") + .next() + .expect("as-built OpenAPI must describe /live before /ready"); + assert!( + live_section.contains("\"503\""), + "GET /live must document HTTP 503 when the process is not live" + ); +} diff --git a/tests/health_process_contract.rs b/tests/health_process_contract.rs new file mode 100644 index 00000000..f109c1db --- /dev/null +++ b/tests/health_process_contract.rs @@ -0,0 +1,348 @@ +//! Process entrypoint contract for operator health probes. +//! +//! A buyer must be able to start one process from listen/store environment +//! variables, point a load balancer at GET `/live` and GET `/ready`, and keep +//! liveness free of store I/O when `PostgreSQL` is down. + +use psychometrics_commons_runtime::health::BacklogHealth; +use psychometrics_commons_runtime::health_http::{HEALTH_LIVE_PATH, HEALTH_READY_PATH}; +use psychometrics_commons_runtime::health_process::{ + bind_health_process, parse_health_process_config, run_health_process, serve_health_process, + HealthProcessConfigError, HealthProcessRunError, HEALTH_BACKLOG_HEALTH_ENV, + HEALTH_DATABASE_URL_ENV, HEALTH_LISTEN_ADDR_ENV, HEALTH_LISTEN_PORT_ENV, +}; +use std::io::{Read, Write}; +use std::net::{IpAddr, Ipv4Addr, SocketAddr, TcpStream}; +use std::sync::mpsc; +use std::thread; +use std::time::Duration; + +fn env_lookup(pairs: &[(&str, &str)]) -> impl Fn(&str) -> Option { + let owned: Vec<(String, String)> = pairs + .iter() + .map(|(key, value)| ((*key).to_owned(), (*value).to_owned())) + .collect(); + move |key| { + owned + .iter() + .find(|(name, _)| name == key) + .map(|(_, value)| value.clone()) + } +} + +fn exchange(addr: SocketAddr, request: &str) -> String { + let mut stream = TcpStream::connect_timeout(&addr, Duration::from_secs(2)).unwrap(); + stream + .set_read_timeout(Some(Duration::from_secs(2))) + .unwrap(); + stream.write_all(request.as_bytes()).unwrap(); + stream.shutdown(std::net::Shutdown::Write).unwrap(); + let mut body = String::new(); + stream.read_to_string(&mut body).unwrap(); + body +} + +#[test] +fn run_health_process_fails_closed_when_the_listen_address_cannot_bind() { + let error = run_health_process(env_lookup(&[(HEALTH_LISTEN_ADDR_ENV, "8.8.8.8:80")])) + .expect_err("binding a foreign address must not start a silent probe process"); + assert!(matches!(error, HealthProcessRunError::Listen(_))); + assert!(error.to_string().contains("HEALTH_LISTEN_ADDR")); +} + +#[test] +fn listen_config_fails_closed_when_no_address_or_port_is_set() { + let error = parse_health_process_config(|_| None) + .expect_err("a process without a listen target must not start"); + assert_eq!(error, HealthProcessConfigError::MissingListenAddress); + assert!(error.to_string().contains("HEALTH_LISTEN_ADDR")); + assert!(error.to_string().contains("PORT")); +} + +#[test] +fn listen_address_rejects_blank_padded_and_unparseable_values() { + for (env_key, value, expected) in [ + ( + HEALTH_LISTEN_ADDR_ENV, + "", + HealthProcessConfigError::InvalidListenAddress, + ), + ( + HEALTH_LISTEN_ADDR_ENV, + " 127.0.0.1:8080", + HealthProcessConfigError::InvalidListenAddress, + ), + ( + HEALTH_LISTEN_ADDR_ENV, + "not-a-socket", + HealthProcessConfigError::InvalidListenAddress, + ), + ( + HEALTH_LISTEN_PORT_ENV, + "", + HealthProcessConfigError::InvalidListenPort, + ), + ( + HEALTH_LISTEN_PORT_ENV, + " 8080", + HealthProcessConfigError::InvalidListenPort, + ), + ( + HEALTH_LISTEN_PORT_ENV, + "65536", + HealthProcessConfigError::InvalidListenPort, + ), + ( + HEALTH_LISTEN_PORT_ENV, + "abc", + HealthProcessConfigError::InvalidListenPort, + ), + ] { + let error = parse_health_process_config(env_lookup(&[(env_key, value)])) + .expect_err("invalid listen configuration must fail closed"); + assert_eq!(error, expected, "{env_key}={value:?}"); + assert!(!error.to_string().is_empty()); + } +} + +#[test] +fn explicit_listen_address_wins_over_platform_port() { + let config = parse_health_process_config(env_lookup(&[ + (HEALTH_LISTEN_ADDR_ENV, "127.0.0.1:0"), + (HEALTH_LISTEN_PORT_ENV, "8080"), + ])) + .expect("an explicit listen address must start the process"); + assert_eq!( + config.listen_addr(), + SocketAddr::new(IpAddr::V4(Ipv4Addr::LOCALHOST), 0) + ); + assert!(config.database_url().is_none()); + assert_eq!(config.backlog_health(), BacklogHealth::Unknown); +} + +#[test] +fn platform_port_binds_all_interfaces() { + let config = parse_health_process_config(env_lookup(&[(HEALTH_LISTEN_PORT_ENV, "8080")])) + .expect("PORT must be enough for a hosted process"); + assert_eq!( + config.listen_addr(), + SocketAddr::new(IpAddr::V4(Ipv4Addr::UNSPECIFIED), 8080) + ); + let ephemeral = parse_health_process_config(env_lookup(&[(HEALTH_LISTEN_PORT_ENV, "0")])) + .expect("PORT=0 must remain a valid ephemeral bind"); + assert_eq!( + ephemeral.listen_addr(), + SocketAddr::new(IpAddr::V4(Ipv4Addr::UNSPECIFIED), 0) + ); +} + +#[test] +fn database_url_and_backlog_fail_closed_on_unknown_semantics() { + for (pairs, expected) in [ + ( + &[ + (HEALTH_LISTEN_ADDR_ENV, "127.0.0.1:0"), + (HEALTH_DATABASE_URL_ENV, ""), + ] as &[(&str, &str)], + HealthProcessConfigError::InvalidDatabaseUrl, + ), + ( + &[ + (HEALTH_LISTEN_ADDR_ENV, "127.0.0.1:0"), + (HEALTH_DATABASE_URL_ENV, " postgres://localhost/db"), + ], + HealthProcessConfigError::InvalidDatabaseUrl, + ), + ( + &[ + (HEALTH_LISTEN_ADDR_ENV, "127.0.0.1:0"), + (HEALTH_DATABASE_URL_ENV, "https://example.test/db"), + ], + HealthProcessConfigError::InvalidDatabaseUrl, + ), + ( + &[ + (HEALTH_LISTEN_ADDR_ENV, "127.0.0.1:0"), + (HEALTH_DATABASE_URL_ENV, "postgres://localhost:65536/db"), + ], + HealthProcessConfigError::InvalidDatabaseUrl, + ), + ( + &[ + (HEALTH_LISTEN_ADDR_ENV, "127.0.0.1:0"), + (HEALTH_BACKLOG_HEALTH_ENV, "green"), + ], + HealthProcessConfigError::InvalidBacklogHealth, + ), + ( + &[ + (HEALTH_LISTEN_ADDR_ENV, "127.0.0.1:0"), + (HEALTH_BACKLOG_HEALTH_ENV, " within_bounds"), + ], + HealthProcessConfigError::InvalidBacklogHealth, + ), + ] { + let owned = pairs.to_vec(); + let error = parse_health_process_config(move |key| { + owned + .iter() + .find(|(name, _)| *name == key) + .map(|(_, value)| (*value).to_string()) + }) + .expect_err("unknown store or backlog semantics must not start the process"); + assert_eq!(error, expected); + assert!(!error.to_string().is_empty()); + } +} + +#[test] +fn postgres_urls_and_named_backlog_values_are_accepted() { + let postgres = parse_health_process_config(env_lookup(&[ + (HEALTH_LISTEN_ADDR_ENV, "127.0.0.1:0"), + ( + HEALTH_DATABASE_URL_ENV, + "postgres://operator:secret@db/product", + ), + (HEALTH_BACKLOG_HEALTH_ENV, "within_bounds"), + ])) + .unwrap(); + assert_eq!( + postgres.database_url(), + Some("postgres://operator:secret@db/product") + ); + assert_eq!(postgres.backlog_health(), BacklogHealth::WithinBounds); + + let postgresql = parse_health_process_config(env_lookup(&[ + (HEALTH_LISTEN_ADDR_ENV, "[::1]:9090"), + (HEALTH_DATABASE_URL_ENV, "postgresql://db/product"), + (HEALTH_BACKLOG_HEALTH_ENV, "stalled"), + ])) + .unwrap(); + assert_eq!(postgresql.database_url(), Some("postgresql://db/product")); + assert_eq!(postgresql.backlog_health(), BacklogHealth::Stalled); + assert_eq!( + postgresql.listen_addr(), + "[::1]:9090".parse::().unwrap() + ); + + let unknown = parse_health_process_config(env_lookup(&[ + (HEALTH_LISTEN_ADDR_ENV, "127.0.0.1:0"), + (HEALTH_BACKLOG_HEALTH_ENV, "unknown"), + ])) + .unwrap(); + assert_eq!(unknown.backlog_health(), BacklogHealth::Unknown); + + let libpq = parse_health_process_config(env_lookup(&[ + (HEALTH_LISTEN_ADDR_ENV, "127.0.0.1:0"), + ( + HEALTH_DATABASE_URL_ENV, + "host=localhost user=postgres password=ci_job_run dbname=psychometrics_commons_test", + ), + ])) + .expect("libpq keyword/value DATABASE_URL used by Runtime CI must start the process"); + assert_eq!( + libpq.database_url(), + Some("host=localhost user=postgres password=ci_job_run dbname=psychometrics_commons_test") + ); +} + +#[test] +fn process_without_database_answers_live_and_fails_ready_closed() { + let config = + parse_health_process_config(env_lookup(&[(HEALTH_LISTEN_ADDR_ENV, "127.0.0.1:0")])) + .unwrap(); + let listener = bind_health_process(&config).unwrap(); + let addr = listener.local_addr().unwrap(); + let stop = listener.try_clone().unwrap(); + let (done_tx, done_rx) = mpsc::channel(); + let server = thread::spawn(move || { + let result = serve_health_process(&listener, &config); + let _ = done_tx.send(result); + }); + + let live = exchange( + addr, + &format!("GET {HEALTH_LIVE_PATH} HTTP/1.1\r\nHost: localhost\r\n\r\n"), + ); + assert!(live.starts_with("HTTP/1.1 200 OK\r\n"), "{live}"); + assert!(live.contains("\"live\":true")); + assert!(!live.contains("postgres")); + assert!(!live.contains("DATABASE_URL")); + + let ready = exchange( + addr, + &format!("GET {HEALTH_READY_PATH} HTTP/1.1\r\nHost: localhost\r\n\r\n"), + ); + assert!( + ready.starts_with("HTTP/1.1 503 Service Unavailable\r\n"), + "a process without store evidence must not advertise readiness: {ready}" + ); + assert!(ready.contains("\"ready\":false")); + assert!(!ready.contains("postgres::")); + assert!(!ready.contains("DbError")); + + stop.set_nonblocking(true) + .expect("the test must be able to stop the process listener"); + let _ = TcpStream::connect_timeout(&addr, Duration::from_millis(200)); + let stopped = done_rx + .recv_timeout(Duration::from_secs(3)) + .expect("serve_health_process must return after accept can no longer block"); + assert!(stopped.is_err(), "{stopped:?}"); + server.join().unwrap(); +} + +#[test] +fn unreachable_database_keeps_liveness_and_hides_driver_errors() { + let config = parse_health_process_config(env_lookup(&[ + (HEALTH_LISTEN_ADDR_ENV, "127.0.0.1:0"), + (HEALTH_DATABASE_URL_ENV, "postgres://127.0.0.1:1/missing"), + (HEALTH_BACKLOG_HEALTH_ENV, "within_bounds"), + ])) + .unwrap(); + let listener = bind_health_process(&config).unwrap(); + let addr = listener.local_addr().unwrap(); + let stop = listener.try_clone().unwrap(); + let (done_tx, done_rx) = mpsc::channel(); + let server = thread::spawn(move || { + let result = serve_health_process(&listener, &config); + let _ = done_tx.send(result); + }); + + let live = exchange( + addr, + &format!("GET {HEALTH_LIVE_PATH} HTTP/1.1\r\nHost: localhost\r\n\r\n"), + ); + assert!( + live.starts_with("HTTP/1.1 200 OK\r\n"), + "a down store must not restart a live process: {live}" + ); + assert!(!live.contains("127.0.0.1")); + assert!(!live.contains("postgres::")); + + let ready = exchange( + addr, + &format!("GET {HEALTH_READY_PATH} HTTP/1.1\r\nHost: localhost\r\n\r\n"), + ); + assert!( + ready.starts_with("HTTP/1.1 503 Service Unavailable\r\n"), + "{ready}" + ); + assert!(!ready.contains("127.0.0.1:1")); + assert!(!ready.contains("DbError")); + assert!(!ready.contains("Connection refused")); + + let named = exchange( + addr, + &format!("GET {HEALTH_READY_PATH}?capability=scoring HTTP/1.1\r\nHost: localhost\r\n\r\n"), + ); + assert!( + named.starts_with("HTTP/1.1 503 Service Unavailable\r\n"), + "{named}" + ); + assert!(!named.contains("Connection refused")); + + stop.set_nonblocking(true).unwrap(); + let _ = TcpStream::connect_timeout(&addr, Duration::from_millis(200)); + let _ = done_rx.recv_timeout(Duration::from_secs(3)); + server.join().unwrap(); +} diff --git a/tests/postgres_health_http_contract.rs b/tests/postgres_health_http_contract.rs new file mode 100644 index 00000000..8896113c --- /dev/null +++ b/tests/postgres_health_http_contract.rs @@ -0,0 +1,290 @@ +//! Wire live `PostgreSQL` operational snapshots into operator HTTP probes. + +use postgres::{Client, NoTls}; +use psychometrics_commons_runtime::health::BacklogHealth; +use psychometrics_commons_runtime::health_http::bind_health_http; +use psychometrics_commons_runtime::health_http::{HEALTH_LIVE_PATH, HEALTH_READY_PATH}; +use psychometrics_commons_runtime::postgres_health::POSTGRES_OPERATIONAL_STORE_CAPABILITY_REF; +use psychometrics_commons_runtime::postgres_health_http::{ + accept_one_postgres_health_http, handle_postgres_health_http_request, + serve_postgres_health_http, +}; +use std::io::{Read, Write}; +use std::net::{IpAddr, Ipv4Addr, SocketAddr, TcpStream}; +use std::sync::mpsc; +use std::thread; +use std::time::Duration; + +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") +} + +fn request(target: &str) -> String { + format!("GET {target} HTTP/1.1\r\nHost: localhost\r\n\r\n") +} + +#[test] +fn writable_store_answers_live_and_postgres_ready_probes() { + let mut client = test_client(); + let live = handle_postgres_health_http_request( + &request(HEALTH_LIVE_PATH), + &mut client, + &["pg_catalog.pg_class"], + BacklogHealth::WithinBounds, + ); + assert_eq!(live.status(), 200); + assert!(live.body().contains("\"live\":true")); + assert!(!live.body().contains("postgres::")); + assert!(!live.body().contains("sql")); + + let ready = handle_postgres_health_http_request( + &request(&format!( + "{HEALTH_READY_PATH}?capability={POSTGRES_OPERATIONAL_STORE_CAPABILITY_REF}" + )), + &mut client, + &["pg_catalog.pg_class"], + BacklogHealth::WithinBounds, + ); + assert_eq!(ready.status(), 200); + assert!(ready.body().contains("\"ready\":true")); + assert!(ready + .body() + .contains(POSTGRES_OPERATIONAL_STORE_CAPABILITY_REF)); +} + +#[test] +fn liveness_probe_does_not_observe_the_store() { + let mut client = test_client(); + let _ = client.batch_execute("SELECT pg_terminate_backend(pg_backend_pid())"); + let live = handle_postgres_health_http_request( + &request(HEALTH_LIVE_PATH), + &mut client, + &["pg_catalog.pg_class"], + BacklogHealth::WithinBounds, + ); + assert_eq!(live.status(), 200); + assert!(live.body().contains("\"live\":true")); + assert!(!live + .body() + .contains(POSTGRES_OPERATIONAL_STORE_CAPABILITY_REF)); + assert!(!live.body().contains("terminate")); + assert!(!live.body().contains("postgres::")); + assert!(!live.body().contains("DbError")); +} + +#[test] +fn bare_ready_probe_fails_closed_when_the_store_cannot_accept_writes() { + let mut client = test_client(); + let mut transaction = client.build_transaction().read_only(true).start().unwrap(); + let ready = handle_postgres_health_http_request( + &request(HEALTH_READY_PATH), + &mut transaction, + &["pg_catalog.pg_class"], + BacklogHealth::WithinBounds, + ); + assert_eq!(ready.status(), 503); + assert!(ready.body().contains("\"live\":true")); + assert!(ready.body().contains("\"ready\":false")); + assert!(ready + .body() + .contains(POSTGRES_OPERATIONAL_STORE_CAPABILITY_REF)); + assert!(!ready.body().contains("sql")); +} + +#[test] +fn missing_required_relation_is_live_but_not_ready() { + let mut client = test_client(); + let ready = handle_postgres_health_http_request( + &request(&format!( + "{HEALTH_READY_PATH}?capability={POSTGRES_OPERATIONAL_STORE_CAPABILITY_REF}" + )), + &mut client, + &["psychometrics_commons_missing_relation"], + BacklogHealth::WithinBounds, + ); + assert_eq!(ready.status(), 503); + assert!(ready.body().contains("\"live\":true")); + assert!(ready.body().contains("\"ready\":false")); + assert!(ready + .body() + .contains("\"data_integrity_health\":\"incompatible\"")); +} + +#[test] +fn probe_failure_fails_readiness_closed_without_driver_text() { + let mut client = test_client(); + let _ = client.batch_execute("SELECT pg_terminate_backend(pg_backend_pid())"); + let ready = handle_postgres_health_http_request( + &request(&format!( + "{HEALTH_READY_PATH}?capability={POSTGRES_OPERATIONAL_STORE_CAPABILITY_REF}" + )), + &mut client, + &["pg_catalog.pg_class"], + BacklogHealth::WithinBounds, + ); + assert_eq!(ready.status(), 503); + assert!(ready.body().contains("\"ready\":false")); + assert!(!ready.body().contains("terminate")); + assert!(!ready.body().contains("postgres::")); + assert!(!ready.body().contains("DbError")); + assert!(!ready.body().contains("sql")); +} + +#[test] +fn stalled_backlog_fails_readiness_even_when_the_store_is_writable() { + let mut client = test_client(); + let ready = handle_postgres_health_http_request( + &request(HEALTH_READY_PATH), + &mut client, + &["pg_catalog.pg_class"], + BacklogHealth::Stalled, + ); + assert_eq!(ready.status(), 503); + assert!(ready.body().contains("\"backlog_health\":\"stalled\"")); +} + +#[test] +fn bound_listener_serves_a_postgres_ready_probe() { + let listener = bind_health_http(SocketAddr::new(IpAddr::V4(Ipv4Addr::LOCALHOST), 0)).unwrap(); + let addr = listener.local_addr().unwrap(); + let server = thread::spawn(move || { + let mut client = test_client(); + accept_one_postgres_health_http( + &listener, + &mut client, + &["pg_catalog.pg_class"], + BacklogHealth::WithinBounds, + ) + .unwrap(); + }); + + let mut stream = TcpStream::connect_timeout(&addr, Duration::from_secs(2)).unwrap(); + stream + .set_read_timeout(Some(Duration::from_secs(2))) + .unwrap(); + stream + .write_all( + format!( + "GET {HEALTH_READY_PATH}?capability={POSTGRES_OPERATIONAL_STORE_CAPABILITY_REF} HTTP/1.1\r\nHost: localhost\r\n\r\n" + ) + .as_bytes(), + ) + .unwrap(); + stream.shutdown(std::net::Shutdown::Write).unwrap(); + let mut body = String::new(); + stream.read_to_string(&mut body).unwrap(); + assert!(body.starts_with("HTTP/1.1 200 OK\r\n")); + assert!(body.contains("\"ready\":true")); + server.join().unwrap(); +} + +#[test] +fn serve_loop_answers_live_then_ready_without_store_io_on_live() { + let listener = bind_health_http(SocketAddr::new(IpAddr::V4(Ipv4Addr::LOCALHOST), 0)).unwrap(); + let addr = listener.local_addr().unwrap(); + let stop = listener.try_clone().unwrap(); + let (done_tx, done_rx) = mpsc::channel(); + let server = thread::spawn(move || { + let mut client = test_client(); + let result = serve_postgres_health_http( + &listener, + &mut client, + &["pg_catalog.pg_class"], + BacklogHealth::WithinBounds, + ); + let _ = done_tx.send(result); + }); + + let mut live_stream = TcpStream::connect_timeout(&addr, Duration::from_secs(2)).unwrap(); + live_stream + .set_read_timeout(Some(Duration::from_secs(2))) + .unwrap(); + live_stream + .write_all(format!("GET {HEALTH_LIVE_PATH} HTTP/1.1\r\nHost: localhost\r\n\r\n").as_bytes()) + .unwrap(); + live_stream.shutdown(std::net::Shutdown::Write).unwrap(); + let mut live = String::new(); + live_stream.read_to_string(&mut live).unwrap(); + assert!(live.starts_with("HTTP/1.1 200 OK\r\n")); + assert!(live.contains("\"live\":true")); + assert!(!live.contains(POSTGRES_OPERATIONAL_STORE_CAPABILITY_REF)); + + let mut ready_stream = TcpStream::connect_timeout(&addr, Duration::from_secs(2)).unwrap(); + ready_stream + .set_read_timeout(Some(Duration::from_secs(2))) + .unwrap(); + ready_stream + .write_all( + format!( + "GET {HEALTH_READY_PATH}?capability={POSTGRES_OPERATIONAL_STORE_CAPABILITY_REF} HTTP/1.1\r\nHost: localhost\r\n\r\n" + ) + .as_bytes(), + ) + .unwrap(); + ready_stream.shutdown(std::net::Shutdown::Write).unwrap(); + let mut ready = String::new(); + ready_stream.read_to_string(&mut ready).unwrap(); + assert!(ready.starts_with("HTTP/1.1 200 OK\r\n")); + assert!(ready.contains("\"ready\":true")); + + stop.set_nonblocking(true).unwrap(); + let _ = TcpStream::connect_timeout(&addr, Duration::from_millis(200)); + let stopped = done_rx + .recv_timeout(Duration::from_secs(3)) + .expect("serve_postgres_health_http must return after accept can no longer block"); + assert!(stopped.is_err()); + server.join().unwrap(); +} + +#[test] +fn serve_loop_keeps_accepting_after_a_client_drops_the_connection() { + let listener = bind_health_http(SocketAddr::new(IpAddr::V4(Ipv4Addr::LOCALHOST), 0)).unwrap(); + let addr = listener.local_addr().unwrap(); + let stop = listener.try_clone().unwrap(); + let (done_tx, done_rx) = mpsc::channel(); + let server = thread::spawn(move || { + let mut client = test_client(); + let result = serve_postgres_health_http( + &listener, + &mut client, + &["pg_catalog.pg_class"], + BacklogHealth::WithinBounds, + ); + let _ = done_tx.send(result); + }); + + { + let mut dropped = TcpStream::connect_timeout(&addr, Duration::from_secs(2)).unwrap(); + dropped + .write_all( + format!("GET {HEALTH_LIVE_PATH} HTTP/1.1\r\nHost: localhost\r\n\r\n").as_bytes(), + ) + .unwrap(); + } + + let mut live_stream = TcpStream::connect_timeout(&addr, Duration::from_secs(2)).unwrap(); + live_stream + .set_read_timeout(Some(Duration::from_secs(2))) + .unwrap(); + live_stream + .write_all(format!("GET {HEALTH_LIVE_PATH} HTTP/1.1\r\nHost: localhost\r\n\r\n").as_bytes()) + .unwrap(); + live_stream.shutdown(std::net::Shutdown::Write).unwrap(); + let mut live = String::new(); + live_stream.read_to_string(&mut live).unwrap(); + assert!( + live.starts_with("HTTP/1.1 200 OK\r\n"), + "a dropped PostgreSQL-backed probe must not stop later GET /live answers: {live}" + ); + assert!(live.contains("\"live\":true")); + assert!(!live.contains(POSTGRES_OPERATIONAL_STORE_CAPABILITY_REF)); + + stop.set_nonblocking(true).unwrap(); + let _ = TcpStream::connect_timeout(&addr, Duration::from_millis(200)); + let _ = done_rx + .recv_timeout(Duration::from_secs(3)) + .expect("serve_postgres_health_http must still return after accept can no longer block"); + server.join().unwrap(); +} diff --git a/tests/postgres_health_process_contract.rs b/tests/postgres_health_process_contract.rs new file mode 100644 index 00000000..a45a093c --- /dev/null +++ b/tests/postgres_health_process_contract.rs @@ -0,0 +1,79 @@ +//! Process entrypoint contract against a reachable `PostgreSQL` operational store. + +use postgres::{Client, NoTls}; +use psychometrics_commons_runtime::health_http::{HEALTH_LIVE_PATH, HEALTH_READY_PATH}; +use psychometrics_commons_runtime::health_process::{ + bind_health_process, parse_health_process_config, serve_health_process, + HEALTH_BACKLOG_HEALTH_ENV, HEALTH_DATABASE_URL_ENV, HEALTH_LISTEN_ADDR_ENV, +}; +use psychometrics_commons_runtime::postgres_health::POSTGRES_OPERATIONAL_STORE_CAPABILITY_REF; +use std::io::{Read, Write}; +use std::net::{SocketAddr, TcpStream}; +use std::sync::mpsc; +use std::thread; +use std::time::Duration; + +fn test_database_url() -> String { + std::env::var("TEST_DATABASE_URL") + .expect("TEST_DATABASE_URL must identify the isolated CI PostgreSQL database") +} + +fn exchange(addr: SocketAddr, request: &str) -> String { + let mut stream = TcpStream::connect_timeout(&addr, Duration::from_secs(2)).unwrap(); + stream + .set_read_timeout(Some(Duration::from_secs(2))) + .unwrap(); + stream.write_all(request.as_bytes()).unwrap(); + stream.shutdown(std::net::Shutdown::Write).unwrap(); + let mut body = String::new(); + stream.read_to_string(&mut body).unwrap(); + body +} + +#[test] +fn process_with_reachable_store_answers_ready_without_exposing_the_url() { + let url = test_database_url(); + Client::connect(&url, NoTls).expect("isolated CI PostgreSQL database must be reachable"); + let owned_url = url.clone(); + let config = parse_health_process_config(move |key| match key { + HEALTH_LISTEN_ADDR_ENV => Some("127.0.0.1:0".to_owned()), + HEALTH_DATABASE_URL_ENV => Some(owned_url.clone()), + HEALTH_BACKLOG_HEALTH_ENV => Some("within_bounds".to_owned()), + _ => None, + }) + .unwrap(); + let listener = bind_health_process(&config).unwrap(); + let addr = listener.local_addr().unwrap(); + let stop = listener.try_clone().unwrap(); + let (done_tx, done_rx) = mpsc::channel(); + let server = thread::spawn(move || { + let result = serve_health_process(&listener, &config); + let _ = done_tx.send(result); + }); + + let live = exchange( + addr, + &format!("GET {HEALTH_LIVE_PATH} HTTP/1.1\r\nHost: localhost\r\n\r\n"), + ); + assert!(live.starts_with("HTTP/1.1 200 OK\r\n"), "{live}"); + assert!(!live.contains(&url)); + + let ready = exchange( + addr, + &format!( + "GET {HEALTH_READY_PATH}?capability={POSTGRES_OPERATIONAL_STORE_CAPABILITY_REF} HTTP/1.1\r\nHost: localhost\r\n\r\n" + ), + ); + assert!( + ready.starts_with("HTTP/1.1 200 OK\r\n"), + "a reachable store with caller-measured backlog must be ready: {ready}" + ); + assert!(ready.contains(POSTGRES_OPERATIONAL_STORE_CAPABILITY_REF)); + assert!(!ready.contains(&url)); + assert!(!ready.contains("postgres::")); + + stop.set_nonblocking(true).unwrap(); + let _ = TcpStream::connect_timeout(&addr, Duration::from_millis(200)); + let _ = done_rx.recv_timeout(Duration::from_secs(3)); + server.join().unwrap(); +} diff --git a/tests/postgres_operational_snapshot.rs b/tests/postgres_operational_snapshot.rs new file mode 100644 index 00000000..e4eab297 --- /dev/null +++ b/tests/postgres_operational_snapshot.rs @@ -0,0 +1,107 @@ +//! Compose `PostgreSQL` probes into one operation-scoped runtime health snapshot. + +use postgres::{Client, NoTls}; +use psychometrics_commons_runtime::health::{BacklogHealth, CapabilityState, DataIntegrityHealth}; +use psychometrics_commons_runtime::postgres_health::{ + observe_postgres_operational_snapshot, POSTGRES_OPERATIONAL_STORE_CAPABILITY_REF, +}; + +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 writable_store_and_verified_relations_are_ready_for_the_postgres_capability() { + let mut client = test_client(); + let snapshot = observe_postgres_operational_snapshot( + &mut client, + &["pg_catalog.pg_class"], + BacklogHealth::WithinBounds, + ); + + assert!(snapshot.is_live()); + assert_eq!(snapshot.backlog_health(), BacklogHealth::WithinBounds); + assert_eq!( + snapshot.data_integrity_health(), + DataIntegrityHealth::Verified + ); + assert!(snapshot.is_ready_for(&[POSTGRES_OPERATIONAL_STORE_CAPABILITY_REF])); + let capability = snapshot + .capability(POSTGRES_OPERATIONAL_STORE_CAPABILITY_REF) + .unwrap(); + assert_eq!(capability.state(), CapabilityState::Available); + assert!(capability.accepts_new_work()); +} + +#[test] +fn missing_required_relation_fails_readiness_closed() { + let mut client = test_client(); + let snapshot = observe_postgres_operational_snapshot( + &mut client, + &["psychometrics_commons_missing_relation"], + BacklogHealth::WithinBounds, + ); + + assert!(snapshot.is_live()); + assert_eq!( + snapshot.data_integrity_health(), + DataIntegrityHealth::Incompatible + ); + assert!(!snapshot.is_ready_for(&[POSTGRES_OPERATIONAL_STORE_CAPABILITY_REF])); +} + +#[test] +fn probe_failure_maps_to_unknown_unready_evidence_without_exposing_the_driver() { + let mut client = test_client(); + let _ = client.batch_execute("SELECT pg_terminate_backend(pg_backend_pid())"); + let snapshot = observe_postgres_operational_snapshot( + &mut client, + &["pg_catalog.pg_class"], + BacklogHealth::WithinBounds, + ); + + assert!(snapshot.is_live()); + assert_eq!( + snapshot.data_integrity_health(), + DataIntegrityHealth::Unknown + ); + let capability = snapshot + .capability(POSTGRES_OPERATIONAL_STORE_CAPABILITY_REF) + .unwrap(); + assert_eq!(capability.state(), CapabilityState::Unknown); + assert!(!capability.accepts_new_work()); + assert!(!snapshot.is_ready_for(&[POSTGRES_OPERATIONAL_STORE_CAPABILITY_REF])); +} + +#[test] +fn caller_supplied_stalled_backlog_fails_readiness_even_when_the_store_is_writable() { + let mut client = test_client(); + let snapshot = observe_postgres_operational_snapshot( + &mut client, + &["pg_catalog.pg_class"], + BacklogHealth::Stalled, + ); + + assert_eq!(snapshot.backlog_health(), BacklogHealth::Stalled); + assert!(!snapshot.is_ready_for(&[POSTGRES_OPERATIONAL_STORE_CAPABILITY_REF])); +} + +#[test] +fn read_only_transaction_is_live_but_not_ready() { + let mut client = test_client(); + let mut transaction = client.build_transaction().read_only(true).start().unwrap(); + let snapshot = observe_postgres_operational_snapshot( + &mut transaction, + &["pg_catalog.pg_class"], + BacklogHealth::WithinBounds, + ); + + assert!(snapshot.is_live()); + let capability = snapshot + .capability(POSTGRES_OPERATIONAL_STORE_CAPABILITY_REF) + .unwrap(); + assert_eq!(capability.state(), CapabilityState::Unavailable); + assert!(!snapshot.is_ready_for(&[POSTGRES_OPERATIONAL_STORE_CAPABILITY_REF])); +}