diff --git a/architecture/recovery.md b/architecture/recovery.md index 83ce92e6..a3db0d51 100644 --- a/architecture/recovery.md +++ b/architecture/recovery.md @@ -34,6 +34,18 @@ position calculations belong in [manifest](../src/source/manifest.rs), [source feedback](../src/source/source_feed.rs), and [retention](../src/ops/retention.rs) +## Archive recovery + +Resolve each archived segment through verified source history. PostgreSQL copies +ancestor prefix into descendant's fork segment; ancestor archive may contain only +`.partial`. Stop replay at each switchpoint before adopting descendant timeline + +When source retention starts beyond a historical fork, verify repeated prefix +against archived descendant WAL before committing branch-aware resume state +Keep lineage, transaction, and replay barriers intact. Resume source replication +only when slot retention reaches requested position. Retry archive gaps with +bounded backoff so delayed uploads can unblock recovery + ## Planned source crossing Pause freezes consumed and received source frontiers while accepted destination diff --git a/src/backfill/backup_backfill.rs b/src/backfill/backup_backfill.rs index fbcc76e1..779fff9c 100644 --- a/src/backfill/backup_backfill.rs +++ b/src/backfill/backup_backfill.rs @@ -30,7 +30,7 @@ //! using commit LSNs, up to each relation's coverage bound //! //! Follow promotions along the branch the stream proved, cross-checked against -//! archived history by [`crate::source::archive_history`] before a gap replay +//! archived history by [`crate::source::archive::verify_history`] before a gap replay //! Reject backups whose redo or finish lies outside that branch //! //! The pre-scan aborts on gap writes that would invalidate the walk: a @@ -82,7 +82,7 @@ use crate::filter::pg_class_decoder::{ use crate::record::{Record, RecordSink, SinkError, WAL_SEG_SIZE, segments_covering_lineage}; use crate::runtime_config::InitialLoadMode; use crate::schema::RelDescriptor; -use crate::source::archive_history; +use crate::source::archive; use crate::source::timeline::TimelineHistory; use crate::ticker::Ticker; use crate::toast::ToastResolver; @@ -200,7 +200,7 @@ async fn run_object_store_pass(ctx: &PassContext, reqs: &[BackupRequest]) -> Res let (patch, gap_segments) = if b_redo < s_max { // Only a replay leg reads the archive's own WAL, so only it needs the // archive to agree on the chain serving those segments - archive_history::verify(settings, &storage, &history) + archive::verify_history(settings, &storage, &history) .await .context("backup_backfill: cross-check archived timeline history")?; let names = diff --git a/src/bin/stream/archive.rs b/src/bin/stream/archive.rs index 434d37d1..ebed5c44 100644 --- a/src/bin/stream/archive.rs +++ b/src/bin/stream/archive.rs @@ -1,137 +1,10 @@ -//! Archive WAL fetch: pull segments the source no longer has from the -//! backup archive, into shadow's `pg_wal`. +//! Hydrate shadow's `pg_wal` from backup archive at bootstrap use std::path::Path; -use std::sync::Arc; -use std::sync::atomic::{AtomicU64, Ordering}; -use std::time::Instant; use anyhow::{Context, Result}; -use futures::{StreamExt, stream as futures_stream}; use walrus::pg::backup::format_pg_lsn; -use walshadow::record::{WAL_SEG_SIZE, segments_covering}; -use walshadow::wal_stream::WalStream; - -/// Fetch archived WAL for source recovery, returning the bytes that begin at -/// exactly `start_lsn`. The archive stores whole 16 MiB segment files, so -/// fetch the single segment containing `start_lsn` (aligned range → one -/// entry) and slice off the already-consumed prefix — the returned bytes line -/// up with `WalStream::next_lsn`, which is byte- not segment-aligned in steady -/// state. -/// -/// Reading whole into memory keeps a prefetch slot off the staging disk, which -/// otherwise costs a 32 MiB round trip per 16 MiB of WAL and leaves a tmp file -/// behind on an aborted leg -pub(crate) async fn fetch_archive_segment( - settings: &walrus::config::Settings, - storage: &walrus::storage::DynStorage, - timeline: u32, - start_lsn: u64, -) -> Result<(String, Vec)> { - let seg_start = WalStream::align_down(start_lsn, WAL_SEG_SIZE); - let name = segments_covering(timeline, seg_start..seg_start + WAL_SEG_SIZE)[0].format(); - let mut bytes = walrus::pg::wal::fetch::read_segment(settings, storage, &name).await?; - if bytes.len() != WAL_SEG_SIZE as usize { - anyhow::bail!( - "archived WAL {name} has {} bytes, expected {WAL_SEG_SIZE}", - bytes.len(), - ); - } - bytes.drain(..(start_lsn - seg_start) as usize); - Ok((name, bytes)) -} - -/// Fetched segment, holding the budget slot it occupies until the pump takes it -type ArchiveSegment = (u64, Vec, tokio::sync::OwnedSemaphorePermit); - -pub(crate) struct ArchiveFeed { - pub(crate) wait_nanos: AtomicU64, - pub(crate) rx: tokio::sync::mpsc::Receiver>, - pub(crate) task: tokio::task::JoinHandle<()>, - pub(crate) fetch_nanos: Arc, -} - -impl ArchiveFeed { - pub(crate) fn spawn( - settings: walrus::config::Settings, - storage: walrus::storage::DynStorage, - timeline: u32, - start: u64, - concurrency: usize, - ) -> Self { - // `buffered` only advances its fetches while the stream is polled, so - // a worker parked on a full channel freezes every download in flight. - // Capacity below `concurrency` caps real depth at that capacity - let (tx, rx) = tokio::sync::mpsc::channel(concurrency); - // Ordered consumption lets a completed fetch sit in `buffered` waiting - // its turn, so slots alone bound nothing. A permit taken before the - // download and released at handoff holds resident segments to - // `concurrency`, plus the one the pump is replaying - let budget = Arc::new(tokio::sync::Semaphore::new(concurrency)); - let fetch_nanos = Arc::new(AtomicU64::new(0)); - let elapsed = fetch_nanos.clone(); - let task = tokio::spawn(async move { - let starts = std::iter::successors(Some(start), |lsn| { - (lsn / WAL_SEG_SIZE + 1).checked_mul(WAL_SEG_SIZE) - }); - let pending = futures_stream::iter(starts) - .map(|lsn| { - let (settings, storage, elapsed) = (&settings, &storage, &elapsed); - let budget = budget.clone(); - async move { - let permit = budget.acquire_owned().await.expect("budget stays open"); - let began = Instant::now(); - let result = fetch_archive_segment(settings, storage, timeline, lsn) - .await - .map(|(_, bytes)| (lsn, bytes, permit)); - elapsed.fetch_add(began.elapsed().as_nanos() as u64, Ordering::Relaxed); - result - } - }) - .buffered(concurrency); - tokio::pin!(pending); - while let Some(result) = pending.next().await { - let failed = result.is_err(); - if tx.send(result).await.is_err() || failed { - break; - } - } - }); - Self { - wait_nanos: AtomicU64::new(0), - rx, - task, - fetch_nanos, - } - } - - pub(crate) async fn next(&mut self) -> Option)>> { - let _elapsed = ArchiveWait { - nanos: &self.wait_nanos, - started: Instant::now(), - }; - let fetched = self.rx.recv().await?; - Some(fetched.map(|(lsn, bytes, _budget)| (lsn, bytes))) - } -} - -struct ArchiveWait<'a> { - nanos: &'a AtomicU64, - started: Instant, -} - -impl Drop for ArchiveWait<'_> { - fn drop(&mut self) { - self.nanos - .fetch_add(self.started.elapsed().as_nanos() as u64, Ordering::Relaxed); - } -} - -impl Drop for ArchiveFeed { - fn drop(&mut self) { - self.task.abort(); - } -} +use walshadow::record::segments_covering; /// Fetch WAL `[start_lsn, end_lsn]` from archive storage into shadow's `pg_wal/`. pub(crate) async fn fetch_wal_into_pg_wal( @@ -202,140 +75,77 @@ pub(crate) async fn fetch_wal_into_pg_wal( #[cfg(test)] mod tests { use super::*; - use std::fs; - - #[tokio::test] - async fn archive_prefetch_preserves_order_and_stops_at_gap() { - let tmp = tempfile::tempdir().unwrap(); - let settings = walrus::config::Settings { - storage: walrus::config::StorageSettings::Fs { - path: tmp.path().join("archive").display().to_string(), - }, - ..Default::default() - }; - let storage = settings.build_storage().unwrap(); - for index in [0u64, 1, 3] { - let name = - segments_covering(1, index * WAL_SEG_SIZE..(index + 1) * WAL_SEG_SIZE)[0].format(); - let path = tmp.path().join(name); - fs::write(&path, vec![index as u8; WAL_SEG_SIZE as usize]).unwrap(); - walrus::pg::wal::push::handle(&settings, storage.clone(), &path) - .await - .unwrap(); - } - let mut reader = ArchiveFeed::spawn(settings, storage, 1, 42, 4); - let (lsn, bytes) = reader.next().await.unwrap().unwrap(); - assert_eq!(lsn, 42); - assert_eq!(bytes.len(), WAL_SEG_SIZE as usize - 42); - assert!(bytes.iter().all(|b| *b == 0)); - let (lsn, bytes) = reader.next().await.unwrap().unwrap(); - assert_eq!(lsn, WAL_SEG_SIZE); - assert!(bytes.iter().all(|b| *b == 1)); - assert!(reader.next().await.unwrap().is_err()); - assert!( - reader.next().await.is_none(), - "must not skip missing segment" - ); - } - - #[tokio::test] - async fn archive_prefetch_drop_cancels_worker() { - let (tx, rx) = tokio::sync::mpsc::channel(1); - let task = tokio::spawn(async move { - let _tx = tx; - std::future::pending::<()>().await; - }); - let abort = task.abort_handle(); - drop(ArchiveFeed { - wait_nanos: AtomicU64::new(0), - rx, - task, - fetch_nanos: Arc::new(AtomicU64::new(0)), - }); - tokio::task::yield_now().await; - assert!(abort.is_finished()); - } + use walshadow::record::WAL_SEG_SIZE; - #[tokio::test] - async fn archive_fetch_reads_exact_segment() { - let tmp = tempfile::tempdir().unwrap(); - let archive = tmp.path().join("archive"); - let segment_path = tmp.path().join("000000010000000000000000"); - fs::write(&segment_path, vec![0; WAL_SEG_SIZE as usize]).unwrap(); - let settings = walrus::config::Settings { - storage: walrus::config::StorageSettings::Fs { - path: archive.display().to_string(), - }, - ..Default::default() - }; - let storage = settings.build_storage().unwrap(); - walrus::pg::wal::push::handle(&settings, storage.clone(), &segment_path) - .await - .unwrap(); - - let (name, bytes) = fetch_archive_segment(&settings, &storage, 1, 0) - .await - .unwrap(); - assert_eq!(name, "000000010000000000000000"); - assert_eq!(bytes.len(), WAL_SEG_SIZE as usize); - } - - #[tokio::test] - async fn archive_fetch_falls_back_across_compressions() { - let tmp = tempfile::tempdir().unwrap(); - let segment_path = tmp.path().join("000000010000000000000000"); - fs::write(&segment_path, vec![7; WAL_SEG_SIZE as usize]).unwrap(); - let storage = walrus::config::StorageSettings::Fs { - path: tmp.path().join("archive").display().to_string(), - }; - let pushed = walrus::config::Settings { - storage: storage.clone(), - compression: walrus::compression::Method::None, - ..Default::default() - }; - let built = pushed.build_storage().unwrap(); - walrus::pg::wal::push::handle(&pushed, built.clone(), &segment_path) + async fn archive(settings: &walrus::config::Settings, dir: &Path, name: &str, body: &[u8]) { + let path = dir.join(name); + tokio::fs::write(&path, body).await.unwrap(); + walrus::pg::wal::push::handle(settings, settings.build_storage().unwrap(), &path) .await .unwrap(); - // A bucket written under another compression must still read back - let reading = walrus::config::Settings { - storage, - ..Default::default() - }; - let (_, bytes) = fetch_archive_segment(&reading, &built, 1, 0).await.unwrap(); - assert_eq!(bytes.len(), WAL_SEG_SIZE as usize); - assert!(bytes.iter().all(|b| *b == 7)); } #[tokio::test] - async fn archive_fetch_slices_from_mid_segment() { - // A mid-segment resume LSN must return the segment's tail beginning at - // that LSN, not the whole segment (which would misalign the replay). + async fn hydrates_every_covering_segment_and_history_when_archived() { let tmp = tempfile::tempdir().unwrap(); - let archive = tmp.path().join("archive"); - let segment_path = tmp.path().join("000000010000000000000000"); - let pattern: Vec = (0..WAL_SEG_SIZE as usize) - .map(|i| (i % 251) as u8) - .collect(); - fs::write(&segment_path, &pattern).unwrap(); let settings = walrus::config::Settings { storage: walrus::config::StorageSettings::Fs { - path: archive.display().to_string(), + path: tmp.path().join("archive").display().to_string(), }, ..Default::default() }; + let segments = segments_covering(2, WAL_SEG_SIZE..3 * WAL_SEG_SIZE); + for seg in &segments { + archive( + &settings, + tmp.path(), + &seg.format(), + &[0; WAL_SEG_SIZE as usize], + ) + .await; + } + let history = walshadow::timeline::history_filename(2); + let shadow = tmp.path().join("shadow"); + + // No history archived yet: segments still land let storage = settings.build_storage().unwrap(); - walrus::pg::wal::push::handle(&settings, storage.clone(), &segment_path) - .await - .unwrap(); + fetch_wal_into_pg_wal( + &settings, + storage.clone(), + &shadow, + WAL_SEG_SIZE + 42, + 2 * WAL_SEG_SIZE, + 2, + ) + .await + .unwrap(); + for seg in &segments { + assert!(shadow.join("pg_wal").join(seg.format()).exists()); + } + assert!(!shadow.join("pg_wal").join(&history).exists()); - let offset = WAL_SEG_SIZE / 2; - let (name, bytes) = fetch_archive_segment(&settings, &storage, 1, offset) + archive( + &settings, + tmp.path(), + &history, + b"1\t0/1000000\tpromotion\n", + ) + .await; + fetch_wal_into_pg_wal( + &settings, + storage.clone(), + &shadow, + WAL_SEG_SIZE, + WAL_SEG_SIZE, + 2, + ) + .await + .unwrap(); + assert!(shadow.join("pg_wal").join(&history).exists()); + + let err = fetch_wal_into_pg_wal(&settings, storage, &shadow, 0, 0, 2) .await - .unwrap(); - // Same segment file, sliced to begin at the mid-segment LSN. - assert_eq!(name, "000000010000000000000000"); - assert_eq!(bytes.len(), (WAL_SEG_SIZE - offset) as usize); - assert_eq!(bytes, pattern[offset as usize..]); + .unwrap_err(); + assert!(format!("{err:#}").contains("fetch WAL"), "{err:#}"); } } diff --git a/src/bin/stream/session.rs b/src/bin/stream/session.rs index b4ea9430..911edd8e 100644 --- a/src/bin/stream/session.rs +++ b/src/bin/stream/session.rs @@ -11,6 +11,7 @@ use tokio::sync::{Mutex, watch}; use tokio_postgres::types::Oid; use tokio_util::sync::CancellationToken; use walrus::pg::backup::format_pg_lsn; +use walshadow::archive::Archive; use walshadow::boundary_hold::{BoundaryGateConfig, BoundaryHoldSink, CatalogBoundaryGate}; use walshadow::ch_emitter::{EmitterConfig, EmitterStats}; use walshadow::config::{ConfigResolver, SourceConn}; @@ -31,7 +32,8 @@ use walshadow::source_db::{DbLink, DbLinkConfig, SourceDb, SourceDbs}; use walshadow::source_feed::{SourceEvent, SourceFeed, StandbyStatus}; use walshadow::timeline::TimelineHistory; use walshadow::transition::{ - CrossingState, ForkGuards, Switchover, TimelineStats, load_boot_history, seed_shadow_branches, + CrossingState, ForkGuards, PrefixOrigin, Switchover, TimelineStats, load_boot_history, + seed_shadow_branches, }; use walshadow::wal_stream::WalStream; use walshadow::xact_buffer::{BufferingDecoderSink, SubxactTracker, XactBuffer, XactBufferConfig}; @@ -59,7 +61,7 @@ use crate::source_db::{ open_source_sql_client, }; use crate::source_recovery::{ - BARRIER_LOG_INTERVAL, FORK_FENCE_DRAIN, PROMOTION_POLL, PromotionGate, Redial, + BARRIER_LOG_INTERVAL, FORK_FENCE_DRAIN, PROMOTION_POLL, PromotionGate, ReconnectBackoff, SOURCE_SWAP_RETRY, SourcePath, SourceRecovery, commit_fork_resume, connect_source_waiting, promotion_gate, resume_manifest, resume_source_feed, stream_branch, swap_reason, }; @@ -267,7 +269,11 @@ pub(crate) async fn run_session( }; let shadow_lifecycle = ShadowLifecycle::spawn(owned, walsender_primary_conninfo(args.walsender_bind)); - let backup_settings = ch_config.as_ref().and_then(|c| c.backup.clone()); + let archive = ch_config + .as_ref() + .and_then(|c| c.backup.clone()) + .map(Archive::open) + .transpose()?; let start_lsn_override: Option> = args .start_lsn .as_deref() @@ -1159,11 +1165,13 @@ pub(crate) async fn run_session( ); } - let source_recovery = SourceRecovery { + let mut source_recovery = SourceRecovery { + system_id: live_identity.system_id, status_interval: Duration::from_secs(args.status_interval), - backup: backup_settings.as_ref(), + backup: archive.as_ref(), floor: &resume_floor, prefetch: usize::from(args.archive_prefetch), + backoff: ReconnectBackoff::default(), }; let mut path = SourcePath::Live; if let Err(e) = feed @@ -1175,15 +1183,8 @@ pub(crate) async fn run_session( .await { path = source_recovery - .recover( - e, - &cfg, - source_conn.slot.as_deref(), - stream_branch(&history, live_identity.system_id, &stream), - stream.next_lsn(), - &mut feed, - ) - .await; + .attempt(Some(e), &source_conn, &history, &stream, &mut feed) + .await?; } let mut segments_shipped = 0u64; @@ -1225,6 +1226,7 @@ pub(crate) async fn run_session( system_id: live_identity.system_id, out_dir: &args.out_dir, shadow_state: &shadow_state, + backup: archive.as_ref(), }; let mut timeline_stats = TimelineStats { // Off the chain, so a restart after a crossing keeps reporting the fork @@ -1238,7 +1240,8 @@ pub(crate) async fn run_session( let mut crossing = CrossingState::default(); let mut barrier_logged: Option = None; let shutdown_reason = 'pump: loop { - if matches!(path, SourcePath::Archive(_)) + // Nothing resumes at a switchpoint, not archive nor redial, only a crossing + if !matches!(path, SourcePath::Live) && history.branch_exhausted(stream.timeline(), stream.next_lsn().get()) { path = SourcePath::Live; @@ -1267,26 +1270,11 @@ pub(crate) async fn run_session( } // Lost source redials whatever `[source]` names now, so a repoint made // during an outage is what the next attempt dials - if let SourcePath::Redial(redial) = &mut path - && let Some(fresh) = source_recovery - .redial( - redial, - &cfg, - source_conn.slot.as_deref(), - stream_branch(&history, live_identity.system_id, &stream), - stream.next_lsn(), - ) - .await? - { - feed = fresh; - path = SourcePath::Live; + if matches!(path, SourcePath::Redial) && source_recovery.backoff.due() { + path = source_recovery + .attempt(None, &source_conn, &history, &stream, &mut feed) + .await?; swap.settled(); - tracing::info!( - target: "walshadow", - endpoint = source_conn.endpoint(), - resume_lsn = %stream.next_lsn(), - "source reconnected — resuming replication", - ); } // Swap between chunks, so the resume point is the byte-contiguous // `next_lsn` and no WalStream state is rebuilt. Old feed stays up @@ -1297,8 +1285,7 @@ pub(crate) async fn run_session( // Not while a crossing is pending: the stream sits at a switchpoint no // branch resumes from, and the crossing dials the live endpoint and slot // itself, so a repoint made mid-crossing lands there instead. - if swap.due(Instant::now()) && !crossing.pending() && !matches!(path, SourcePath::Redial(_)) - { + if swap.due(Instant::now()) && !crossing.pending() && !matches!(path, SourcePath::Redial) { match resume_source_feed( &cfg, source_conn.slot.as_deref(), @@ -1477,12 +1464,9 @@ pub(crate) async fn run_session( result = async { path.archive().expect("guarded by arm").next().await }, if matches!(path, SourcePath::Archive(_)) && !paused && !crossing.pending() => { match result { - Some(Ok((start_lsn, mut bytes))) => { + Some(Ok((start_lsn, bytes))) => { + source_recovery.backoff.reset(); anyhow::ensure!(start_lsn == stream.next_lsn().get(), "archive WAL discontinuity"); - if let Some(fork) = history.switchpoint_of(stream.timeline()) { - anyhow::ensure!(start_lsn < fork, "archive read past timeline fork"); - bytes.truncate((fork - start_lsn).min(bytes.len() as u64) as usize); - } archived_bytes = bytes; archived_segment = true; Some(walshadow::source_feed::WalChunk { @@ -1491,14 +1475,13 @@ pub(crate) async fn run_session( data: &archived_bytes, }) } - result => { - let reason = match result { - Some(Err(e)) => format!("{e:#}"), - None => "archive reader stopped".to_string(), - Some(Ok(_)) => unreachable!(), - }; + ended => { + let reason = ended.and_then(Result::err).map_or_else( + || "archive reader stopped".to_string(), + |e| format!("{e:#}"), + ); tracing::info!(target: "walshadow", reason, "archive ended, reconnecting source"); - path = SourcePath::Redial(Redial::now(reason)); + path = SourcePath::Redial; None } } @@ -1548,15 +1531,8 @@ pub(crate) async fn run_session( } }; path = source_recovery - .recover( - err, - &cfg, - source_conn.slot.as_deref(), - stream_branch(&history, live_identity.system_id, &stream), - stream.next_lsn(), - &mut feed, - ) - .await; + .attempt(Some(err), &source_conn, &history, &stream, &mut feed) + .await?; // Recovery dials the live endpoint, so a queued swap is done swap.settled(); None @@ -1799,6 +1775,10 @@ pub(crate) async fn run_session( slot = source_conn.slot.as_deref(), "crossed source timeline", ); + if crossed.prefix_origin == PrefixOrigin::Archive { + source_recovery.backoff.reset(); + path = SourcePath::Redial; + } history = crossed.history; history_tx.send_replace(Arc::new(history.clone())); crossing.committed(); diff --git a/src/bin/stream/source_recovery.rs b/src/bin/stream/source_recovery.rs index 8f51515c..65fad216 100644 --- a/src/bin/stream/source_recovery.rs +++ b/src/bin/stream/source_recovery.rs @@ -7,6 +7,7 @@ use std::time::{Duration, Instant}; use anyhow::{Context, Result}; use tokio_postgres::types::PgLsn; use walrus::pg::replication::conn::PgConfig; +use walshadow::archive::{Archive, ArchiveFeed}; use walshadow::config::SourceConn; use walshadow::manifest; use walshadow::pos::{Floor, Monotone, Pos}; @@ -16,7 +17,6 @@ use walshadow::timeline::TimelineHistory; use walshadow::transition::{TransitionError, source_history}; use walshadow::wal_stream::WalStream; -use crate::archive::ArchiveFeed; use crate::args::{Args, cli_base}; /// How long the fork proofs wait for the pump-side queue to drain. Past it the @@ -411,7 +411,9 @@ pub(crate) fn swap_reason(err: &anyhow::Error) -> &'static str { pub(crate) enum SourcePath { Live, Archive(ArchiveFeed), - Redial(Redial), + /// Source lost with no archive to read, redialed each due pump iteration so + /// the loop keeps publishing, pausing and applying `[source]` repoints + Redial, } impl SourcePath { @@ -423,66 +425,82 @@ impl SourcePath { } } -/// Source lost with no archive to read, redialed each due pump iteration so -/// the loop keeps publishing, pausing and applying `[source]` repoints -pub(crate) struct Redial { - pub(crate) retry_at: Instant, - pub(crate) backoff: Duration, - /// Why the archive could not stand in, named if the source cannot serve - pub(crate) archive_error: String, +/// Redial pacing, held across archive legs so an archive gap waits out the +/// delay its failed dial set instead of redialing hot +pub(crate) struct ReconnectBackoff { + retry_at: Instant, + delay: Duration, } -impl Redial { - pub(crate) const MIN_BACKOFF: Duration = Duration::from_millis(200); - pub(crate) const MAX_BACKOFF: Duration = Duration::from_secs(10); - - pub(crate) fn now(archive_error: String) -> Self { +impl Default for ReconnectBackoff { + fn default() -> Self { Self { retry_at: Instant::now(), - backoff: Self::MIN_BACKOFF, - archive_error, + delay: Self::MIN, } } } +impl ReconnectBackoff { + const MIN: Duration = Duration::from_millis(200); + const MAX: Duration = Duration::from_secs(10); + + pub(crate) fn due(&self) -> bool { + Instant::now() >= self.retry_at + } + + pub(crate) fn reset(&mut self) { + *self = Self::default(); + } + + fn failed(&mut self) { + self.retry_at = Instant::now() + self.delay; + self.delay = (self.delay * 2).min(Self::MAX); + } +} + pub(crate) struct SourceRecovery<'a> { + pub(crate) system_id: u64, pub(crate) status_interval: Duration, - pub(crate) backup: Option<&'a walrus::config::Settings>, + pub(crate) backup: Option<&'a Archive>, /// Published resume floor, which is what a slot on the far end has to still /// reach — the reconnect's own `resume_lsn` sits above it pub(crate) floor: &'a Monotone, pub(crate) prefetch: usize, + pub(crate) backoff: ReconnectBackoff, } impl SourceRecovery<'_> { - /// Try source, otherwise start bounded archive fetches for normal pump, - /// otherwise redial. `cfg`, `slot`, and `branch` are the live endpoint, - /// slot name, and proved branch, passed per call rather than held, so a - /// recovery that starts after a `[source]` reload or a crossing dials the - /// new address under the new name and asks for the descendant, with the - /// archive read under its segment names. - pub(crate) async fn recover( - &self, - source_error: anyhow::Error, - cfg: &PgConfig, - slot: Option<&str>, - branch: SourceBranch, - resume_lsn: Pos, + /// Dial source, otherwise start bounded archive fetches for normal pump, + /// otherwise redial. `source` and `history` are the live endpoint and + /// proved chain, passed per call rather than held, so a recovery that + /// starts after a `[source]` reload or a crossing dials the new address + /// under the new slot and asks for the descendant, with the archive read + /// under its segment names. + /// + /// `lost` is what ended a live feed, starting a fresh outage. A removed-WAL + /// (58P01) loss means source genuinely can't serve resume point, so skip + /// straight to archive + pub(crate) async fn attempt( + &mut self, + lost: Option, + source: &SourceConn, + history: &TimelineHistory, + stream: &WalStream, feed: &mut SourceFeed, - ) -> SourcePath { - // Source first (primary_conninfo analog): a plain drop is usually - // transient, so try the source again at the exact resume point before - // reaching for the archive. A removed-WAL (58P01) error means the - // source genuinely can't serve it — skip straight to the archive. - let source_missing = walshadow::source_feed::is_wal_segment_removed(&source_error); - let reason = if source_missing { - source_error - } else { - match resume_source_feed( - cfg, - slot, + ) -> Result { + if lost.is_some() { + self.backoff.reset(); + } + let resume_lsn = stream.next_lsn(); + // Source first (primary_conninfo analog), plain drop is usually transient + let error = match lost { + Some(e) if walshadow::source_feed::is_wal_segment_removed(&e) => e, + _ => match resume_source_feed( + &source.to_pg_config(), + source.slot.as_deref(), resume_lsn, - branch, + stream_branch(history, self.system_id, stream), self.floor.get(), self.status_interval, ) @@ -490,88 +508,57 @@ impl SourceRecovery<'_> { { Ok(fresh) => { *feed = fresh; + self.backoff.reset(); tracing::info!( target: "walshadow", + endpoint = source.endpoint(), resume_lsn = %resume_lsn, "source reconnected — resuming replication", ); - return SourcePath::Live; + return Ok(SourcePath::Live); } - Err(retry_error) => retry_error, - } + Err(e) => e, + }, }; tracing::warn!( target: "walshadow", - error = %reason, + error = %format!("{error:#}"), + endpoint = source.endpoint(), resume_lsn = %resume_lsn, - source_missing, + retry_in_ms = self.backoff.delay.as_millis() as u64, "source cannot serve the resume point — trying archive", ); - // Archive fallback (restore_command analog). Redial covers every "no - // archive": a transient error retries the source with backoff, a - // removed-WAL error surfaces the operator-action message. - let archive_error = match self.backup.map(|s| (s, s.build_storage())) { - None => "no [backup] archive configured".to_string(), - Some((_, Err(e))) => format!("build archive storage: {e:#}"), - Some((settings, Ok(storage))) => { + self.fall_back(error, history, stream.timeline(), resume_lsn) + } + + /// Archive fallback (restore_command analog). Without an archive, removed + /// WAL needs an operator, any other failure redials after backoff + fn fall_back( + &mut self, + error: anyhow::Error, + history: &TimelineHistory, + timeline: u32, + resume_lsn: Pos, + ) -> Result { + self.backoff.failed(); + match self.backup { + Some(archive) => { tracing::info!(target: "walshadow", resume_lsn = %resume_lsn, prefetch = self.prefetch, "starting archive recovery"); - return SourcePath::Archive(ArchiveFeed::spawn( - settings.clone(), - storage, - branch.timeline, + Ok(SourcePath::Archive(archive.feed( + history.clone(), + timeline, resume_lsn.get(), self.prefetch, - )); + ))) } - }; - SourcePath::Redial(Redial::now(archive_error)) - } - - /// One attempt once `redial` is due. Removed WAL needs an operator; any - /// other failure backs off for the next iteration. Every attempt goes - /// through [`resume_source_feed`]'s proofs - pub(crate) async fn redial( - &self, - redial: &mut Redial, - cfg: &PgConfig, - slot: Option<&str>, - branch: SourceBranch, - resume_lsn: Pos, - ) -> Result> { - if Instant::now() < redial.retry_at { - return Ok(None); - } - match resume_source_feed( - cfg, - slot, - resume_lsn, - branch, - self.floor.get(), - self.status_interval, - ) - .await - { - Ok(feed) => Ok(Some(feed)), - Err(e) if walshadow::source_feed::is_wal_segment_removed(&e) => { - Err(e.context(format!( - "source cannot serve WAL at {resume_lsn}; {}; \ + None if walshadow::source_feed::is_wal_segment_removed(&error) => { + Err(error.context(format!( + "source cannot serve WAL at {resume_lsn}; no [backup] archive configured; \ base-backup refresh requires operator action", - redial.archive_error, ))) } - Err(e) => { - tracing::warn!( - target: "walshadow", - error = %e, - endpoint = %format!("{}:{}", cfg.host, cfg.port), - retry_in_ms = redial.backoff.as_millis() as u64, - "source reconnect failed — retrying", - ); - redial.retry_at = Instant::now() + redial.backoff; - redial.backoff = (redial.backoff * 2).min(Redial::MAX_BACKOFF); - Ok(None) - } + None => Ok(SourcePath::Redial), } } } @@ -580,6 +567,105 @@ impl SourceRecovery<'_> { mod tests { use super::*; + fn removed_wal() -> anyhow::Error { + walshadow::source_feed::WalSegmentRemoved { + error_code: walshadow::source_feed::SQLSTATE_UNDEFINED_FILE.to_string(), + message: "requested WAL segment has already been removed".to_string(), + } + .into() + } + + fn slot_too_new() -> anyhow::Error { + TransitionError::Slot(walshadow::source_feed::SlotError::TooNew { + slot: "walshadow".to_string(), + restart_lsn: 0x1_0E00_0490, + resume_lsn: 0x6200_0000, + }) + .into() + } + + fn recovery<'a>(backup: Option<&'a Archive>, floor: &'a Monotone) -> SourceRecovery<'a> { + SourceRecovery { + system_id: 1, + status_interval: Duration::from_secs(10), + backup, + floor, + prefetch: 1, + backoff: ReconnectBackoff::default(), + } + } + + #[tokio::test] + async fn archive_retries_after_slot_refusal_and_delayed_upload() { + let tmp = tempfile::tempdir().unwrap(); + let settings = walrus::config::Settings { + storage: walrus::config::StorageSettings::Fs { + path: tmp.path().join("archive").display().to_string(), + }, + ..Default::default() + }; + let floor = Monotone::new(Pos::new(0x6200_0000)); + let archive = Archive::open(settings.clone()).unwrap(); + let mut recovery = recovery(Some(&archive), &floor); + let history = TimelineHistory::root(1); + let resume = Pos::new(0x6200_002A); + for attempt in 0..8 { + let delay = recovery.backoff.delay; + let before = Instant::now(); + let error = if attempt % 2 == 0 { + slot_too_new() + } else { + removed_wal() + }; + let mut path = recovery.fall_back(error, &history, 1, resume).unwrap(); + let error = path.archive().unwrap().next().await.unwrap().unwrap_err(); + assert!(format!("{error:#}").contains("000000010000000000000062")); + assert!(recovery.backoff.retry_at >= before + delay); + assert_eq!( + recovery.backoff.delay, + (delay * 2).min(ReconnectBackoff::MAX) + ); + } + assert_eq!(recovery.backoff.delay, ReconnectBackoff::MAX); + + let segment = tmp.path().join("000000010000000000000062"); + std::fs::write(&segment, vec![0x5a; WAL_SEG_SIZE as usize]).unwrap(); + walrus::pg::wal::push::handle(&settings, settings.build_storage().unwrap(), &segment) + .await + .unwrap(); + let mut path = recovery + .fall_back(slot_too_new(), &history, 1, resume) + .unwrap(); + let (lsn, bytes) = path.archive().unwrap().next().await.unwrap().unwrap(); + assert_eq!(lsn, resume.get()); + assert_eq!(bytes.len(), WAL_SEG_SIZE as usize - 42); + assert!(bytes.iter().all(|b| *b == 0x5a)); + let error = path.archive().unwrap().next().await.unwrap().unwrap_err(); + assert!(format!("{error:#}").contains("000000010000000000000063")); + } + + #[test] + fn removed_wal_without_archive_requires_operator() { + let floor = Monotone::new(Pos::new(0x6200_0000)); + let mut recovery = recovery(None, &floor); + let history = TimelineHistory::root(1); + let Err(error) = recovery.fall_back(removed_wal(), &history, 1, floor.get()) else { + panic!("removed WAL without archive must fail"); + }; + assert!( + error + .to_string() + .contains("base-backup refresh requires operator action") + ); + assert!(error.to_string().contains("no [backup] archive configured")); + assert!(walshadow::source_feed::is_wal_segment_removed(&error)); + let path = recovery + .fall_back(slot_too_new(), &history, 1, floor.get()) + .unwrap(); + assert!(matches!(path, SourcePath::Redial)); + assert!(!recovery.backoff.due()); + } + /// Two standbys of one primary, promoted independently, are both timeline 2 /// under one system identifier. The chain places either one, so only where /// the branch begins refuses the wrong one diff --git a/src/lib.rs b/src/lib.rs index 3e1366ab..7287af10 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -75,7 +75,7 @@ pub use ops::{ }; #[doc(hidden)] pub use source::{ - archive_history, boundary_hold, catalog_capture, manifest, queueing_record_sink, segment_sink, + archive, boundary_hold, catalog_capture, manifest, queueing_record_sink, segment_sink, shadow_stream, source_feed, timeline, transition, wal_stream, }; #[doc(hidden)] diff --git a/src/source/archive.rs b/src/source/archive.rs new file mode 100644 index 00000000..93413156 --- /dev/null +++ b/src/source/archive.rs @@ -0,0 +1,420 @@ +//! Read archived WAL using source timeline history. + +use std::sync::Arc; +use std::sync::atomic::{AtomicU64, Ordering}; +use std::time::Instant; + +use anyhow::{Context, Result}; +use futures::{StreamExt, stream as futures_stream}; +use walrus::config::Settings; +use walrus::storage::DynStorage; + +use crate::record::{WAL_SEG_SIZE, segments_covering}; +use crate::source::timeline::{TimelineHistory, history_filename}; + +/// Initialize archive storage once so invalid `[backup]` settings fail at startup. +#[derive(Clone)] +pub struct Archive { + settings: Settings, + storage: DynStorage, +} + +impl Archive { + pub fn open(settings: Settings) -> Result { + let storage = settings.build_storage().context("build archive storage")?; + Ok(Self { settings, storage }) + } + + /// Read archived WAL from `lsn` to end of its segment. + /// + /// At a timeline fork, PostgreSQL copies earlier WAL into a segment named + /// for its new timeline (`XLogInitNewTimeline`). Read that file because + /// its parent timeline may have archived only a `.partial` file. + /// + /// Read into memory to avoid writing and reading 16 MiB of WAL on disk, + /// or leaving temporary files behind when recovery is interrupted. + pub async fn read_segment( + &self, + history: &TimelineHistory, + lsn: u64, + ) -> Result<(String, Vec)> { + let seg_start = lsn / WAL_SEG_SIZE * WAL_SEG_SIZE; + let timeline = history + .tli_of_segment(seg_start, WAL_SEG_SIZE) + .context("archive position outside source history")?; + let name = segments_covering(timeline, seg_start..seg_start + WAL_SEG_SIZE)[0].format(); + let mut bytes = + walrus::pg::wal::fetch::read_segment(&self.settings, &self.storage, &name).await?; + anyhow::ensure!( + bytes.len() == WAL_SEG_SIZE as usize, + "archived WAL {name} has {} bytes, expected {WAL_SEG_SIZE}", + bytes.len(), + ); + bytes.drain(..(lsn - seg_start) as usize); + Ok((name, bytes)) + } + + /// Prefetch WAL from `start` until `timeline` ends so replay can switch timelines. + pub fn feed( + &self, + history: TimelineHistory, + timeline: u32, + start: u64, + concurrency: usize, + ) -> ArchiveFeed { + let end = history.switchpoint_of(timeline).unwrap_or(u64::MAX); + let archive = self.clone(); + // `buffered` advances downloads only while polled. A full channel + // pauses all downloads, so allow room for `concurrency` segments. + let (tx, rx) = tokio::sync::mpsc::channel(concurrency); + // Completed downloads can wait in `buffered` as well as in this channel. + // Hold a permit until replay receives each segment to limit memory use + // to `concurrency` segments plus one being replayed. + let budget = Arc::new(tokio::sync::Semaphore::new(concurrency)); + let fetch_nanos = Arc::new(AtomicU64::new(0)); + let elapsed = fetch_nanos.clone(); + let task = tokio::spawn(async move { + let starts = std::iter::successors(Some(start), |lsn| { + (lsn / WAL_SEG_SIZE + 1).checked_mul(WAL_SEG_SIZE) + }) + .take_while(|lsn| *lsn < end); + let pending = futures_stream::iter(starts) + .map(|lsn| { + let (archive, history, elapsed) = (&archive, &history, &elapsed); + let budget = budget.clone(); + async move { + let permit = budget.acquire_owned().await.expect("budget stays open"); + let began = Instant::now(); + let result = + archive + .read_segment(history, lsn) + .await + .map(|(_, mut bytes)| { + bytes.truncate(bytes.len().min((end - lsn) as usize)); + (lsn, bytes, permit) + }); + elapsed.fetch_add(began.elapsed().as_nanos() as u64, Ordering::Relaxed); + result + } + }) + .buffered(concurrency); + tokio::pin!(pending); + while let Some(result) = pending.next().await { + let failed = result.is_err(); + if tx.send(result).await.is_err() || failed { + break; + } + } + }); + ArchiveFeed { + wait_nanos: AtomicU64::new(0), + rx, + task, + fetch_nanos, + } + } +} + +/// Check that archive and source have matching timeline histories, since replay +/// uses source history to choose archived segment names. +pub async fn verify_history( + settings: &Settings, + storage: &DynStorage, + source: &TimelineHistory, +) -> Result<()> { + let name = history_filename(source.target()); + let raw = walrus::pg::wal::fetch::read_segment(settings, storage, &name) + .await + .with_context(|| format!("fetch {name}"))?; + let archived = + TimelineHistory::parse(source.target(), &raw).with_context(|| format!("parse {name}"))?; + anyhow::ensure!( + archived.entries() == source.entries(), + "archived {name} disagrees with source timeline history", + ); + Ok(()) +} + +/// Keep a downloaded segment's memory permit until replay receives it. +type ArchiveSegment = (u64, Vec, tokio::sync::OwnedSemaphorePermit); + +pub struct ArchiveFeed { + pub wait_nanos: AtomicU64, + rx: tokio::sync::mpsc::Receiver>, + task: tokio::task::JoinHandle<()>, + pub fetch_nanos: Arc, +} + +impl ArchiveFeed { + pub async fn next(&mut self) -> Option)>> { + let _elapsed = ArchiveWait { + nanos: &self.wait_nanos, + started: Instant::now(), + }; + let fetched = self.rx.recv().await?; + Some(fetched.map(|(lsn, bytes, _budget)| (lsn, bytes))) + } +} + +struct ArchiveWait<'a> { + nanos: &'a AtomicU64, + started: Instant, +} + +impl Drop for ArchiveWait<'_> { + fn drop(&mut self) { + self.nanos + .fetch_add(self.started.elapsed().as_nanos() as u64, Ordering::Relaxed); + } +} + +impl Drop for ArchiveFeed { + fn drop(&mut self) { + self.task.abort(); + } +} + +#[cfg(test)] +mod tests { + use std::path::Path; + + use super::*; + use std::fs; + + /// Match `archive_command`: store history uncompressed under `wal_005/`. + async fn archive_history_file(settings: &Settings, storage: &DynStorage, tli: u32, body: &str) { + let dir = tempfile::tempdir().unwrap(); + let path = dir.path().join(history_filename(tli)); + tokio::fs::write(&path, body).await.unwrap(); + walrus::pg::wal::push::handle(settings, storage.clone(), &path) + .await + .unwrap(); + } + + fn fs_archive(root: &Path) -> (Settings, DynStorage) { + let settings = Settings { + storage: walrus::config::StorageSettings::Fs { + path: root.display().to_string(), + }, + ..Default::default() + }; + let storage = settings.build_storage().unwrap(); + (settings, storage) + } + + #[tokio::test] + async fn reject_an_archive_missing_the_sources_own_history() { + let tmp = tempfile::tempdir().unwrap(); + let (settings, storage) = fs_archive(&tmp.path().join("archive")); + archive_history_file(&settings, &storage, 2, "1\t0/3000000\tpromotion\n").await; + let source = TimelineHistory::parse(3, b"1\t0/3000000\tpromotion\n").unwrap(); + let err = verify_history(&settings, &storage, &source) + .await + .unwrap_err(); + assert!(err.to_string().contains("00000003.history"), "{err:#}"); + } + + #[tokio::test] + async fn archived_ids_need_not_be_consecutive() { + let tmp = tempfile::tempdir().unwrap(); + let (settings, storage) = fs_archive(&tmp.path().join("archive")); + archive_history_file(&settings, &storage, 4, "1\t0/3000000\tpromotion\n").await; + let source = TimelineHistory::parse(4, b"1\t0/3000000\tpromotion\n").unwrap(); + verify_history(&settings, &storage, &source).await.unwrap(); + } + + #[tokio::test] + async fn reject_archived_history_with_different_ancestry_or_switchpoint() { + for body in ["1\t0/4000000\tpromotion\n", "2\t0/3000000\tpromotion\n"] { + let tmp = tempfile::tempdir().unwrap(); + let (settings, storage) = fs_archive(&tmp.path().join("archive")); + archive_history_file(&settings, &storage, 3, body).await; + let source = TimelineHistory::parse(3, b"1\t0/3000000\tpromotion\n").unwrap(); + let err = verify_history(&settings, &storage, &source) + .await + .unwrap_err(); + assert!(err.to_string().contains("disagrees"), "{err:#}"); + } + } + + #[tokio::test] + async fn archive_prefetch_preserves_order_and_stops_at_gap() { + let tmp = tempfile::tempdir().unwrap(); + let settings = walrus::config::Settings { + storage: walrus::config::StorageSettings::Fs { + path: tmp.path().join("archive").display().to_string(), + }, + ..Default::default() + }; + let storage = settings.build_storage().unwrap(); + for index in [0u64, 1, 3] { + let name = + segments_covering(1, index * WAL_SEG_SIZE..(index + 1) * WAL_SEG_SIZE)[0].format(); + let path = tmp.path().join(name); + fs::write(&path, vec![index as u8; WAL_SEG_SIZE as usize]).unwrap(); + walrus::pg::wal::push::handle(&settings, storage.clone(), &path) + .await + .unwrap(); + } + let archive = Archive { settings, storage }; + let mut reader = archive.feed(TimelineHistory::root(1), 1, 42, 4); + let (lsn, bytes) = reader.next().await.unwrap().unwrap(); + assert_eq!(lsn, 42); + assert_eq!(bytes.len(), WAL_SEG_SIZE as usize - 42); + assert!(bytes.iter().all(|b| *b == 0)); + let (lsn, bytes) = reader.next().await.unwrap().unwrap(); + assert_eq!(lsn, WAL_SEG_SIZE); + assert!(bytes.iter().all(|b| *b == 1)); + assert!(reader.next().await.unwrap().is_err()); + assert!( + reader.next().await.is_none(), + "must not skip missing segment" + ); + } + + #[tokio::test] + async fn archive_prefetch_reads_fork_segment_under_descendant_timeline() { + let tmp = tempfile::tempdir().unwrap(); + let wal = tmp.path().join("wal_005"); + fs::create_dir(&wal).unwrap(); + fs::write( + wal.join("000000010000000000000000"), + vec![1; WAL_SEG_SIZE as usize], + ) + .unwrap(); + fs::write( + wal.join("000000010000000000000001.partial"), + vec![9; WAL_SEG_SIZE as usize], + ) + .unwrap(); + fs::write( + wal.join("000000020000000000000001"), + vec![2; WAL_SEG_SIZE as usize], + ) + .unwrap(); + let settings = walrus::config::Settings { + storage: walrus::config::StorageSettings::Fs { + path: tmp.path().display().to_string(), + }, + ..Default::default() + }; + let history = TimelineHistory::parse(2, b"1\t0/10000A0\n").unwrap(); + let mut feed = Archive::open(settings) + .unwrap() + .feed(history, 1, WAL_SEG_SIZE - 42, 2); + let (lsn, bytes) = feed.next().await.unwrap().unwrap(); + assert_eq!(lsn, WAL_SEG_SIZE - 42); + assert_eq!(bytes, vec![1; 42]); + let (lsn, bytes) = feed.next().await.unwrap().unwrap(); + assert_eq!(lsn, WAL_SEG_SIZE); + assert_eq!(bytes, vec![2; 0xA0], "cut at ancestor's switchpoint"); + assert!(feed.next().await.is_none(), "feed ends at switchpoint"); + } + + #[tokio::test] + async fn archive_prefetch_drop_cancels_worker() { + let (tx, rx) = tokio::sync::mpsc::channel(1); + let task = tokio::spawn(async move { + let _tx = tx; + std::future::pending::<()>().await; + }); + let abort = task.abort_handle(); + drop(ArchiveFeed { + wait_nanos: AtomicU64::new(0), + rx, + task, + fetch_nanos: Arc::new(AtomicU64::new(0)), + }); + tokio::task::yield_now().await; + assert!(abort.is_finished()); + } + + #[tokio::test] + async fn archive_fetch_reads_exact_segment() { + let tmp = tempfile::tempdir().unwrap(); + let archive = tmp.path().join("archive"); + let segment_path = tmp.path().join("000000010000000000000000"); + fs::write(&segment_path, vec![0; WAL_SEG_SIZE as usize]).unwrap(); + let settings = walrus::config::Settings { + storage: walrus::config::StorageSettings::Fs { + path: archive.display().to_string(), + }, + ..Default::default() + }; + let storage = settings.build_storage().unwrap(); + walrus::pg::wal::push::handle(&settings, storage.clone(), &segment_path) + .await + .unwrap(); + + let (name, bytes) = Archive { settings, storage } + .read_segment(&TimelineHistory::root(1), 0) + .await + .unwrap(); + assert_eq!(name, "000000010000000000000000"); + assert_eq!(bytes.len(), WAL_SEG_SIZE as usize); + } + + #[tokio::test] + async fn archive_fetch_falls_back_across_compressions() { + let tmp = tempfile::tempdir().unwrap(); + let segment_path = tmp.path().join("000000010000000000000000"); + fs::write(&segment_path, vec![7; WAL_SEG_SIZE as usize]).unwrap(); + let storage = walrus::config::StorageSettings::Fs { + path: tmp.path().join("archive").display().to_string(), + }; + let pushed = walrus::config::Settings { + storage: storage.clone(), + compression: walrus::compression::Method::None, + ..Default::default() + }; + let built = pushed.build_storage().unwrap(); + walrus::pg::wal::push::handle(&pushed, built.clone(), &segment_path) + .await + .unwrap(); + // Read archives written with different compression settings. + let reading = walrus::config::Settings { + storage, + ..Default::default() + }; + let (_, bytes) = Archive { + settings: reading, + storage: built, + } + .read_segment(&TimelineHistory::root(1), 0) + .await + .unwrap(); + assert_eq!(bytes.len(), WAL_SEG_SIZE as usize); + assert!(bytes.iter().all(|b| *b == 7)); + } + + #[tokio::test] + async fn archive_fetch_slices_from_mid_segment() { + // Resume at requested LSN to keep replay aligned. + let tmp = tempfile::tempdir().unwrap(); + let archive = tmp.path().join("archive"); + let segment_path = tmp.path().join("000000010000000000000000"); + let pattern: Vec = (0..WAL_SEG_SIZE as usize) + .map(|i| (i % 251) as u8) + .collect(); + fs::write(&segment_path, &pattern).unwrap(); + let settings = walrus::config::Settings { + storage: walrus::config::StorageSettings::Fs { + path: archive.display().to_string(), + }, + ..Default::default() + }; + let storage = settings.build_storage().unwrap(); + walrus::pg::wal::push::handle(&settings, storage.clone(), &segment_path) + .await + .unwrap(); + + let offset = WAL_SEG_SIZE / 2; + let (name, bytes) = Archive { settings, storage } + .read_segment(&TimelineHistory::root(1), offset) + .await + .unwrap(); + assert_eq!(name, "000000010000000000000000"); + assert_eq!(bytes.len(), (WAL_SEG_SIZE - offset) as usize); + assert_eq!(bytes, pattern[offset as usize..]); + } +} diff --git a/src/source/archive_history.rs b/src/source/archive_history.rs deleted file mode 100644 index ae369c8a..00000000 --- a/src/source/archive_history.rs +++ /dev/null @@ -1,86 +0,0 @@ -//! Cross-check archived timeline history against the branch the source proved - -use anyhow::{Context, Result}; -use walrus::config::Settings; -use walrus::storage::DynStorage; - -use crate::source::timeline::{TimelineHistory, history_filename}; - -/// Gap replay names segments off `source`'s chain, so the archive serving them -/// has to record that same chain -pub async fn verify( - settings: &Settings, - storage: &DynStorage, - source: &TimelineHistory, -) -> Result<()> { - let name = history_filename(source.target()); - let raw = walrus::pg::wal::fetch::read_segment(settings, storage, &name) - .await - .with_context(|| format!("fetch {name}"))?; - let archived = - TimelineHistory::parse(source.target(), &raw).with_context(|| format!("parse {name}"))?; - anyhow::ensure!( - archived.entries() == source.entries(), - "archived {name} disagrees with source timeline history", - ); - Ok(()) -} - -#[cfg(test)] -mod tests { - use std::path::Path; - - use super::*; - - /// Match `archive_command`: store history uncompressed under `wal_005/` - async fn archive_history_file(settings: &Settings, storage: &DynStorage, tli: u32, body: &str) { - let dir = tempfile::tempdir().unwrap(); - let path = dir.path().join(history_filename(tli)); - tokio::fs::write(&path, body).await.unwrap(); - walrus::pg::wal::push::handle(settings, storage.clone(), &path) - .await - .unwrap(); - } - - fn fs_archive(root: &Path) -> (Settings, DynStorage) { - let settings = Settings { - storage: walrus::config::StorageSettings::Fs { - path: root.display().to_string(), - }, - ..Default::default() - }; - let storage = settings.build_storage().unwrap(); - (settings, storage) - } - - #[tokio::test] - async fn reject_an_archive_missing_the_sources_own_history() { - let tmp = tempfile::tempdir().unwrap(); - let (settings, storage) = fs_archive(&tmp.path().join("archive")); - archive_history_file(&settings, &storage, 2, "1\t0/3000000\tpromotion\n").await; - let source = TimelineHistory::parse(3, b"1\t0/3000000\tpromotion\n").unwrap(); - let err = verify(&settings, &storage, &source).await.unwrap_err(); - assert!(err.to_string().contains("00000003.history"), "{err:#}"); - } - - #[tokio::test] - async fn archived_ids_need_not_be_consecutive() { - let tmp = tempfile::tempdir().unwrap(); - let (settings, storage) = fs_archive(&tmp.path().join("archive")); - archive_history_file(&settings, &storage, 4, "1\t0/3000000\tpromotion\n").await; - let source = TimelineHistory::parse(4, b"1\t0/3000000\tpromotion\n").unwrap(); - verify(&settings, &storage, &source).await.unwrap(); - } - - #[tokio::test] - async fn reject_archived_history_with_different_ancestry_or_switchpoint() { - for body in ["1\t0/4000000\tpromotion\n", "2\t0/3000000\tpromotion\n"] { - let tmp = tempfile::tempdir().unwrap(); - let (settings, storage) = fs_archive(&tmp.path().join("archive")); - archive_history_file(&settings, &storage, 3, body).await; - let source = TimelineHistory::parse(3, b"1\t0/3000000\tpromotion\n").unwrap(); - let err = verify(&settings, &storage, &source).await.unwrap_err(); - assert!(err.to_string().contains("disagrees"), "{err:#}"); - } - } -} diff --git a/src/source/mod.rs b/src/source/mod.rs index 82d0f0fc..dbe72bf5 100644 --- a/src/source/mod.rs +++ b/src/source/mod.rs @@ -1,4 +1,4 @@ -pub mod archive_history; +pub mod archive; pub mod boundary_hold; pub mod catalog_capture; pub mod manifest; diff --git a/src/source/transition.rs b/src/source/transition.rs index eb53442b..6f19e9ad 100644 --- a/src/source/transition.rs +++ b/src/source/transition.rs @@ -41,6 +41,7 @@ use walrus::pg::backup::format_pg_lsn; use crate::pos::{Drain, FilterDurable, Floor, Pos, ResumeSafe, ShadowReplay, Switchpoint}; use crate::record::{RecordSink, SegmentSink}; +use crate::source::archive::Archive; use crate::source::shadow_stream::ShadowStreamState; use crate::source::source_feed::{SlotError, SourceEvent, SourceFeed, StandbyStatus}; use crate::source::timeline::{HistoryError, TimelineHistory, history_filename}; @@ -148,6 +149,16 @@ impl TransitionError { | Self::Slot(SlotError::Query(_)) ) } + + /// Source no longer holds fork segment or slot sits past it, both of which + /// archive can stand in for + fn archive_can_serve(&self) -> bool { + match self { + Self::Slot(SlotError::TooNew { .. }) => true, + Self::Source(e) => crate::source_feed::is_wal_segment_removed(e), + _ => false, + } + } } /// Cumulative crossing counters, read by the metrics snapshot. @@ -338,6 +349,15 @@ pub struct ForkResume { pub switch_lsn: Pos, } +/// Where a crossing read descendant's copy of fork prefix +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum PrefixOrigin { + /// Feed is left streaming descendant + Live, + /// Feed never reached descendant, so source must be redialed + Archive, +} + /// One crossing's outcome. The history comes back with it: the floor timeline /// is resolved against this chain, and a caller still holding the pre-fork /// history would pin the floor on the ancestor forever. @@ -349,6 +369,7 @@ pub struct Crossing { pub switch_lsn: u64, pub prefix_bytes: u64, pub history: TimelineHistory, + pub prefix_origin: PrefixOrigin, } pub struct Switchover<'a> { @@ -362,6 +383,7 @@ pub struct Switchover<'a> { /// probe can find them. pub out_dir: &'a Path, pub shadow_state: &'a Arc>, + pub backup: Option<&'a Archive>, } impl Switchover<'_> { @@ -577,28 +599,23 @@ impl Switchover<'_> { to: switch_lsn, }); } - // The promotion target's slot has to already reach the position about to - // be committed. `START_REPLICATION` would answer for the request alone, - // leaving a slot that pins nothing below it to be found at the next - // restart (architecture/recovery.md) - if let Some(name) = slot { - let restart_lsn = feed - .prove_physical_slot(name, Pos::new(seg_start), Pos::new(seg_start)) - .await?; - tracing::info!( - target: "walshadow", - slot = name, - restart_lsn = restart_lsn.map(|l| walrus::pg::backup::format_pg_lsn(l).to_string()), - resume_lsn = %walrus::pg::backup::format_pg_lsn(seg_start), - "target slot reaches the fork resume position", - ); - } - feed.start_physical_replication(slot, seg_start, next_tli) + let (prefix_bytes, past_fork, prefix_origin) = match self + .verify_prefix_live(feed, slot, seg_start, fork, status, ancestor) .await - .map_err(source)?; - let (prefix_bytes, past_fork) = self - .verify_prefix(feed, switch_lsn, status, ancestor) - .await?; + { + Ok((prefix_bytes, past_fork)) => (prefix_bytes, past_fork, PrefixOrigin::Live), + Err(e) + if e.archive_can_serve() + && let Some(archive) = self.backup => + { + let (prefix_bytes, past_fork) = + archived_fork_prefix(archive, fork, ancestor).await?; + tracing::info!(target: "walshadow", error = %e, + next_timeline = next_tli, "crossing source timeline through archive"); + (prefix_bytes, past_fork, PrefixOrigin::Archive) + } + Err(e) => return Err(e), + }; // The hinge. Above this line the stream is still the ancestor's at `F` // and every failure re-crosses from there; below it the resume position // is the fork segment's start on the descendant @@ -641,9 +658,43 @@ impl Switchover<'_> { switch_lsn, prefix_bytes, history: fork.live_history().clone(), + prefix_origin, }) } + /// Stream descendant from fork segment's start and verify its prefix copy. + /// Promotion target's slot has to already reach the position about to be + /// committed. `START_REPLICATION` would answer for the request alone, + /// leaving a slot that pins nothing below it to be found at the next + /// restart (architecture/recovery.md) + async fn verify_prefix_live( + &self, + feed: &mut SourceFeed, + slot: Option<&str>, + seg_start: u64, + fork: &ForkPoint, + status: StandbyStatus, + ancestor: ForkPrefix, + ) -> Result<(u64, Vec), TransitionError> { + if let Some(name) = slot { + let restart_lsn = feed + .prove_physical_slot(name, Pos::new(seg_start), Pos::new(seg_start)) + .await?; + tracing::info!( + target: "walshadow", + slot = name, + restart_lsn = restart_lsn.map(|l| format_pg_lsn(l).to_string()), + resume_lsn = %format_pg_lsn(seg_start), + "target slot reaches the fork resume position", + ); + } + feed.start_physical_replication(slot, seg_start, fork.next_tli) + .await + .map_err(source)?; + self.verify_prefix(feed, fork.switch_lsn, status, ancestor) + .await + } + /// Always restart the descendant at the fork segment's start, matching /// `pg_receivewal`: the repeated `[from, switch_lsn)` bytes are the verbatim /// copy PostgreSQL made of the ancestor prefix (`XLogInitNewTimeline`), so @@ -687,16 +738,48 @@ impl Switchover<'_> { // until its copy of the prefix has answered for itself past_fork.extend_from_slice(&chunk.data[overlap..]); } - if crc != ancestor.crc { - return Err(TransitionError::ForkPrefixMismatch { - from: ancestor.from, - to: switch_lsn, - }); - } + check_prefix(ancestor, crc, switch_lsn)?; Ok((switch_lsn - ancestor.from, past_fork)) } } +/// Archive's copy of fork segment, verified against ancestor prefix like +/// [`Switchover::verify_prefix`], cut at descendant's own switchpoint when both +/// forks share the segment +async fn archived_fork_prefix( + archive: &Archive, + fork: &ForkPoint, + ancestor: ForkPrefix, +) -> Result<(u64, Vec), TransitionError> { + let history = fork.live_history(); + let (_, mut bytes) = archive + .read_segment(history, ancestor.from) + .await + .map_err(source)?; + let prefix = fork.switch_lsn - ancestor.from; + check_prefix( + ancestor, + crc32c::crc32c(&bytes[..prefix as usize]), + fork.switch_lsn, + )?; + if let Some(next_fork) = history.switchpoint_of(fork.next_tli) { + bytes.truncate(bytes.len().min((next_fork - ancestor.from) as usize)); + } + bytes.drain(..prefix as usize); + Ok((prefix, bytes)) +} + +/// Descendant's copy of `[ancestor.from, switch_lsn)` digests to what ancestor fed +fn check_prefix(ancestor: ForkPrefix, crc: u32, switch_lsn: u64) -> Result<(), TransitionError> { + if crc != ancestor.crc { + return Err(TransitionError::ForkPrefixMismatch { + from: ancestor.from, + to: switch_lsn, + }); + } + Ok(()) +} + fn source(e: anyhow::Error) -> TransitionError { TransitionError::Source(e) } @@ -968,6 +1051,60 @@ pub async fn seed_shadow_branches( mod tests { use super::*; + #[tokio::test] + async fn archived_prefix_checks_identity_and_stops_at_next_fork() { + let tmp = tempfile::tempdir().unwrap(); + let settings = walrus::config::Settings { + storage: walrus::config::StorageSettings::Fs { + path: tmp.path().display().to_string(), + }, + ..Default::default() + }; + let archive = Archive::open(settings).unwrap(); + let seg_size = crate::record::WAL_SEG_SIZE; + let history = TimelineHistory::parse(3, b"1\t0/10000A0\n2\t0/1000140\n").unwrap(); + let fork = ForkPoint { + finished_tli: 1, + next_tli: 2, + live_tli: 3, + switch_lsn: seg_size + 160, + histories: vec![history], + }; + let bytes = vec![0x5a; seg_size as usize]; + let ancestor = ForkPrefix { + from: seg_size, + through: fork.switch_lsn, + crc: crc32c::crc32c(&bytes[..160]), + }; + let wal = tmp.path().join("wal_005"); + std::fs::create_dir(&wal).unwrap(); + let segment = wal.join("000000030000000000000001"); + let error = archived_fork_prefix(&archive, &fork, ancestor) + .await + .unwrap_err(); + assert!(error.retryable()); + std::fs::write(&segment, &bytes).unwrap(); + let (prefix, tail) = archived_fork_prefix(&archive, &fork, ancestor) + .await + .unwrap(); + assert_eq!(prefix, 160); + assert_eq!(tail, vec![0x5a; 160]); + let wrong = ForkPrefix { + crc: ancestor.crc ^ 1, + ..ancestor + }; + let error = archived_fork_prefix(&archive, &fork, wrong) + .await + .unwrap_err(); + assert!(matches!(error, TransitionError::ForkPrefixMismatch { .. })); + assert!(!error.retryable()); + std::fs::write(&segment, &bytes[..320]).unwrap(); + let error = archived_fork_prefix(&archive, &fork, ancestor) + .await + .unwrap_err(); + assert!(error.to_string().contains("320 bytes")); + } + #[test] fn every_reason_has_a_label() { let errs = [ diff --git a/tests/backfill_gap_across_promotion.rs b/tests/backfill_gap_across_promotion.rs index 83459e62..572d5a7e 100644 --- a/tests/backfill_gap_across_promotion.rs +++ b/tests/backfill_gap_across_promotion.rs @@ -23,7 +23,7 @@ use walrus::pg::wal::segment::SegmentName; use walrus::pg::walparser::Oid; use walrus::storage::DynStorage; use walrus::storage::fs::FsStorage; -use walshadow::archive_history; +use walshadow::archive; use walshadow::backfill_bootstrap::seed_catalog_from_source; use walshadow::backup_backfill::fetch_segments; use walshadow::heap_decoder::{ColumnValue, decode_heap_record}; @@ -235,7 +235,7 @@ async fn gap_replay_crosses_a_promotion() { 2, "source pins the branch, not the archive" ); - archive_history::verify(&settings, &storage, &history) + archive::verify_history(&settings, &storage, &history) .await .expect("archive records the chain the source serves"); drop(feed); diff --git a/tests/common/pgext.rs b/tests/common/pgext.rs index 4be6de8a..43d031ba 100644 --- a/tests/common/pgext.rs +++ b/tests/common/pgext.rs @@ -465,7 +465,7 @@ impl Cluster { let deadline = Instant::now() + Duration::from_secs(30); loop { if let Some(pid) = self.worker_pid() - && self.bridge_path().exists() + && UnixStream::connect(self.bridge_path()).is_ok() { return pid; } diff --git a/tests/control_plane_e2e.rs b/tests/control_plane_e2e.rs index 6622cca1..03d8b9c7 100644 --- a/tests/control_plane_e2e.rs +++ b/tests/control_plane_e2e.rs @@ -19,6 +19,7 @@ use std::process::{Child, Command, Stdio}; use std::time::{Duration, Instant}; use anyhow::{Context, Result, bail, ensure}; +use walrus::pg::wal::segment::SegmentName; use walshadow::pg::parse_pg_lsn; use walshadow::record::WAL_SEG_SIZE; use walshadow::shadow::{Shadow, ShadowConfig}; @@ -1402,6 +1403,89 @@ async fn switchover_crosses_fork_and_keeps_every_row() { } } +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] +async fn archive_crosses_forks_when_ancestor_segments_are_partial_and_slot_is_ahead() { + if !gated() { + return; + } + let mut h = Harness::up(&fx::Ports::alloc()) + .await + .expect("bring up harness"); + let target = h.promotion_target().expect("build promotion target"); + let result = async { + pause_and_stop_writes(&h, &target).await?; + h.stop_daemon(); + promote(&target)?; + let fork = parse_pg_lsn(&fork_switch_lsn(&target, 2)?)?; + let seg_size = walshadow::record::WAL_SEG_SIZE; + let fork_start = fork / seg_size * seg_size; + ensure!(fork > fork_start, "need a fork inside a segment"); + target.psql_one("UPDATE demo.users SET email = 'archive-descendant@x' WHERE id = 1")?; + target.stop()?; + target.write_standby_signal()?; + target.start()?; + promote(&target)?; + let history = walshadow::timeline::TimelineHistory::parse(3, + &fs::read(target.config().data_dir.join("pg_wal/00000003.history"))?)?; + let second_fork = history.switchpoint_of(2).context("second fork missing")?; + ensure!(second_fork / seg_size == fork / seg_size, "forks must share a segment"); + target.psql_one("SELECT pg_switch_wal()")?; + target.psql_one("CHECKPOINT")?; + target.psql_one("SELECT pg_create_physical_replication_slot('archive_resume', true)")?; + let retained = parse_pg_lsn(&target.psql_one( + "SELECT restart_lsn FROM pg_replication_slots WHERE slot_name = 'archive_resume'" + )?)?; + ensure!(retained > fork, "slot must sit beyond historical fork"); + target.psql_one("SELECT pg_switch_wal()")?; + let current = target.psql_one("SELECT pg_walfile_name(pg_current_wal_insert_lsn())")?; + target.stop()?; + + let archive = h.tmp.path().join("archive"); + let wal = archive.join("wal_005"); + fs::create_dir_all(&wal)?; + let ancestor_name = walshadow::record::segments_covering(1, fork_start..fork_start + seg_size)[0].format(); + let current = SegmentName::parse(current.trim())?.start_lsn(seg_size); + let pg_wal = target.config().data_dir.join("pg_wal"); + for entry in fs::read_dir(&pg_wal)? { + let entry = entry?; + let Ok(seg) = SegmentName::parse(&entry.file_name().to_string_lossy()) else { + continue; + }; + let ancestor = seg.timeline < 3; + let start = seg.start_lsn(seg_size); + if start < current { + let archived_name = if ancestor && start == fork_start { format!("{}.partial", seg.format()) } else { seg.format() }; + fs::copy(entry.path(), wal.join(archived_name))?; + } + if ancestor { + fs::remove_file(entry.path())?; + } + } + ensure!(!wal.join(&ancestor_name).exists(), "ancestor fork segment must be absent"); + target.start()?; + fs::write(&h.frag_path, format!( + "[source]\nhost = \"{}\"\nport = {TARGET_PORT}\nslot = \"archive_resume\"\n\n[stream]\npaused = false\n\n[backup]\narchive = \"file://{}\"\n", + target.config().socket_dir.display(), archive.display(), + ))?; + h.start_daemon(Duration::from_secs(60)).await?; + h.wait_ch(SECOND_EMAIL, "below-fork@x", Duration::from_secs(60)).await?; + h.wait_ch(USER_EMAIL, "archive-descendant@x", Duration::from_secs(60)).await?; + h.wait_log("crossing source timeline through archive", Duration::from_secs(60)).await?; + h.wait_log("source reconnected", Duration::from_secs(60)).await?; + target.psql_one("UPDATE demo.users SET email = 'live-after-archive@x' WHERE id = 1")?; + h.wait_ch(USER_EMAIL, "live-after-archive@x", Duration::from_secs(60)).await?; + ensure!(h.metric("walshadow_timeline_switches_total")? == 2, "expected two forks"); + ensure!(h.metric("walshadow_timeline_prefix_bytes_verified_total")? > 0, "prefix was not verified"); + ensure!(h.alive(), "daemon exited after archive recovery"); + Ok::<(), anyhow::Error>(()) + }.await; + let stderr = h.teardown(); + let _ = target.stop(); + if let Err(e) = result { + panic!("{e:#}\n--- daemon stderr ---\n{stderr}"); + } +} + /// A crossing has to survive a restart, on both sides of the fork /// (architecture/recovery.md). Same drill as /// `switchover_crosses_fork_and_keeps_every_row`, then two restarts: