diff --git a/Cargo.lock b/Cargo.lock index f39da570..a325d2c4 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -3033,7 +3033,7 @@ dependencies = [ [[package]] name = "walshadow" -version = "0.1.3" +version = "0.1.4" dependencies = [ "ahash", "anyhow", @@ -3082,7 +3082,7 @@ dependencies = [ [[package]] name = "walshadow-bench" -version = "0.1.3" +version = "0.1.4" dependencies = [ "ab_glyph", "anyhow", diff --git a/Cargo.toml b/Cargo.toml index c45ea493..136db714 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -3,7 +3,7 @@ members = ["bench"] [package] name = "walshadow" -version = "0.1.3" +version = "0.1.4" edition = "2024" description = "Schema-only Postgres + WAL replay catalog mirror for CDC to ClickHouse" license = "AGPL-3.0-only" diff --git a/bench/Cargo.toml b/bench/Cargo.toml index 06579e7d..002f2271 100644 --- a/bench/Cargo.toml +++ b/bench/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "walshadow-bench" -version = "0.1.3" +version = "0.1.4" edition = "2024" description = "Replication-latency benchmarks (Postgres → ClickHouse / PG standby), engine + CLIs" license = "AGPL-3.0-only" diff --git a/src/bin/stream/shadow_proc.rs b/src/bin/stream/shadow_proc.rs index cb59a44c..5158f89d 100644 --- a/src/bin/stream/shadow_proc.rs +++ b/src/bin/stream/shadow_proc.rs @@ -91,6 +91,7 @@ pub(crate) async fn start_owned_shadow( .context("wait for running shadow startup")? { s.validate_running().context("validate running shadow")?; + warn_database_absent(&s); tracing::info!(target: "walshadow::shadow", "reusing running shadow"); return Ok(()); } @@ -117,6 +118,7 @@ pub(crate) async fn start_owned_shadow( "shadow caught up to bootstrap end_lsn", ); } + warn_database_absent(&s); Ok(()) }) .await @@ -320,6 +322,20 @@ impl Drop for ShadowLifecycle { } } +/// A shadow missing the applied `[source] dbname` cannot serve its bridge, but +/// a database the source created after the base backup arrives with replay — +/// so this is a warning at startup, and the bridge connect is what refuses +fn warn_database_absent(shadow: &walshadow::shadow::Shadow) { + if let Ok(false) = shadow.has_database() { + tracing::warn!( + target: "walshadow::shadow", + dbname = shadow.config().dbname.as_str(), + "shadow holds no such database yet; replay has to create it before the \ + bridge can attach. A data dir provisioned for another source never will", + ); + } +} + #[cfg(test)] mod tests { use super::*; diff --git a/src/catalog/shadow.rs b/src/catalog/shadow.rs index b9ded3cd..4d497c74 100644 --- a/src/catalog/shadow.rs +++ b/src/catalog/shadow.rs @@ -286,11 +286,17 @@ pub struct HealthReport { pub struct Shadow { config: ShadowConfig, + /// Database the cluster-level probes connect to, resolved once. See + /// [`Shadow::probe_one`] + probe_db: std::sync::OnceLock, } impl Shadow { pub fn new(config: ShadowConfig) -> Self { - Self { config } + Self { + config, + probe_db: std::sync::OnceLock::new(), + } } pub fn config(&self) -> &ShadowConfig { @@ -302,20 +308,20 @@ impl Shadow { if !self.is_in_recovery()? { return Err(ShadowError::NotInRecovery); } - let actual = self.psql_one("SHOW data_directory")?; + let actual = self.probe_one("SHOW data_directory")?; if fs::canonicalize(actual)? != fs::canonicalize(&self.config.data_dir)? { return Err(ShadowError::Incompatible("data_directory differs".into())); } let restore = format!("cp {}/%f %p", self.config.filter_out_dir.display()); - if self.psql_one("SHOW restore_command")? != restore { + if self.probe_one("SHOW restore_command")? != restore { return Err(ShadowError::Incompatible("restore_command differs".into())); } if let Some(bridge) = &self.config.bridge { - if self.psql_one("SHOW walshadow.socket_path")? != bridge.socket_path.to_string_lossy() + if self.probe_one("SHOW walshadow.socket_path")? != bridge.socket_path.to_string_lossy() { return Err(ShadowError::Incompatible("bridge socket differs".into())); } - let workers = self.psql_one("SHOW walshadow.bridge_workers")?; + let workers = self.probe_one("SHOW walshadow.bridge_workers")?; if workers .parse::() .ok() @@ -468,7 +474,7 @@ impl Shadow { /// gain a line per boot pub fn point_at_walsender(&self, conninfo: &str) -> Result<()> { if self - .psql_one("SHOW primary_conninfo") + .probe_one("SHOW primary_conninfo") .is_ok_and(|cur| cur == conninfo) { return Ok(()); @@ -667,6 +673,30 @@ impl Shadow { /// Trimmed stdout of one statement. `-tAXq` keeps each row /// unadorned. pub fn psql_one(&self, sql: &str) -> Result { + self.psql_in(self.config.dbname.as_str(), sql) + } + + /// `SHOW`, recovery state, and shared-catalog reads are cluster-wide, so + /// they must not need the source's database to be present: a shadow + /// provisioned for another `dbname` would otherwise fail every probe with + /// `database "..." does not exist` rather than the mismatch itself. + /// Prefers `postgres`, falling back to the configured database + pub fn probe_one(&self, sql: &str) -> Result { + if let Some(db) = self.probe_db.get() { + return self.psql_in(db, sql); + } + // `template1` is deliberately not a candidate: a session on it blocks + // `CREATE DATABASE` replay, and on a standby that stalls recovery + let db = match self.psql_in("postgres", "SELECT 1") { + Ok(_) => "postgres".to_owned(), + Err(_) => self.config.dbname.clone(), + }; + let db = self.probe_db.get_or_init(|| db); + self.psql_in(db, sql) + } + + /// Every cluster-level probe reads the shadow through here + fn psql_in(&self, dbname: &str, sql: &str) -> Result { let out = Command::new(self.config.bin("psql")) .args([ "-h", @@ -676,7 +706,7 @@ impl Shadow { "-U", self.config.user.as_str(), "-d", - self.config.dbname.as_str(), + dbname, "-tAXq", "-c", sql, @@ -691,7 +721,7 @@ impl Shadow { /// binary-upgrade dump can actually build here. pub fn available_extensions(&self) -> Result> { let raw = self - .psql_one("SELECT coalesce(string_agg(name, ','), '') FROM pg_available_extensions")?; + .probe_one("SELECT coalesce(string_agg(name, ','), '') FROM pg_available_extensions")?; Ok(raw .split(',') .map(str::trim) @@ -701,7 +731,7 @@ impl Shadow { } pub fn is_in_recovery(&self) -> Result { - match self.psql_one("SELECT pg_is_in_recovery()")?.as_str() { + match self.probe_one("SELECT pg_is_in_recovery()")?.as_str() { "t" => Ok(true), "f" => Ok(false), other => Err(ShadowError::PsqlParse(format!( @@ -713,7 +743,7 @@ impl Shadow { /// `pg_last_wal_replay_lsn()`; `None` on `NULL` (normal-mode /// cluster, or standby that has not replayed anything). pub fn last_replay_lsn(&self) -> Result> { - let s = self.psql_one("SELECT pg_last_wal_replay_lsn()")?; + let s = self.probe_one("SELECT pg_last_wal_replay_lsn()")?; if s.is_empty() { return Ok(None); } @@ -723,6 +753,16 @@ impl Shadow { } /// Block until in recovery and replay LSN ≥ `target`, or `timeout`. + pub fn has_database(&self) -> Result { + // `standard_conforming_strings` is on by default, so doubling the + // quote is the whole escape + let literal = self.config.dbname.replace('\'', "''"); + let present = self.probe_one(&format!( + "SELECT EXISTS (SELECT 1 FROM pg_database WHERE datname = '{literal}')" + ))?; + Ok(present == "t") + } + pub fn wait_for_replay(&self, target: u64, timeout: Duration) -> Result { let start = Instant::now(); loop { @@ -754,7 +794,7 @@ impl Shadow { /// value from `pg_control`. An operator `pg_wal_replay_pause` (or a /// recovery-target pause) leaves floor equal to running, so it holds pub fn try_pg_wal_replay_resume(&self) -> Result { - if self.psql_one("SELECT pg_get_wal_replay_pause_state()")? == "not paused" { + if self.probe_one("SELECT pg_get_wal_replay_pause_state()")? == "not paused" { return Ok(ResumeOutcome::NotPaused); } if !self @@ -763,7 +803,7 @@ impl Shadow { { return Ok(ResumeOutcome::PausedForeign); } - self.psql_one("SELECT pg_wal_replay_resume()")?; + self.probe_one("SELECT pg_wal_replay_resume()")?; Ok(ResumeOutcome::ResumedForFloor) } @@ -771,7 +811,7 @@ impl Shadow { /// [`control_guc_floor`](Self::control_guc_floor) to tell a floor-raise /// pause apart from other replay pauses fn running_guc_floor(&self) -> Result { - let row = self.psql_one( + let row = self.probe_one( "SELECT current_setting('max_connections'), \ current_setting('max_worker_processes'), \ current_setting('max_wal_senders'), \ @@ -799,12 +839,12 @@ impl Shadow { pub fn health(&self) -> Result { let in_recovery = self.is_in_recovery()?; let replay_lsn = self.last_replay_lsn()?; - let count_s = self.psql_one("SELECT count(*) FROM pg_class")?; + let count_s = self.probe_one("SELECT count(*) FROM pg_class")?; let pg_class_count = count_s .parse::() .map_err(|e| ShadowError::PsqlParse(format!("count(*) FROM pg_class: {e}")))?; let pg_proc_relname = - self.psql_one("SELECT relname FROM pg_class WHERE oid = 'pg_proc'::regclass")?; + self.probe_one("SELECT relname FROM pg_class WHERE oid = 'pg_proc'::regclass")?; Ok(HealthReport { in_recovery, replay_lsn, diff --git a/src/source_db.rs b/src/source_db.rs index 7b7e212e..40061a4a 100644 --- a/src/source_db.rs +++ b/src/source_db.rs @@ -61,7 +61,11 @@ impl DbLink { .await .with_context(|| { format!( - "connect bridge for database {} at {}", + "connect bridge for database {} at {}. The socket is registered by \ + the shadow's bridge worker for that database, so it is absent \ + while the shadow has no such database — replay creates one the \ + source made after the base backup, but a data dir provisioned for \ + another source never gains it and needs a reset", cfg.name, cfg.bridge_path.display() ) diff --git a/tests/shadow_dbname_guard.rs b/tests/shadow_dbname_guard.rs new file mode 100644 index 00000000..42869c17 --- /dev/null +++ b/tests/shadow_dbname_guard.rs @@ -0,0 +1,62 @@ +//! A shadow data dir provisioned for another source database must name itself. +//! +//! The resume path adopts any dir holding `PG_VERSION`, and the shadow follows +//! the applied `[source] dbname`, so a dir built for a different database (or +//! by an older `--shadow-dbname`) is reused verbatim. Before this guard the +//! first symptom was `psql: FATAL: database "..." does not exist` from the +//! replay wait, then a missing bridge socket, restarting forever. + +#![cfg(target_os = "linux")] + +#[path = "common/inproc_harness.rs"] +mod fx; + +use walshadow::shadow::Shadow; + +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] +async fn a_shadow_without_the_applied_database_is_named_not_probed_to_death() { + if !fx::tools::requirements_available() { + return; + } + let slot = fx::Ports::alloc(); + let tmp = tempfile::tempdir().unwrap(); + let ( + fx::BootstrappedClusters { + source, + shadow, + shadow_filter_dir: _, + }, + _stream, + ) = fx::bootstrap_clusters(&tmp, "", slot.source, slot.shadow, slot.walsender).await; + let _src_stop = fx::StopOnDrop { sh: &source }; + let _shd_stop = fx::StopOnDrop { sh: &shadow }; + + assert!( + shadow + .has_database() + .expect("probe the shadow's database list"), + "the shadow holds the database it was provisioned for", + ); + + let mut foreign = shadow.config().clone(); + foreign.dbname = "a_database_the_shadow_never_had".to_owned(); + let foreign = Shadow::new(foreign); + + // Cluster-level probes read the same postmaster, so they must not depend + // on the source's database being present + assert!( + foreign.is_in_recovery().is_ok(), + "recovery probe must not need the source database", + ); + assert!( + foreign.last_replay_lsn().is_ok(), + "replay probe must not need the source database", + ); + + assert!( + !foreign + .has_database() + .expect("probe a database the shadow lacks"), + "a database the shadow never had reads as absent", + ); +}