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
5 changes: 3 additions & 2 deletions plans/bootstrap.md
Original file line number Diff line number Diff line change
Expand Up @@ -61,8 +61,9 @@ for rendered diagram. Five clusters top→bottom:
4. **Shadow handoff** — `BootstrapOutcome { start, end }` returned;
daemon writes `standby.signal` and calls `materialize_conf` to
replace shadow's config files. Config includes walshadow settings,
minimum GUC values from `pg_control`, `restore_command`, and
`primary_conninfo`. Daemon empties `postgresql.auto.conf`, starts
minimum GUC values from `pg_control`, and `restore_command`, but not
`primary_conninfo`, which daemon adds after walsender starts
(see [shadow.md](shadow.md)). Daemon empties `postgresql.auto.conf`, starts
shadow with `start_with_floor_retry`, waits for `end_lsn` with
`wait_for_replay`, then supervises it (see [shadow.md](shadow.md))
5. **Manifest + WAL pump start** — the ack atomic seeds at `end_lsn`
Expand Down
18 changes: 11 additions & 7 deletions plans/shadow.md
Original file line number Diff line number Diff line change
Expand Up @@ -38,9 +38,9 @@ otherwise compose at daemon level
dir. Never turn standby recovery failure into automatic rebootstrap
2. Run `write_standby_signal`. Standby signal keeps shadow in recovery
while it receives continuous WAL stream
3. Run `control_guc_floor` and
`materialize_conf(floor, primary_conninfo)`. They read five minimum
GUC values from shadow's `pg_control` with `pg_controldata`.
3. Run `control_guc_floor` and `materialize_conf(floor, None)`. They
read five minimum GUC values from shadow's `pg_control` with
`pg_controldata`.
`LC_ALL=C` keeps output labels stable. PostgreSQL
checks these values against `pg_control`, so reading them locally
matches WAL being replayed and avoids querying source. Current
Expand All @@ -51,8 +51,9 @@ otherwise compose at daemon level
`autovacuum = off`, `fsync = on`, `hot_standby = on`,
`wal_level = replica`, `listen_addresses = ''`),
`restore_command = 'cp <filter_dir>/%f %p'`,
`recovery_target_timeline = 'latest'`, and `primary_conninfo =
'<walsender>'`. Empty `postgresql.auto.conf` to remove source
and `recovery_target_timeline = 'latest'`. Omit `primary_conninfo`
so shadow reads only archived WAL until step 5
Empty `postgresql.auto.conf` to remove source
`ALTER SYSTEM` settings included by BASE_BACKUP. Write
socket-only `pg_hba.conf` using trust authentication and empty
`pg_ident.conf`. Do not use config files from backup because Debian
Expand All @@ -67,7 +68,10 @@ otherwise compose at daemon level
`pg_ctl` only reports "could not start server" while log includes
required value. After a fresh bootstrap, run
`wait_for_replay(end_lsn, timeout)` against WAL included in backup
5. Call `is_running` every 2 s. If postmaster stops, restart it with
5. Once walsender listens, run `point_at_walsender(conninfo)` to append
`primary_conninfo` and reload PostgreSQL. This avoids waiting for
`wal_retrieve_retry_interval` after a failed connection
6. Call `is_running` every 2 s. If postmaster stops, restart it with
backoff and read minimum GUC values again. Hot standby can pause
replay when WAL requires higher value. Detect a pause with
`pg_get_wal_replay_pause_state()`, then confirm the cause: a floor
Expand All @@ -78,7 +82,7 @@ otherwise compose at daemon level
running settings (operator `pg_wal_replay_pause`, recovery target)
holds untouched. On daemon exit, run `pg_ctl stop -m fast` so data
dir is ready for next startup
6. Run `health` to check recovery state, replay LSN, `pg_class` count,
7. Run `health` to check recovery state, replay LSN, `pg_class` count,
and `pg_proc` lookup in one corruption probe

