From e035b8b7f27f402b56cb996e0535635b6a6f1371 Mon Sep 17 00:00:00 2001 From: fylorn <249551762+fylorn@users.noreply.github.com> Date: Mon, 5 Oct 2026 18:39:47 +0800 Subject: [PATCH 1/5] fix(startup): set up the schema one instance at a time Every instance applies db/schema.sql, the seeds and the guard-settings conversion when it starts. Two or more starting together against one database - Helm replicaCount above 1, a rolling upgrade, an autoscaler - ran them side by side, and the losers exited with "Database migration failed". schema.sql goes to Postgres as one multi-statement query, so one implicit transaction, and on an up-to-date database two of them deadlock on api_keys: each holds a lock that the other's DROP TRIGGER IF EXISTS / CREATE TRIGGER trg_api_keys_surfaces_valid needs (the relation on one side, the trigger's pg_trigger row on the other). On an empty database they instead race to create the same objects, and all but the first fail with a unique violation on pg_extension_name_index (CREATE EXTENSION IF NOT EXISTS pgcrypto). Reproduced with psql and with four run_migrations at once; 3 of 4 failed in both cases. run_migrations now takes a session-level advisory lock (SCHEMA_LOCK, "twschema") and runs the schema, the seeds and the conversion on the connection that holds it. The next instance waits on pg_advisory_lock, logging that it does, and then finds everything applied: the schema is idempotent, the seeds are ON CONFLICT DO NOTHING, and the conversion finds its marker. The conversion keeps its own transaction-level lock, which a 3.0 or 3.1 instance restarting during the rollout also takes. The lock lives on a pooled connection marked close-on-drop, and release closes it, so a connection still holding the lock never goes back to the pool, whether setup succeeds, fails or is cancelled; an instance that dies holding it loses the session and Postgres releases the lock. (A transaction-level lock releases itself too, but the ClickHouse lock below would then be a transaction left idle for as long as the backfill runs.) An instance of an earlier version does not take the lock, so one of those restarting at the very moment a new one migrates can still collide with it. ClickHouse setup had the same kind of race, without an error to show for it: ensure_clickhouse_tables backfills cost_rollup_hourly, provider_health_5m and mcp_server_call_counts from the logs when it finds them empty, so instances starting together could each find them empty and each copy the logs in - the SummingMergeTree rollups then counted every request once per instance. It now runs under a second lock (CLICKHOUSE_SETUP_LOCK, "twchinit"), taken per attempt rather than across the retry backoff: with ClickHouse down, instances taking turns through each other's ~22 s of retries would outlast the chart's startup probe. Nothing else at start-up writes shared state: the setup wizard already serializes on its own advisory lock and is not run at boot, there is no initial-admin bootstrap, and the TTL and bucket-lifecycle reconciles set values every instance computes the same way. tests/concurrent_startup.rs starts four setups at once on a fresh database, again on the result (a restart), and on a database put back to an earlier schema with the old guard settings, and compares the schema and seeds with a database one instance migrated. It also checks that a holder that dies or drops the lock releases it, and that four ClickHouse setups at once backfill a rollup once. Before this change the Postgres cases failed each time they ran, the ClickHouse one in one run of three (100 requests counted for 25). Co-Authored-By: Claude Opus 5.5 --- CHANGELOG.md | 14 + crates/common/src/db.rs | 86 ++++- crates/common/src/guard_policy/legacy.rs | 12 +- crates/server/src/init.rs | 81 +++-- .../src/services/mcp_store_repository.rs | 3 + crates/test-support/src/lib.rs | 35 +- crates/test-support/src/pg.rs | 17 +- .../test-support/tests/concurrent_startup.rs | 339 ++++++++++++++++++ db/schema.sql | 5 +- 9 files changed, 536 insertions(+), 56 deletions(-) create mode 100644 crates/test-support/tests/concurrent_startup.rs diff --git a/CHANGELOG.md b/CHANGELOG.md index 675ab46a..37c7f4b4 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -75,6 +75,20 @@ target. file with that CA, which is then trusted alone; the Helm chart sets it from a Secret given in `redis.caSecret`. The chart's README describes both. +- **Several instances starting at once.** Server instances starting together + against one database — a Helm `replicaCount` above 1, a rolling upgrade, an + autoscaler adding pods — applied the schema side by side, and all but one + could exit with `Database migration failed: apply db/schema.sql: … deadlock + detected` (on an empty database: `duplicate key value violates unique + constraint "pg_extension_name_index"`). With ClickHouse, the rollups that an + instance fills from the logs when it finds them empty (`cost_rollup_hourly`, + `provider_health_5m`, `mcp_server_call_counts`) could be filled by each of + them, counting every request once per instance on the cost pages, the + dashboard and the MCP server list. Instances now set up Postgres and + ClickHouse one at a time, under Postgres advisory locks: the others wait, + logging `Another instance is setting up the database schema; waiting for it + to finish`, then find it done. An instance that dies holding a lock releases + it with its connection. ## [3.1.0] — 2026-10-05 diff --git a/crates/common/src/db.rs b/crates/common/src/db.rs index 25e70127..f92a8828 100644 --- a/crates/common/src/db.rs +++ b/crates/common/src/db.rs @@ -1,5 +1,6 @@ -use sqlx::PgPool; +use sqlx::pool::PoolConnection; use sqlx::postgres::PgPoolOptions; +use sqlx::{Executor, PgConnection, PgPool, Postgres}; pub async fn create_pool(database_url: &str) -> anyhow::Result { if !database_url.contains("sslmode=") { @@ -28,6 +29,64 @@ pub async fn create_pool(database_url: &str) -> anyhow::Result { Ok(pool) } +/// Advisory-lock key held while an instance applies the schema, the +/// seeds and the conversions ([`run_migrations`]): "twschema" in ASCII. +pub const SCHEMA_LOCK: i64 = i64::from_be_bytes(*b"twschema"); + +/// Advisory-lock key held while an instance sets up the ClickHouse tables +/// and backfills the rollups: "twchinit" in ASCII. +pub const CLICKHOUSE_SETUP_LOCK: i64 = i64::from_be_bytes(*b"twchinit"); + +/// A Postgres session-level advisory lock on a connection of its own. +/// +/// Start-up work that every instance runs, but that two instances must +/// not run at once, holds one: an instance that starts while another +/// holds the lock waits for it, then finds the work done. +/// +/// The lock lasts as long as the session. The connection never goes back +/// to the pool — [`release`](Self::release) closes it, and so does +/// dropping the guard on an error — so no later query can find itself +/// holding the lock. An instance that dies holding it loses its +/// connection, and Postgres releases the lock with the session. +pub struct StartupLock { + conn: PoolConnection, +} + +impl StartupLock { + /// Take the lock `key`, waiting while another instance holds it. + /// `work` says what the holder is doing, for the log line a waiting + /// instance writes. + pub async fn acquire(pool: &PgPool, key: i64, work: &str) -> anyhow::Result { + let mut conn = pool.acquire().await?; + conn.close_on_drop(); + let free: bool = sqlx::query_scalar("SELECT pg_try_advisory_lock($1)") + .bind(key) + .fetch_one(&mut *conn) + .await?; + if !free { + tracing::info!("Another instance is {work}; waiting for it to finish"); + sqlx::query("SELECT pg_advisory_lock($1)") + .bind(key) + .execute(&mut *conn) + .await?; + } + Ok(Self { conn }) + } + + /// The connection holding the lock. The work runs on it. + pub fn conn(&mut self) -> &mut PgConnection { + &mut self.conn + } + + /// Release the lock by ending the session. + pub async fn release(self) { + if let Err(e) = self.conn.close().await { + // The session is gone either way, and the lock with it. + tracing::debug!("closing the start-up lock's connection: {e}"); + } + } +} + /// Bring the database up to the schema declared in `db/schema.sql`, /// then idempotent-apply seeds from `db/seeds.sql`. Both files are /// embedded at compile time, so the running binary never needs disk @@ -45,20 +104,35 @@ pub async fn create_pool(database_url: &str) -> anyhow::Result { /// to carry them out by itself, a conversion that runs here in one /// transaction and finds nothing left to do on the next boot (the guard /// settings, [`crate::guard_policy::legacy`]). +/// +/// **One instance at a time.** Several replicas start together — a +/// first install with more than one, a rolling upgrade, an autoscaler — +/// and each runs this. Run side by side, the statements deadlock (each +/// session holding a lock on `api_keys` or its trigger that another +/// needs) or collide creating the same object on an empty database, and +/// the instances that lose exit. So all of it runs under +/// [`SCHEMA_LOCK`], on the connection that holds it: the next instance +/// waits, then applies a schema that is already there, seeds rows that +/// already exist, and finds nothing left to convert. pub async fn run_migrations(pool: &PgPool) -> anyhow::Result<()> { + let mut lock = + StartupLock::acquire(pool, SCHEMA_LOCK, "setting up the database schema").await?; + let conn = lock.conn(); let schema = include_str!("../../../db/schema.sql"); - sqlx::raw_sql(schema) - .execute(pool) + // `Executor::execute` rather than `RawSql::execute`: with a borrowed + // connection, the compiler cannot prove the latter's future `Send`, + // and a caller that spawns this would not compile. + conn.execute(sqlx::raw_sql(schema)) .await .map_err(|e| anyhow::anyhow!("apply db/schema.sql: {e}"))?; let seeds = include_str!("../../../db/seeds.sql"); - sqlx::raw_sql(seeds) - .execute(pool) + conn.execute(sqlx::raw_sql(seeds)) .await .map_err(|e| anyhow::anyhow!("apply db/seeds.sql: {e}"))?; - crate::guard_policy::legacy::upgrade(pool) + crate::guard_policy::legacy::upgrade(conn) .await .map_err(|e| anyhow::anyhow!("convert the previous guard settings: {e}"))?; + lock.release().await; tracing::info!("Database schema + seeds applied"); Ok(()) } diff --git a/crates/common/src/guard_policy/legacy.rs b/crates/common/src/guard_policy/legacy.rs index 7fe61d1f..568fdcd8 100644 --- a/crates/common/src/guard_policy/legacy.rs +++ b/crates/common/src/guard_policy/legacy.rs @@ -42,14 +42,16 @@ //! **Once, in one transaction, safe to run again.** The old keys and the //! old column go in the same transaction that writes what replaces them, //! so the next boot finds nothing to convert, and a failure leaves -//! everything as it was. An advisory lock keeps two replicas booting at -//! once from both converting. +//! everything as it was. Two replicas booting at once don't both convert: +//! it runs under the schema lock [`crate::db::run_migrations`] holds, and +//! under a lock of its own, which an instance of a version from before +//! the schema lock takes too. use std::collections::{BTreeMap, BTreeSet, HashSet}; use serde::Deserialize; use serde_json::Value; -use sqlx::PgPool; +use sqlx::{Connection, PgConnection}; use tw_guard::policy::{ ContentAction, ContentMatch, ContentPolicy, CustomContentRule, CustomRedactRule, CustomToolRule, DEFAULT_LABEL, Guard, LABEL_MAX, Mode, RedactPolicy, ToolAction, ToolPolicy, @@ -94,8 +96,8 @@ pub const MARKER: &str = "security.legacy_converted"; /// a rollback, or an old replica restarting, whose seeds write their /// defaults back. Converting them would overwrite the policies in force /// with those defaults, so they are only removed, with a warning. -pub async fn upgrade(pool: &PgPool) -> anyhow::Result<()> { - let mut tx = pool.begin().await?; +pub async fn upgrade(conn: &mut PgConnection) -> anyhow::Result<()> { + let mut tx = conn.begin().await?; sqlx::query("SELECT pg_advisory_xact_lock($1)") .bind(LOCK) .execute(&mut *tx) diff --git a/crates/server/src/init.rs b/crates/server/src/init.rs index aac66a66..c9dd5dfa 100644 --- a/crates/server/src/init.rs +++ b/crates/server/src/init.rs @@ -18,6 +18,7 @@ use sqlx::PgPool; use think_watch_auth::oidc::OidcManager; use think_watch_common::audit::{self, AuditConfig, AuditLogger}; use think_watch_common::config::AppConfig; +use think_watch_common::db; use think_watch_common::dynamic_config::{self, DynamicConfig}; use think_watch_common::tasks::supervise; @@ -43,33 +44,8 @@ pub async fn init_state( .await .context("persisted rate-limit / weight rows fail validation")?; - // ClickHouse tables. Same bounded retry as production but without - // the metrics counter (recorder is not installed in tests). - if ch_client.is_some() { - let mut attempt = 0u32; - loop { - match audit::ensure_clickhouse_tables(&ch_client).await { - Ok(()) => break, - Err(e) if attempt < 4 => { - let backoff_ms = 1_500u64 * 2u64.pow(attempt); - tracing::warn!( - attempt = attempt + 1, - backoff_ms, - "ClickHouse table init failed, retrying: {e}" - ); - tokio::time::sleep(std::time::Duration::from_millis(backoff_ms)).await; - attempt += 1; - } - Err(e) => { - tracing::error!( - "ClickHouse table init failed after {} attempts: {e}", - attempt + 1 - ); - break; - } - } - } - } + // ClickHouse tables, when ClickHouse is configured. + ensure_clickhouse_schema(&pool, &ch_client).await?; let dynamic_config = Arc::new(DynamicConfig::load(pool.clone()).await?); @@ -174,6 +150,57 @@ pub async fn init_state( Ok(state) } +/// Create or bring up to date the ClickHouse tables, and backfill the +/// rollups that are still empty. Retried with a backoff; a ClickHouse +/// that stays unreachable is logged, not fatal. +/// +/// One instance at a time, under [`db::CLICKHOUSE_SETUP_LOCK`] in +/// Postgres: a rollup is backfilled when it is found empty, and two +/// instances starting together both found it empty and each copied the +/// whole log into it, counting every request twice. The lock is held for +/// an attempt, not across the backoff, so that instances starting while +/// ClickHouse is down don't wait out each other's retries — longer than +/// the chart's startup probe allows. +pub async fn ensure_clickhouse_schema( + pool: &PgPool, + ch_client: &Option, +) -> anyhow::Result<()> { + if ch_client.is_none() { + return Ok(()); + } + let mut attempt = 0u32; + loop { + let lock = db::StartupLock::acquire( + pool, + db::CLICKHOUSE_SETUP_LOCK, + "setting up the ClickHouse tables", + ) + .await?; + let result = audit::ensure_clickhouse_tables(ch_client).await; + lock.release().await; + match result { + Ok(()) => return Ok(()), + Err(e) if attempt < 4 => { + let backoff_ms = 1_500u64 * 2u64.pow(attempt); + tracing::warn!( + attempt = attempt + 1, + backoff_ms, + "ClickHouse table init failed, retrying: {e}" + ); + tokio::time::sleep(std::time::Duration::from_millis(backoff_ms)).await; + attempt += 1; + } + Err(e) => { + tracing::error!( + "ClickHouse table init failed after {} attempts: {e}", + attempt + 1 + ); + return Ok(()); + } + } + } +} + fn audit_config(config: &AppConfig) -> AuditConfig { config.audit_config() } diff --git a/crates/server/src/services/mcp_store_repository.rs b/crates/server/src/services/mcp_store_repository.rs index f445d384..9afeae54 100644 --- a/crates/server/src/services/mcp_store_repository.rs +++ b/crates/server/src/services/mcp_store_repository.rs @@ -21,6 +21,9 @@ use uuid::Uuid; /// Reserved advisory lock keys (keep this list current): /// * `MCP_STORE_INSTALL_LOCK_KEY` (here): template-install /// serialization in `create_server` when `template_slug` is set. +/// * `think_watch_common::db::SCHEMA_LOCK` ("twschema") and +/// `CLICKHOUSE_SETUP_LOCK` ("twchinit"): one instance at a time +/// through the start-up schema setup. const MCP_STORE_INSTALL_LOCK_KEY: i64 = 0x6D637053746F7265; /// A category and how many templates are in it. diff --git a/crates/test-support/src/lib.rs b/crates/test-support/src/lib.rs index fa8db92a..fc0791f9 100644 --- a/crates/test-support/src/lib.rs +++ b/crates/test-support/src/lib.rs @@ -151,9 +151,7 @@ impl TestApp { pub async fn try_spawn_with(opts: SpawnOptions) -> anyhow::Result { init_test_tracing(); - let base_url = std::env::var("TEST_DATABASE_BASE_URL").unwrap_or_else(|_| { - "postgres://thinkwatch:7c3fe6307d00fe3f2f29f534e806ac71@localhost:5432".into() - }); + let base_url = database_base_url(); // Default to logical DB 1 so we never trample the dev Redis // (DB 0). Override via env when CI uses a dedicated instance. let redis_url = std::env::var("TEST_REDIS_URL").unwrap_or_else(|_| { @@ -195,14 +193,7 @@ impl TestApp { // Per-test ClickHouse — only when the test asked for it. let (ch_owner, ch_client, ch_url, ch_db, ch_user, ch_password) = if opts.clickhouse { - let url = std::env::var("TEST_CLICKHOUSE_URL") - .unwrap_or_else(|_| "http://localhost:8123".into()); - let user = std::env::var("TEST_CLICKHOUSE_USER") - .ok() - .or_else(|| Some("thinkwatch".into())); - let password = std::env::var("TEST_CLICKHOUSE_PASSWORD") - .ok() - .or_else(|| Some("c693ded3da8388c7b6a4288dac91a2ad".into())); + let (url, user, password) = clickhouse_env(); let owner = IsolatedClickHouseDatabase::create(&url, user.as_deref(), password.as_deref()) .await @@ -523,6 +514,28 @@ fn init_test_tracing() { }); } +/// The Postgres server tests create their databases on +/// (`TEST_DATABASE_BASE_URL`), without a database name. +pub fn database_base_url() -> String { + std::env::var("TEST_DATABASE_BASE_URL").unwrap_or_else(|_| { + "postgres://thinkwatch:7c3fe6307d00fe3f2f29f534e806ac71@localhost:5432".into() + }) +} + +/// The ClickHouse server tests create their databases on: URL, user and +/// password (`TEST_CLICKHOUSE_URL`, `_USER`, `_PASSWORD`). +pub fn clickhouse_env() -> (String, Option, Option) { + let url = + std::env::var("TEST_CLICKHOUSE_URL").unwrap_or_else(|_| "http://localhost:8123".into()); + let user = std::env::var("TEST_CLICKHOUSE_USER") + .ok() + .or_else(|| Some("thinkwatch".into())); + let password = std::env::var("TEST_CLICKHOUSE_PASSWORD") + .ok() + .or_else(|| Some("c693ded3da8388c7b6a4288dac91a2ad".into())); + (url, user, password) +} + /// Convenience re-exports so test files only need one `use`. pub mod prelude { pub use crate::SpawnOptions; diff --git a/crates/test-support/src/pg.rs b/crates/test-support/src/pg.rs index ab77687d..77d37707 100644 --- a/crates/test-support/src/pg.rs +++ b/crates/test-support/src/pg.rs @@ -24,6 +24,17 @@ impl IsolatedDatabase { /// including) the database name — e.g. /// `postgres://user:pwd@localhost:5432`. pub async fn create(base_url: &str) -> anyhow::Result { + let db = Self::create_empty(base_url).await?; + // Use the workspace migrator from common so the schema stays + // in lockstep with production. + think_watch_common::db::run_migrations(db.pool()) + .await + .context("run migrations into per-test DB")?; + Ok(db) + } + + /// A fresh database with nothing in it: what a first start finds. + pub async fn create_empty(base_url: &str) -> anyhow::Result { let admin_url = format!("{}/postgres", base_url.trim_end_matches('/')); let name = format!("test_tw_{}", Uuid::new_v4().simple()); @@ -45,12 +56,6 @@ impl IsolatedDatabase { .await .context("connect to per-test DB")?; - // Use the workspace migrator from common so the schema stays - // in lockstep with production. - think_watch_common::db::run_migrations(&pool) - .await - .context("run migrations into per-test DB")?; - Ok(Self { base_url: base_url.trim_end_matches('/').to_string(), name, diff --git a/crates/test-support/tests/concurrent_startup.rs b/crates/test-support/tests/concurrent_startup.rs new file mode 100644 index 00000000..ca8d6057 --- /dev/null +++ b/crates/test-support/tests/concurrent_startup.rs @@ -0,0 +1,339 @@ +//! Several instances starting at the same moment against one database: +//! a Helm `replicaCount` above 1, a rolling upgrade, an autoscaler adding +//! pods. Every instance applies `db/schema.sql`, the seeds and the +//! guard-settings conversion, then sets up the ClickHouse tables and +//! backfills the rollups. Run side by side, the schema statements of one +//! instance deadlocked with another's, and one of them exited at its +//! migration; two backfills of the same empty rollup each copied the +//! whole log into it. +//! +//! Each round here starts the setup on [`INSTANCES`] tasks at once, each +//! with a connection pool of its own as separate processes would have: +//! every one succeeds, and the database ends up as one instance alone +//! leaves it. + +use std::sync::Arc; +use std::time::Duration; + +use serde_json::Value; +use sqlx::PgPool; +use sqlx::postgres::PgPoolOptions; +use think_watch_common::db::{SCHEMA_LOCK, StartupLock, run_migrations}; +use think_watch_test_support::{ + IsolatedClickHouseDatabase, IsolatedDatabase, clickhouse_env, database_base_url, +}; +use tokio::sync::Barrier; + +const INSTANCES: usize = 4; + +/// One pool per instance, against the same database. +async fn instance_pools(db: &IsolatedDatabase) -> Vec { + let mut pools = Vec::with_capacity(INSTANCES); + for _ in 0..INSTANCES { + pools.push( + PgPoolOptions::new() + .max_connections(4) + .connect(db.url()) + .await + .expect("connect an instance's pool"), + ); + } + pools +} + +/// Start the schema setup on every instance at the same moment and +/// return what each one returned. +async fn migrate_at_once(db: &IsolatedDatabase) -> Vec> { + let barrier = Arc::new(Barrier::new(INSTANCES)); + let mut tasks = Vec::with_capacity(INSTANCES); + for pool in instance_pools(db).await { + let barrier = barrier.clone(); + tasks.push(tokio::spawn(async move { + barrier.wait().await; + let result = run_migrations(&pool).await.map_err(|e| format!("{e:#}")); + pool.close().await; + result + })); + } + let mut results = Vec::with_capacity(INSTANCES); + for task in tasks { + results.push(task.await.expect("an instance's setup panicked")); + } + results +} + +fn assert_all_started(results: &[Result<(), String>]) { + let failed: Vec<&String> = results.iter().filter_map(|r| r.as_ref().err()).collect(); + assert!( + failed.is_empty(), + "{} of {INSTANCES} instances failed their migration: {failed:#?}", + failed.len() + ); +} + +/// Every object the schema setup creates, one line each, sorted. +async fn schema_of(pool: &PgPool) -> Vec { + sqlx::query_scalar( + "SELECT 'column ' || table_name || '.' || column_name || ' ' || data_type + || ' nullable=' || is_nullable || ' default=' || coalesce(column_default, '') + FROM information_schema.columns WHERE table_schema = 'public' + UNION ALL + SELECT 'index ' || indexdef FROM pg_indexes WHERE schemaname = 'public' + UNION ALL + SELECT 'constraint ' || conrelid::regclass::text || ' ' || conname || ' ' + || pg_get_constraintdef(oid) + FROM pg_constraint WHERE connamespace = 'public'::regnamespace + UNION ALL + SELECT 'trigger ' || pg_get_triggerdef(oid) FROM pg_trigger WHERE NOT tgisinternal + UNION ALL + SELECT 'function ' || oid::regprocedure::text + FROM pg_proc WHERE pronamespace = 'public'::regnamespace + UNION ALL + SELECT 'extension ' || extname FROM pg_extension + ORDER BY 1", + ) + .fetch_all(pool) + .await + .expect("read the schema") +} + +/// Every row the seeds write, one line each, sorted. The conversion's +/// record says when it was written, so only its presence is compared. +async fn seeds_of(pool: &PgPool) -> Vec { + sqlx::query_scalar( + "SELECT 'role ' || name || ' ' || policy_document::text FROM rbac_roles + UNION ALL + SELECT 'surface ' || name FROM api_key_surface_kinds + UNION ALL + SELECT 'pricing ' || id::text FROM platform_pricing + UNION ALL + SELECT 'setting ' || key || ' ' + || CASE WHEN key = 'security.legacy_converted' THEN 'present' ELSE value::text END + FROM system_settings + UNION ALL + SELECT 'template ' || slug FROM mcp_store_templates + ORDER BY 1", + ) + .fetch_all(pool) + .await + .expect("read the seeded rows") +} + +async fn setting(pool: &PgPool, key: &str) -> Option { + sqlx::query_scalar("SELECT value FROM system_settings WHERE key = $1") + .bind(key) + .fetch_optional(pool) + .await + .unwrap() +} + +#[ignore = "integration test — run via `make test-it`"] +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] +async fn replicas_starting_together_on_a_fresh_database_all_start() { + let base = database_base_url(); + let alone = IsolatedDatabase::create(&base).await.unwrap(); + let db = IsolatedDatabase::create_empty(&base).await.unwrap(); + + // The first start of a deployment: an empty database. + assert_all_started(&migrate_at_once(&db).await); + assert_eq!(schema_of(db.pool()).await, schema_of(alone.pool()).await); + assert_eq!(seeds_of(db.pool()).await, seeds_of(alone.pool()).await); + + // All of them restarting at once (a rollout of the same version): + // everything is there already, and nothing changes. + assert_all_started(&migrate_at_once(&db).await); + assert_eq!(schema_of(db.pool()).await, schema_of(alone.pool()).await); + assert_eq!(seeds_of(db.pool()).await, seeds_of(alone.pool()).await); +} + +#[ignore = "integration test — run via `make test-it`"] +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] +async fn replicas_upgrading_together_apply_the_upgrade_once() { + let base = database_base_url(); + let alone = IsolatedDatabase::create(&base).await.unwrap(); + let db = IsolatedDatabase::create(&base).await.unwrap(); + + // The database as an earlier version leaves it: columns and an index + // the schema has added since are missing, and the guard settings are + // still the old ones, with nothing converted yet. + sqlx::raw_sql( + r#"ALTER TABLE models DROP COLUMN max_output_tokens; + ALTER TABLE model_routes DROP COLUMN upstream_protocol; + DROP INDEX idx_api_keys_user_not_deleted; + DELETE FROM system_settings WHERE key = 'security.legacy_converted'; + INSERT INTO system_settings (key, value, category, description) VALUES + ('security.hidden_text', '"block"', 'security', 'old'), + ('security.tool_inspection', + '{"mode": "enforce", "disabled": ["chmod-777"], "actions": {}, "custom": []}', + 'security', 'old'); + ALTER TABLE models ADD COLUMN output_guardrails JSONB NOT NULL DEFAULT '[]'::jsonb;"#, + ) + .execute(db.pool()) + .await + .unwrap(); + + assert_all_started(&migrate_at_once(&db).await); + + // The schema is complete, the old column is gone... + assert_eq!(schema_of(db.pool()).await, schema_of(alone.pool()).await); + // ...and the old settings were converted, by one of them. + for key in ["security.hidden_text", "security.tool_inspection"] { + assert_eq!(setting(db.pool(), key).await, None, "{key}"); + } + let marker = setting(db.pool(), "security.legacy_converted") + .await + .expect("the conversion was recorded"); + let mut converted: Vec<&str> = marker["converted"] + .as_array() + .expect("{marker}") + .iter() + .filter_map(Value::as_str) + .collect(); + converted.sort_unstable(); + assert_eq!( + converted, + ["security.hidden_text", "security.tool_inspection"], + "{marker}" + ); + assert_eq!( + setting(db.pool(), "security.inspect_tools").await, + Some(serde_json::json!({"mode": "enforce", "disable": ["chmod-777"]})), + ); +} + +/// Advisory locks held in `pool`'s database: `(granted, waiting)`. +async fn advisory_locks(pool: &PgPool) -> (i64, i64) { + sqlx::query_as( + "SELECT count(*) FILTER (WHERE granted), count(*) FILTER (WHERE NOT granted) + FROM pg_locks + WHERE locktype = 'advisory' + AND database = (SELECT oid FROM pg_database WHERE datname = current_database())", + ) + .fetch_one(pool) + .await + .unwrap() +} + +/// Wait until `pool`'s database shows `want` advisory locks. +async fn until_advisory_locks(pool: &PgPool, want: (i64, i64), what: &str) { + tokio::time::timeout(Duration::from_secs(10), async { + while advisory_locks(pool).await != want { + tokio::time::sleep(Duration::from_millis(20)).await; + } + }) + .await + .unwrap_or_else(|_| panic!("{what}: the advisory locks never came to {want:?}")); +} + +#[ignore = "integration test — run via `make test-it`"] +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] +async fn the_schema_lock_never_outlives_its_holder() { + let db = IsolatedDatabase::create(&database_base_url()) + .await + .unwrap(); + + // An instance that dies holding the lock: its session ends, and the + // instance waiting for it goes ahead. + let mut holder = StartupLock::acquire(db.pool(), SCHEMA_LOCK, "testing") + .await + .unwrap(); + let pid: i32 = sqlx::query_scalar("SELECT pg_backend_pid()") + .fetch_one(holder.conn()) + .await + .unwrap(); + let pool = db.pool().clone(); + let waiter = tokio::spawn(async move { run_migrations(&pool).await }); + until_advisory_locks(db.pool(), (1, 1), "the second instance waits").await; + sqlx::query("SELECT pg_terminate_backend($1)") + .bind(pid) + .execute(db.pool()) + .await + .unwrap(); + tokio::time::timeout(Duration::from_secs(30), waiter) + .await + .expect("the waiting instance got the lock once its holder was gone") + .unwrap() + .unwrap(); + drop(holder); + until_advisory_locks(db.pool(), (0, 0), "released after the migration").await; + + // A setup that fails drops the lock without releasing it: its + // connection is closed, not handed back to the pool still holding it. + let lock = StartupLock::acquire(db.pool(), SCHEMA_LOCK, "testing") + .await + .unwrap(); + assert_eq!(advisory_locks(db.pool()).await, (1, 0)); + drop(lock); + until_advisory_locks(db.pool(), (0, 0), "released on drop").await; +} + +#[ignore = "integration test — run via `make test-it`"] +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] +async fn replicas_starting_together_backfill_the_rollups_once() { + let pg = IsolatedDatabase::create(&database_base_url()) + .await + .unwrap(); + let (url, user, password) = clickhouse_env(); + let ch = IsolatedClickHouseDatabase::create(&url, user.as_deref(), password.as_deref()) + .await + .unwrap(); + let client = ch.client().clone(); + + // Requests logged before the rollups existed: the rows are in + // gateway_logs, and the rollups a new version adds start empty. + client + .query( + "INSERT INTO gateway_logs (id, model_id, provider, status_code, latency_ms, \ + input_tokens, output_tokens) \ + SELECT toString(number), 'm', 'p', 200, 10, 3, 4 FROM numbers(25)", + ) + .execute() + .await + .unwrap(); + for table in ["cost_rollup_hourly", "provider_health_5m"] { + client + .query(&format!("TRUNCATE TABLE {table}")) + .execute() + .await + .unwrap(); + } + + let barrier = Arc::new(Barrier::new(INSTANCES)); + let mut tasks = Vec::with_capacity(INSTANCES); + for pool in instance_pools(&pg).await { + let barrier = barrier.clone(); + let ch = Some(client.clone()); + tasks.push(tokio::spawn(async move { + barrier.wait().await; + let result = think_watch_server::init::ensure_clickhouse_schema(&pool, &ch) + .await + .map_err(|e| format!("{e:#}")); + pool.close().await; + result + })); + } + let mut results = Vec::with_capacity(INSTANCES); + for task in tasks { + results.push(task.await.expect("an instance's setup panicked")); + } + assert_all_started(&results); + + let requests: u64 = client + .query("SELECT toUInt64(sum(request_count)) FROM cost_rollup_hourly") + .fetch_one() + .await + .unwrap(); + assert_eq!(requests, 25, "cost_rollup_hourly counts each request once"); + let tokens: i64 = client + .query("SELECT toInt64(sum(output_tokens)) FROM cost_rollup_hourly") + .fetch_one() + .await + .unwrap(); + assert_eq!(tokens, 100); + let health: u64 = client + .query("SELECT toUInt64(sum(total_requests)) FROM provider_health_5m") + .fetch_one() + .await + .unwrap(); + assert_eq!(health, 25, "provider_health_5m counts each request once"); +} diff --git a/db/schema.sql b/db/schema.sql index 0245ea58..877519f4 100644 --- a/db/schema.sql +++ b/db/schema.sql @@ -5,7 +5,10 @@ -- in place when the schema changes; the application calls -- `sqlx::raw_sql(include_str!("../../../db/schema.sql"))` on every -- boot, and every statement here is wrapped in `IF NOT EXISTS` / --- `OR REPLACE` so a re-run is a no-op on an up-to-date DB. +-- `OR REPLACE` so a re-run is a no-op on an up-to-date DB. Replicas +-- starting together apply it one at a time, under an advisory lock +-- (crates/common/src/db.rs::run_migrations): run side by side, its DDL +-- deadlocks. -- -- Limits of declarative apply: -- * column rename, type narrowing, or DROP COLUMN need an explicit From 70b811701cb183e095ee41addaa9c9f40c5c3f85 Mon Sep 17 00:00:00 2001 From: fylorn <249551762+fylorn@users.noreply.github.com> Date: Mon, 5 Oct 2026 18:39:58 +0800 Subject: [PATCH 2/5] fix(helm): let the network policy follow the database URLs' ports With networkPolicy.enabled the server's egress allowed Postgres on 5432, Redis on 6379 and ClickHouse on 8123 whatever postgres/redis/clickhouse .externalUrl said. A database on any other port was blocked: Azure Cache for Redis over TLS (6380), ClickHouse Cloud (8443), a managed Postgres on a port of its own. The server then could not start, or for ClickHouse started without writing its logs. The allowed ports now come from the configured endpoint. A bundled service keeps its fixed port. For an external URL, tw.urlPorts takes every port the URL names - each host, comma-separated hosts included, each node= of a Redis Cluster or Sentinel URL, and a Postgres ?port= - and, for a host without one, the default the client uses for the scheme: 5432 for postgres, 6379 for redis and rediss (the server's Redis client doesn't change it for TLS, so Azure's 6380 has to be in the URL anyway), 26379 for a Sentinel plus 6379 for the primary it points to, 80 for http and 443 for https (not ClickHouse's 8123/8443: the HTTP client goes by the URL). So the policy allows the ports the server dials for what the URL names, with no second setting that could disagree with the URL, as an explicit redis.port could. The URL is read with regexes, not urlParse, which panics on a URL it can't parse: a malformed one still renders, with the default port. What no URL can say - the ports Redis Cluster nodes announce beyond the seeds, a Sentinel's primary on another port, an upstream or MCP server on a port other than 443 - goes in networkPolicy.extraEgress, rules appended to the server's egress as written. Rendered with Helm 3.22.0 and 4.3.0: with bundled databases and with the README's external values the policy is unchanged apart from its comments, and nothing outside networkpolicy.yaml changes. Co-Authored-By: Claude Opus 5.5 --- CHANGELOG.md | 12 ++++ deploy/helm/think-watch/README.md | 40 +++++++++++- .../helm/think-watch/templates/_helpers.tpl | 65 +++++++++++++++++++ .../think-watch/templates/networkpolicy.yaml | 25 +++---- deploy/helm/think-watch/values.yaml | 19 ++++++ 5 files changed, 147 insertions(+), 14 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 37c7f4b4..d07c4c58 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -89,6 +89,18 @@ target. logging `Another instance is setting up the database schema; waiting for it to finish`, then find it done. An instance that dies holding a lock releases it with its connection. +- **Helm network policy and databases on other ports.** With + `networkPolicy.enabled`, the server could reach PostgreSQL only on `5432`, + Redis on `6379` and ClickHouse on `8123`, whatever their `externalUrl` said, + so a database on another port was blocked — Azure Cache for Redis over TLS + (`6380`), ClickHouse Cloud (`8443`), a managed Postgres on a port of its own: + the server could not start, or started without writing to ClickHouse. The + allowed ports now follow `postgres.externalUrl`, `redis.externalUrl` and + `clickhouse.externalUrl`: every port a URL names, and the client's default + for its scheme where it names none. `networkPolicy.extraEgress` adds egress + rules as written, for ports no URL names (Redis Cluster nodes announcing + other ports, an upstream or MCP server on a port other than `443`). The + chart's README describes both. ## [3.1.0] — 2026-10-05 diff --git a/deploy/helm/think-watch/README.md b/deploy/helm/think-watch/README.md index 5efc586d..a5720a56 100644 --- a/deploy/helm/think-watch/README.md +++ b/deploy/helm/think-watch/README.md @@ -86,7 +86,9 @@ redis: - Every node must be reachable from the server pods at the address it announces to the cluster (`cluster-announce-ip` / `-port`): the server - follows the cluster's redirects to it. + follows the cluster's redirects to it. With `networkPolicy.enabled`, + ports the URL doesn't name go in `networkPolicy.extraEgress` (see + [Network policy](#network-policy)). - A cluster has only database 0, so the URL names no `/`. - Nothing else is needed. Every script the server runs keeps its keys in one hash slot — the counters of one user's request share the tag @@ -119,8 +121,8 @@ redis: IP addresses that their certificates do not name cannot be used over TLS. - The port is whatever the service uses for TLS (Azure Cache for Redis: - `6380`). With `networkPolicy.enabled`, the chart's egress rule lets - the server reach Redis on `6379` only. + `6380`). Write it in the URL: a `rediss://` URL without one means + `6379`, as `redis://` does. A self-hosted Redis whose certificate a private CA signed needs that CA. Put its PEM certificate in a Secret and name it; the server then @@ -146,6 +148,38 @@ server pods after changing the Secret. Outside the chart, set (mutual TLS) are not supported: give such a Redis `tls-auth-clients no` and authenticate with the password. +## Network policy + +`networkPolicy.enabled` limits what the server pods may reach: DNS, port +`443` (upstreams, the OIDC provider), `9000` (ClickHouse's native +protocol), and PostgreSQL, Redis and +ClickHouse on the ports the server connects to them on: + +- A bundled database: its service port (`5432`, `6379`, `8123`). +- An external one: every port its `externalUrl` names, so a database on + another port needs no setting of its own. That is the port of each host + in the URL, of each `node=` of a Redis Cluster or Sentinel URL, and a + Postgres `?port=`. +- A URL without a port: the client's default for the scheme, which is + `5432` for `postgres://`, `6379` for `redis://` and `rediss://` (TLS + does not change it), `26379` for a Sentinel and `6379` for the primary it + points to, `80` for `http://` and `443` for `https://`. + +What no URL names goes in `networkPolicy.extraEgress`, rules added to the +server's egress as written: Redis Cluster nodes that announce ports the +URL doesn't list, a Sentinel's primary on a port other than `6379`, an +upstream or MCP server on a port other than `443`. + +```yaml +networkPolicy: + enabled: true + extraEgress: + - ports: + - port: 7000 + endPort: 7005 + protocol: TCP +``` + ## Rotating secrets `-secrets` is kept on `helm uninstall`. To rotate passwords: diff --git a/deploy/helm/think-watch/templates/_helpers.tpl b/deploy/helm/think-watch/templates/_helpers.tpl index 6735a849..1bb7c69c 100644 --- a/deploy/helm/think-watch/templates/_helpers.tpl +++ b/deploy/helm/think-watch/templates/_helpers.tpl @@ -24,3 +24,68 @@ Resource names for bundled databases — stable across upgrades. {{- define "tw.redis.name" -}}{{ .Release.Name }}-redis{{- end -}} {{- define "tw.clickhouse.name" -}}{{ .Release.Name }}-clickhouse{{- end -}} {{- define "tw.secret.name" -}}{{ .Release.Name }}-secrets{{- end -}} + +{{/* +The TCP ports a URL reaches, space-separated: the port of each host in it +(comma-separated hosts included), of each `node=` in its query (the seed +nodes of a Redis Cluster or Sentinel) and of a `port=` in its query (which +Postgres honours); `default` for a host written without a port. Never +fails on a URL it cannot read: the default stands. +Usage: {{ include "tw.urlPorts" (dict "url" $url "default" 6379) }} +*/}} +{{- define "tw.urlPorts" -}} +{{- $rest := regexReplaceAll "^[A-Za-z][A-Za-z0-9+.-]*://" .url "" -}} +{{- $hosts := regexReplaceAll "^.*@" (regexFind "^[^/?#]*" $rest) "" -}} +{{- $query := "" -}} +{{- if contains "?" $rest -}} +{{- $query = regexReplaceAll "#.*$" (regexReplaceAll "^[^?]*[?]" $rest "") "" -}} +{{- end -}} +{{- $ports := list -}} +{{- range $host := splitList "," $hosts -}} +{{- $port := trimPrefix ":" (regexFind ":[0-9]+$" (regexReplaceAll "\\[[^]]*\\]" $host "")) -}} +{{- $ports = append $ports (default (toString $.default) $port) -}} +{{- end -}} +{{- range $param := splitList "&" $query -}} +{{- if hasPrefix "node=" $param -}} +{{- $node := trimPrefix "node=" $param | replace "%3A" ":" | replace "%3a" ":" -}} +{{- $port := trimPrefix ":" (regexFind ":[0-9]+$" (regexReplaceAll "\\[[^]]*\\]" $node "")) -}} +{{- if $port -}}{{- $ports = append $ports $port -}}{{- end -}} +{{- else if hasPrefix "port=" $param -}} +{{- $port := regexFind "^[0-9]+$" (trimPrefix "port=" $param) -}} +{{- if $port -}}{{- $ports = append $ports $port -}}{{- end -}} +{{- end -}} +{{- end -}} +{{- $ports | uniq | join " " -}} +{{- end -}} + +{{/* +The ports the server reaches each database on, for the network policy: +the bundled service's, or those of its externalUrl, with the client's +default for the scheme where the URL names none (Postgres 5432; Redis +6379, a Sentinel 26379 and 6379 for its primary; ClickHouse over http 80, +over https 443). +*/}} +{{- define "tw.postgres.ports" -}} +{{- if .Values.postgres.bundled -}}5432 +{{- else -}}{{ include "tw.urlPorts" (dict "url" .Values.postgres.externalUrl "default" 5432) }} +{{- end -}} +{{- end -}} + +{{- define "tw.redis.ports" -}} +{{- if .Values.redis.bundled -}}6379 +{{- else -}} +{{- $url := .Values.redis.externalUrl -}} +{{- $sentinel := or (regexMatch "^[A-Za-z]+-sentinel://" $url) (contains "sentinelServiceName=" $url) -}} +{{- include "tw.urlPorts" (dict "url" $url "default" (ternary 26379 6379 $sentinel)) -}} +{{- /* The primary a Sentinel points to: the URL can't say where. */ -}} +{{- if $sentinel }} 6379{{ end -}} +{{- end -}} +{{- end -}} + +{{- define "tw.clickhouse.ports" -}} +{{- if .Values.clickhouse.bundled -}}8123 +{{- else -}} +{{- $url := .Values.clickhouse.externalUrl -}} +{{- include "tw.urlPorts" (dict "url" $url "default" (ternary 443 80 (hasPrefix "https://" (lower $url)))) -}} +{{- end -}} +{{- end -}} diff --git a/deploy/helm/think-watch/templates/networkpolicy.yaml b/deploy/helm/think-watch/templates/networkpolicy.yaml index d2c77e56..15d269ef 100644 --- a/deploy/helm/think-watch/templates/networkpolicy.yaml +++ b/deploy/helm/think-watch/templates/networkpolicy.yaml @@ -1,6 +1,7 @@ {{- if .Values.networkPolicy.enabled }} # Server pod: allow ingress on gateway (3000) and console (3001) ports, -# and egress to PostgreSQL, Redis, ClickHouse, and DNS. +# and egress to PostgreSQL, Redis, ClickHouse, DNS, HTTPS and +# networkPolicy.extraEgress. apiVersion: networking.k8s.io/v1 kind: NetworkPolicy metadata: @@ -37,18 +38,16 @@ spec: protocol: UDP - port: 53 protocol: TCP - # PostgreSQL + # PostgreSQL, Redis and ClickHouse (HTTP), on the ports the server + # connects to them on (tw..ports in _helpers.tpl) + {{- range $db := list "postgres" "redis" "clickhouse" }} + # {{ $db }} - ports: - - port: 5432 - protocol: TCP - # Redis - - ports: - - port: 6379 - protocol: TCP - # ClickHouse HTTP - - ports: - - port: 8123 + {{- range splitList " " (include (printf "tw.%s.ports" $db) $) }} + - port: {{ . }} protocol: TCP + {{- end }} + {{- end }} # ClickHouse native - ports: - port: 9000 @@ -57,6 +56,10 @@ spec: - ports: - port: 443 protocol: TCP + {{- with .Values.networkPolicy.extraEgress }} + # networkPolicy.extraEgress + {{- toYaml . | nindent 4 }} + {{- end }} --- # Web pod: allow ingress on port 80, egress only to server console port and DNS. apiVersion: networking.k8s.io/v1 diff --git a/deploy/helm/think-watch/values.yaml b/deploy/helm/think-watch/values.yaml index 38b9509a..b7b9bc5e 100644 --- a/deploy/helm/think-watch/values.yaml +++ b/deploy/helm/think-watch/values.yaml @@ -135,6 +135,25 @@ clickhouse: networkPolicy: enabled: false + # The server may reach PostgreSQL, Redis and ClickHouse on the ports it + # connects to them on: a bundled service's port, or every port its + # externalUrl names (each host, each `node=` of a Redis Cluster or + # Sentinel URL, a Postgres `port=`). Where the URL names no port, the + # client's default for its scheme: postgres:// 5432; redis:// and + # rediss:// 6379; a Sentinel 26379, plus 6379 for the primary it points + # to; http:// 80; https:// 443. A database on another port, such as Azure + # Cache for Redis over TLS (rediss://...:6380), needs nothing more than + # the port in its URL. + # + # Egress rules added as written, for ports that no URL names: Redis + # Cluster nodes announcing ports the URL doesn't list, a Sentinel's + # primary on a port other than 6379, an upstream or MCP server on a port + # other than 443. + extraEgress: [] + # - ports: + # - port: 7000 + # endPort: 7005 + # protocol: TCP autoscaling: enabled: false From 9da239d64c612c3deda6a92961eeb016b4f4b950 Mon Sep 17 00:00:00 2001 From: fylorn <249551762+fylorn@users.noreply.github.com> Date: Mon, 5 Oct 2026 18:54:31 +0800 Subject: [PATCH 3/5] fix(clickhouse): stop resetting the body-column TTL on every start ensure_clickhouse_tables runs deploy/clickhouse/initdb.d/01_init.sql on every start, and the file ended its gateway_logs and mcp_logs sections with ALTER TABLE ... MODIFY COLUMN TTL ... 30 DAY, for request_body, response_body, tool_arguments and tool_result. The configured audit.body_retention_days is put back afterwards, by reconcile_clickhouse_ttls in main.rs, but ClickHouse materializes a TTL on the data already stored as soon as it is set (materialize_ttl_after_modify), so a deployment keeping bodies longer than 30 days lost the ones older than that on each restart; the rows stayed, their bodies became NULL. Reproduced against ClickHouse 26.3: body TTL set to 90 days as the reconcile leaves it, a row from 60 days ago with its bodies, then the start-up sequence. Right after the table setup SHOW CREATE TABLE showed toIntervalDay(30) on all four columns, and once the mutations had run the row's bodies were NULL - also when the 90-day TTL was restored immediately after, as a real start does. The TTL now sits in the ADD COLUMN IF NOT EXISTS that creates each column, and the MODIFY statements are gone: a new column gets the 30-day default (a new column has nothing to lose), an existing one keeps what the server last set. Reading the configured value before the table setup would have needed Postgres in ensure_clickhouse_tables and a second place deciding the TTL; the reconcile stays the only one. The other TTLs at start-up don't have the pattern. Each log table's TTL is only in its CREATE TABLE IF NOT EXISTS, which leaves an existing table alone, and the reconcile sets the configured values. (It re-issues them on every start even when unchanged, which costs a TTL mutation per table but loses nothing.) The chart's own init SQL for the bundled ClickHouse runs once, on an empty data directory, and sets no body TTL. tests/clickhouse_retention_restart.rs runs a deployment with 90-day bodies and 365-day logs through a restart, step by step, and checks every TTL after each step and that the old bodies and rows are still there. Before the change it failed after the table setup (30 instead of 90) and, without that check, on the bodies (NULL). Co-Authored-By: Claude Opus 5.5 --- CHANGELOG.md | 11 ++ .../tests/clickhouse_retention_restart.rs | 166 ++++++++++++++++++ deploy/clickhouse/initdb.d/01_init.sql | 36 ++-- 3 files changed, 194 insertions(+), 19 deletions(-) create mode 100644 crates/test-support/tests/clickhouse_retention_restart.rs diff --git a/CHANGELOG.md b/CHANGELOG.md index d07c4c58..032b3c2e 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -89,6 +89,17 @@ target. logging `Another instance is setting up the database schema; waiting for it to finish`, then find it done. An instance that dies holding a lock releases it with its connection. +- **Captured bodies kept as long as configured.** With ClickHouse and + `audit.body_retention_days` above 30, every server start could clear the + captured request and response bodies older than 30 days + (`gateway_logs.request_body` / `response_body`, `mcp_logs.tool_arguments` + / `tool_result`; the rows themselves stayed). The start-up table setup set + those columns' TTL to 30 days each time, and ClickHouse applies a TTL to + the data already stored as soon as it is set, before the server put the + configured TTL back a moment later. The setup now gives these columns a + TTL only when it creates them, so a restart leaves the configured one in + place. Bodies already cleared cannot be recovered. The log tables' own + TTLs (`data.retention_days_*`) were not affected. - **Helm network policy and databases on other ports.** With `networkPolicy.enabled`, the server could reach PostgreSQL only on `5432`, Redis on `6379` and ClickHouse on `8123`, whatever their `externalUrl` said, diff --git a/crates/test-support/tests/clickhouse_retention_restart.rs b/crates/test-support/tests/clickhouse_retention_restart.rs new file mode 100644 index 00000000..bbef5c5a --- /dev/null +++ b/crates/test-support/tests/clickhouse_retention_restart.rs @@ -0,0 +1,166 @@ +//! A restart never takes a ClickHouse retention below the configured one. +//! +//! The server sets up the ClickHouse tables on every start, then applies +//! the configured retentions: `data.retention_days_*` as each log table's +//! TTL, `audit.body_retention_days` as the TTL of the columns holding +//! request and response bodies. The table setup used to give the body +//! columns a 30-day TTL each time, and ClickHouse acts on a TTL as soon +//! as it is set: with a longer body retention configured, every restart +//! cleared the bodies older than 30 days, although the configured TTL was +//! back a moment later. + +use think_watch_server::handlers::admin::reconcile_clickhouse_ttls; +use think_watch_server::init::ensure_clickhouse_schema; +use think_watch_test_support::prelude::*; + +const BODY_COLUMNS: [(&str, &str); 4] = [ + ("gateway_logs", "request_body"), + ("gateway_logs", "response_body"), + ("mcp_logs", "tool_arguments"), + ("mcp_logs", "tool_result"), +]; + +/// The days in a `toIntervalDay(N)` on `line`. +fn interval_days(line: &str) -> Option { + let rest = &line[line.find("toIntervalDay(")? + "toIntervalDay(".len()..]; + rest[..rest.find(')')?].parse().ok() +} + +/// Every TTL that matters here, in days: each log table's and each body +/// column's, as `SHOW CREATE TABLE` gives them. +async fn ttls(ch: &clickhouse::Client) -> Vec<(String, Option)> { + let mut out = Vec::new(); + for table in ["gateway_logs", "mcp_logs"] { + let create: String = ch + .query(&format!("SHOW CREATE TABLE {table}")) + .fetch_one() + .await + .unwrap(); + let table_ttl = create + .lines() + .find(|l| l.starts_with("TTL ")) + .and_then(interval_days); + out.push((table.to_string(), table_ttl)); + for (_, column) in BODY_COLUMNS.iter().filter(|(t, _)| *t == table) { + let line = create + .lines() + .find(|l| l.trim_start().starts_with(&format!("`{column}` "))) + .unwrap_or_else(|| panic!("no {table}.{column} in {create}")); + out.push((format!("{table}.{column}"), interval_days(line))); + } + } + out +} + +fn configured() -> Vec<(String, Option)> { + vec![ + ("gateway_logs".into(), Some(365)), + ("gateway_logs.request_body".into(), Some(90)), + ("gateway_logs.response_body".into(), Some(90)), + ("mcp_logs".into(), Some(365)), + ("mcp_logs.tool_arguments".into(), Some(90)), + ("mcp_logs.tool_result".into(), Some(90)), + ] +} + +/// Wait until ClickHouse has carried out every mutation of this database, +/// the TTL it materializes after an `ALTER … TTL` among them. +async fn mutations_done(ch: &clickhouse::Client) { + tokio::time::timeout(std::time::Duration::from_secs(60), async { + loop { + let pending: u64 = ch + .query( + "SELECT count() FROM system.mutations \ + WHERE database = currentDatabase() AND NOT is_done", + ) + .fetch_one() + .await + .unwrap(); + if pending == 0 { + return; + } + tokio::time::sleep(std::time::Duration::from_millis(100)).await; + } + }) + .await + .expect("ClickHouse finished its mutations"); +} + +#[ignore = "integration test — run via `make test-it`"] +#[tokio::test] +async fn a_restart_keeps_retentions_longer_than_the_defaults() { + let app = TestApp::spawn_with_clickhouse().await; + let ch = app.state.clickhouse.clone().expect("ClickHouse"); + + // A deployment that keeps bodies 90 days and logs a year, running. + for (key, days) in [ + ("audit.body_retention_days", 90), + ("data.retention_days_gateway", 365), + ("data.retention_days_mcp", 365), + ] { + fixtures::set_setting(&app.db, key, json!(days)) + .await + .unwrap(); + } + app.state.dynamic_config.reload().await.unwrap(); + reconcile_clickhouse_ttls(&app.state).await; + mutations_done(&ch).await; + assert_eq!(ttls(&ch).await, configured(), "as the deployment runs"); + + // Requests from 60 days ago, their bodies within the retention, and + // from 200 days ago, within the logs'. + ch.query( + "INSERT INTO gateway_logs (id, model_id, request_body, response_body, created_at) VALUES \ + ('60d', 'm', 'request', 'response', now64(3) - INTERVAL 60 DAY)", + ) + .execute() + .await + .unwrap(); + ch.query( + "INSERT INTO gateway_logs (id, model_id, created_at) VALUES \ + ('200d', 'm', now64(3) - INTERVAL 200 DAY)", + ) + .execute() + .await + .unwrap(); + ch.query( + "INSERT INTO mcp_logs (id, tool_arguments, tool_result, created_at) VALUES \ + ('60d', 'arguments', 'result', now64(3) - INTERVAL 60 DAY)", + ) + .execute() + .await + .unwrap(); + + // The restart, step by step as the server takes it: the tables, then + // the configured retentions. No step may shorten one. + ensure_clickhouse_schema(&app.db, &app.state.clickhouse) + .await + .unwrap(); + assert_eq!(ttls(&ch).await, configured(), "after the table setup"); + reconcile_clickhouse_ttls(&app.state).await; + assert_eq!(ttls(&ch).await, configured(), "after the retentions"); + + mutations_done(&ch).await; + let gateway: Vec = ch + .query( + "SELECT concat(id, ' ', ifNull(request_body, '-'), ' ', ifNull(response_body, '-')) \ + FROM gateway_logs ORDER BY id", + ) + .fetch_all() + .await + .unwrap(); + assert_eq!( + gateway, + ["200d - -", "60d request response"], + "the 60-day-old bodies and the 200-day-old row are still there" + ); + let mcp: Vec = ch + .query( + "SELECT concat(id, ' ', ifNull(tool_arguments, '-'), ' ', ifNull(tool_result, '-')) \ + FROM mcp_logs", + ) + .fetch_all() + .await + .unwrap(); + assert_eq!(mcp, ["60d arguments result"]); +} diff --git a/deploy/clickhouse/initdb.d/01_init.sql b/deploy/clickhouse/initdb.d/01_init.sql index 20c153b0..a876095f 100644 --- a/deploy/clickhouse/initdb.d/01_init.sql +++ b/deploy/clickhouse/initdb.d/01_init.sql @@ -185,24 +185,24 @@ ALTER TABLE gateway_logs ADD PROJECTION IF NOT EXISTS proj_by_latency ( -- compression dominates the cold-storage footprint. Sit AFTER -- session_id so the existing column order is preserved and new -- deployments + upgraded ones converge on the same shape. -ALTER TABLE gateway_logs ADD COLUMN IF NOT EXISTS request_body Nullable(String) CODEC(ZSTD(6)) AFTER session_id; -ALTER TABLE gateway_logs ADD COLUMN IF NOT EXISTS response_body Nullable(String) CODEC(ZSTD(6)) AFTER request_body; +-- +-- Body columns get a SHORTER TTL than the row-level retention: when it +-- fires, the value is reset to NULL while the row stays around for the +-- full table TTL, so metadata queries remain whole after bodies have +-- aged out. The 30 days here is only the TTL a column is created with. +-- The server sets the operator's `audit.body_retention_days` once it is +-- up (`apply_body_column_ttls`, crates/server/src/handlers/admin/ +-- retention.rs), and this file runs on every start, so it must never +-- set the TTL of a column that exists: a 30-day TTL applied, even for +-- the moment until the server restores a longer one, makes ClickHouse +-- clear every body older than 30 days. +ALTER TABLE gateway_logs ADD COLUMN IF NOT EXISTS request_body Nullable(String) CODEC(ZSTD(6)) TTL toDateTime(created_at) + INTERVAL 30 DAY AFTER session_id; +ALTER TABLE gateway_logs ADD COLUMN IF NOT EXISTS response_body Nullable(String) CODEC(ZSTD(6)) TTL toDateTime(created_at) + INTERVAL 30 DAY AFTER request_body; ALTER TABLE gateway_logs ADD COLUMN IF NOT EXISTS request_body_bytes Nullable(UInt32) AFTER response_body; ALTER TABLE gateway_logs ADD COLUMN IF NOT EXISTS response_body_bytes Nullable(UInt32) AFTER request_body_bytes; -- 'captured' | 'truncated' | 'disabled' | 'from_cache' | 'error' ALTER TABLE gateway_logs ADD COLUMN IF NOT EXISTS body_capture_status LowCardinality(Nullable(String)) AFTER response_body_bytes; --- Body columns get a SHORTER TTL than the row-level retention. The --- Rust side (`apply_body_column_ttls` in handlers/admin.rs) re-issues --- these at startup against the operator-configurable --- `audit.body_retention_days` setting (default 30); this seed-default --- exists so a CH bootstrap that happens before the server ever runs --- still has the right shape. When the column TTL fires, the value is --- reset to NULL while the row stays around for the full table TTL — --- so metadata queries remain whole even after bodies have aged out. -ALTER TABLE gateway_logs MODIFY COLUMN request_body TTL toDateTime(created_at) + INTERVAL 30 DAY; -ALTER TABLE gateway_logs MODIFY COLUMN response_body TTL toDateTime(created_at) + INTERVAL 30 DAY; - -- Substring search across captured bodies — auditors searching -- "which conversations mentioned API key XYZ" or "which tool calls -- referenced /etc/passwd" otherwise have to scan every row. tokenbf @@ -254,16 +254,14 @@ ALTER TABLE mcp_logs ADD PROJECTION IF NOT EXISTS proj_by_duration ( -- by sanitize_detail); promote both arguments and the upstream result -- to first-class columns so audit queries don't have to JSON-parse on -- every row. Same ZSTD(6) trade-off as gateway_logs. -ALTER TABLE mcp_logs ADD COLUMN IF NOT EXISTS tool_arguments Nullable(String) CODEC(ZSTD(6)) AFTER detail; -ALTER TABLE mcp_logs ADD COLUMN IF NOT EXISTS tool_result Nullable(String) CODEC(ZSTD(6)) AFTER tool_arguments; +-- Their TTL as gateway_logs' body columns above: set when the column +-- is created, never again by this file. +ALTER TABLE mcp_logs ADD COLUMN IF NOT EXISTS tool_arguments Nullable(String) CODEC(ZSTD(6)) TTL toDateTime(created_at) + INTERVAL 30 DAY AFTER detail; +ALTER TABLE mcp_logs ADD COLUMN IF NOT EXISTS tool_result Nullable(String) CODEC(ZSTD(6)) TTL toDateTime(created_at) + INTERVAL 30 DAY AFTER tool_arguments; ALTER TABLE mcp_logs ADD COLUMN IF NOT EXISTS arguments_bytes Nullable(UInt32) AFTER tool_result; ALTER TABLE mcp_logs ADD COLUMN IF NOT EXISTS result_bytes Nullable(UInt32) AFTER arguments_bytes; ALTER TABLE mcp_logs ADD COLUMN IF NOT EXISTS body_capture_status LowCardinality(Nullable(String)) AFTER result_bytes; --- Same body-column TTL story as gateway_logs above; see comment there. -ALTER TABLE mcp_logs MODIFY COLUMN tool_arguments TTL toDateTime(created_at) + INTERVAL 30 DAY; -ALTER TABLE mcp_logs MODIFY COLUMN tool_result TTL toDateTime(created_at) + INTERVAL 30 DAY; - -- Same substring-search rationale as gateway_logs above. ALTER TABLE mcp_logs ADD INDEX IF NOT EXISTS idx_tool_arguments ifNull(tool_arguments, '') TYPE tokenbf_v1(512, 3, 0) GRANULARITY 4; ALTER TABLE mcp_logs ADD INDEX IF NOT EXISTS idx_tool_result ifNull(tool_result, '') TYPE tokenbf_v1(512, 3, 0) GRANULARITY 4; From 5f4c06ba09dc79c605068a86922cc472ea6a573b Mon Sep 17 00:00:00 2001 From: fylorn <249551762+fylorn@users.noreply.github.com> Date: Mon, 5 Oct 2026 18:56:01 +0800 Subject: [PATCH 4/5] fix(helm): stop allowing ClickHouse's native port in the network policy The server's egress rules allowed port 9000, labelled "ClickHouse native", to any destination. The server never speaks ClickHouse's native protocol: the clickhouse crate it uses talks HTTP, on the port in CLICKHOUSE_URL, which the policy now allows by itself. Nothing else the chart configures uses 9000 either - body offload to S3 needs S3_BUCKET, which the chart doesn't set - so the rule only opened a port. An S3 endpoint on 9000 (RustFS, MinIO) set up outside the chart goes in networkPolicy.extraEgress; values.yaml and the README say so. helm lint and helm template, Helm 3.22.0 and 4.3.0, on the eleven value sets used for the previous change: the server's egress is DNS, 443 and the database ports, and nothing outside networkpolicy.yaml renders differently from dev. Co-Authored-By: Claude Opus 5.5 --- CHANGELOG.md | 5 ++++- deploy/helm/think-watch/README.md | 9 ++++++--- deploy/helm/think-watch/templates/networkpolicy.yaml | 4 ---- deploy/helm/think-watch/values.yaml | 5 +++-- 4 files changed, 13 insertions(+), 10 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 032b3c2e..c6f1198f 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -111,7 +111,10 @@ target. for its scheme where it names none. `networkPolicy.extraEgress` adds egress rules as written, for ports no URL names (Redis Cluster nodes announcing other ports, an upstream or MCP server on a port other than `443`). The - chart's README describes both. + chart's README describes both. Port `9000`, ClickHouse's native protocol, + is no longer allowed: the server reaches ClickHouse over HTTP only. An S3 + endpoint on `9000` (RustFS, MinIO) configured outside the chart needs a + rule in `networkPolicy.extraEgress`. ## [3.1.0] — 2026-10-05 diff --git a/deploy/helm/think-watch/README.md b/deploy/helm/think-watch/README.md index a5720a56..8d6b62d4 100644 --- a/deploy/helm/think-watch/README.md +++ b/deploy/helm/think-watch/README.md @@ -151,8 +151,7 @@ and authenticate with the password. ## Network policy `networkPolicy.enabled` limits what the server pods may reach: DNS, port -`443` (upstreams, the OIDC provider), `9000` (ClickHouse's native -protocol), and PostgreSQL, Redis and +`443` (upstreams, the OIDC provider), and PostgreSQL, Redis and ClickHouse on the ports the server connects to them on: - A bundled database: its service port (`5432`, `6379`, `8123`). @@ -165,10 +164,14 @@ ClickHouse on the ports the server connects to them on: does not change it), `26379` for a Sentinel and `6379` for the primary it points to, `80` for `http://` and `443` for `https://`. +The server talks to ClickHouse over HTTP only, so ClickHouse's native port +(`9000`) is not allowed. + What no URL names goes in `networkPolicy.extraEgress`, rules added to the server's egress as written: Redis Cluster nodes that announce ports the URL doesn't list, a Sentinel's primary on a port other than `6379`, an -upstream or MCP server on a port other than `443`. +upstream, MCP server or S3 endpoint on a port other than `443` (RustFS and +MinIO listen on `9000`). ```yaml networkPolicy: diff --git a/deploy/helm/think-watch/templates/networkpolicy.yaml b/deploy/helm/think-watch/templates/networkpolicy.yaml index 15d269ef..06955d37 100644 --- a/deploy/helm/think-watch/templates/networkpolicy.yaml +++ b/deploy/helm/think-watch/templates/networkpolicy.yaml @@ -48,10 +48,6 @@ spec: protocol: TCP {{- end }} {{- end }} - # ClickHouse native - - ports: - - port: 9000 - protocol: TCP # OIDC provider (HTTPS) - ports: - port: 443 diff --git a/deploy/helm/think-watch/values.yaml b/deploy/helm/think-watch/values.yaml index b7b9bc5e..c82367fc 100644 --- a/deploy/helm/think-watch/values.yaml +++ b/deploy/helm/think-watch/values.yaml @@ -147,8 +147,9 @@ networkPolicy: # # Egress rules added as written, for ports that no URL names: Redis # Cluster nodes announcing ports the URL doesn't list, a Sentinel's - # primary on a port other than 6379, an upstream or MCP server on a port - # other than 443. + # primary on a port other than 6379, an upstream, MCP server or S3 + # endpoint on a port other than 443 (RustFS and MinIO: 9000). ClickHouse + # is reached over HTTP only, so its native port 9000 is not allowed. extraEgress: [] # - ports: # - port: 7000 From cd274e91c7c749c43213582c6e204d84f1b7a4cb Mon Sep 17 00:00:00 2001 From: fylorn <249551762+fylorn@users.noreply.github.com> Date: Mon, 5 Oct 2026 18:56:26 +0800 Subject: [PATCH 5/5] docs(postgres): no transaction-mode pooler; upgrade with an ordinary rollout The start-up schema setup now holds a Postgres session-level advisory lock (SCHEMA_LOCK), and a session lock needs the session. Behind a pooler in transaction mode - PgBouncer's pool_mode = transaction, or a managed pooler's transaction port - each statement outside a transaction may go to a different server connection, and PgBouncer runs no reset query when a client leaves in that mode, so the lock can stay held on a server connection the pooler keeps: every instance started afterwards would wait on it with no end. The server has no separate URL for its migrations; it runs them on the connections of DATABASE_URL. So the chart's README (under the external databases, where Postgres is set up), the postgres .externalUrl comment in values.yaml and .env.example say to reach Postgres directly or through a pooler in session mode. The CHANGELOG's "Read before upgrading" gains that, and a sentence on the rollout: instances of earlier versions don't take the lock, so an old instance restarting while the first new one sets up the schema can still collide with it, which a normal rolling update doesn't do. Co-Authored-By: Claude Opus 5.5 --- .env.example | 3 +++ CHANGELOG.md | 8 ++++++++ deploy/helm/think-watch/README.md | 13 +++++++++++++ deploy/helm/think-watch/values.yaml | 4 +++- 4 files changed, 27 insertions(+), 1 deletion(-) diff --git a/.env.example b/.env.example index 988a7efb..389320a2 100644 --- a/.env.example +++ b/.env.example @@ -26,6 +26,9 @@ DB_MAX_CONNECTIONS=10 DB_PASSWORD=__SECRET_HEX_16__ # dev: DATABASE_URL=postgres://${DB_USER}:${DB_PASSWORD}@localhost:5432/${DB_NAME} # prod: DATABASE_URL=postgres://${DB_USER}:${DB_PASSWORD}@postgres:5432/${DB_NAME}?sslmode=disable +# Postgres directly, or a pooler in session mode: never one in transaction +# mode (PgBouncer pool_mode = transaction). Start-up schema setup holds a +# session-level advisory lock (deploy/helm/think-watch/README.md). # --- Redis --- REDIS_PASSWORD=__SECRET_HEX_16__ diff --git a/CHANGELOG.md b/CHANGELOG.md index c6f1198f..801c5488 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -29,6 +29,14 @@ target. narrow what its owner may do through that key, never widen it. A key given a higher limit than its owner to give it more room needs the owner's limit raised instead. +- **Upgrade with an ordinary rollout.** Earlier versions don't take the lock + that now makes instances set up the schema one at a time, so don't restart + instances of the old version while the first one of this version starts. +- **No transaction-mode pooler in front of Postgres.** Schema setup now holds a + Postgres session-level advisory lock, which a pooler in transaction mode + (PgBouncer `pool_mode = transaction`) can leave held, and every later start + then waits for it: point `DATABASE_URL` at Postgres itself or at a pooler in + session mode (the Helm chart's README has the details). ### Fixed diff --git a/deploy/helm/think-watch/README.md b/deploy/helm/think-watch/README.md index 8d6b62d4..1cb9f0fc 100644 --- a/deploy/helm/think-watch/README.md +++ b/deploy/helm/think-watch/README.md @@ -72,6 +72,19 @@ helm upgrade --install thinkwatch deploy/helm/think-watch \ When `bundled=false` and `externalUrl` is empty the chart fails at install-time with an explicit message — no silent broken Secret. +### PostgreSQL behind a connection pooler + +Each server instance sets up the schema when it starts, holding a +Postgres session-level advisory lock so that instances starting together +take turns. The lock belongs to one Postgres session, so `externalUrl` +must reach Postgres directly or through a pooler in session mode, never +one in transaction mode (PgBouncer `pool_mode = transaction`, or the +transaction-mode port of a managed pooler such as Supabase's): there, +the lock can stay held on a server connection the pooler keeps after the +instance is done with it, and every later start waits for it for good. +The server runs its schema setup on the connections of `DATABASE_URL`; +there is no separate URL for it. + ### Redis Cluster The external Redis can be a Redis Cluster. Give its URL the diff --git a/deploy/helm/think-watch/values.yaml b/deploy/helm/think-watch/values.yaml index c82367fc..73ba061d 100644 --- a/deploy/helm/think-watch/values.yaml +++ b/deploy/helm/think-watch/values.yaml @@ -76,7 +76,9 @@ postgres: image: postgres:18-alpine user: thinkwatch database: think_watch - # Used only when bundled=false; e.g. postgres://user:pass@host:5432/db + # Used only when bundled=false; e.g. postgres://user:pass@host:5432/db. + # Postgres itself or a pooler in session mode, never transaction mode + # (see README.md, "PostgreSQL behind a connection pooler"). externalUrl: "" storage: 10Gi storageClassName: ""