Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 3 additions & 0 deletions .env.example
Original file line number Diff line number Diff line change
Expand Up @@ -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__
Expand Down
48 changes: 48 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down Expand Up @@ -75,6 +83,46 @@ 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.
- **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,
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. 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

Expand Down
86 changes: 80 additions & 6 deletions crates/common/src/db.rs
Original file line number Diff line number Diff line change
@@ -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<PgPool> {
if !database_url.contains("sslmode=") {
Expand Down Expand Up @@ -28,6 +29,64 @@ pub async fn create_pool(database_url: &str) -> anyhow::Result<PgPool> {
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<Postgres>,
}

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<Self> {
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
Expand All @@ -45,20 +104,35 @@ pub async fn create_pool(database_url: &str) -> anyhow::Result<PgPool> {
/// 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(())
}
12 changes: 7 additions & 5 deletions crates/common/src/guard_policy/legacy.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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)
Expand Down
81 changes: 54 additions & 27 deletions crates/server/src/init.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;

Expand All @@ -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?);

Expand Down Expand Up @@ -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<clickhouse::Client>,
) -> 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()
}
Expand Down
3 changes: 3 additions & 0 deletions crates/server/src/services/mcp_store_repository.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
35 changes: 24 additions & 11 deletions crates/test-support/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -151,9 +151,7 @@ impl TestApp {
pub async fn try_spawn_with(opts: SpawnOptions) -> anyhow::Result<Self> {
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(|_| {
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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<String>, Option<String>) {
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;
Expand Down
Loading
Loading