After bootstrap marker clears, every later start is standby recovery.
Expand Down
32 changes: 17 additions & 15 deletions src/bin/stream.rs
Original file line number Diff line number Diff line change
Expand Up @@ -960,25 +960,25 @@ async fn run_session(
} else {
None
};
// Regenerate config because walsender address and port may change
// Read minimum GUC values from shadow pg_control
// Regenerate config because shadow's port, socket, and GUC floor may change
// Keep shadow alive until pipeline teardown finishes
let shadow_lifecycle: Option<ShadowLifecycle> = match &shadow_start {
ShadowStart::External => None,
ShadowStart::Bootstrap(dir) | ShadowStart::Resume(dir) => {
let shadow = Arc::new(build_owned_shadow(args, dir.clone()));
let conninfo = walsender_primary_conninfo(args.walsender_bind);
shadow
.write_standby_signal()
.context("write standby.signal")?;
start_owned_shadow(
&shadow,
conninfo.clone(),
bootstrap_end_lsn,
Duration::from_secs(args.bootstrap_shadow_replay_timeout),
)
.await?;
Some(ShadowLifecycle::spawn(shadow, conninfo))
Some(ShadowLifecycle::spawn(
shadow,
walsender_primary_conninfo(args.walsender_bind),
))
}
};
let backup_settings = ch_config.as_ref().and_then(|c| c.backup.clone());
Expand Down Expand Up @@ -1151,10 +1151,7 @@ async fn run_session(
.collect();

let mut stream = WalStream::new(start_timeline, WAL_SEG_SIZE, aligned)?;
// Bind walsender listener BEFORE shadow's walreceiver can connect.
// Without an active sink, the catalog gate inside `BufferingDecoderSink`
// deadlocks: shadow's replay LSN never advances since segment-sink fires
// after per-record dispatch in the current ordering.
// Shadow must attach to this listener before catalog replay can advance
let mut shadow_boot = walshadow::shadow_stream::ShadowStreamState::new(
history.shadow_boot_branch(stored_timeline, aligned.get(), start_timeline),
ident.sysid.clone(),
Expand Down Expand Up @@ -1197,6 +1194,14 @@ async fn run_session(
stream.set_bytes_sink(Box::new(walshadow::shadow_stream::ShadowStreamSink::new(
shadow_state.clone(),
)));
// Set address after bind so first connection succeeds
// Supervisor restarts a shadow that is down, with the address in its conf
if let (Some(lifecycle), Some(conninfo)) = (
&shadow_lifecycle,
walsender_primary_conninfo(args.walsender_bind),
) {
probe_blocking(&lifecycle.shadow, move |s| s.point_at_walsender(&conninfo)).await;
}

// Seed catalog tracker from source's current pg_class before
// START_REPLICATION. Closes the "source rotated a mapped catalog above
Expand Down Expand Up @@ -4655,15 +4660,13 @@ fn walsender_primary_conninfo(bind: SocketAddr) -> Option<String> {
})
}

/// Start daemon-owned shadow with current walsender address and minimum
/// GUC values from its `pg_control`
/// Start daemon-owned shadow using archived WAL
/// After fresh bootstrap, wait for backup `end_lsn`; direct mode includes
/// required WAL in `base.tar`
/// Restart a postmaster left alive by an unclean prior exit so it binds
/// this daemon's port, socket, and walsender address
/// this daemon's port and socket
async fn start_owned_shadow(
shadow: &Arc<Shadow>,
conninfo: Option<String>,
replay_target: Option<u64>,
replay_timeout: Duration,
) -> Result<()> {
Expand All @@ -4681,8 +4684,7 @@ async fn start_owned_shadow(
s.stop().context("stop stale shadow before restart")?;
}
s.clear_stale_pid().context("clear stale postmaster.pid")?;
s.start_with_floor_retry(conninfo.as_deref())
.context("shadow start")?;
s.start_with_floor_retry(None).context("shadow start")?;
if let Some(target) = replay_target {
let lsn = s
.wait_for_replay(target, replay_timeout)
Expand Down
18 changes: 15 additions & 3 deletions src/catalog/shadow.rs
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@
//! Daemon path: [`control_guc_floor`](Shadow::control_guc_floor),
//! [`materialize_conf`](Shadow::materialize_conf),
//! [`write_standby_signal`](Shadow::write_standby_signal),
//! [`point_at_walsender`](Shadow::point_at_walsender),
//! [`clear_stale_pid`](Shadow::clear_stale_pid),
//! [`start_with_floor_retry`](Shadow::start_with_floor_retry),
//! [`wait_for_replay`](Shadow::wait_for_replay),
Expand Down Expand Up @@ -360,9 +361,7 @@ impl Shadow {
max_locks_per_transaction = floor.max_locks_per_transaction,
);
if let Some(conninfo) = primary_conninfo {
// Escape single quotes in PostgreSQL config string
let escaped = conninfo.replace('\'', "''");
conf.push_str(&format!("primary_conninfo = '{escaped}'\n"));
conf.push_str(&primary_conninfo_line(conninfo));
}
conf.push_str(&self.config.bridge_conf());
let d = &self.config.data_dir;
Expand All @@ -381,6 +380,15 @@ impl Shadow {
Ok(())
}

/// Set `primary_conninfo` and reload running shadow
pub fn point_at_walsender(&self, conninfo: &str) -> Result<()> {
let conf_path = self.config.data_dir.join("postgresql.conf");
let mut f = fs::OpenOptions::new().append(true).open(&conf_path)?;
f.write_all(primary_conninfo_line(conninfo).as_bytes())?;
self.run("pg_ctl", ["-D", self.config.data_str(), "reload"])?;
Ok(())
}

/// Create empty `standby.signal`
pub fn write_standby_signal(&self) -> Result<()> {
fs::write(self.config.data_dir.join("standby.signal"), b"")?;
Expand Down Expand Up @@ -711,6 +719,10 @@ impl Shadow {
}
}

fn primary_conninfo_line(conninfo: &str) -> String {
format!("primary_conninfo = '{}'\n", conninfo.replace('\'', "''"))
}

/// Parse required GUC values from `pg_controldata`
/// `LC_ALL=C` keeps labels stable; two labels use abbreviated `xact`
fn parse_controldata_floor(text: &str) -> Result<SourceGucFloor> {
Expand Down
7 changes: 2 additions & 5 deletions tests/common/bootstrap_ch_fixture.rs
Original file line number Diff line number Diff line change
Expand Up @@ -271,16 +271,13 @@ pub fn create_ch_dest_table(ch: &ChServer, database: &str, table: &str) -> Resul
Ok(())
}

/// Append `wal_level=logical` + `max_wal_senders` to the source PG's
/// postgresql.conf so the daemon's preflight (which insists on
/// `logical`) clears and BASE_BACKUP can attach. Mirrors the
/// `append_source_conf` helpers in the bootstrap drill tests.
/// Configure bootstrap source with spare sender slots for switchover
pub fn append_source_conf(sh: &Shadow) -> Result<()> {
let path = sh.config().data_dir.join("postgresql.conf");
let mut f = fs::OpenOptions::new().append(true).open(&path)?;
writeln!(f, "\n# walshadow bootstrap-CH source overrides")?;
writeln!(f, "wal_level = logical")?;
writeln!(f, "max_wal_senders = 4")?;
writeln!(f, "max_wal_senders = 8")?;
Ok(())
}

Expand Down