diff --git a/plans/bootstrap.md b/plans/bootstrap.md index 594d4b4a..24084864 100644 --- a/plans/bootstrap.md +++ b/plans/bootstrap.md @@ -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` diff --git a/plans/shadow.md b/plans/shadow.md index 91c4604f..fd08ecf9 100644 --- a/plans/shadow.md +++ b/plans/shadow.md @@ -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 @@ -51,8 +51,9 @@ otherwise compose at daemon level `autovacuum = off`, `fsync = on`, `hot_standby = on`, `wal_level = replica`, `listen_addresses = ''`), `restore_command = 'cp /%f %p'`, - `recovery_target_timeline = 'latest'`, and `primary_conninfo = - ''`. 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 @@ -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 @@ -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. diff --git a/src/bin/stream.rs b/src/bin/stream.rs index b4b750d2..7eb259b2 100644 --- a/src/bin/stream.rs +++ b/src/bin/stream.rs @@ -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 = 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()); @@ -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(), @@ -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 @@ -4655,15 +4660,13 @@ fn walsender_primary_conninfo(bind: SocketAddr) -> Option { }) } -/// 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, - conninfo: Option, replay_target: Option, replay_timeout: Duration, ) -> Result<()> { @@ -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) diff --git a/src/catalog/shadow.rs b/src/catalog/shadow.rs index 2d5825cf..4d812f80 100644 --- a/src/catalog/shadow.rs +++ b/src/catalog/shadow.rs @@ -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), @@ -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; @@ -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"")?; @@ -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 { diff --git a/tests/common/bootstrap_ch_fixture.rs b/tests/common/bootstrap_ch_fixture.rs index e120c788..2a9f7179 100644 --- a/tests/common/bootstrap_ch_fixture.rs +++ b/tests/common/bootstrap_ch_fixture.rs @@ -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(()) }