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
4 changes: 2 additions & 2 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

2 changes: 1 addition & 1 deletion Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down
2 changes: 1 addition & 1 deletion bench/Cargo.toml
Original file line number Diff line number Diff line change
@@ -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"
Expand Down
16 changes: 16 additions & 0 deletions src/bin/stream/shadow_proc.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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(());
}
Expand All @@ -117,6 +118,7 @@ pub(crate) async fn start_owned_shadow(
"shadow caught up to bootstrap end_lsn",
);
}
warn_database_absent(&s);
Ok(())
})
.await
Expand Down Expand Up @@ -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::*;
Expand Down
70 changes: 55 additions & 15 deletions src/catalog/shadow.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<String>,
}

impl Shadow {
pub fn new(config: ShadowConfig) -> Self {
Self { config }
Self {
config,
probe_db: std::sync::OnceLock::new(),
}
}

pub fn config(&self) -> &ShadowConfig {
Expand All @@ -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::<usize>()
.ok()
Expand Down Expand Up @@ -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(());
Expand Down Expand Up @@ -667,6 +673,30 @@ impl Shadow {
/// Trimmed stdout of one statement. `-tAXq` keeps each row
/// unadorned.
pub fn psql_one(&self, sql: &str) -> Result<String> {
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<String> {
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<String> {
let out = Command::new(self.config.bin("psql"))
.args([
"-h",
Expand All @@ -676,7 +706,7 @@ impl Shadow {
"-U",
self.config.user.as_str(),
"-d",
self.config.dbname.as_str(),
dbname,
"-tAXq",
"-c",
sql,
Expand All @@ -691,7 +721,7 @@ impl Shadow {
/// binary-upgrade dump can actually build here.
pub fn available_extensions(&self) -> Result<ahash::HashSet<String>> {
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)
Expand All @@ -701,7 +731,7 @@ impl Shadow {
}

pub fn is_in_recovery(&self) -> Result<bool> {
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!(
Expand All @@ -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<Option<u64>> {
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);
}
Expand All @@ -723,6 +753,16 @@ impl Shadow {
}

/// Block until in recovery and replay LSN ≥ `target`, or `timeout`.
pub fn has_database(&self) -> Result<bool> {
// `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<u64> {
let start = Instant::now();
loop {
Expand Down Expand Up @@ -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<ResumeOutcome> {
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
Expand All @@ -763,15 +803,15 @@ 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)
}

/// Read running GUC values with `current_setting`. Compare to
/// [`control_guc_floor`](Self::control_guc_floor) to tell a floor-raise
/// pause apart from other replay pauses
fn running_guc_floor(&self) -> Result<SourceGucFloor> {
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'), \
Expand Down Expand Up @@ -799,12 +839,12 @@ impl Shadow {
pub fn health(&self) -> Result<HealthReport> {
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::<u64>()
.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,
Expand Down
6 changes: 5 additions & 1 deletion src/source_db.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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()
)
Expand Down
62 changes: 62 additions & 0 deletions tests/shadow_dbname_guard.rs
Original file line number Diff line number Diff line change
@@ -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",
);
}
Loading