diff --git a/Cargo.lock b/Cargo.lock index c02c2ad3e37..3631dfb6fbb 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -4448,6 +4448,24 @@ dependencies = [ "wat", ] +[[package]] +name = "ironclaw_hooks_postgres" +version = "0.1.0" +dependencies = [ + "async-trait", + "blake3", + "chrono", + "deadpool-postgres", + "ironclaw_hooks", + "ironclaw_host_api", + "rust_decimal", + "thiserror 2.0.18", + "tokio", + "tokio-postgres", + "tracing", + "uuid", +] + [[package]] name = "ironclaw_host_api" version = "0.1.0" diff --git a/Cargo.toml b/Cargo.toml index f2a49191f24..0b23fa499c7 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -1,5 +1,5 @@ [workspace] -members = [".", "crates/ironclaw_common", "crates/ironclaw_host_api", "crates/ironclaw_filesystem", "crates/ironclaw_memory", "crates/ironclaw_events", "crates/ironclaw_event_projections", "crates/ironclaw_event_streams", "crates/ironclaw_reborn_event_store", "crates/ironclaw_extensions", "crates/ironclaw_processes", "crates/ironclaw_dispatcher", "crates/ironclaw_scripts", "crates/ironclaw_process_sandbox", "crates/ironclaw_mcp", "crates/ironclaw_wasm", "crates/ironclaw_wasm_sandbox_core", "crates/ironclaw_wasm_limiter", "crates/ironclaw_capabilities", "crates/ironclaw_secrets", "crates/ironclaw_network", "crates/ironclaw_host_runtime", "crates/ironclaw_runtime_policy", "crates/ironclaw_authorization", "crates/ironclaw_run_state", "crates/ironclaw_approvals", "crates/ironclaw_resources", "crates/ironclaw_auth", "crates/ironclaw_trust", "crates/ironclaw_turns", "crates/ironclaw_agent_loop", "crates/ironclaw_threads", "crates/ironclaw_prompt_envelope", "crates/ironclaw_hooks", "crates/ironclaw_loop_support", "crates/ironclaw_reborn", "crates/ironclaw_reborn_config", "crates/ironclaw_reborn_composition", "crates/ironclaw_first_party_extensions", "crates/ironclaw_reborn_cli", "crates/ironclaw_reborn_traces", "crates/ironclaw_reborn_webui_ingress", "crates/ironclaw_conversations", "crates/ironclaw_product_adapters", "crates/ironclaw_product_workflow", "crates/ironclaw_product_workflow_storage", "crates/ironclaw_product_adapter_registry", "crates/ironclaw_wasm_product_adapters", "crates/ironclaw_telegram_v2_adapter", "crates/ironclaw_slack_v2_adapter", "crates/ironclaw_outbound", "crates/ironclaw_triggers", "crates/ironclaw_architecture", "crates/ironclaw_safety", "crates/ironclaw_skills", "crates/ironclaw_oauth", "crates/ironclaw_llm", "crates/ironclaw_embeddings", "crates/ironclaw_engine", "crates/ironclaw_gateway", "crates/ironclaw_tui", "crates/ironclaw_webui_v2", "crates/ironclaw_webui_v2_static"] +members = [".", "crates/ironclaw_common", "crates/ironclaw_host_api", "crates/ironclaw_filesystem", "crates/ironclaw_memory", "crates/ironclaw_events", "crates/ironclaw_event_projections", "crates/ironclaw_event_streams", "crates/ironclaw_reborn_event_store", "crates/ironclaw_extensions", "crates/ironclaw_processes", "crates/ironclaw_dispatcher", "crates/ironclaw_scripts", "crates/ironclaw_process_sandbox", "crates/ironclaw_mcp", "crates/ironclaw_wasm", "crates/ironclaw_wasm_sandbox_core", "crates/ironclaw_wasm_limiter", "crates/ironclaw_capabilities", "crates/ironclaw_secrets", "crates/ironclaw_network", "crates/ironclaw_host_runtime", "crates/ironclaw_runtime_policy", "crates/ironclaw_authorization", "crates/ironclaw_run_state", "crates/ironclaw_approvals", "crates/ironclaw_resources", "crates/ironclaw_auth", "crates/ironclaw_trust", "crates/ironclaw_turns", "crates/ironclaw_agent_loop", "crates/ironclaw_threads", "crates/ironclaw_prompt_envelope", "crates/ironclaw_hooks", "crates/ironclaw_hooks_postgres", "crates/ironclaw_loop_support", "crates/ironclaw_reborn", "crates/ironclaw_reborn_config", "crates/ironclaw_reborn_composition", "crates/ironclaw_first_party_extensions", "crates/ironclaw_reborn_cli", "crates/ironclaw_reborn_traces", "crates/ironclaw_reborn_webui_ingress", "crates/ironclaw_conversations", "crates/ironclaw_product_adapters", "crates/ironclaw_product_workflow", "crates/ironclaw_product_workflow_storage", "crates/ironclaw_product_adapter_registry", "crates/ironclaw_wasm_product_adapters", "crates/ironclaw_telegram_v2_adapter", "crates/ironclaw_slack_v2_adapter", "crates/ironclaw_outbound", "crates/ironclaw_triggers", "crates/ironclaw_architecture", "crates/ironclaw_safety", "crates/ironclaw_skills", "crates/ironclaw_oauth", "crates/ironclaw_llm", "crates/ironclaw_embeddings", "crates/ironclaw_engine", "crates/ironclaw_gateway", "crates/ironclaw_tui", "crates/ironclaw_webui_v2", "crates/ironclaw_webui_v2_static"] exclude = [ "channels-src/discord", "channels-src/feishu", diff --git a/crates/ironclaw_hooks_postgres/AGENTS.md b/crates/ironclaw_hooks_postgres/AGENTS.md new file mode 100644 index 00000000000..9ca6d6b6b08 --- /dev/null +++ b/crates/ironclaw_hooks_postgres/AGENTS.md @@ -0,0 +1,56 @@ +# Agent Map — ironclaw_hooks_postgres + +## Start Here + +- No crate-local `CLAUDE.md` exists yet; use this map plus the contracts below. +- Read `src/lib.rs` first — it documents why this is a separate crate, the + dual-backend split (Postgres here, libSQL in `ironclaw_hooks_libsql`, parity + in `ironclaw_hooks_parity`), and the shared two-table typed schema. +- Read `Cargo.toml` for actual dependencies and the `postgres` feature gate. +- The trait contract this crate implements lives in + `ironclaw_hooks::predicate_state` (`PredicateStateBackend`); the shared + contract harness is exposed via that crate's `contract-tests` feature. + +## What This Crate Owns + +- The durable PostgreSQL `PredicateStateBackend` implementation: the atomic + record-and-read transaction body, per-key advisory-lock serialization, the + fail-closed per-key sample cap, and the per-scope LRU quota enforcement. +- The Postgres predicate-state schema (`hooks_predicate_invocations` / + `hooks_predicate_values`) and its migration SQL. +- Crate-local public API, tests, and fixtures needed to prove that ownership. + +## Do Not Move In Here + +- The `PredicateStateBackend` trait itself, the in-memory backend, or the + libSQL backend (those live in `ironclaw_hooks` / `ironclaw_hooks_libsql`). +- Evaluator policy, hook-framework wiring, or backend selection. +- Secrets, raw connection strings, host names, schema details, or raw DB + error text in errors, events, logs, or docs — `PredicateBackendError` + payloads are sanitized (`DB_UNAVAILABLE_MSG` and the quota fail-closed + constants); keep the raw error behind `tracing` only. + +## Validation + +- Fast local check: `cargo test -p ironclaw_hooks_postgres` +- Postgres-backed integration/adversarial tests are env-gated on + `IRONCLAW_HOOKS_POSTGRES_URL` / `DATABASE_URL` and skip (passing) when no + DB is reachable; run them against a live Postgres to exercise the advisory + locks, cap fail-closed, and LRU eviction paths. +- If production persistence behavior changes, keep PostgreSQL/libSQL parity: + update the libSQL counterpart (`ironclaw_hooks_libsql`) and the + cross-backend parity suite (`ironclaw_hooks_parity`) in lockstep. + +## Agent Notes + +- Keep edits inside this crate unless the trait contract in `ironclaw_hooks` + explicitly requires a neighboring crate change. +- The advisory-lock protocol is load-bearing: same-key writers serialize on a + per-key `(int4,int4)` lock, scope-quota passes serialize on a per-scope + `(int8)` lock folded with a kind byte, and victim eviction uses the + NON-blocking `pg_try_advisory_xact_lock` to stay deadlock-free. Do not change + lock derivation, isolation level, or fail-closed posture without updating the + module-level docs and the disjointness/lock-equality unit tests. +- Cap and quota are fail-closed by contract: overflow returns + `WindowOverflow`, an unenforceable scope quota returns `Unavailable`. Never + silently drop-oldest or commit an over-quota scope. diff --git a/crates/ironclaw_hooks_postgres/Cargo.toml b/crates/ironclaw_hooks_postgres/Cargo.toml new file mode 100644 index 00000000000..9bd93222c80 --- /dev/null +++ b/crates/ironclaw_hooks_postgres/Cargo.toml @@ -0,0 +1,37 @@ +[package] +name = "ironclaw_hooks_postgres" +version = "0.1.0" +edition = "2024" +publish = false +description = "Durable PostgreSQL-backed PredicateStateBackend for the reborn hook framework (durable backend PR 2/4)." +authors = ["NEAR AI "] +license = "MIT OR Apache-2.0" +homepage = "https://github.com/nearai/ironclaw" +repository = "https://github.com/nearai/ironclaw" + +[features] +default = [] +# The Postgres backend itself. Mirrors the per-crate `postgres` feature gate +# used by `ironclaw_reborn_event_store` / `ironclaw_filesystem`: nothing in +# this crate compiles a DB dependency unless `postgres` is on. +postgres = ["dep:deadpool-postgres", "dep:tokio-postgres"] + +[dependencies] +async-trait = "0.1" +blake3 = "1" +chrono = { version = "0.4", features = ["serde"] } +deadpool-postgres = { version = "0.14", optional = true } +ironclaw_hooks = { path = "../ironclaw_hooks" } +ironclaw_host_api = { path = "../ironclaw_host_api" } +rust_decimal = { version = "1", features = ["serde", "serde-with-str", "db-tokio-postgres"] } +thiserror = "2" +tokio-postgres = { version = "0.7", optional = true, features = ["with-chrono-0_4"] } +tracing = "0.1" + +[dev-dependencies] +# `contract-tests` exposes the trait-level contract harness from +# `ironclaw_hooks` so this crate's Postgres impl runs the SAME suite the +# in-memory backend runs (proven-by-construction, not per-impl). +ironclaw_hooks = { path = "../ironclaw_hooks", features = ["contract-tests"] } +tokio = { version = "1", features = ["macros", "rt", "rt-multi-thread", "sync"] } +uuid = { version = "1", features = ["v4"] } diff --git a/crates/ironclaw_hooks_postgres/migrations/V1__predicate_state.sql b/crates/ironclaw_hooks_postgres/migrations/V1__predicate_state.sql new file mode 100644 index 00000000000..73548c145ee --- /dev/null +++ b/crates/ironclaw_hooks_postgres/migrations/V1__predicate_state.sql @@ -0,0 +1,93 @@ +-- Durable predicate sliding-window state for the reborn hook framework. +-- +-- This crate owns its own schema (per-crate pattern, like +-- ironclaw_reborn_event_store / ironclaw_filesystem) rather than going +-- through the legacy main-binary refinery `migrations/` directory. The +-- DDL is embedded verbatim into `schema.rs` via `include_str!` and applied +-- as an idempotent `CREATE TABLE IF NOT EXISTS` batch by `run_migrations()`; +-- this file is the single human-reviewable canonical source for that schema. +-- +-- ## Canonical typed two-table shape (cross-backend invariant) +-- +-- The two durable backends (Postgres + libSQL) share ONE logical schema: +-- two typed tables — one for invocation-count samples, one for +-- numeric-value samples — with identical column names and semantics. The +-- storage TYPES differ per backend (Postgres uses native TIMESTAMPTZ + +-- NUMERIC; libSQL uses epoch-ms INTEGER + TEXT) but the table count, column +-- names, primary keys, and eviction/dedup/quota semantics are identical. +-- The cross-backend parity suite (ironclaw_hooks_parity) proves they are +-- behaviorally interchangeable. Replacing the earlier single +-- `hook_predicate_counters(kind CHAR(1), …)` table with two explicit typed +-- tables removes the `kind` discriminator AND the `value: Option` +-- "count smuggled through a NUMERIC" abstraction the shared record path used. +-- +-- ## Hash columns +-- +-- `scope_hash` and `key_hash` are blake3 digests (32 raw bytes, BYTEA): +-- scope_hash = blake3(len-prefixed tenant_id) -- tenant grain +-- key_hash = blake3(map-discriminant ++ hook_id ++ tenant_id +-- ++ capability [++ field]) -- full bucket +-- `scope_hash` is the trust boundary + per-tenant LRU-quota grain; `key_hash` +-- is the full bucket identity (the dedup + count/sum grain). BYTEA (not TEXT) +-- keeps the index keys fixed-width and avoids collation surprises. +-- +-- ## event_id column type (codex #3635 finding) +-- +-- The replay-dedup id is a `PredicateEventId` — an opaque host-assigned +-- string whose canonical synth shape is a 64-char blake3 hex digest, but +-- callers may stamp other formats. Postgres `uuid` is a fixed 128-bit type +-- and will REJECT a 64-char hex digest, so `event_id` is `TEXT`, NOT `uuid`. +-- (#3635 docs pinned a 64-char id while the old schema said uuid; TEXT +-- resolves that contradiction, and matches the libSQL sibling.) +-- +-- ## Window-clock basis +-- +-- `occurred_at` is the wall-clock timestamp passed by the caller +-- (TIMESTAMPTZ). Window comparisons are performed against a caller-supplied +-- cutoff computed from the same clock the in-memory backend uses, so the trim +-- semantics (`occurred_at < cutoff`, entry at exact cutoff retained) match the +-- in-memory backend bit-for-bit. See `schema.rs` for the DB-clock rationale. + +-- Invocation-count samples. One row per recorded invocation event; the +-- in-window COUNT(*) is the invocation count. +CREATE TABLE IF NOT EXISTS hooks_predicate_invocations ( + scope_hash BYTEA NOT NULL, + key_hash BYTEA NOT NULL, + event_id TEXT NOT NULL, + occurred_at TIMESTAMPTZ NOT NULL, + PRIMARY KEY (key_hash, event_id) +); + +-- Per-key window-trim + COUNT scan: every record_invocation prunes and +-- aggregates over (key_hash, occurred_at). +CREATE INDEX IF NOT EXISTS hooks_predicate_invocations_key_ts_idx + ON hooks_predicate_invocations (key_hash, occurred_at); +-- Per-scope (tenant) distinct-key LRU eviction. enforce_scope_quota runs +-- COUNT(DISTINCT key_hash) and ranks victims by MIN(occurred_at) per key, +-- both scoped by scope_hash; the (scope_hash, key_hash, occurred_at) cover +-- lets those run as index-only scans for tenants with many recorded keys. +CREATE INDEX IF NOT EXISTS hooks_predicate_invocations_scope_idx + ON hooks_predicate_invocations (scope_hash, key_hash, occurred_at); +-- Operator reaper (`evict_older_than`) deletes globally by age. +CREATE INDEX IF NOT EXISTS hooks_predicate_invocations_ts_idx + ON hooks_predicate_invocations (occurred_at); + +-- Numeric-value samples. One row per recorded value event; the in-window +-- SUM(value) is the running sum. `value` is NOT NULL here (the typed tables +-- make the count-vs-sum distinction explicit, so no nullable double-duty). +CREATE TABLE IF NOT EXISTS hooks_predicate_values ( + scope_hash BYTEA NOT NULL, + key_hash BYTEA NOT NULL, + event_id TEXT NOT NULL, + occurred_at TIMESTAMPTZ NOT NULL, + value NUMERIC NOT NULL, + PRIMARY KEY (key_hash, event_id) +); + +CREATE INDEX IF NOT EXISTS hooks_predicate_values_key_ts_idx + ON hooks_predicate_values (key_hash, occurred_at); +-- Same index-only-scan cover for the value table's scope-quota pass. +CREATE INDEX IF NOT EXISTS hooks_predicate_values_scope_idx + ON hooks_predicate_values (scope_hash, key_hash, occurred_at); +CREATE INDEX IF NOT EXISTS hooks_predicate_values_ts_idx + ON hooks_predicate_values (occurred_at); diff --git a/crates/ironclaw_hooks_postgres/src/backend.rs b/crates/ironclaw_hooks_postgres/src/backend.rs new file mode 100644 index 00000000000..b20b936177f --- /dev/null +++ b/crates/ironclaw_hooks_postgres/src/backend.rs @@ -0,0 +1,987 @@ +//! [`PostgresPredicateStateBackend`] — durable, cross-host-consistent +//! implementation of [`PredicateStateBackend`]. +//! +//! # Atomic record-and-read (the load-bearing property) +//! +//! Every `record_*` call runs inside a single `READ COMMITTED` +//! transaction guarded by a transaction-scoped advisory lock on the +//! bucket (and, for new-key eviction, a second advisory lock on the +//! scope). The advisory lock — not a high isolation level — is what +//! serializes concurrent writers to the same bucket: it makes the second +//! writer *block* until the first commits rather than aborting it with a +//! serialization failure (which would force a caller retry loop). The +//! transaction body: +//! +//! 1. **Trims** rows for the key with `occurred_at < cutoff` (out of +//! window). This also frees the dedup `event_id` of any trimmed row, so +//! an `event_id` that aged out of the window can be re-recorded — +//! matching the in-memory backend, whose dedup memory is exactly the +//! in-window entry set. +//! 2. **Dedup-checks** the incoming `event_id` against the in-window rows +//! for the key. A match is a replay — a no-op that short-circuits before +//! the cap check and returns the unchanged aggregate. This is the +//! cross-host replay defense (a row written by host A blocks host B's +//! re-insert of the same id regardless of clock skew) AND the property +//! that lets a replay survive the cap boundary. +//! 3. **Caps fail-closed**: if the `event_id` is new (not a replay) and the +//! in-window count is already at [`MAX_SAMPLES_PER_KEY`], the call +//! returns [`PredicateBackendError::WindowOverflow`] WITHOUT inserting — +//! matching the in-memory backend's `if !dedup && len >= cap { Err }` +//! contract exactly. Silently evicting the oldest sample would weaken +//! cap enforcement and break replay refusal (the evicted id would leave +//! the dedup set while still logically in-window), so we fail closed. +//! 4. **Inserts** the new row `ON CONFLICT (key_hash, event_id) DO NOTHING` +//! into the kind-specific table (`hooks_predicate_invocations` for +//! counts, `hooks_predicate_values` for numeric sums — explicit typed +//! tables, not a generic `kind` discriminator column). +//! 5. **Evicts** the scope's oldest-front key — the key whose oldest +//! retained sample (`MIN(occurred_at)`) is oldest — when the scope's +//! distinct-key count exceeds [`MAX_KEYS_PER_TENANT`] (the durable +//! analogue of the in-memory per-tenant LRU quota, which ranks buckets by +//! their oldest entry). Each victim's rows +//! are deleted only after acquiring that victim key's per-key advisory +//! lock with the NON-blocking `pg_try_advisory_xact_lock`, so eviction +//! obeys the same per-bucket serialization as a recorder and can never +//! deadlock against (or tear the aggregate of) a concurrent write to the +//! victim bucket. Eviction is an explicit quota OUTCOME: it requeries and +//! keeps evicting until the per-scope cap is met, or FAILS CLOSED +//! ([`PredicateBackendError::Unavailable`]) rather than committing an +//! over-quota scope — it is never silently best-effort (see +//! [`PostgresPredicateStateBackend::enforce_scope_quota`]). +//! 6. **Aggregates** the in-window `COUNT(*)` (invocation table) / `SUM(value)` +//! (value table) and returns it. +//! +//! Steps 1-6 share one transaction under the bucket advisory lock, so two +//! concurrent writers can never both observe "1 under cap" and both +//! proceed — the second blocks until the first commits. This is the codex +//! Critical atomicity requirement from PR #3635. +//! +//! # Cap semantics — fail-closed, NOT drop-oldest +//! +//! When a key's in-window sample count reaches [`MAX_SAMPLES_PER_KEY`] and +//! a NEW distinct id arrives, the backend returns +//! [`PredicateBackendError::WindowOverflow`] rather than evicting the +//! oldest sample to make room. This matches the in-memory backend and the +//! trait contract (PR #3635 followup / #3929): the evaluator maps the error +//! to a restrictive DENY/PauseApproval, so overflow surfaces as a refusal, +//! never a silent Allow. Replay of an already-recorded in-window id still +//! short-circuits to a no-op before the cap check (step 2), so replay +//! refusal survives the cap boundary. + +use std::time::Duration; + +use async_trait::async_trait; +use chrono::{DateTime, Utc}; +use deadpool_postgres::Pool; +use ironclaw_hooks::predicate_state::{ + InvocationKey, MAX_KEYS_PER_TENANT, MAX_SAMPLES_PER_KEY, PredicateBackendError, + PredicateEventId, PredicateStateBackend, ValueKey, +}; +use rust_decimal::Decimal; +use std::sync::atomic::{AtomicU64, Ordering}; +use tokio_postgres::IsolationLevel; + +use crate::hashing::{Digest, invocation_key_hash, scope_hash, value_key_hash}; +use crate::schema::{INVOCATIONS_TABLE, VALUES_TABLE}; + +/// A fully-resolved, typed record request for one of the two predicate +/// tables. This replaces the old `Bucket { kind } + value: Option` +/// pair: the count path is [`RecordPlan::Invocation`] (no value field can +/// exist) and the sum path is [`RecordPlan::Value`] which *carries* the +/// `Decimal` in the variant. There is no nullable mode flag and no +/// `debug_assert_eq` invariant tying a separate `value` arg to a separate +/// `kind` — the type makes "value present iff value table" structural. +/// +/// All four common fields (`scope`, `key`, `label`) live in the shared +/// [`RecordPlan::common`] accessor so the lock/trim/dedup/quota helpers stay +/// table-agnostic while the INSERT column list and the final aggregate are +/// driven by the variant. +enum RecordPlan { + /// `hooks_predicate_invocations`; aggregate is `COUNT(*)`. + Invocation(PlanCommon), + /// `hooks_predicate_values`; aggregate is `SUM(value)`. The recorded + /// numeric lives HERE, not in a nullable side-channel. + Value { common: PlanCommon, value: Decimal }, +} + +/// Bucket identity shared by both [`RecordPlan`] variants. +struct PlanCommon { + /// Scope (tenant) digest — the per-tenant LRU-quota grain. + scope: Digest, + /// Full bucket-identity digest — the dedup + count/sum grain and the + /// per-key advisory-lock key. + key: Digest, + /// Human-readable label for the bucket, used only in the + /// `WindowOverflow` error. Mirrors the in-memory backend's format: + /// `{tenant}/{capability}` for invocations and + /// `{tenant}/{capability}#{field}` for values. + label: String, +} + +impl RecordPlan { + /// The shared bucket identity, regardless of variant. + fn common(&self) -> &PlanCommon { + match self { + RecordPlan::Invocation(c) => c, + RecordPlan::Value { common, .. } => common, + } + } + + /// One-byte discriminant folded into the per-scope advisory-lock key so + /// the invocation and value scope-quota passes lock independently (they + /// touch disjoint tables, so they must not serialize against each other). + fn lock_tag(&self) -> &'static [u8] { + match self { + RecordPlan::Invocation(_) => b"i", + RecordPlan::Value { .. } => b"v", + } + } +} + +/// Per-table SQL statements with the table name already interpolated. +/// +/// The table name folded into every `record()` statement is fixed per kind +/// (`hooks_predicate_invocations` for counts, `hooks_predicate_values` for +/// sums), so the statement strings are built ONCE at backend construction +/// rather than re-`format!`'d on every hot-path `record()` call (which +/// allocated 5-6 throwaway `String`s per invocation — henrypark perf +/// finding). `record()` selects the matching set via [`RecordPlan::table`]. +struct TableStatements { + trim: String, + dedup_precount: String, + insert: String, + aggregate_count: String, + aggregate_sum: String, + scope_distinct: String, + scope_candidates: String, + evict_victim: String, + reap: String, +} + +impl TableStatements { + /// Statements for the invocation (count) table. The invocation table has + /// no `value` column, so its INSERT omits it; the aggregate is `COUNT(*)`. + fn for_invocations(table: &str) -> Self { + Self::build( + table, + format!( + "INSERT INTO {table} (scope_hash, key_hash, event_id, occurred_at) + VALUES ($1, $2, $3, $4) + ON CONFLICT (key_hash, event_id) DO NOTHING" + ), + ) + } + + /// Statements for the value (sum) table. The value table carries a NOT + /// NULL `value` column the invocation table lacks, so its INSERT includes + /// it; the aggregate is `SUM(value)`. + fn for_values(table: &str) -> Self { + Self::build( + table, + format!( + "INSERT INTO {table} (scope_hash, key_hash, event_id, occurred_at, value) + VALUES ($1, $2, $3, $4, $5) + ON CONFLICT (key_hash, event_id) DO NOTHING" + ), + ) + } + + /// Pre-format every table-agnostic statement around the table-specific + /// `insert` SQL the named constructors above supply. The trim / dedup / + /// aggregate / scope-quota / reap statements are identical in shape across + /// the two typed tables (only the table name is interpolated), so they are + /// built once here. + fn build(table: &str, insert: String) -> Self { + Self { + trim: format!("DELETE FROM {table} WHERE key_hash = $1 AND occurred_at < $2"), + dedup_precount: format!( + "SELECT COUNT(*)::BIGINT AS cnt, + BOOL_OR(event_id = $3) AS dup + FROM {table} + WHERE key_hash = $1 AND occurred_at >= $2" + ), + insert, + aggregate_count: format!( + "SELECT COUNT(*)::BIGINT FROM {table} + WHERE key_hash = $1 AND occurred_at >= $2" + ), + aggregate_sum: format!( + "SELECT COALESCE(SUM(value), 0)::NUMERIC FROM {table} + WHERE key_hash = $1 AND occurred_at >= $2" + ), + scope_distinct: format!( + "SELECT COUNT(DISTINCT key_hash)::BIGINT + FROM {table} + WHERE scope_hash = $1" + ), + scope_candidates: format!( + "SELECT key_hash FROM ( + SELECT key_hash, MIN(occurred_at) AS oldest_ts + FROM {table} + WHERE scope_hash = $1 + AND key_hash <> $2 + GROUP BY key_hash + ORDER BY oldest_ts ASC + LIMIT $3 + ) victims" + ), + evict_victim: format!("DELETE FROM {table} WHERE scope_hash = $1 AND key_hash = $2"), + reap: format!("DELETE FROM {table} WHERE occurred_at < $1"), + } + } +} + +/// Durable PostgreSQL [`PredicateStateBackend`]. Holds a `deadpool` +/// connection pool; construct one pool per process and share the backend +/// behind an `Arc`. +pub struct PostgresPredicateStateBackend { + pool: Pool, + /// Local mirror of LRU evictions performed by THIS process instance, + /// matching the in-memory backend's `evictions_observed()` contract + /// (a process-local monitoring counter, not a global DB total). + evictions: AtomicU64, + /// Pre-formatted SQL for the invocation (count) table. + invocation_sql: TableStatements, + /// Pre-formatted SQL for the value (sum) table. + value_sql: TableStatements, +} + +impl PostgresPredicateStateBackend { + /// Wrap a `deadpool` pool. Call [`Self::run_migrations`] once before + /// first use to ensure the schema exists. + pub fn new(pool: Pool) -> Self { + Self { + pool, + evictions: AtomicU64::new(0), + invocation_sql: TableStatements::for_invocations(INVOCATIONS_TABLE), + value_sql: TableStatements::for_values(VALUES_TABLE), + } + } + + /// Select the pre-formatted statement set matching a plan's typed table. + fn statements(&self, plan: &RecordPlan) -> &TableStatements { + match plan { + RecordPlan::Invocation(_) => &self.invocation_sql, + RecordPlan::Value { .. } => &self.value_sql, + } + } + + /// Apply the idempotent schema. Safe to call repeatedly and + /// concurrently (`CREATE … IF NOT EXISTS`). + pub async fn run_migrations(&self) -> Result<(), PredicateBackendError> { + let client = self.client().await?; + client + .batch_execute(crate::schema::POSTGRES_PREDICATE_SCHEMA) + .await + .map_err(map_pg)?; + Ok(()) + } + + async fn client(&self) -> Result { + self.pool.get().await.map_err(map_pool) + } + + /// Compute the wall-clock cutoff `now - window`, saturating to `now` + /// for windows beyond chrono's range (nothing trimmed — conservative + /// for a rate/value cap). Mirrors `predicate_state::window_cutoff`. + fn cutoff(now: DateTime, window: Duration) -> DateTime { + match chrono::Duration::from_std(window) { + Ok(d) => now.checked_sub_signed(d).unwrap_or(now), + Err(_) => now, + } + } + + /// Shared transaction body for both record paths. The trim / dedup / cap + /// / quota steps are identical across the two typed tables, so they are + /// driven generically off [`RecordPlan::common`]; only the INSERT column + /// list and the final aggregate are variant-specific, dispatched on the + /// [`RecordPlan`] enum (NOT on a nullable `value: Option` mode + /// flag — the count path can no longer carry a value at all, and the sum + /// path carries its `Decimal` inside the variant). The caller-facing + /// `record_invocation` / `record_value` methods build the typed plan and + /// map the returned aggregate to their typed return. + async fn record( + &self, + plan: RecordPlan, + event_id: &PredicateEventId, + now: DateTime, + window: Duration, + ) -> Result { + let PlanCommon { scope, key, label } = plan.common(); + let (scope, key, label) = (*scope, *key, label.clone()); + let sql = self.statements(&plan); + let cutoff = Self::cutoff(now, window); + let mut client = self.client().await?; + // READ COMMITTED + a transaction-scoped advisory lock keyed on the + // bucket. The advisory lock serializes ALL writers to the same + // key (the durable analogue of the in-memory backend's single + // Mutex), so the trim / insert / aggregate steps see a consistent + // view and two concurrent writers can never both observe "under + // cap" and both proceed (codex Critical atomicity requirement). + // + // We deliberately do NOT use REPEATABLE READ here: under that + // level concurrent same-key writers abort with + // `could not serialize access`, forcing a caller retry loop. The + // advisory lock instead makes the second writer *block* until the + // first commits — same correctness, no spurious aborts. Writers to + // DIFFERENT keys take different advisory locks and proceed fully + // concurrently. + let tx = client + .build_transaction() + .isolation_level(IsolationLevel::ReadCommitted) + .start() + .await + .map_err(map_pg)?; + + let scope_ref: &[u8] = &scope; + let key_ref: &[u8] = &key; + + // Serialize same-key writers. The advisory lock key is two i32s + // derived from the bucket's key_hash; pg_advisory_xact_lock is + // released automatically at commit/rollback. Collisions across + // distinct keys (same 64-bit lock key) only cost extra + // serialization, never correctness. + let lock_key = advisory_lock_key_from_bytes(&key); + tx.execute( + "SELECT pg_advisory_xact_lock($1, $2)", + &[&lock_key.0, &lock_key.1], + ) + .await + .map_err(map_pg)?; + + // (1) Trim out-of-window rows for this key. Doing this BEFORE the + // insert frees the dedup id of any row that aged out of the + // window, so a re-used id whose original entry is no longer + // in-window records fresh — matching the in-memory backend, whose + // dedup memory is exactly the in-window entry set. + tx.execute(&sql.trim, &[&key_ref, &cutoff]) + .await + .map_err(map_pg)?; + + // (2) Replay-dedup check + pre-insert count, computed atomically in + // one statement under the advisory lock. `cnt` is the in-window + // sample count BEFORE this call's insert; `dup` is whether the + // incoming id is already recorded in-window for this key. A replay + // (`dup = true`) short-circuits to a no-op below — matching the + // in-memory backend's `if !dedup_ids.contains(event_id)` guard. + let pre_row = tx + .query_one( + &sql.dedup_precount, + &[&key_ref, &cutoff, &event_id.as_str()], + ) + .await + .map_err(map_pg)?; + let pre_count: i64 = pre_row.get("cnt"); + // BOOL_OR over an empty set is NULL; treat NULL as "no duplicate". + let is_replay: bool = pre_row.get::<_, Option>("dup").unwrap_or(false); + + if is_replay { + // Replay refusal: the id is already in-window for this key, so + // this is a no-op against the count/sum. Short-circuit BEFORE + // the cap check so a replay at the cap dedups rather than + // overflowing — matching the in-memory contract. Aggregate and + // return the unchanged state. + let agg = self.aggregate(&tx, &plan, key_ref, &cutoff).await?; + tx.commit().await.map_err(map_pg)?; + return Ok(agg); + } + + // (3) Per-key sample cap — FAIL CLOSED. The id is new (not a + // replay). If the in-window count is already at the cap, refuse to + // insert and return `WindowOverflow` (PR #3635 followup / #3929). + // Silently dropping the oldest sample to make room would weaken cap + // enforcement and break replay refusal — so we fail closed, + // matching the in-memory backend's `if !dedup && len >= cap { Err }`. + if pre_count.max(0) as usize >= MAX_SAMPLES_PER_KEY { + // Roll back so the trim above (which freed aged-out dedup ids) + // is not committed independently of a rejected record; the + // caller observes a clean no-write overflow. + drop(tx); + return Err(PredicateBackendError::WindowOverflow { + key: label, + cap: MAX_SAMPLES_PER_KEY, + }); + } + + // (4) Insert the new row, deduping on the PRIMARY KEY + // (key_hash, event_id) as belt-and-suspenders against a concurrent + // racer that inserted the same id between our dedup check and here + // (the advisory lock serializes same-key writers, so this conflict + // is not expected, but ON CONFLICT keeps it a no-op if it occurs). + // The column list differs per typed table: the invocation table has + // no `value` column; the value table's `value` is NOT NULL. The + // value is read straight off the typed `RecordPlan::Value` variant — + // there is no nullable side-channel to unwrap. + match &plan { + RecordPlan::Invocation(_) => { + tx.execute( + &sql.insert, + &[&scope_ref, &key_ref, &event_id.as_str(), &now], + ) + .await + .map_err(map_pg)?; + } + RecordPlan::Value { value, .. } => { + tx.execute( + &sql.insert, + &[&scope_ref, &key_ref, &event_id.as_str(), &now, value], + ) + .await + .map_err(map_pg)?; + } + } + + // In-window sample count AFTER this call's insert. We DERIVE it as + // `pre_count + 1` rather than re-querying: we are under this key's + // per-key advisory lock (so no concurrent writer can add or remove a + // sample for this key), `is_replay` was false and the cap gate + // passed, and the `INSERT ... ON CONFLICT DO NOTHING` therefore added + // exactly one new in-window row (the dedup check already proved the + // id was absent, so the ON CONFLICT no-op branch cannot fire here). + // This eliminates a COUNT round trip on the record() hot path that + // would return the identical value (codex/henrypark perf finding). + let in_window_count: i64 = pre_count.max(0) + 1; + + // (5) Per-scope distinct-key LRU quota. Only scan when this key is + // newly material (count == 1 after insert — equivalently + // `pre_count == 0` — means we may have just created the scope's Nth + // key). Distinct keys are counted by key_hash within the scope (one + // typed table per kind, so no kind filter is needed); if over quota, + // evict the least-recently-active key's rows entirely. Scope-LRU + // eviction only ever touches OTHER keys, so it cannot change this + // key's aggregate and does not require a re-read. + let evicted = if in_window_count == 1 { + self.enforce_scope_quota(&tx, &plan, scope_ref, key_ref) + .await? + } else { + 0 + }; + + // Final returned aggregate. For the invocation table the aggregate is + // exactly the in-window COUNT, which we already hold as the derived + // `in_window_count` — so return it directly and skip a third COUNT + // round trip. The value table's aggregate is a `SUM(value)`, which is + // NOT derivable from the sample count, so it still issues one query. + let agg = match &plan { + RecordPlan::Invocation(_) => Decimal::from(in_window_count.max(0) as u64), + RecordPlan::Value { .. } => self.aggregate(&tx, &plan, key_ref, &cutoff).await?, + }; + + tx.commit().await.map_err(map_pg)?; + + if evicted > 0 { + // Mirror only on a successful commit so the monitoring counter + // never advances for a rolled-back eviction. + self.evictions.fetch_add(evicted, Ordering::Relaxed); + } + + Ok(agg) + } + + /// In-window aggregate for a key: `COUNT(*)` (as a `Decimal`) for the + /// invocation table, `SUM(value)` for the value table. Centralizes the + /// per-table aggregate SQL so the invocation table — which has no `value` + /// column — is never asked to `SUM(value)`. + async fn aggregate( + &self, + tx: &deadpool_postgres::Transaction<'_>, + plan: &RecordPlan, + key_ref: &[u8], + cutoff: &DateTime, + ) -> Result { + let sql = self.statements(plan); + match plan { + RecordPlan::Invocation(_) => { + let count: i64 = tx + .query_one(&sql.aggregate_count, &[&key_ref, &cutoff]) + .await + .map_err(map_pg)? + .get(0); + Ok(Decimal::from(count.max(0) as u64)) + } + RecordPlan::Value { .. } => { + let total: Decimal = tx + .query_one(&sql.aggregate_sum, &[&key_ref, &cutoff]) + .await + .map_err(map_pg)? + .get(0); + Ok(total) + } + } + } + + /// Enforce [`MAX_KEYS_PER_TENANT`] distinct keys per scope+kind. + /// Returns the number of keys evicted (0 or more). Eviction drops the + /// key whose OLDEST retained sample is oldest (`MIN(ts)` per key) — + /// the "oldest-front" victim selection, matching the in-memory backend + /// (which ranks buckets by their front/oldest entry) and the libSQL + /// backend. It never touches the key we just inserted. + /// + /// # Lock discipline for victim eviction (deadlock + race fix) + /// + /// Deleting a victim key's rows is itself a write to that bucket, so it + /// MUST participate in the per-key advisory-lock serialization just like + /// a `record` call would — otherwise this pass could delete rows for + /// key B while another transaction is concurrently recording B under + /// B's own per-key lock, producing a torn aggregate (the recorder's + /// COUNT/SUM straddling a delete it never serialized against) and a + /// deadlock: + /// + /// - Txn V (recording victim key B): holds B's per-key lock (it trimmed + /// B's out-of-window rows), then *blocks* waiting for the scope lock + /// inside its own `enforce_scope_quota`. + /// - Txn L (this LRU pass): holds the scope lock, then *blocks* waiting + /// on B's row locks to DELETE them. + /// - Cycle: V waits on the scope lock held by L; L waits on B's row + /// locks held by V → deadlock. + /// + /// The fix: before deleting a victim's rows, `pg_try_advisory_xact_lock` + /// that victim's per-key lock. The *try* variant returns immediately + /// (never blocks), so the cycle above can never form — L observes B's + /// lock held by V and skips B instead of waiting. Skipped (in-flight) + /// victims are passed over for the next-staleest candidate. + /// + /// # Quota is an explicit OUTCOME, never silently best-effort + /// + /// An earlier revision over-fetched a fixed candidate set, skipped + /// locked victims, and returned `Ok(evicted)` even if the scope was + /// STILL above [`MAX_KEYS_PER_TENANT`] — making the per-scope bound an + /// accident of how many victims happened to be in-flight (serrrfirat + /// BLOCKER on #3933). This loop closes that: it requeries the distinct + /// count + a fresh candidate batch each pass and keeps evicting until + /// **the cap is actually met**. If a whole pass makes zero progress + /// (every remaining stale candidate is locked by an in-flight recorder) + /// while the scope is still over cap, it FAILS CLOSED with + /// [`PredicateBackendError::Unavailable`] rather than committing an + /// over-quota scope. The caller's transaction rolls back, the record is + /// not applied, and the evaluator maps the error to a restrictive + /// outcome — the same fail-closed posture the per-key cap uses. The pass + /// budget bounds worst-case work so a pathological scope cannot spin. + async fn enforce_scope_quota( + &self, + tx: &deadpool_postgres::Transaction<'_>, + plan: &RecordPlan, + scope_ref: &[u8], + current_key: &[u8], + ) -> Result { + let sql = self.statements(plan); + // Serialize quota enforcement within the scope. Concurrent inserts + // of DISTINCT new keys in the same scope each reach this path with + // `count == 1`, but under READ COMMITTED neither sees the other's + // just-inserted row, so each would under-count `distinct` and + // under-evict — leaving the scope above the cap. A scope-level + // advisory lock makes the eviction check serial per scope, so the + // count is exact. It is taken in the SINGLE-arg `(int8)` advisory + // space, disjoint from the per-key `(int4,int4)` lock space. Hot-path + // same-key writes never reach here (only newly-material keys do), so + // this does not serialize steady-state traffic. + let scope_lock = scope_advisory_lock_key(scope_ref, plan.lock_tag()); + tx.execute("SELECT pg_advisory_xact_lock($1)", &[&scope_lock]) + .await + .map_err(map_pg)?; + + // Per-pass victim batch size. Bounded so a pathological scope can't + // pull an unbounded candidate set into memory in one query; the + // outer loop requeries for more if a pass exhausts its batch while + // still over quota. + const VICTIM_BATCH: i64 = 64; + // Worst-case pass budget. Each productive pass evicts at least one + // key, so the cap is met within `over_quota` productive passes; the + // budget additionally tolerates passes that make no progress because + // every candidate is momentarily locked, after which we fail closed + // rather than spin. This bound keeps the transaction from looping + // unboundedly under sustained contention. + const MAX_PASSES: usize = 1_024; + + let mut evicted = 0u64; + for _pass in 0..MAX_PASSES { + // Recompute the live distinct-key count under the scope lock. On + // the first pass this is the authoritative over-quota measure; + // on later passes it reflects the rows this loop already deleted, + // so the loop terminates exactly when the cap is met. + let distinct: i64 = tx + .query_one(&sql.scope_distinct, &[&scope_ref]) + .await + .map_err(map_pg)? + .get(0); + + if distinct as usize <= MAX_KEYS_PER_TENANT { + // Cap is met (either it never was over, or we evicted enough). + return Ok(evicted); + } + // Exact deficit to clear THIS pass — we must evict precisely this + // many keys, never the whole candidate batch, or we would + // over-evict below the cap (which would, on a flood, reset more of + // the tenant's surviving counters than the quota requires). + let to_evict = distinct as usize - MAX_KEYS_PER_TENANT; + + // Victim candidates: rank keys in this scope by their OLDEST + // retained sample (MIN(ts)) and evict the key whose oldest sample + // is oldest — "oldest-front" selection, matching the in-memory and + // libSQL backends. The in-memory backend ranks buckets by their + // front (oldest) entry's timestamp (`entries.front()` + + // `min_by_key`), so the durable analogue is MIN(ts) per key, NOT + // MAX(ts). Using MAX(ts) here would diverge: a key with one ancient + // sample and one fresh sample would be ranked by the fresh sample + // and spared, while the in-memory backend ranks it by the ancient + // sample and evicts it. The single-sample-per-key parity matrix + // masks this (MIN == MAX), but multi-sample keys would evict + // different keys across backends. + // + // Exclude the key we just inserted so a flood can never evict + // itself and mask the new entry. + let candidate_rows = tx + .query( + &sql.scope_candidates, + &[&scope_ref, ¤t_key, &VICTIM_BATCH], + ) + .await + .map_err(map_pg)?; + + // Evict candidates one at a time, each under its own per-key + // advisory lock taken with the NON-blocking + // `pg_try_advisory_xact_lock`. A victim whose lock is already held + // (a concurrent `record` is mid-flight against that bucket) is + // skipped — never waited on — so this pass cannot deadlock against + // a recorder, and it never deletes rows out from under a + // transaction that did not serialize against us. + let mut progressed = false; + let mut evicted_this_pass = 0usize; + for row in &candidate_rows { + if evicted_this_pass >= to_evict { + // Cleared this pass's deficit; stop before over-evicting. + break; + } + let victim_key: Vec = row.get(0); + let lock_key = advisory_lock_key_from_bytes(&victim_key); + let got_lock: bool = tx + .query_one( + "SELECT pg_try_advisory_xact_lock($1, $2)", + &[&lock_key.0, &lock_key.1], + ) + .await + .map_err(map_pg)? + .get(0); + if !got_lock { + // In-flight under its own per-key lock; skip and try the + // next-staleest candidate rather than block (deadlock-free). + continue; + } + tx.execute(&sql.evict_victim, &[&scope_ref, &victim_key]) + .await + .map_err(map_pg)?; + evicted += 1; + evicted_this_pass += 1; + progressed = true; + } + + if !progressed { + // The scope is still over quota AND every stale candidate is + // locked by an in-flight recorder, so we cannot bring the + // scope under the cap on this transaction's watch. Refuse to + // commit an over-quota scope: fail closed. The caller's + // transaction rolls back (this record is not applied) and the + // evaluator maps the error restrictively — same posture as the + // per-key cap's `WindowOverflow`. A retry (or a concurrent + // recorder finishing) clears the contention. The operational + // detail (the quota constant, the contention state) stays in + // the debug log; the caller-facing message is the sanitized + // constant, matching the `DB_UNAVAILABLE_MSG` contract so the + // evaluator only observes the error type, not the payload. + tracing::debug!( + max_keys_per_tenant = MAX_KEYS_PER_TENANT, + "scope quota enforcement contended: every stale eviction \ + candidate is locked by an in-flight recorder; failing \ + closed rather than committing an over-quota scope" + ); + return Err(PredicateBackendError::Unavailable( + QUOTA_CONTENDED_MSG.to_string(), + )); + } + } + + // Exhausted the pass budget while still over quota: treat the same as + // an unenforceable cap and fail closed rather than commit over-quota. + tracing::debug!( + max_keys_per_tenant = MAX_KEYS_PER_TENANT, + max_passes = MAX_PASSES, + "scope quota not met after eviction-pass budget exhausted; failing closed" + ); + Err(PredicateBackendError::Unavailable( + QUOTA_BUDGET_MSG.to_string(), + )) + } +} + +#[async_trait] +impl PredicateStateBackend for PostgresPredicateStateBackend { + async fn record_invocation( + &self, + key: &InvocationKey, + event_id: &PredicateEventId, + now: DateTime, + window: Duration, + ) -> Result { + let plan = RecordPlan::Invocation(PlanCommon { + scope: scope_hash(key.tenant_id.as_str()), + key: invocation_key_hash(key), + label: format!("{}/{}", key.tenant_id.as_str(), key.capability), + }); + let count = self.record(plan, event_id, now, window).await?; + // `record` returns the invocation count as a Decimal (COUNT(*), + // capped at MAX_SAMPLES_PER_KEY). Narrow to u32 via the integer + // value; the cap (4_096) guarantees it fits. + use rust_decimal::prelude::ToPrimitive; + let n = count.to_u32().unwrap_or(u32::MAX); + Ok(n) + } + + async fn record_value( + &self, + key: &ValueKey, + event_id: &PredicateEventId, + now: DateTime, + value: Decimal, + window: Duration, + ) -> Result { + let plan = RecordPlan::Value { + common: PlanCommon { + scope: scope_hash(key.tenant_id.as_str()), + key: value_key_hash(key), + label: format!( + "{}/{}#{}", + key.tenant_id.as_str(), + key.capability, + key.field + ), + }, + value, + }; + self.record(plan, event_id, now, window).await + } + + fn evictions_observed(&self) -> u64 { + self.evictions.load(Ordering::Relaxed) + } + + async fn evict_older_than(&self, cutoff: DateTime) -> Result { + let mut client = self.client().await?; + // Reap both typed tables atomically. Run both DELETEs inside one + // transaction so the reap is all-or-nothing: if the value DELETE + // fails after the invocation DELETE ran, the transaction rolls back + // (no commit) and BOTH tables are left untouched, so a retry reaps a + // consistent snapshot rather than finding invocation rows already + // gone and value rows still present. READ COMMITTED is sufficient — + // the reaper does not interleave with the per-key/scope advisory-lock + // protocol the record path uses (it only deletes already-stale rows), + // and the two DELETEs touch disjoint tables. Return the total rows + // deleted across both. + let tx = client + .build_transaction() + .isolation_level(IsolationLevel::ReadCommitted) + .start() + .await + .map_err(map_pg)?; + let inv = tx + .execute(&self.invocation_sql.reap, &[&cutoff]) + .await + .map_err(map_pg)?; + let val = tx + .execute(&self.value_sql.reap, &[&cutoff]) + .await + .map_err(map_pg)?; + tx.commit().await.map_err(map_pg)?; + Ok(inv + val) + } +} + +/// Derive the two-`i32` advisory-lock key from a bucket's `key_hash` bytes. +/// `pg_advisory_xact_lock(int4, int4)` namespaces the lock by the pair, so +/// we feed the first four bytes as the classifier and the next four as the +/// object id. A hash collision across distinct keys merely serializes two +/// unrelated buckets — a (rare) throughput cost, never a correctness bug. +/// +/// Both the recording path (which holds the typed [`Digest`], passed via +/// slice coercion) and the scope-LRU eviction path (which reads candidate +/// victims' `key_hash` back from the `BYTEA` column as raw bytes) call this +/// over the SAME bytes, so they derive the identical `(i32, i32)` lock key +/// for the same bucket — otherwise the victim try-lock would guard a +/// different lock than the recorder holds and the serialization would be +/// defeated. `key_hash` is always a 32-byte blake3 digest, so the first 8 +/// bytes are present; a shorter slice (never expected) is zero-padded so the +/// function is total rather than panicking on an out-of-range index. +fn advisory_lock_key_from_bytes(key: &[u8]) -> (i32, i32) { + let mut buf = [0u8; 8]; + let n = key.len().min(8); + buf[..n].copy_from_slice(&key[..n]); + let a = i32::from_le_bytes([buf[0], buf[1], buf[2], buf[3]]); + let b = i32::from_le_bytes([buf[4], buf[5], buf[6], buf[7]]); + (a, b) +} + +/// Derive the single-`i64` scope advisory-lock key. Uses the `(int8)` +/// advisory space, which Postgres keeps disjoint from the `(int4,int4)` +/// space used for per-key locks, so a key lock and a scope lock can never +/// alias each other. Folds in the `kind` byte so the invocation and value +/// maps lock independently. +fn scope_advisory_lock_key(scope: &[u8], kind: &[u8]) -> i64 { + let mut hasher = blake3::Hasher::new(); + hasher.update(scope); + hasher.update(kind); + let d = hasher.finalize(); + let bytes = d.as_bytes(); + i64::from_le_bytes([ + bytes[0], bytes[1], bytes[2], bytes[3], bytes[4], bytes[5], bytes[6], bytes[7], + ]) +} + +/// Sanitized message returned to callers for any database-layer failure. +/// The raw error (which can embed connection strings, host names, schema +/// details, or SQL fragments) is logged at `warn` for operators but NOT +/// surfaced through [`PredicateBackendError::Unavailable`], whose payload +/// can reach the evaluator/caller (henrypark security finding). The +/// evaluator only needs to know the backend is unavailable to fail closed; +/// it does not need the raw DB error text. +const DB_UNAVAILABLE_MSG: &str = "predicate state backend unavailable (database error)"; + +/// Sanitized message for the scope-quota contention fail-closed path. Mirrors +/// the `DB_UNAVAILABLE_MSG` contract: the operational detail (the quota +/// constant `MAX_KEYS_PER_TENANT`, the lock-contention state) is logged at +/// `debug` for operators but NOT surfaced through +/// [`PredicateBackendError::Unavailable`], whose payload can reach the +/// evaluator/caller. The evaluator only needs the error type to fail closed. +const QUOTA_CONTENDED_MSG: &str = + "predicate state backend unavailable (quota enforcement contended)"; + +/// Sanitized message for the scope-quota pass-budget-exhaustion fail-closed +/// path. Same sanitization posture as [`QUOTA_CONTENDED_MSG`]. +const QUOTA_BUDGET_MSG: &str = + "predicate state backend unavailable (quota enforcement budget exhausted)"; + +fn map_pg(e: tokio_postgres::Error) -> PredicateBackendError { + tracing::warn!(error = %e, "postgres predicate backend error"); + PredicateBackendError::Unavailable(DB_UNAVAILABLE_MSG.to_string()) +} + +fn map_pool(e: deadpool_postgres::PoolError) -> PredicateBackendError { + tracing::warn!(error = %e, "postgres predicate backend pool error"); + PredicateBackendError::Unavailable(DB_UNAVAILABLE_MSG.to_string()) +} + +#[cfg(test)] +mod tests { + use super::*; + + /// The scope-LRU eviction path try-locks each victim key over the + /// `key_hash` bytes it read back from the DB as an owned `Vec`, while + /// the recording path locks over the typed `Digest` via slice coercion. + /// Both go through `advisory_lock_key_from_bytes`, so the derivation is + /// the same function — but the inputs reach it by different paths + /// (`&digest[..]` slice coercion vs. an owned `Vec`). If those ever + /// produced different keys the eviction would guard a DIFFERENT advisory + /// lock than a concurrent recorder holds, defeating the per-bucket + /// serialization the fix depends on. This pins them equal for the full + /// 32-byte digest — a provable-by-inspection guard for the + /// lock-acquisition invariant that does not need a live Postgres. + #[test] + fn eviction_and_record_derive_identical_per_key_lock() { + let digest: Digest = { + let mut d = [0u8; 32]; + for (i, b) in d.iter_mut().enumerate() { + *b = (i as u8).wrapping_mul(7).wrapping_add(3); + } + d + }; + // Recorder path: a typed `Digest` coerced to `&[u8]`. + let from_digest = advisory_lock_key_from_bytes(&digest); + // Eviction path: the same bytes read back from the DB as an owned vec. + let victim_bytes: Vec = digest.to_vec(); + let from_bytes = advisory_lock_key_from_bytes(&victim_bytes); + assert_eq!( + from_digest, from_bytes, + "victim try-lock key must equal the recorder's lock key for the same bucket" + ); + } + + /// The `Err(_) => now` arm of `cutoff` fires when `window` exceeds + /// chrono's maximum `Duration` (e.g. `Duration::MAX`). In that case + /// `cutoff` saturates to `now`, so the entire window is treated as + /// in-scope and nothing is trimmed — the conservative trim-nothing posture + /// for a rate/value cap. This is a pure `fn cutoff(now, window)` on the + /// struct, exercisable via `use super::*` with no live Postgres. + #[test] + fn cutoff_with_overflow_window_saturates_to_now() { + let now = DateTime::from_timestamp(1_700_000_000, 0).expect("static timestamp is in range"); + // `Duration::MAX` exceeds chrono's max i64-nanosecond `Duration`, so + // `chrono::Duration::from_std` returns `Err` and `cutoff` saturates. + assert_eq!( + PostgresPredicateStateBackend::cutoff(now, Duration::MAX), + now + ); + } + + /// Distinct buckets must (almost always) map to distinct per-key lock + /// keys; a single differing leading byte must change the derived lock. + #[test] + fn distinct_digests_yield_distinct_lock_keys() { + let mut a = [0u8; 32]; + let mut b = [0u8; 32]; + a[0] = 1; + b[0] = 2; + assert_ne!( + advisory_lock_key_from_bytes(&a), + advisory_lock_key_from_bytes(&b) + ); + } + + /// The invocation (`b"i"`) and value (`b"v"`) scope-quota passes MUST take + /// distinct scope advisory locks for the same tenant, or every value-table + /// scope-quota pass would serialize behind invocation-table passes for that + /// tenant (they touch disjoint tables and must not block each other). The + /// `lock_tag` byte folded into `scope_advisory_lock_key` is what keeps them + /// disjoint. This pins that invariant on a pure function with no live + /// Postgres: the two tags over the same scope bytes derive different keys. + #[test] + fn invocation_and_value_scope_lock_keys_are_distinct() { + let scope: [u8; 32] = { + let mut s = [0u8; 32]; + for (i, b) in s.iter_mut().enumerate() { + *b = (i as u8).wrapping_mul(11).wrapping_add(5); + } + s + }; + let inv = RecordPlan::Invocation(PlanCommon { + scope, + key: [0u8; 32], + label: String::new(), + }); + let val = RecordPlan::Value { + common: PlanCommon { + scope, + key: [0u8; 32], + label: String::new(), + }, + value: Decimal::ZERO, + }; + assert_eq!(inv.lock_tag(), b"i"); + assert_eq!(val.lock_tag(), b"v"); + assert_ne!( + scope_advisory_lock_key(&scope, inv.lock_tag()), + scope_advisory_lock_key(&scope, val.lock_tag()), + "invocation and value scope-quota passes must derive disjoint scope \ + advisory locks for the same tenant" + ); + } + + /// `advisory_lock_key_from_bytes` is total: a short slice (never + /// expected from a 32-byte `key_hash`, but defensive) zero-pads rather + /// than panicking on an out-of-range index. + #[test] + fn short_slice_zero_pads_without_panic() { + assert_eq!(advisory_lock_key_from_bytes(&[]), (0, 0)); + assert_eq!( + advisory_lock_key_from_bytes(&[0xFF]), + (i32::from_le_bytes([0xFF, 0, 0, 0]), 0) + ); + } +} diff --git a/crates/ironclaw_hooks_postgres/src/hashing.rs b/crates/ironclaw_hooks_postgres/src/hashing.rs new file mode 100644 index 00000000000..fbd2a25dd7d --- /dev/null +++ b/crates/ironclaw_hooks_postgres/src/hashing.rs @@ -0,0 +1,127 @@ +//! Canonical hashing of predicate bucket identities into fixed-width +//! BYTEA index keys. +//! +//! Two digests are derived per bucket: +//! +//! - `scope_hash` = blake3 over the length-prefixed `tenant_id`. This is +//! the trust boundary (one tenant's counters never affect another's) +//! and the grain at which the distinct-key LRU quota is enforced. +//! - `key_hash` = blake3 over a length-prefixed canonical serialization +//! of the *whole* bucket identity, including a one-byte map +//! discriminant so an invocation key and a value key that share +//! `(hook, tenant, capability)` never collide. +//! +//! Length-prefixing every field (8-byte big-endian length ++ bytes) +//! makes the serialization injective: `("ab", "c")` and `("a", "bc")` +//! produce distinct digests, closing the classic concatenation-collision +//! hole that a naive `a ++ b` would leave open. + +use ironclaw_hooks::predicate_state::{InvocationKey, ValueKey}; + +/// Map discriminants folded into `key_hash` so the invocation and value +/// maps share a table without cross-contaminating dedup. +const KIND_INVOCATION: u8 = b'i'; +const KIND_VALUE: u8 = b'v'; + +/// 32-byte blake3 digest stored as `BYTEA`. +pub(crate) type Digest = [u8; 32]; + +fn feed(hasher: &mut blake3::Hasher, field: &[u8]) { + // 8-byte length prefix makes the field boundary unambiguous; u64 makes + // the usize->len conversion infallible on all supported platforms (no + // saturation corner that could alias two fields differing only beyond + // u32::MAX), keeping the serialization strictly injective. + let len = field.len() as u64; + hasher.update(&len.to_be_bytes()); + hasher.update(field); +} + +/// `scope_hash` for a tenant — the LRU-quota and trust grain. +pub(crate) fn scope_hash(tenant_id: &str) -> Digest { + let mut hasher = blake3::Hasher::new(); + feed(&mut hasher, tenant_id.as_bytes()); + *hasher.finalize().as_bytes() +} + +/// `key_hash` for an invocation-counter bucket. +pub(crate) fn invocation_key_hash(key: &InvocationKey) -> Digest { + let mut hasher = blake3::Hasher::new(); + hasher.update(&[KIND_INVOCATION]); + feed(&mut hasher, key.hook_id.as_bytes()); + feed(&mut hasher, key.tenant_id.as_str().as_bytes()); + feed(&mut hasher, key.capability.as_bytes()); + *hasher.finalize().as_bytes() +} + +/// `key_hash` for a numeric-value-sum bucket. +pub(crate) fn value_key_hash(key: &ValueKey) -> Digest { + let mut hasher = blake3::Hasher::new(); + hasher.update(&[KIND_VALUE]); + feed(&mut hasher, key.hook_id.as_bytes()); + feed(&mut hasher, key.tenant_id.as_str().as_bytes()); + feed(&mut hasher, key.capability.as_bytes()); + feed(&mut hasher, key.field.as_bytes()); + *hasher.finalize().as_bytes() +} + +#[cfg(test)] +mod tests { + use super::*; + use ironclaw_hooks::identity::{ExtensionId, HookId, HookLocalId, HookVersion}; + use ironclaw_host_api::TenantId; + + fn hook() -> HookId { + HookId::derive( + &ExtensionId::new("ext").unwrap(), + "1.0", + &HookLocalId::new("h").unwrap(), + HookVersion::ONE, + ) + } + + fn inv(tenant: &str, capability: &str) -> InvocationKey { + InvocationKey { + hook_id: hook(), + tenant_id: TenantId::new(tenant).unwrap(), + capability: capability.to_string(), + } + } + + fn val(tenant: &str, capability: &str, field: &str) -> ValueKey { + ValueKey { + hook_id: hook(), + tenant_id: TenantId::new(tenant).unwrap(), + capability: capability.to_string(), + field: field.to_string(), + } + } + + #[test] + fn distinct_tenants_have_distinct_scope_hashes() { + assert_ne!(scope_hash("alpha"), scope_hash("beta")); + } + + #[test] + fn invocation_and_value_keys_never_collide() { + // Same hook/tenant/capability across the two maps must hash apart + // because of the map discriminant. + let i = invocation_key_hash(&inv("t", "cap.x")); + let v = value_key_hash(&val("t", "cap.x", "cap.x")); + assert_ne!(i, v); + } + + #[test] + fn field_boundary_is_injective() { + // Length-prefixing prevents ("ab","c") and ("a","bc") aliasing. + let a = value_key_hash(&val("t", "ab", "c")); + let b = value_key_hash(&val("t", "a", "bc")); + assert_ne!(a, b); + } + + #[test] + fn capability_boundary_is_injective_for_invocations() { + let a = invocation_key_hash(&inv("t", "abc")); + let b = invocation_key_hash(&inv("tabc", "")); + assert_ne!(a, b); + } +} diff --git a/crates/ironclaw_hooks_postgres/src/lib.rs b/crates/ironclaw_hooks_postgres/src/lib.rs new file mode 100644 index 00000000000..091fcbe4e18 --- /dev/null +++ b/crates/ironclaw_hooks_postgres/src/lib.rs @@ -0,0 +1,72 @@ +//! Durable PostgreSQL-backed [`PredicateStateBackend`]. +//! +//! This is durable-backend PR 2/4 in the predicate-state split. It +//! implements the *exact same* trait contract as the in-memory backend +//! ([`ironclaw_hooks::predicate_state::InMemoryPredicateStateBackend`]) +//! and is proven against the shared contract harness (see +//! `tests/predicate_state_postgres_contract.rs`). All eight contract +//! functions plus adversarial multi-host tests run against this impl. +//! +//! # Why a separate crate +//! +//! The trait was widened to `pub` in PR 1/4 specifically so durable +//! backends live *out of crate* and depend on `ironclaw_hooks` with the +//! `contract-tests` feature — exactly the way +//! [`ironclaw_reborn_event_store`] is a separate per-domain durable crate +//! rather than living in `ironclaw_events`. Keeping the Postgres +//! dependency surface out of `ironclaw_hooks` keeps the hook framework +//! itself DB-free. +//! +//! # Dual-backend compliance (libSQL counterpart) +//! +//! The repo rule "new persistence must support both PostgreSQL and +//! libSQL" is satisfied across the staged durable-backend series, not in +//! this crate alone. This crate is the **Postgres** half; the **libSQL** +//! counterpart lives in the sibling crate `ironclaw_hooks_libsql` (PR +//! #3936), and behavioral interchangeability between the two is enforced +//! by the cross-backend parity suite `ironclaw_hooks_parity` (PR #3937), +//! which runs the identical contract assertions against both backends. +//! Merge ordering: this crate (PR 2/4, #3933) lands first, then the +//! libSQL crate (#3936), then the parity suite (#3937) gates that the two +//! durable backends are drop-in equivalent. Both backends share the one +//! logical two-table typed schema (see `migrations/V1__predicate_state.sql`); +//! only the storage column types differ (Postgres TIMESTAMPTZ + NUMERIC +//! vs. libSQL epoch-ms INTEGER + TEXT). +//! +//! [`PredicateStateBackend`]: +//! ironclaw_hooks::predicate_state::PredicateStateBackend +//! [`ironclaw_reborn_event_store`]: https://docs.rs/ironclaw_reborn_event_store + +#[cfg(feature = "postgres")] +mod backend; +#[cfg(feature = "postgres")] +mod hashing; +#[cfg(feature = "postgres")] +mod schema; + +#[cfg(feature = "postgres")] +pub use backend::PostgresPredicateStateBackend; +#[cfg(feature = "postgres")] +pub use schema::POSTGRES_PREDICATE_SCHEMA; + +/// Test-only accessors for the crate-internal bucket hashing, used by the +/// adversarial integration tests to compute the `key_hash` / `scope_hash` +/// bytes a bucket maps to so a test can query rows directly. Not part of the +/// public API surface (this crate is `publish = false`); kept out of the +/// rendered docs. Production code must never depend on this module. +#[cfg(feature = "postgres")] +#[doc(hidden)] +pub mod test_support { + use ironclaw_hooks::predicate_state::InvocationKey; + + /// `key_hash` bytes for an invocation bucket — see + /// `crate::hashing::invocation_key_hash`. + pub fn invocation_key_hash_bytes(key: &InvocationKey) -> [u8; 32] { + crate::hashing::invocation_key_hash(key) + } + + /// `scope_hash` bytes for a tenant — see `crate::hashing::scope_hash`. + pub fn scope_hash_bytes(tenant_id: &str) -> [u8; 32] { + crate::hashing::scope_hash(tenant_id) + } +} diff --git a/crates/ironclaw_hooks_postgres/src/schema.rs b/crates/ironclaw_hooks_postgres/src/schema.rs new file mode 100644 index 00000000000..ebeb6dd8d10 --- /dev/null +++ b/crates/ironclaw_hooks_postgres/src/schema.rs @@ -0,0 +1,57 @@ +//! Embedded idempotent schema for the durable predicate backend. +//! +//! The DDL is the *single source of truth* file +//! `migrations/V1__predicate_state.sql`, pulled in verbatim at compile +//! time via `include_str!`. There is no hand-maintained second copy to +//! drift against. It is applied via +//! [`PostgresPredicateStateBackend::run_migrations`] using +//! a single `batch_execute` (which tolerates the file's `--` SQL +//! comments), the same per-crate pattern +//! `ironclaw_filesystem::PostgresRootFilesystem::run_migrations` uses. We +//! deliberately do NOT route through the legacy main-binary refinery +//! `migrations/` directory: that system is scoped to `src/db/` and the +//! reborn durable crates each own their schema. +//! +//! # DB-clock decision (cross-host correctness) +//! +//! The trait passes `now: DateTime`. There are two candidate clocks +//! for the *window comparison basis*: +//! +//! 1. The caller's `now` (stored in `ts`, compared against a +//! caller-computed `cutoff`). +//! 2. The database's `NOW()`. +//! +//! We use **the caller's `now`** as the comparison basis — `ts < cutoff` +//! where `cutoff = now - window` is computed host-side exactly as the +//! in-memory backend does. This is the choice that makes the Postgres +//! backend a *drop-in* for the in-memory backend under the shared +//! contract harness: the contract tests drive a deterministic fixed +//! clock (`at(0)`, `at(60)`, …) and assert exact counts at the window +//! boundary. If we substituted `NOW()` for the comparison basis those +//! tests could not pin a deterministic result, and a host whose clock +//! the operator already trusts (the same `Utc::now()` the in-memory +//! backend trusts) would silently disagree with the DB clock. +//! +//! The trade-off this accepts: cross-host window correctness now depends +//! on the hosts' wall clocks being roughly synchronized (NTP), the same +//! assumption the rest of the system makes for `occurred_at` timestamps. +//! The load-bearing cross-host property — *replay dedup* — does NOT +//! depend on clock agreement: it is enforced by the +//! `PRIMARY KEY (key_hash, event_id)` constraint and `ON CONFLICT DO NOTHING`, +//! which is exact regardless of clock skew. Atomicity is enforced by +//! running prune + dedup-check + insert + aggregate inside one +//! `READ COMMITTED` transaction guarded by a per-key advisory lock, also +//! clock-independent. + +/// Idempotent schema applied by `run_migrations()`. Sourced directly from +/// `migrations/V1__predicate_state.sql` via `include_str!` so the file +/// is the only copy — no embedded duplicate can drift out of sync. +pub const POSTGRES_PREDICATE_SCHEMA: &str = include_str!("../migrations/V1__predicate_state.sql"); + +/// Table holding invocation-count samples (one row per recorded invocation +/// event; the in-window `COUNT(*)` is the invocation count). +pub(crate) const INVOCATIONS_TABLE: &str = "hooks_predicate_invocations"; + +/// Table holding numeric-value samples (one row per recorded value event; +/// the in-window `SUM(value)` is the running sum). +pub(crate) const VALUES_TABLE: &str = "hooks_predicate_values"; diff --git a/crates/ironclaw_hooks_postgres/tests/predicate_state_postgres_adversarial.rs b/crates/ironclaw_hooks_postgres/tests/predicate_state_postgres_adversarial.rs new file mode 100644 index 00000000000..fd79bccd142 --- /dev/null +++ b/crates/ironclaw_hooks_postgres/tests/predicate_state_postgres_adversarial.rs @@ -0,0 +1,895 @@ +//! Adversarial / multi-host tests for [`PostgresPredicateStateBackend`]. +//! +//! These prove the durable backend's cross-host correctness properties +//! that the in-memory backend explicitly does NOT provide (its dedup is +//! process-local). Each "host" is a separate `deadpool` pool over the +//! same database, simulating distinct processes pointing at one Postgres. +//! +//! Gated on a reachable Postgres via `IRONCLAW_HOOKS_POSTGRES_URL` / +//! `DATABASE_URL`; skipped (passing) otherwise — same env-gate pattern as +//! the contract suite. Serialized behind a process-global lock because +//! they share fixed keys against one table. + +#![cfg(feature = "postgres")] + +use std::sync::Arc; +use std::time::Duration; + +use chrono::{DateTime, Utc}; +use deadpool_postgres::Pool; +use ironclaw_hooks::identity::{ExtensionId, HookId, HookLocalId, HookVersion}; +use ironclaw_hooks::predicate_state::{ + InvocationKey, MAX_KEYS_PER_TENANT, MAX_SAMPLES_PER_KEY, PredicateBackendError, + PredicateEventId, PredicateStateBackend, ValueKey, +}; +use ironclaw_hooks_postgres::PostgresPredicateStateBackend; +use ironclaw_host_api::TenantId; +use rust_decimal::Decimal; + +static TEST_LOCK: std::sync::Mutex<()> = std::sync::Mutex::new(()); + +fn db_url() -> Option { + std::env::var("IRONCLAW_HOOKS_POSTGRES_URL") + .or_else(|_| std::env::var("DATABASE_URL")) + .ok() +} + +/// Dedicated schema so this binary cannot collide with the contract-test +/// binary that `cargo test` runs in parallel against the same database. +const TEST_SCHEMA: &str = "hooks_predicate_adversarial_test"; + +fn build_pool(url: &str) -> Option { + let config = url.parse::().ok()?; + let manager = deadpool_postgres::Manager::new(config, tokio_postgres::NoTls); + deadpool_postgres::Pool::builder(manager) + .max_size(16) + .post_create(deadpool_postgres::Hook::async_fn(|client, _| { + Box::pin(async move { + client + .batch_execute(&format!("SET search_path TO {TEST_SCHEMA}")) + .await + .map_err(|e| deadpool_postgres::HookError::message(e.to_string()))?; + Ok(()) + }) + })) + .build() + .ok() +} + +/// Build N independent backends ("hosts") over the same DB, ensure schema, +/// and truncate once so the table starts empty. +async fn hosts(url: &str, n: usize) -> Vec> { + // Ensure the isolated schema exists before any pooled connection sets + // its search_path to it. + { + let (client, conn) = tokio_postgres::connect(url, tokio_postgres::NoTls) + .await + .expect("connect"); + tokio::spawn(conn); + client + .batch_execute(&format!("CREATE SCHEMA IF NOT EXISTS {TEST_SCHEMA}")) + .await + .expect("create schema"); + } + let mut out = Vec::with_capacity(n); + for i in 0..n { + let pool = build_pool(url).expect("pool"); + let backend = PostgresPredicateStateBackend::new(pool.clone()); + backend.run_migrations().await.expect("migrate"); + if i == 0 { + let client = pool.get().await.expect("client"); + client + .batch_execute("TRUNCATE TABLE hooks_predicate_invocations, hooks_predicate_values") + .await + .expect("truncate"); + } + out.push(Arc::new(backend)); + } + out +} + +fn hook() -> HookId { + HookId::derive( + &ExtensionId::new("ext").unwrap(), + "1.0", + &HookLocalId::new("h").unwrap(), + HookVersion::ONE, + ) +} + +fn inv_key(tenant: &str, capability: &str) -> InvocationKey { + InvocationKey { + hook_id: hook(), + tenant_id: TenantId::new(tenant).unwrap(), + capability: capability.to_string(), + } +} + +fn val_key(tenant: &str, capability: &str, field: &str) -> ValueKey { + ValueKey { + hook_id: hook(), + tenant_id: TenantId::new(tenant).unwrap(), + capability: capability.to_string(), + field: field.to_string(), + } +} + +fn ev(s: &str) -> PredicateEventId { + PredicateEventId::new(s).expect("valid event id") +} + +fn base() -> DateTime { + DateTime::from_timestamp(1_700_000_000, 0).unwrap() +} + +macro_rules! guarded { + () => {{ + let Some(url) = db_url() else { + eprintln!("skipping postgres adversarial test: no DB URL set"); + return; + }; + // Lock recovered-on-poison; serialize across per-test runtimes. + let guard = TEST_LOCK.lock().unwrap_or_else(|p| p.into_inner()); + (url, guard) + }}; +} + +/// Two hosts hammering the SAME key with distinct event ids must produce +/// a count equal to the total number of distinct ids — no lost-update +/// desync. This exercises the single-transaction atomic record-and-read +/// across two connection pools. +#[tokio::test] +#[allow(clippy::await_holding_lock)] +async fn two_hosts_write_storm_no_count_desync() { + let (url, _guard) = guarded!(); + let hs = hosts(&url, 2).await; + let key = inv_key("storm-tenant", "cap.storm"); + let window = Duration::from_secs(3600); + let now = base(); + + const PER_HOST: usize = 100; + let mut handles = Vec::new(); + for (h, backend) in hs.iter().enumerate() { + for i in 0..PER_HOST { + let backend = Arc::clone(backend); + let key = key.clone(); + let id = ev(&format!("h{h}-e{i}")); + handles.push(tokio::spawn(async move { + backend + .record_invocation(&key, &id, now, window) + .await + .expect("record ok") + })); + } + } + for handle in handles { + handle.await.expect("join"); + } + + // Final count observed via a duplicate-id no-op read on host 0. + let final_count = hs[0] + .record_invocation(&key, &ev("h0-e0"), now, window) + .await + .expect("read ok"); + assert_eq!( + final_count as usize, + 2 * PER_HOST, + "every distinct-id write across both hosts must be counted exactly once" + ); +} + +/// Two hosts each record the SAME event id for the same key. Cross-host +/// replay dedup (the PRIMARY KEY + ON CONFLICT) must count it once. +#[tokio::test] +#[allow(clippy::await_holding_lock)] +async fn cross_host_replay_counts_once() { + let (url, _guard) = guarded!(); + let hs = hosts(&url, 2).await; + let key = inv_key("replay-tenant", "cap.replay"); + let window = Duration::from_secs(3600); + let now = base(); + let id = ev("shared-event-X"); + + let c_a = hs[0] + .record_invocation(&key, &id, now, window) + .await + .expect("host A"); + let c_b = hs[1] + .record_invocation(&key, &id, now, window) + .await + .expect("host B"); + + assert_eq!(c_a, 1, "host A records the id fresh"); + assert_eq!( + c_b, 1, + "host B replaying the same id must NOT double-count (cross-host dedup)" + ); +} + +/// Cross-host replay on the value path: the running sum must reflect a +/// single contribution even though two hosts recorded the same id. +#[tokio::test] +#[allow(clippy::await_holding_lock)] +async fn cross_host_value_replay_sums_once() { + let (url, _guard) = guarded!(); + let hs = hosts(&url, 2).await; + let key = val_key("replay-tenant", "cap.spend", "amount"); + let window = Duration::from_secs(3600); + let now = base(); + let id = ev("shared-value-X"); + + let s_a = hs[0] + .record_value(&key, &id, now, Decimal::from(50), window) + .await + .expect("host A"); + let s_b = hs[1] + .record_value(&key, &id, now, Decimal::from(50), window) + .await + .expect("host B"); + + assert_eq!(s_a, Decimal::from(50)); + assert_eq!( + s_b, + Decimal::from(50), + "duplicate id from a second host must not double the sum" + ); +} + +/// Per-key sample cap under a flood: filling a key to `MAX_SAMPLES_PER_KEY` +/// with distinct in-window ids succeeds, and the next distinct id FAILS +/// CLOSED with `WindowOverflow` rather than silently dropping the oldest +/// sample (PR #3635 followup / #3929). A replay of an already-recorded +/// in-window id at the cap still dedups to a no-op. This mirrors the +/// in-memory `record_invocation_overflow_is_fail_closed` contract against +/// the durable backend. +#[tokio::test] +#[allow(clippy::await_holding_lock)] +async fn per_key_sample_cap_fails_closed_under_flood() { + let (url, _guard) = guarded!(); + let hs = hosts(&url, 1).await; + let backend = &hs[0]; + let key = inv_key("flood-tenant", "cap.hot"); + let window = Duration::from_secs(86_400); + + // Fill exactly to the cap with distinct in-window ids — all succeed. + for i in 0..MAX_SAMPLES_PER_KEY { + let ts = base() + chrono::Duration::milliseconds(i as i64); + let count = backend + .record_invocation(&key, &ev(&format!("flood-{i}")), ts, window) + .await + .expect("inserts up to the cap succeed"); + assert_eq!(count as usize, i + 1); + } + + // The next distinct in-window id must fail closed, not silent-evict. + let overflow_ts = base() + chrono::Duration::milliseconds(MAX_SAMPLES_PER_KEY as i64); + let result = backend + .record_invocation(&key, &ev("flood-overflow"), overflow_ts, window) + .await; + assert!( + matches!(result, Err(PredicateBackendError::WindowOverflow { .. })), + "hitting the per-key cap must fail closed, got {result:?}" + ); + + // A replay of an in-window id at the cap dedups to a no-op rather than + // overflowing — replay refusal survives the cap boundary. + let replay_ts = base() + chrono::Duration::milliseconds(MAX_SAMPLES_PER_KEY as i64 + 1); + let replay = backend + .record_invocation(&key, &ev("flood-0"), replay_ts, window) + .await + .expect("replay of an in-window id must dedup, not overflow"); + assert_eq!( + replay as usize, MAX_SAMPLES_PER_KEY, + "replay at the cap is a no-op against the count" + ); +} + +/// Per-key sample cap under a flood on the VALUE path: the cap-reject branch +/// in the shared `record()` body runs identically for `record_value`, but the +/// value-table INSERT and `SUM(value)` aggregate are variant-specific, so a +/// regression there (e.g. the value INSERT failing to dedup, or `aggregate_sum` +/// miscounting) could silently admit over-cap values while the invocation path +/// stays correct. This mirrors `per_key_sample_cap_fails_closed_under_flood` +/// against `record_value`: fill exactly to the cap with distinct in-window +/// ids, assert the next distinct id at `MAX_SAMPLES_PER_KEY+1` fails closed +/// with `WindowOverflow`, and assert an in-window replay at the cap still +/// dedups (sum unchanged) rather than overflowing. +/// +/// Gated on a reachable Postgres via `IRONCLAW_HOOKS_POSTGRES_URL` / +/// `DATABASE_URL`; skipped (passing) otherwise. +#[tokio::test] +#[allow(clippy::await_holding_lock)] +async fn per_value_key_sample_cap_fails_closed_under_flood() { + let (url, _guard) = guarded!(); + let hs = hosts(&url, 1).await; + let backend = &hs[0]; + let key = val_key("flood-value-tenant", "cap.spend", "amount"); + let window = Duration::from_secs(86_400); + + // Each sample contributes a fixed amount so the running sum is a direct + // multiple of the in-window sample count — lets us assert the sum tracks + // the count exactly as we fill to the cap. + let per_sample = Decimal::from(2); + for i in 0..MAX_SAMPLES_PER_KEY { + let ts = base() + chrono::Duration::milliseconds(i as i64); + let sum = backend + .record_value(&key, &ev(&format!("vflood-{i}")), ts, per_sample, window) + .await + .expect("value inserts up to the cap succeed"); + assert_eq!( + sum, + per_sample * Decimal::from(i + 1), + "running sum must track the in-window sample count up to the cap" + ); + } + + // The next distinct in-window id must fail closed, not silent-evict. + let overflow_ts = base() + chrono::Duration::milliseconds(MAX_SAMPLES_PER_KEY as i64); + let result = backend + .record_value( + &key, + &ev("vflood-overflow"), + overflow_ts, + per_sample, + window, + ) + .await; + assert!( + matches!(result, Err(PredicateBackendError::WindowOverflow { .. })), + "hitting the per-value-key cap must fail closed, got {result:?}" + ); + + // A replay of an in-window id at the cap dedups to a no-op against the + // sum rather than overflowing — replay refusal survives the cap boundary. + let replay_ts = base() + chrono::Duration::milliseconds(MAX_SAMPLES_PER_KEY as i64 + 1); + let replay = backend + .record_value(&key, &ev("vflood-0"), replay_ts, per_sample, window) + .await + .expect("replay of an in-window value id must dedup, not overflow"); + assert_eq!( + replay, + per_sample * Decimal::from(MAX_SAMPLES_PER_KEY), + "replay at the cap is a no-op against the sum" + ); +} + +/// `evictions_observed()` increments only AFTER `tx.commit()`. A per-key cap +/// overflow rolls the transaction back via `drop(tx)` (it never commits), so +/// the monitoring counter must NOT advance on that path. A bug crediting an +/// eviction on rollback would emit spurious telemetry. This drives a single +/// key to `MAX_SAMPLES_PER_KEY+1` (forcing the `WindowOverflow` rollback) and +/// asserts `evictions_observed()` is unchanged from the pre-overflow snapshot. +/// +/// Gated on a reachable Postgres via `IRONCLAW_HOOKS_POSTGRES_URL` / +/// `DATABASE_URL`; skipped (passing) otherwise. +#[tokio::test] +#[allow(clippy::await_holding_lock)] +async fn evictions_counter_unchanged_on_window_overflow() { + let (url, _guard) = guarded!(); + let hs = hosts(&url, 1).await; + let backend = &hs[0]; + let key = inv_key("overflow-evict-tenant", "cap.hot"); + let window = Duration::from_secs(86_400); + + // Fill exactly to the cap — these are single-sample-per-key-free inserts + // against ONE key, so no scope-LRU eviction fires (one distinct key, well + // under MAX_KEYS_PER_TENANT). + for i in 0..MAX_SAMPLES_PER_KEY { + let ts = base() + chrono::Duration::milliseconds(i as i64); + backend + .record_invocation(&key, &ev(&format!("ovf-{i}")), ts, window) + .await + .expect("inserts up to the cap succeed"); + } + + // Snapshot the eviction counter immediately before the overflow. + let before = backend.evictions_observed(); + + // Drive the key one past the cap: this must roll back via `drop(tx)`. + let overflow_ts = base() + chrono::Duration::milliseconds(MAX_SAMPLES_PER_KEY as i64); + let result = backend + .record_invocation(&key, &ev("ovf-overflow"), overflow_ts, window) + .await; + assert!( + matches!(result, Err(PredicateBackendError::WindowOverflow { .. })), + "the over-cap insert must fail closed, got {result:?}" + ); + + let after = backend.evictions_observed(); + assert_eq!( + before, after, + "evictions_observed() must NOT advance when a per-key cap overflow \ + rolls the transaction back (counter increments only after commit)" + ); +} + +/// Per-scope (tenant) LRU quota under concurrent insert pressure across +/// two hosts: a single tenant's distinct-key footprint must be bounded at +/// `MAX_KEYS_PER_TENANT`, and the eviction counter must advance. +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] +#[allow(clippy::await_holding_lock)] +async fn per_scope_lru_eviction_bounds_distinct_keys() { + let (url, _guard) = guarded!(); + let hs = hosts(&url, 2).await; + let window = Duration::from_secs(3600); + + // Drive distinct keys past the per-tenant quota from two hosts at + // once. Each key gets one row; strictly increasing ts so LRU victim + // selection is deterministic. + let total = MAX_KEYS_PER_TENANT + 50; + let mut handles = Vec::new(); + for i in 0..total { + let backend = Arc::clone(&hs[i % 2]); + let key = inv_key("lru-tenant", &format!("cap.{i}")); + let ts = base() + chrono::Duration::milliseconds(i as i64); + handles.push(tokio::spawn(async move { + backend + .record_invocation(&key, &ev(&format!("lru-e{i}")), ts, window) + .await + .expect("record") + })); + } + for handle in handles { + handle.await.expect("join"); + } + + // Assert the per-scope bound holds for EVERY scope present. Other + // test binaries may share this database concurrently, so we check the + // maximum distinct-key count across all scopes rather than a global + // total — the LRU quota is per-scope, and no scope may exceed it. + let pool = build_pool(&url).expect("pool"); + let client = pool.get().await.expect("client"); + let row = client + .query_one( + "SELECT COALESCE(MAX(kc), 0)::BIGINT FROM ( + SELECT COUNT(DISTINCT key_hash) AS kc + FROM hooks_predicate_invocations + GROUP BY scope_hash + ) per_scope", + &[], + ) + .await + .expect("count"); + let max_per_scope: i64 = row.get(0); + assert!( + max_per_scope as usize <= MAX_KEYS_PER_TENANT, + "per-scope LRU must bound distinct keys at MAX_KEYS_PER_TENANT for every scope; \ + worst scope had {max_per_scope}" + ); + let evictions: u64 = hs.iter().map(|h| h.evictions_observed()).sum(); + assert!( + evictions >= 1, + "LRU eviction counter must advance when the per-scope quota is exceeded" + ); +} + +/// Regression: scope-LRU eviction must take each victim key's per-key +/// advisory lock (via the non-blocking `pg_try_advisory_xact_lock`) before +/// deleting its rows. Before the fix, `enforce_scope_quota` deleted victim +/// rows holding ONLY the scope lock, which produced two failures: +/// +/// 1. **Deadlock** — a transaction recording the victim key holds that +/// key's per-key lock and then waits for the scope lock inside its own +/// quota pass, while the LRU transaction holds the scope lock and waits +/// on the victim's row locks. Cycle → deadlock. +/// 2. **Torn aggregate** — the LRU pass could delete rows for a key while +/// another transaction was actively aggregating that same key under its +/// per-key lock, so the recorder's COUNT/SUM straddled a delete it +/// never serialized against. +/// +/// This test drives both pressures at once against ONE scope: a flood of +/// fresh distinct keys (each triggering an eviction pass) concurrently with +/// a flood of distinct-id records against a single "hot" key that is a prime +/// eviction candidate. The fix makes every task complete (no deadlock) under +/// a hard timeout, and keeps the hot key's reported aggregate consistent with +/// its surviving rows (no torn/lost update). Run against a real Postgres via +/// `IRONCLAW_HOOKS_POSTGRES_URL` / `DATABASE_URL`; skipped (passing) otherwise. +#[tokio::test(flavor = "multi_thread", worker_threads = 8)] +#[allow(clippy::await_holding_lock)] +async fn scope_lru_eviction_serializes_against_victim_writes_no_deadlock() { + let (url, _guard) = guarded!(); + let hs = hosts(&url, 2).await; + let window = Duration::from_secs(86_400); + let tenant = "lru-race-tenant"; + + // Seed the scope to exactly the quota with distinct keys, all OLDER than + // the activity we drive below, so the LRU pass has plenty of stale + // victims to choose from. The "hot" key is seeded oldest so it is a top + // eviction candidate while we also hammer it concurrently. + let hot = inv_key(tenant, "cap.hot"); + hs[0] + .record_invocation(&hot, &ev("hot-seed"), base(), window) + .await + .expect("seed hot key"); + for i in 0..MAX_KEYS_PER_TENANT { + let k = inv_key(tenant, &format!("seed.{i}")); + let ts = base() + chrono::Duration::seconds(1 + i as i64); + hs[0] + .record_invocation(&k, &ev(&format!("seed-e{i}")), ts, window) + .await + .expect("seed key"); + } + + // Concurrent pressure: fresh keys that force eviction passes, plus + // distinct-id records against the hot key (which holds the hot key's + // per-key lock during its own transaction). If eviction ignored the + // victim's per-key lock these would deadlock. + const FRESH: usize = 60; + const HOT_HITS: usize = 60; + let mut handles = Vec::new(); + for i in 0..FRESH { + let backend = Arc::clone(&hs[i % 2]); + let k = inv_key(tenant, &format!("fresh.{i}")); + let ts = base() + chrono::Duration::seconds(10_000 + i as i64); + handles.push(tokio::spawn(async move { + // RETURN the backend result so the join below can assert it. A + // discarded result would let a deadlock-detected/serialization + // DB error pass silently as long as the task returned before the + // timeout (serrrfirat regression on #3933). The result is checked + // against the allowed-outcome set after the join. + backend + .record_invocation(&k, &ev(&format!("fresh-e{i}")), ts, window) + .await + })); + } + for i in 0..HOT_HITS { + let backend = Arc::clone(&hs[i % 2]); + let hot = hot.clone(); + // All hot-key records share one in-window instant region; distinct + // ids so each is a real insert (subject to the per-key sample cap). + let ts = base() + chrono::Duration::seconds(20_000 + i as i64); + handles.push(tokio::spawn(async move { + backend + .record_invocation(&hot, &ev(&format!("hot-e{i}")), ts, window) + .await + })); + } + + // Hard timeout is the deadlock detector: pre-fix this join would hang + // (or surface a Postgres deadlock-detected error) instead of completing. + // Collect every task's backend result so we can assert the allowed + // outcomes EXPLICITLY rather than discarding them. + let join_all = async { + let mut results = Vec::with_capacity(handles.len()); + for handle in handles { + results.push(handle.await.expect("task did not panic")); + } + results + }; + let results = tokio::time::timeout(Duration::from_secs(60), join_all) + .await + .expect("all record tasks completed without deadlock/hang"); + + // Assert the production failure mode this test documents actually fails + // the test: each record must be either `Ok` (recorded / replayed) or a + // benign quota/window outcome. A `Unavailable(..)` whose message looks + // like a Postgres deadlock-detected or serialization-failure error means + // the eviction lock discipline regressed — that must surface as a test + // failure, not a silent pass. + for (i, r) in results.iter().enumerate() { + match r { + Ok(_) => {} + Err(PredicateBackendError::WindowOverflow { .. }) => { + // A fresh/hot key can legitimately hit the per-key sample cap. + } + Err(PredicateBackendError::Unavailable(msg)) => { + let lower = msg.to_lowercase(); + assert!( + !(lower.contains("deadlock") || lower.contains("serialize")), + "task {i} surfaced a deadlock/serialization DB error \ + (eviction lock discipline regressed): {msg}" + ); + // A non-deadlock `Unavailable` here is the explicit fail-closed + // quota outcome (every stale victim momentarily locked) — that + // is an allowed outcome of the new quota-enforcement contract. + } + } + } + + // Consistency: whatever survived for the hot key, the backend's reported + // aggregate must equal its actual surviving IN-WINDOW row count — no torn + // aggregate (the bug would let an LRU delete straddle the recorder's + // COUNT). We read the backend's aggregate FIRST via a no-op replay of a + // known id, then compare against a window-matched direct COUNT computed + // with the SAME `now`/cutoff the replay used, so the two are apples-to- + // apples (an all-rows COUNT would spuriously differ from the in-window + // aggregate). Reading the backend first also pins the row set: the replay + // is a no-op (no insert/delete for the hot key), so the direct count that + // follows observes exactly the rows the replay aggregated. + let read_now = base() + chrono::Duration::seconds(100_000); + let reported = hs[0] + .record_invocation(&hot, &ev("hot-seed"), read_now, window) + .await + .expect("hot read"); + let cutoff = read_now - chrono::Duration::from_std(window).unwrap(); + let pool = build_pool(&url).expect("pool"); + let client = pool.get().await.expect("client"); + let hot_hash = ironclaw_hooks_postgres::test_support::invocation_key_hash_bytes(&hot); + let rows: i64 = client + .query_one( + "SELECT COUNT(*)::BIGINT FROM hooks_predicate_invocations \ + WHERE key_hash = $1 AND occurred_at >= $2", + &[&&hot_hash[..], &cutoff], + ) + .await + .expect("count hot rows") + .get(0); + assert_eq!( + reported as i64, rows, + "hot key's reported aggregate must match its surviving in-window row count (no torn update)" + ); + + // The per-scope quota bound must still hold for this scope. + let scope = ironclaw_hooks_postgres::test_support::scope_hash_bytes(tenant); + let distinct: i64 = client + .query_one( + "SELECT COUNT(DISTINCT key_hash)::BIGINT FROM hooks_predicate_invocations \ + WHERE scope_hash = $1", + &[&&scope[..]], + ) + .await + .expect("distinct count") + .get(0); + assert!( + distinct as usize <= MAX_KEYS_PER_TENANT, + "per-scope LRU bound must hold after the concurrent race; had {distinct}" + ); +} + +/// Stateful proof that the per-scope quota is ACTUALLY enforced, not merely +/// best-effort. Drives a single uncontended host sequentially well past +/// [`MAX_KEYS_PER_TENANT`] distinct keys, then asserts the scope holds at +/// EXACTLY the cap — no overshoot. This is the regression guard for the +/// "silently best-effort `Ok(evicted)`" BLOCKER (serrrfirat on #3933): the +/// old code could commit an over-quota scope if victims happened to be +/// locked; with a single sequential writer no victim is ever locked, so the +/// quota-enforcement loop must drive the scope exactly to the cap on every +/// newly-material insert. A regression that under-evicts would leave +/// `distinct > MAX_KEYS_PER_TENANT` and fail here. Run against a real +/// Postgres via `IRONCLAW_HOOKS_POSTGRES_URL` / `DATABASE_URL`; skipped +/// (passing) otherwise. +#[tokio::test] +#[allow(clippy::await_holding_lock)] +async fn scope_quota_is_enforced_exactly_not_best_effort() { + let (url, _guard) = guarded!(); + let hs = hosts(&url, 1).await; + let backend = &hs[0]; + let window = Duration::from_secs(86_400); + let tenant = "quota-exact-tenant"; + + // Insert OVERSHOOT keys beyond the cap. Each key gets one in-window + // sample with a strictly increasing timestamp so the oldest-front LRU + // victim selection is deterministic (key N is staler than key N+1). + const OVERSHOOT: usize = 200; + let total = MAX_KEYS_PER_TENANT + OVERSHOOT; + for i in 0..total { + let k = inv_key(tenant, &format!("k.{i}")); + let ts = base() + chrono::Duration::seconds(i as i64); + backend + .record_invocation(&k, &ev(&format!("e{i}")), ts, window) + .await + .expect("record under uncontended single writer must succeed"); + } + + // The scope must hold at EXACTLY the cap — eviction kept pace with every + // over-cap insert. Verify via a direct distinct-key count. + let pool = build_pool(&url).expect("pool"); + let client = pool.get().await.expect("client"); + let scope = ironclaw_hooks_postgres::test_support::scope_hash_bytes(tenant); + let distinct: i64 = client + .query_one( + "SELECT COUNT(DISTINCT key_hash)::BIGINT FROM hooks_predicate_invocations \ + WHERE scope_hash = $1", + &[&&scope[..]], + ) + .await + .expect("distinct count") + .get(0); + assert_eq!( + distinct as usize, MAX_KEYS_PER_TENANT, + "uncontended sequential flood must leave the scope at EXACTLY the cap \ + (quota enforced, not best-effort); had {distinct}" + ); + + // And the surviving keys must be the MOST-RECENT ones: the oldest-front + // victims (k.0 .. k.OVERSHOOT-1) were evicted, so k.0 must be gone and + // the newest key (k.{total-1}) must survive. + let oldest = inv_key(tenant, "k.0"); + let oldest_hash = ironclaw_hooks_postgres::test_support::invocation_key_hash_bytes(&oldest); + let oldest_rows: i64 = client + .query_one( + "SELECT COUNT(*)::BIGINT FROM hooks_predicate_invocations WHERE key_hash = $1", + &[&&oldest_hash[..]], + ) + .await + .expect("count oldest") + .get(0); + assert_eq!( + oldest_rows, 0, + "the oldest-front key must have been evicted under the quota" + ); + let newest = inv_key(tenant, &format!("k.{}", total - 1)); + let newest_hash = ironclaw_hooks_postgres::test_support::invocation_key_hash_bytes(&newest); + let newest_rows: i64 = client + .query_one( + "SELECT COUNT(*)::BIGINT FROM hooks_predicate_invocations WHERE key_hash = $1", + &[&&newest_hash[..]], + ) + .await + .expect("count newest") + .get(0); + assert_eq!(newest_rows, 1, "the most-recent key must survive eviction"); +} + +/// Deterministic reproduction of the eviction deadlock cycle at the raw-SQL +/// lock level — independent of the timing luck a concurrency stress test +/// relies on. This replays the exact advisory-lock + row-lock sequence the +/// production code takes and proves the fix's protocol (try-lock the victim +/// key BEFORE deleting its rows) cannot deadlock, whereas the old protocol +/// (blocking on the victim's rows while holding the scope lock) deadlocks. +/// +/// Setup mirrors production: +/// * Txn V ("victim recorder"): takes key B's per-key advisory lock, then +/// writes a row for B (holding a row lock) — exactly what a `record` +/// call for B does before it reaches its own quota pass. +/// * Txn L ("LRU pass"): takes the scope advisory lock, then must evict B. +/// +/// The OLD code would `DELETE` B's rows here and BLOCK on V's row lock; if V +/// then waited on the scope lock (its quota pass) the two would deadlock. The +/// FIXED code instead runs `pg_try_advisory_xact_lock` on B's key first — +/// which returns FALSE because V holds it — and skips B without blocking. +/// We assert that non-blocking outcome directly: the try-lock from L returns +/// false while V holds B's key lock, so L never waits on V and the cycle +/// cannot form. This is the load-bearing invariant of the fix. +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] +#[allow(clippy::await_holding_lock)] +async fn eviction_try_lock_does_not_block_on_in_flight_victim() { + let (url, _guard) = guarded!(); + // Build the schema/table without disturbing other tests' data. + let _hs = hosts(&url, 1).await; + + // Two independent connections = two independent transactions. + let (mut client_v, conn_v) = tokio_postgres::connect(&url, tokio_postgres::NoTls) + .await + .expect("connect V"); + tokio::spawn(conn_v); + let (mut client_l, conn_l) = tokio_postgres::connect(&url, tokio_postgres::NoTls) + .await + .expect("connect L"); + tokio::spawn(conn_l); + for c in [&client_v, &client_l] { + c.batch_execute(&format!("SET search_path TO {TEST_SCHEMA}")) + .await + .expect("search_path"); + } + + // A distinct victim-key lock pair unlikely to collide with other tests. + let (lk_a, lk_b): (i32, i32) = (0x7EED_1234u32 as i32, 0x0BAD_5678u32 as i32); + + // Txn V: take key B's per-key advisory lock (the recorder's first act). + let tx_v = client_v.transaction().await.expect("begin V"); + let got_v: bool = tx_v + .query_one("SELECT pg_try_advisory_xact_lock($1, $2)", &[&lk_a, &lk_b]) + .await + .expect("V lock") + .get(0); + assert!(got_v, "V must acquire the victim key lock first"); + + // Txn L: while V holds B's key lock, L (the eviction pass) must NOT block + // on it — the fix uses the non-blocking try-lock and skips B. Pre-fix, L + // would issue a blocking DELETE on B's rows here and stall. + let tx_l = client_l.transaction().await.expect("begin L"); + let got_l: bool = tokio::time::timeout( + Duration::from_secs(5), + tx_l.query_one("SELECT pg_try_advisory_xact_lock($1, $2)", &[&lk_a, &lk_b]), + ) + .await + .expect("try-lock must return promptly, never block (deadlock-free)") + .expect("L try-lock query") + .get(0); + assert!( + !got_l, + "eviction try-lock on an in-flight victim key MUST fail (so the pass \ + skips it) rather than block — this is what breaks the deadlock cycle" + ); + + // Clean up both transactions. + tx_l.rollback().await.expect("rollback L"); + tx_v.rollback().await.expect("rollback V"); +} + +/// `evict_older_than` is the time-based reaper (distinct from the per-scope +/// LRU eviction): it deletes every row with `occurred_at < cutoff` from BOTH +/// typed tables and returns the total rows removed. This proves the Postgres +/// override actually reaps both tables and spares in-window rows — the two +/// `DELETE` statements had no integration coverage (henrypark tests finding). +/// +/// Gated on a reachable Postgres via `IRONCLAW_HOOKS_POSTGRES_URL` / +/// `DATABASE_URL`; skipped (passing) otherwise. +#[tokio::test] +#[allow(clippy::await_holding_lock)] +async fn evict_older_than_removes_stale_rows_from_both_tables() { + let (url, _guard) = guarded!(); + let hs = hosts(&url, 1).await; + let backend = &hs[0]; + let window = Duration::from_secs(86_400); + + // Unique tenant so this test's rows are isolated from sibling tests that + // share the truncated-once schema. + let tenant = "reaper-tenant"; + let inv = inv_key(tenant, "cap.reap"); + let val = val_key(tenant, "cap.reap", "amount"); + + // Two stale rows (well before the cutoff) and one fresh row (after it), + // on EACH table — so we can prove the reaper hits both tables and spares + // the fresh rows. + let stale_a = base(); + let stale_b = base() + chrono::Duration::seconds(10); + let fresh = base() + chrono::Duration::seconds(1_000); + + backend + .record_invocation(&inv, &ev("inv-stale-a"), stale_a, window) + .await + .expect("inv stale a"); + backend + .record_invocation(&inv, &ev("inv-stale-b"), stale_b, window) + .await + .expect("inv stale b"); + backend + .record_invocation(&inv, &ev("inv-fresh"), fresh, window) + .await + .expect("inv fresh"); + + backend + .record_value(&val, &ev("val-stale-a"), stale_a, Decimal::from(5), window) + .await + .expect("val stale a"); + backend + .record_value(&val, &ev("val-fresh"), fresh, Decimal::from(7), window) + .await + .expect("val fresh"); + + // Cutoff between the stale and fresh rows: reaps 2 invocation + 1 value + // stale rows = 3 total, leaving 1 invocation + 1 value fresh row. + let cutoff = base() + chrono::Duration::seconds(100); + let removed = backend + .evict_older_than(cutoff) + .await + .expect("evict_older_than"); + assert_eq!( + removed, 3, + "reaper must delete all 3 stale rows across both tables, got {removed}" + ); + + // Verify directly: only the fresh rows survive in each table. + let pool = build_pool(&url).expect("pool"); + let client = pool.get().await.expect("client"); + let inv_hash = ironclaw_hooks_postgres::test_support::invocation_key_hash_bytes(&inv); + let inv_remaining: i64 = client + .query_one( + "SELECT COUNT(*)::BIGINT FROM hooks_predicate_invocations WHERE key_hash = $1", + &[&&inv_hash[..]], + ) + .await + .expect("count inv") + .get(0); + assert_eq!( + inv_remaining, 1, + "only the fresh invocation row must survive the reap" + ); + + // The value table shares the key derivation; query its surviving rows by + // the value key's hash directly via a fresh dedup record-and-read would + // mutate state, so count rows for this tenant's value key instead. We + // re-derive the value key's hash the same way the backend does by + // recording a replay of the fresh id (a no-op that returns the live sum). + let live_sum = backend + .record_value(&val, &ev("val-fresh"), fresh, Decimal::from(7), window) + .await + .expect("replay fresh value"); + assert_eq!( + live_sum, + Decimal::from(7), + "only the fresh value row must survive the reap (stale value row gone)" + ); +} diff --git a/crates/ironclaw_hooks_postgres/tests/predicate_state_postgres_contract.rs b/crates/ironclaw_hooks_postgres/tests/predicate_state_postgres_contract.rs new file mode 100644 index 00000000000..3316fd6105a --- /dev/null +++ b/crates/ironclaw_hooks_postgres/tests/predicate_state_postgres_contract.rs @@ -0,0 +1,179 @@ +//! `PostgresPredicateStateBackend` run against the shared trait-level +//! contract harness from `ironclaw_hooks::predicate_state::contract`. +//! +//! All nine contract functions are exercised. The harness is the same +//! one the in-memory backend is wired through (PR 1/4), so this proves +//! the durable backend honors the identical isolation / dedup / window / +//! atomicity invariants by construction. +//! +//! # DB gating +//! +//! These tests need a reachable Postgres. Set +//! `IRONCLAW_HOOKS_POSTGRES_URL` (or `DATABASE_URL`) to a libpq URL. With +//! no URL the tests print a skip notice and pass, matching the +//! env-gated-skip pattern used by `ironclaw_reborn_event_store` and +//! `ironclaw_filesystem` (no testcontainers dependency). +//! +//! # Isolation +//! +//! The contract functions use fixed tenant/key names ("alpha", "beta"), +//! so concurrent runs against a shared table would collide. We serialize +//! these tests behind a process-global async mutex and `TRUNCATE` the +//! table at the start of each, giving every contract a fresh-empty +//! backend exactly as the in-memory factory does. + +#![cfg(feature = "postgres")] + +use std::sync::Arc; + +use deadpool_postgres::Pool; +use ironclaw_hooks_postgres::PostgresPredicateStateBackend; + +/// Process-global serialization lock. The contract suite reuses fixed +/// keys ("alpha", "beta"), so tests must not interleave against a shared +/// table. This is a `std::sync::Mutex` (not tokio) because each +/// `#[tokio::test]` runs on its OWN runtime — a tokio Mutex / pool tied +/// to one test's runtime cannot be reused from another's reactor (that +/// caused spurious `kind: Closed` connection errors). We hold the lock +/// across the whole test via a guard. +static TEST_LOCK: std::sync::Mutex<()> = std::sync::Mutex::new(()); + +fn db_url() -> Option { + std::env::var("IRONCLAW_HOOKS_POSTGRES_URL") + .or_else(|_| std::env::var("DATABASE_URL")) + .ok() +} + +/// Dedicated Postgres schema for THIS test binary so it cannot collide +/// with other integration-test binaries (e.g. the adversarial suite) that +/// `cargo test` runs in parallel against the same database. Every pooled +/// connection sets `search_path` to this schema via a `post_create` hook, +/// so the backend's unqualified `hooks_predicate_*` tables resolve here. +const TEST_SCHEMA: &str = "hooks_predicate_contract_test"; + +/// Build a pool ON THE CURRENT runtime pinned to the dedicated schema, +/// ensure the schema + table, truncate, and return a fresh backend. Each +/// test calls this so the pool's connections are bound to that test's own +/// reactor. +async fn fresh_backend(url: &str) -> Option> { + // Create the isolated schema first via a one-off connection. + { + let (client, conn) = tokio_postgres::connect(url, tokio_postgres::NoTls) + .await + .ok()?; + tokio::spawn(conn); + client + .batch_execute(&format!("CREATE SCHEMA IF NOT EXISTS {TEST_SCHEMA}")) + .await + .ok()?; + } + + let config = url.parse::().ok()?; + let manager = deadpool_postgres::Manager::new(config, tokio_postgres::NoTls); + let pool: Pool = deadpool_postgres::Pool::builder(manager) + .max_size(8) + .post_create(deadpool_postgres::Hook::async_fn(|client, _| { + Box::pin(async move { + client + .batch_execute(&format!("SET search_path TO {TEST_SCHEMA}")) + .await + .map_err(|e| deadpool_postgres::HookError::message(e.to_string()))?; + Ok(()) + }) + })) + .build() + .ok()?; + let backend = PostgresPredicateStateBackend::new(pool.clone()); + backend.run_migrations().await.ok()?; + let client = pool.get().await.ok()?; + client + .batch_execute("TRUNCATE TABLE hooks_predicate_invocations, hooks_predicate_values") + .await + .ok()?; + Some(Arc::new(backend)) +} + +/// Macro: one `#[tokio::test]` per contract function. Each acquires the +/// global lock, truncates, then drives the shared contract with a factory +/// returning a clone of the freshly-prepared backend handle. +macro_rules! pg_contract { + ($name:ident) => { + #[tokio::test] + // The std Mutex guard is intentionally held across awaits to + // serialize tests that share a fixed-key table across separate + // per-test runtimes; a tokio Mutex would be runtime-bound. + #[allow(clippy::await_holding_lock)] + async fn $name() { + let Some(url) = db_url() else { + eprintln!( + "skipping postgres predicate contract `{}`: \ + IRONCLAW_HOOKS_POSTGRES_URL / DATABASE_URL not set", + stringify!($name) + ); + return; + }; + // Serialize across tests' separate runtimes. Recover from a + // poisoned lock (a panicking test still released the table via + // the next test's TRUNCATE). + let _guard = TEST_LOCK.lock().unwrap_or_else(|p| p.into_inner()); + let backend = fresh_backend(&url).await.expect("postgres setup"); + ironclaw_hooks::predicate_state::contract::$name(move || { + // Factory is called once by the contract; hand back a + // clone over the same isolated (truncated) table. + PgBackendHandle(backend.clone()) + }) + .await; + } + }; +} + +/// Newtype wrapper so the contract's `B: PredicateStateBackend` bound is +/// satisfied by a cheaply-cloneable `Arc` handle. Delegates every method +/// to the inner backend. +#[derive(Clone)] +struct PgBackendHandle(Arc); + +#[async_trait::async_trait] +impl ironclaw_hooks::predicate_state::PredicateStateBackend for PgBackendHandle { + async fn record_invocation( + &self, + key: &ironclaw_hooks::predicate_state::InvocationKey, + event_id: &ironclaw_hooks::predicate_state::PredicateEventId, + now: chrono::DateTime, + window: std::time::Duration, + ) -> Result { + self.0.record_invocation(key, event_id, now, window).await + } + + async fn record_value( + &self, + key: &ironclaw_hooks::predicate_state::ValueKey, + event_id: &ironclaw_hooks::predicate_state::PredicateEventId, + now: chrono::DateTime, + value: rust_decimal::Decimal, + window: std::time::Duration, + ) -> Result { + self.0.record_value(key, event_id, now, value, window).await + } + + fn evictions_observed(&self) -> u64 { + self.0.evictions_observed() + } + + async fn evict_older_than( + &self, + cutoff: chrono::DateTime, + ) -> Result { + self.0.evict_older_than(cutoff).await + } +} + +pg_contract!(invocation_counts_within_window); +pg_contract!(invocation_trims_outside_window); +pg_contract!(value_sums_within_window); +pg_contract!(tenant_isolation); +pg_contract!(duplicate_event_id_is_noop_for_invocations); +pg_contract!(duplicate_event_id_is_noop_for_values); +pg_contract!(invocation_retains_entry_at_exact_window_cutoff); +pg_contract!(event_id_dedup_isolated_across_maps); +pg_contract!(record_invocation_overflow_is_fail_closed);