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
12 changes: 12 additions & 0 deletions architecture/recovery.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
6 changes: 3 additions & 3 deletions src/backfill/backup_backfill.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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;
Expand Down Expand Up @@ -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 =
Expand Down
308 changes: 59 additions & 249 deletions src/bin/stream/archive.rs
Original file line number Diff line number Diff line change
@@ -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<u8>)> {
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<u8>, tokio::sync::OwnedSemaphorePermit);

pub(crate) struct ArchiveFeed {
pub(crate) wait_nanos: AtomicU64,
pub(crate) rx: tokio::sync::mpsc::Receiver<Result<ArchiveSegment>>,
pub(crate) task: tokio::task::JoinHandle<()>,
pub(crate) fetch_nanos: Arc<AtomicU64>,
}

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<Result<(u64, Vec<u8>)>> {
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(
Expand Down Expand Up @@ -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<u8> = (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:#}");
}
}
Loading
Loading