use std::path::{Path, PathBuf};
use std::sync::Arc;
use std::time::{Duration, Instant};
use anyhow::{Context, Result};
use object_store::{ObjectStore, ObjectStoreExt};
use crate::claim::{self, ClaimLost, ClaimOutcome, ClaimRecord};
use crate::dedup::{self, DedupOutcome};
use crate::hydrate::{subject_path, Tier};
use crate::snapshot::{self, BackupTarget, SnapshotOutcome};
use crate::stream::{
self, CoreWalSeam, PreconditionSupport, PreflightStage, SourceFingerprint, StreamConfig,
StreamOutcome,
};
const MAX_BASE_ATTEMPTS: u32 = 3;
pub struct TailRequest<'a> {
pub store: Arc<dyn ObjectStore>,
pub store_prefix: &'a str,
pub volume_root: &'a Path,
pub subjects: &'a [String],
pub tier: Tier,
pub owner: &'a str,
pub page_size: usize,
pub rpo_target: Option<Duration>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum TailRefusal {
SinkNotFenced { stage: PreflightStage },
ClaimLost { current: ClaimRecord },
}
impl TailRefusal {
pub fn headline(&self) -> String {
match self {
TailRefusal::SinkNotFenced { stage } => format!(
"the sink does not enforce conditional puts ({}), so an ownership claim on it \
cannot fence anybody; refusing to stream rather than ship bytes behind a fence \
that is not there",
stage.as_str()
),
TailRefusal::ClaimLost { current } => format!(
"lost the ownership claim race to epoch {} (owner {})",
current.epoch, current.owner
),
}
}
}
pub struct TailSession {
epoch: u64,
displaced: Option<ClaimRecord>,
workload_target: BackupTarget,
subjects: Vec<SubjectState>,
tier: Tier,
owner: String,
page_size: usize,
rpo_target: Option<Duration>,
}
struct SubjectState {
name: String,
path: PathBuf,
target: BackupTarget,
base_snapshot_key: Option<String>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum SubjectOutcome {
Absent,
SourceTooHot { detail: String },
Snapshot(SnapshotOutcome),
Dedup(DedupOutcome),
Stream {
base_snapshot_key: String,
base_published: bool,
base_attempts: u32,
outcome: StreamOutcome,
},
}
#[derive(Debug, Clone, PartialEq)]
pub struct SubjectReport {
pub subject: String,
pub outcome: SubjectOutcome,
pub seconds: f64,
}
#[derive(Debug, Clone, PartialEq)]
pub enum RoundOutcome {
Backed(Vec<SubjectReport>),
Fenced { detail: String },
}
pub async fn start(req: TailRequest<'_>) -> Result<std::result::Result<TailSession, TailRefusal>> {
let workload_target = BackupTarget {
store: req.store.clone(),
prefix: req.store_prefix.to_string(),
};
if let PreconditionSupport::Degraded { stage } =
stream::probe_conditional_puts(&workload_target).await?
{
return Ok(Err(TailRefusal::SinkNotFenced { stage }));
}
let (epoch, displaced) = match claim::acquire(&workload_target, req.owner).await? {
ClaimOutcome::Granted { claim, previous } => (claim.epoch, previous),
ClaimOutcome::Lost { current } => return Ok(Err(TailRefusal::ClaimLost { current })),
};
let prefix = req.store_prefix.trim_matches('/');
let mut subjects = Vec::with_capacity(req.subjects.len());
for name in req.subjects {
subjects.push(SubjectState {
path: subject_path(req.volume_root, name)?,
target: BackupTarget {
store: req.store.clone(),
prefix: if prefix.is_empty() {
name.clone()
} else {
format!("{prefix}/{name}")
},
},
name: name.clone(),
base_snapshot_key: None,
});
}
Ok(Ok(TailSession {
epoch,
displaced,
workload_target,
subjects,
tier: req.tier,
owner: req.owner.to_string(),
page_size: req.page_size,
rpo_target: req.rpo_target,
}))
}
impl TailSession {
pub fn epoch(&self) -> u64 {
self.epoch
}
pub fn displaced(&self) -> Option<&ClaimRecord> {
self.displaced.as_ref()
}
pub fn subjects(&self) -> impl Iterator<Item = &str> {
self.subjects.iter().map(|s| s.name.as_str())
}
}
pub async fn round(session: &mut TailSession) -> Result<RoundOutcome> {
if let Err(lost) = claim::assert_holds(&session.workload_target, session.epoch).await? {
return Ok(RoundOutcome::Fenced {
detail: fenced_detail(&lost),
});
}
let tier = session.tier;
let page_size = session.page_size;
let owner = session.owner.clone();
let epoch = session.epoch;
let rpo_target = session.rpo_target;
let workload_target = BackupTarget {
store: session.workload_target.store.clone(),
prefix: session.workload_target.prefix.clone(),
};
let mut reports = Vec::with_capacity(session.subjects.len());
for subject in &mut session.subjects {
let started = Instant::now();
if !subject.path.exists() {
reports.push(SubjectReport {
subject: subject.name.clone(),
outcome: SubjectOutcome::Absent,
seconds: started.elapsed().as_secs_f64(),
});
continue;
}
let outcome = back_up_subject(
subject,
&workload_target,
tier,
page_size,
&owner,
epoch,
rpo_target,
)
.await
.with_context(|| format!("backing up subject {}", subject.name))?;
if let SubjectOutcome::Stream {
outcome: StreamOutcome::Fenced { current_epoch, our_epoch, .. },
..
} = &outcome
{
return Ok(RoundOutcome::Fenced {
detail: format!(
"the sink bounced subject {}: it is stamped with epoch {current_epoch} and we \
hold {our_epoch}",
subject.name
),
});
}
reports.push(SubjectReport {
subject: subject.name.clone(),
outcome,
seconds: started.elapsed().as_secs_f64(),
});
}
Ok(RoundOutcome::Backed(reports))
}
fn fenced_detail(lost: &ClaimLost) -> String {
match lost {
ClaimLost::Superseded { .. } => format!("the ownership claim moved: {lost}"),
ClaimLost::Vanished { .. } => format!("{lost}; treating that as fenced"),
}
}
async fn back_up_subject(
subject: &mut SubjectState,
workload_target: &BackupTarget,
tier: Tier,
page_size: usize,
owner: &str,
epoch: u64,
rpo_target: Option<Duration>,
) -> Result<SubjectOutcome> {
let path = subject
.path
.to_str()
.with_context(|| format!("subject path {} is not UTF-8", subject.path.display()))?
.to_string();
let path = path.as_str();
match tier {
Tier::Snapshot => {
let image = match live_image(path, page_size).await? {
Ok(image) => image,
Err(too_hot) => return Ok(too_hot),
};
Ok(SubjectOutcome::Snapshot(
snapshot::upload_snapshot_image(&subject.target, &image).await?,
))
}
Tier::Dedup => {
let image = match live_image(path, page_size).await? {
Ok(image) => image,
Err(too_hot) => return Ok(too_hot),
};
Ok(SubjectOutcome::Dedup(
dedup::snapshot_dedup_image(&image, &subject.target).await?,
))
}
Tier::Stream => {
stream_subject(subject, workload_target, path, page_size, owner, epoch, rpo_target)
.await
}
}
}
async fn live_image(
path: &str,
page_size: usize,
) -> Result<std::result::Result<Vec<u8>, SubjectOutcome>> {
match stream::raw_consistent_copy_live(path, page_size).await {
Ok(image) => Ok(Ok(image)),
Err(e) if is_source_too_hot(&e) => Ok(Err(SubjectOutcome::SourceTooHot {
detail: format!("{e:#}"),
})),
Err(e) => Err(e),
}
}
fn is_source_too_hot(e: &anyhow::Error) -> bool {
format!("{e:#}").contains("the source moved under every one of")
}
async fn stream_subject(
subject: &mut SubjectState,
workload_target: &BackupTarget,
path: &str,
page_size: usize,
owner: &str,
epoch: u64,
rpo_target: Option<Duration>,
) -> Result<SubjectOutcome> {
if subject.base_snapshot_key.is_none() {
subject.base_snapshot_key = snapshot::latest_snapshot_key(&subject.target).await?;
}
if let Some(base) = subject.base_snapshot_key.clone() {
let outcome = tail_once(subject, path, &base, page_size, owner, epoch, rpo_target).await?;
if matches!(outcome, StreamOutcome::Restarted { .. }) {
return rebase(
subject,
workload_target,
path,
page_size,
owner,
epoch,
rpo_target,
)
.await;
}
return Ok(SubjectOutcome::Stream {
base_snapshot_key: base,
base_published: false,
base_attempts: 0,
outcome,
});
}
let mut last_generation = None;
for attempt in 1..=MAX_BASE_ATTEMPTS {
let before = SourceFingerprint::read(path)?;
let image = match live_image(path, page_size).await? {
Ok(image) => image,
Err(too_hot) => return Ok(too_hot),
};
let base = snapshot::upload_base_snapshot(&subject.target, &image).await?;
let outcome = tail_once(subject, path, &base, page_size, owner, epoch, rpo_target).await?;
let after = SourceFingerprint::read(path)?;
if before.stable_across(&after) {
subject.base_snapshot_key = Some(base.clone());
return Ok(SubjectOutcome::Stream {
base_snapshot_key: base,
base_published: true,
base_attempts: attempt,
outcome,
});
}
last_generation = Some((before, after));
}
let (before, after) = last_generation.expect("MAX_BASE_ATTEMPTS is non-zero");
anyhow::bail!(
"could not anchor a tier-2 base for {path}: the source WAL was folded during each of \
{MAX_BASE_ATTEMPTS} attempts (last: {before:?} -> {after:?}). The base and the frames \
would describe different points in time, so nothing was anchored; the next round retries. \
An application checkpointing faster than its database can be copied needs a longer \
wal_autocheckpoint, not a longer retry loop."
);
}
async fn rebase(
subject: &mut SubjectState,
workload_target: &BackupTarget,
path: &str,
page_size: usize,
owner: &str,
epoch: u64,
rpo_target: Option<Duration>,
) -> Result<SubjectOutcome> {
if let Err(lost) = claim::assert_holds(workload_target, epoch).await? {
anyhow::bail!(
"refusing to re-anchor {path} after a WAL restart: {lost}. Rebasing rewrites the \
chain, and doing that without the claim would destroy the real owner's"
);
}
let image = match live_image(path, page_size).await? {
Ok(image) => image,
Err(too_hot) => return Ok(too_hot),
};
let base = snapshot::upload_base_snapshot(&subject.target, &image).await?;
delete_generation_manifests(&subject.target).await?;
subject
.target
.store
.delete(&subject.target.watermark_key())
.await
.or_else(ignore_absent)
.with_context(|| format!("clearing the stream watermark for {path}"))?;
let outcome = tail_once(subject, path, &base, page_size, owner, epoch, rpo_target).await?;
subject.base_snapshot_key = Some(base.clone());
Ok(SubjectOutcome::Stream {
base_snapshot_key: base,
base_published: true,
base_attempts: 1,
outcome,
})
}
async fn delete_generation_manifests(target: &BackupTarget) -> Result<()> {
let prefix = target.prefix.trim_matches('/');
let dir = object_store::path::Path::from(if prefix.is_empty() {
"generations".to_string()
} else {
format!("{prefix}/generations")
});
let listing = target
.store
.list_with_delimiter(Some(&dir))
.await
.with_context(|| format!("listing generation manifests under {dir}"))?;
for object in listing.objects {
target
.store
.delete(&object.location)
.await
.or_else(ignore_absent)
.with_context(|| format!("deleting stale generation manifest {}", object.location))?;
}
Ok(())
}
fn ignore_absent(e: object_store::Error) -> std::result::Result<(), object_store::Error> {
match e {
object_store::Error::NotFound { .. } => Ok(()),
other => Err(other),
}
}
async fn tail_once(
subject: &SubjectState,
path: &str,
base_snapshot_key: &str,
page_size: usize,
owner: &str,
epoch: u64,
rpo_target: Option<Duration>,
) -> Result<StreamOutcome> {
let cfg = StreamConfig {
base_snapshot_key,
page_size,
backpressure: Default::default(),
rpo_target,
epoch,
owner: Some(owner),
pointer_generation: 0,
};
let seam = CoreWalSeam::open_reader(path)
.with_context(|| format!("opening a read-only WAL seam on {path}"))?;
stream::tail_frames(&seam, &subject.target, &cfg)
.await
.with_context(|| format!("tailing {path} at epoch {epoch}"))
}
#[cfg(test)]
mod tests {
use super::*;
use object_store::memory::InMemory;
struct Volume(PathBuf);
impl Volume {
fn new(tag: &str) -> Self {
let dir = std::env::temp_dir().join(format!(
"turso-backup-tail-{tag}-{}-{:?}",
std::process::id(),
std::thread::current().id()
));
let _ = std::fs::remove_dir_all(&dir);
std::fs::create_dir_all(&dir).unwrap();
Self(dir)
}
fn path(&self) -> &Path {
&self.0
}
fn subject(&self, name: &str) -> String {
self.0.join(name).to_str().unwrap().to_string()
}
}
impl Drop for Volume {
fn drop(&mut self) {
let _ = std::fs::remove_dir_all(&self.0);
}
}
async fn seed(path: &str, start: i64, count: i64) {
let db = turso::Builder::new_local(path).build().await.unwrap();
let conn = db.connect().unwrap();
conn.execute("CREATE TABLE IF NOT EXISTS t (id INTEGER PRIMARY KEY, v TEXT)", ())
.await
.unwrap();
conn.execute("BEGIN", ()).await.unwrap();
for i in start..start + count {
conn.execute("INSERT INTO t (id, v) VALUES (?, ?)", (i, format!("v{i}")))
.await
.unwrap();
}
conn.execute("COMMIT", ()).await.unwrap();
}
async fn count_rows(path: &str) -> i64 {
let db = turso::Builder::new_local(path).build().await.unwrap();
let conn = db.connect().unwrap();
let mut r = conn.query("SELECT COUNT(*) FROM t", ()).await.unwrap();
r.next().await.unwrap().unwrap().get::<i64>(0).unwrap()
}
fn store() -> Arc<dyn ObjectStore> {
Arc::new(InMemory::new())
}
fn request<'a>(
store: Arc<dyn ObjectStore>,
vol: &'a Volume,
subjects: &'a [String],
tier: Tier,
owner: &'a str,
) -> TailRequest<'a> {
TailRequest {
store,
store_prefix: "workloads/acct",
volume_root: vol.path(),
subjects,
tier,
owner,
page_size: 4096,
rpo_target: None,
}
}
#[tokio::test]
async fn a_tail_then_a_hydrate_round_trips_the_declared_state() {
let vol = Volume::new("roundtrip");
let subjects = vec!["accounts.db".to_string()];
seed(&vol.subject("accounts.db"), 0, 120).await;
let store = store();
let mut session =
match start(request(store.clone(), &vol, &subjects, Tier::Stream, "node-a"))
.await
.unwrap()
{
Ok(s) => s,
Err(r) => panic!("start refused: {}", r.headline()),
};
assert_eq!(session.epoch(), claim::FIRST_EPOCH);
assert!(session.displaced().is_none(), "a virgin prefix displaces nobody");
let RoundOutcome::Backed(reports) = round(&mut session).await.unwrap() else {
panic!("the first round should not be fenced");
};
assert_eq!(reports.len(), 1);
assert!(
matches!(&reports[0].outcome, SubjectOutcome::Stream { base_published: true, .. }),
"the first round publishes the base: {:?}",
reports[0].outcome
);
seed(&vol.subject("accounts.db"), 120, 80).await;
let RoundOutcome::Backed(reports) = round(&mut session).await.unwrap() else {
panic!("the second round should not be fenced");
};
assert!(
matches!(&reports[0].outcome, SubjectOutcome::Stream { base_published: false, .. }),
"the base is published once, not per round: {:?}",
reports[0].outcome
);
let dest = Volume::new("roundtrip-dest");
let outcome = crate::hydrate::hydrate(crate::hydrate::HydrateRequest {
store,
store_prefix: "workloads/acct",
volume_root: dest.path(),
subjects: &subjects,
tier: Tier::Stream,
owner: "node-b",
})
.await
.unwrap();
assert!(
matches!(outcome, crate::hydrate::HydrateOutcome::Hydrated { .. }),
"the tail's output must be hydratable: {outcome:?}"
);
assert_eq!(
count_rows(&dest.subject("accounts.db")).await,
200,
"every row the application committed before the last round must come back"
);
}
#[tokio::test]
async fn every_declared_subject_is_backed_up_not_just_the_first() {
let vol = Volume::new("three");
let subjects = vec![
"accounts.db".to_string(),
"passkeys.db".to_string(),
"sessions.db".to_string(),
];
for (i, s) in subjects.iter().enumerate() {
seed(&vol.subject(s), 0, 10 * (i as i64 + 1)).await;
}
let store = store();
let mut session = start(request(store.clone(), &vol, &subjects, Tier::Stream, "node-a"))
.await
.unwrap()
.expect("start");
let RoundOutcome::Backed(reports) = round(&mut session).await.unwrap() else {
panic!("not fenced");
};
assert_eq!(reports.len(), 3);
for r in &reports {
assert!(
matches!(r.outcome, SubjectOutcome::Stream { .. }),
"{} was not streamed: {:?}",
r.subject,
r.outcome
);
}
let dest = Volume::new("three-dest");
crate::hydrate::hydrate(crate::hydrate::HydrateRequest {
store,
store_prefix: "workloads/acct",
volume_root: dest.path(),
subjects: &subjects,
tier: Tier::Stream,
owner: "node-b",
})
.await
.unwrap();
for (i, s) in subjects.iter().enumerate() {
assert_eq!(
count_rows(&dest.subject(s)).await,
10 * (i as i64 + 1),
"subject {s} did not come back"
);
}
}
#[tokio::test]
async fn a_source_that_never_holds_still_is_too_hot_not_a_failure() {
let vol = Volume::new("hot");
let path = vol.subject("accounts.db");
seed(&path, 0, 10).await;
let err = stream::validated_against_source(&path, "live copy", || async {
let mut f = std::fs::OpenOptions::new().append(true).open(&path).unwrap();
std::io::Write::write_all(&mut f, &vec![0u8; 4096]).unwrap();
Ok(Vec::<u8>::new())
})
.await
.expect_err("a source that moves under every attempt must be refused");
assert!(
is_source_too_hot(&err),
"the too-hot refusal was not recognised, so a busy appliance's tail would die \
instead of retrying next round: {err:#}"
);
let missing = stream::raw_consistent_copy_live(&vol.subject("nope.db"), 4096)
.await
.expect_err("a missing database is a failure");
assert!(!is_source_too_hot(&missing), "{missing:#}");
}
#[tokio::test]
async fn a_wal_restart_re_anchors_the_chain_instead_of_breaking_the_restore() {
let vol = Volume::new("refold");
let subjects = vec!["accounts.db".to_string()];
let db = vol.subject("accounts.db");
seed(&db, 0, 50).await;
let store = store();
let mut session = start(request(store.clone(), &vol, &subjects, Tier::Stream, "node-a"))
.await
.unwrap()
.expect("start");
assert!(matches!(round(&mut session).await.unwrap(), RoundOutcome::Backed(_)));
{
let d = turso::Builder::new_local(&db).build().await.unwrap();
let c = d.connect().unwrap();
let mut rows = c.query("PRAGMA wal_checkpoint(TRUNCATE)", ()).await.unwrap();
while rows.next().await.unwrap().is_some() {}
}
seed(&db, 50, 25).await;
let RoundOutcome::Backed(reports) = round(&mut session).await.unwrap() else {
panic!("not fenced")
};
match &reports[0].outcome {
SubjectOutcome::Stream { base_published, .. } => assert!(
base_published,
"a WAL restart must re-anchor onto a fresh base, not keep streaming onto the \
old one: {:?}",
reports[0].outcome
),
other => panic!("{other:?}"),
}
let dest = Volume::new("refold-dest");
crate::hydrate::hydrate(crate::hydrate::HydrateRequest {
store,
store_prefix: "workloads/acct",
volume_root: dest.path(),
subjects: &subjects,
tier: Tier::Stream,
owner: "node-b",
})
.await
.expect("a prefix that survived a WAL restart must still be restorable");
assert_eq!(
count_rows(&dest.subject("accounts.db")).await,
75,
"the restore came back at the wrong point in time"
);
}
#[tokio::test]
async fn an_absent_subject_is_reported_rather_than_failing_the_round() {
let vol = Volume::new("absent");
let subjects = vec!["accounts.db".to_string(), "not-yet.db".to_string()];
seed(&vol.subject("accounts.db"), 0, 5).await;
let store = store();
let mut session = start(request(store, &vol, &subjects, Tier::Stream, "node-a"))
.await
.unwrap()
.expect("start");
let RoundOutcome::Backed(reports) = round(&mut session).await.unwrap() else {
panic!("not fenced");
};
assert!(matches!(reports[0].outcome, SubjectOutcome::Stream { .. }));
assert_eq!(reports[1].outcome, SubjectOutcome::Absent);
}
#[tokio::test]
async fn a_displaced_tail_is_fenced_on_its_next_round() {
let vol = Volume::new("fence");
let subjects = vec!["accounts.db".to_string()];
seed(&vol.subject("accounts.db"), 0, 20).await;
let store = store();
let mut first = start(request(store.clone(), &vol, &subjects, Tier::Stream, "node-a"))
.await
.unwrap()
.expect("start a");
assert!(matches!(round(&mut first).await.unwrap(), RoundOutcome::Backed(_)));
let second = start(request(store, &vol, &subjects, Tier::Stream, "node-b"))
.await
.unwrap()
.expect("start b");
assert_eq!(second.epoch(), first.epoch() + 1, "a takeover is a monotonic bump");
assert_eq!(
second.displaced().map(|c| c.owner.as_str()),
Some("node-a"),
"the takeover records who it displaced"
);
match round(&mut first).await.unwrap() {
RoundOutcome::Fenced { detail } => {
assert!(detail.contains("claim"), "unhelpful fence detail: {detail}")
}
other => panic!("the displaced tail kept writing: {other:?}"),
}
}
#[tokio::test]
async fn a_tier_1_tail_is_fenced_too_even_though_its_writes_carry_no_epoch() {
let vol = Volume::new("fence-t1");
let subjects = vec!["accounts.db".to_string()];
seed(&vol.subject("accounts.db"), 0, 20).await;
let store = store();
let mut first = start(request(store.clone(), &vol, &subjects, Tier::Snapshot, "node-a"))
.await
.unwrap()
.expect("start a");
assert!(matches!(round(&mut first).await.unwrap(), RoundOutcome::Backed(_)));
let _second = start(request(store, &vol, &subjects, Tier::Snapshot, "node-b"))
.await
.unwrap()
.expect("start b");
assert!(
matches!(round(&mut first).await.unwrap(), RoundOutcome::Fenced { .. }),
"a tier-1a tail that lost the claim must stop writing"
);
}
#[tokio::test]
async fn an_idle_tier_1a_subject_stops_uploading() {
let vol = Volume::new("idle");
let subjects = vec!["accounts.db".to_string()];
seed(&vol.subject("accounts.db"), 0, 30).await;
let store = store();
let mut session = start(request(store, &vol, &subjects, Tier::Snapshot, "node-a"))
.await
.unwrap()
.expect("start");
let RoundOutcome::Backed(first) = round(&mut session).await.unwrap() else {
panic!("not fenced")
};
assert!(
matches!(first[0].outcome, SubjectOutcome::Snapshot(SnapshotOutcome::Uploaded { .. })),
"{:?}",
first[0].outcome
);
let RoundOutcome::Backed(second) = round(&mut session).await.unwrap() else {
panic!("not fenced")
};
assert!(
matches!(
second[0].outcome,
SubjectOutcome::Snapshot(SnapshotOutcome::Deduplicated { .. })
),
"an untouched database was re-uploaded: {:?}",
second[0].outcome
);
}
#[tokio::test]
async fn tier_1b_backs_up_a_live_database_and_hydrate_reads_it() {
let vol = Volume::new("dedup");
let subjects = vec!["accounts.db".to_string()];
seed(&vol.subject("accounts.db"), 0, 60).await;
let store = store();
let mut session = start(request(store.clone(), &vol, &subjects, Tier::Dedup, "node-a"))
.await
.unwrap()
.expect("start");
let RoundOutcome::Backed(reports) = round(&mut session).await.unwrap() else {
panic!("not fenced")
};
assert!(
matches!(reports[0].outcome, SubjectOutcome::Dedup(DedupOutcome::Snapshotted { .. })),
"{:?}",
reports[0].outcome
);
let dest = Volume::new("dedup-dest");
crate::hydrate::hydrate(crate::hydrate::HydrateRequest {
store,
store_prefix: "workloads/acct",
volume_root: dest.path(),
subjects: &subjects,
tier: Tier::Dedup,
owner: "node-b",
})
.await
.unwrap();
assert_eq!(count_rows(&dest.subject("accounts.db")).await, 60);
}
#[tokio::test]
async fn a_restarted_tail_adopts_the_existing_base_instead_of_republishing() {
let vol = Volume::new("restart");
let subjects = vec!["accounts.db".to_string()];
seed(&vol.subject("accounts.db"), 0, 40).await;
let store = store();
let mut first = start(request(store.clone(), &vol, &subjects, Tier::Stream, "node-a"))
.await
.unwrap()
.expect("start a");
assert!(matches!(round(&mut first).await.unwrap(), RoundOutcome::Backed(_)));
drop(first);
let mut second = start(request(store, &vol, &subjects, Tier::Stream, "node-a"))
.await
.unwrap()
.expect("restart a");
let RoundOutcome::Backed(reports) = round(&mut second).await.unwrap() else {
panic!("not fenced")
};
match &reports[0].outcome {
SubjectOutcome::Stream { base_published, base_snapshot_key, .. } => {
assert!(!base_published, "a restart re-copied the whole database for nothing");
assert!(base_snapshot_key.contains("snapshots/"), "{base_snapshot_key}");
}
other => panic!("{other:?}"),
}
}
}