use std::path::{Path, PathBuf};
use std::sync::Arc;
use std::time::Instant;
use anyhow::{bail, Context, Result};
use object_store::ObjectStore;
use crate::claim::{self, ClaimLost, ClaimOutcome, ClaimRecord};
use crate::snapshot::BackupTarget;
use crate::stream::{PreconditionSupport, PreflightStage};
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum Tier {
Snapshot,
Dedup,
Stream,
}
impl Tier {
pub fn as_str(&self) -> &'static str {
match self {
Tier::Snapshot => "snapshot",
Tier::Dedup => "dedup",
Tier::Stream => "stream",
}
}
}
impl std::fmt::Display for Tier {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.write_str(self.as_str())
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum SubjectState {
Absent,
Populated { bytes: u64 },
}
impl SubjectState {
fn is_absent(&self) -> bool {
matches!(self, SubjectState::Absent)
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum Readiness {
Hydrate,
AlreadyPopulated,
TornVolume {
populated: Vec<String>,
absent: Vec<String>,
},
}
pub fn assess(states: &[(String, SubjectState)]) -> Result<Readiness> {
if states.is_empty() {
bail!("no subjects to hydrate — a tier that ships bytes must name at least one database");
}
let (absent, populated): (Vec<_>, Vec<_>) =
states.iter().partition(|(_, s)| s.is_absent());
match (absent.is_empty(), populated.is_empty()) {
(true, false) => Ok(Readiness::AlreadyPopulated),
(false, true) => Ok(Readiness::Hydrate),
(false, false) => Ok(Readiness::TornVolume {
populated: populated.iter().map(|(n, _)| n.clone()).collect(),
absent: absent.iter().map(|(n, _)| n.clone()).collect(),
}),
(true, true) => unreachable!("states was checked non-empty"),
}
}
pub fn inspect_volume(
volume_root: &Path,
subjects: &[String],
) -> Result<Vec<(String, SubjectState)>> {
subjects
.iter()
.map(|subject| {
let path = subject_path(volume_root, subject)?;
let state = match std::fs::metadata(&path) {
Ok(m) if m.len() > 0 => SubjectState::Populated { bytes: m.len() },
Ok(_) => SubjectState::Absent,
Err(e) if e.kind() == std::io::ErrorKind::NotFound => SubjectState::Absent,
Err(e) => {
return Err(e).with_context(|| format!("stat-ing subject {}", path.display()))
}
};
Ok((subject.clone(), state))
})
.collect()
}
pub fn subject_path(volume_root: &Path, subject: &str) -> Result<PathBuf> {
if subject.is_empty() {
bail!("empty subject path");
}
let p = Path::new(subject);
if p.is_absolute() {
bail!("subject {subject:?} is absolute; subjects are relative to the volume root");
}
if p.components()
.any(|c| matches!(c, std::path::Component::ParentDir | std::path::Component::CurDir))
{
bail!("subject {subject:?} contains a \".\" or \"..\" component");
}
Ok(volume_root.join(p))
}
#[derive(Debug, Clone, PartialEq)]
pub struct SubjectRestore {
pub subject: String,
pub source: String,
pub bytes: u64,
pub seconds: f64,
pub frames_replayed: Option<u64>,
}
#[derive(Debug, Clone, PartialEq)]
pub enum HydrateOutcome {
Hydrated {
epoch: u64,
subjects: Vec<SubjectRestore>,
seconds: f64,
},
AlreadyPopulated,
NothingInTheStore { epoch: u64, subjects: Vec<String> },
Refused(HydrateRefusal),
}
#[derive(Debug, Clone, PartialEq)]
pub enum HydrateRefusal {
TornVolume {
populated: Vec<String>,
absent: Vec<String>,
},
ClaimLost { current: ClaimRecord },
FencedMidRestore {
lost: ClaimLost,
restored: Vec<SubjectRestore>,
},
SinkNotFenced { stage: PreflightStage },
}
impl HydrateRefusal {
pub fn headline(&self) -> String {
match self {
HydrateRefusal::TornVolume { populated, absent } => format!(
"volume is torn: {} present ({}), {} missing ({}). Restoring only the missing \
ones would rebuild them at a different point in time from the ones already \
here; decide which copy is authoritative and clear or complete the volume by \
hand",
populated.len(),
populated.join(", "),
absent.len(),
absent.join(", ")
),
HydrateRefusal::ClaimLost { current } => format!(
"another node holds this workload's state: epoch {} (owner {}). It is being \
placed somewhere else right now",
current.epoch, current.owner
),
HydrateRefusal::FencedMidRestore { lost, restored } => format!(
"{lost} while restoring {} subject(s); the bytes are on disk but this node no \
longer owns them — do not start the workload",
restored.len()
),
HydrateRefusal::SinkNotFenced { stage } => format!(
"the object store does not enforce {} — the ownership claim would degrade to \
last-write-wins and two nodes could both believe they own this workload's \
state. Check that the bucket supports conditional writes (R2 and MinIO do; \
some S3-compatible backends do not)",
stage.as_str()
),
}
}
}
pub struct HydrateRequest<'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,
}
impl HydrateRequest<'_> {
fn workload_target(&self) -> BackupTarget {
BackupTarget {
store: self.store.clone(),
prefix: self.store_prefix.to_string(),
}
}
fn subject_target(&self, subject: &str) -> BackupTarget {
let prefix = self.store_prefix.trim_matches('/');
BackupTarget {
store: self.store.clone(),
prefix: if prefix.is_empty() {
subject.to_string()
} else {
format!("{prefix}/{subject}")
},
}
}
}
pub async fn hydrate(req: HydrateRequest<'_>) -> Result<HydrateOutcome> {
let states = inspect_volume(req.volume_root, req.subjects)?;
match assess(&states)? {
Readiness::AlreadyPopulated => return Ok(HydrateOutcome::AlreadyPopulated),
Readiness::TornVolume { populated, absent } => {
return Ok(HydrateOutcome::Refused(HydrateRefusal::TornVolume {
populated,
absent,
}))
}
Readiness::Hydrate => {}
}
let started = Instant::now();
let workload_target = req.workload_target();
if let PreconditionSupport::Degraded { stage } =
crate::stream::probe_conditional_puts(&workload_target).await?
{
return Ok(HydrateOutcome::Refused(HydrateRefusal::SinkNotFenced {
stage,
}));
}
let epoch = match claim::acquire(&workload_target, req.owner).await? {
ClaimOutcome::Granted { claim, .. } => claim.epoch,
ClaimOutcome::Lost { current } => {
return Ok(HydrateOutcome::Refused(HydrateRefusal::ClaimLost {
current,
}))
}
};
let mut restored = Vec::new();
let mut missing = Vec::new();
for subject in req.subjects {
let dest = subject_path(req.volume_root, subject)?;
if let Some(parent) = dest.parent() {
std::fs::create_dir_all(parent)
.with_context(|| format!("creating {} for subject {subject}", parent.display()))?;
}
let dest_str = dest
.to_str()
.with_context(|| format!("subject path {} is not UTF-8", dest.display()))?;
let target = req.subject_target(subject);
let subject_started = Instant::now();
let Some((source, frames_replayed)) =
restore_subject(&target, req.tier, dest_str).await?
else {
missing.push(subject.clone());
continue;
};
let bytes = std::fs::metadata(&dest)
.with_context(|| format!("stat-ing restored subject {}", dest.display()))?
.len();
restored.push(SubjectRestore {
subject: subject.clone(),
source,
bytes,
seconds: subject_started.elapsed().as_secs_f64(),
frames_replayed,
});
}
if !missing.is_empty() && !restored.is_empty() {
return Ok(HydrateOutcome::Refused(HydrateRefusal::TornVolume {
populated: restored.into_iter().map(|r| r.subject).collect(),
absent: missing,
}));
}
if let Err(lost) = claim::assert_holds(&workload_target, epoch).await? {
return Ok(HydrateOutcome::Refused(HydrateRefusal::FencedMidRestore {
lost,
restored,
}));
}
if restored.is_empty() {
return Ok(HydrateOutcome::NothingInTheStore {
epoch,
subjects: missing,
});
}
Ok(HydrateOutcome::Hydrated {
epoch,
subjects: restored,
seconds: started.elapsed().as_secs_f64(),
})
}
async fn restore_subject(
target: &BackupTarget,
tier: Tier,
dest: &str,
) -> Result<Option<(String, Option<u64>)>> {
match tier {
Tier::Snapshot => match crate::snapshot::restore_latest(target, dest).await {
Ok(key) => Ok(Some((key, None))),
Err(e) if is_absent_source(&e) => Ok(None),
Err(e) => Err(e),
},
Tier::Dedup => match crate::dedup::load_latest_manifest(target).await? {
None => Ok(None),
Some((key, manifest)) => {
crate::dedup::restore_from_manifest(target, &manifest, dest).await?;
Ok(Some((key, None)))
}
},
Tier::Stream => {
let manifests = crate::stream::list_and_parse_generation_manifests(target).await?;
if manifests.is_empty() {
return match crate::snapshot::restore_latest(target, dest).await {
Ok(key) => Ok(Some((key, Some(0)))),
Err(e) if is_absent_source(&e) => Ok(None),
Err(e) => Err(e),
};
}
let outcome =
crate::stream::restore_stream_from_manifests(target, dest, &manifests).await?;
Ok(Some((
outcome.base_snapshot_key.clone(),
Some(outcome.frames_replayed),
)))
}
}
}
fn is_absent_source(e: &anyhow::Error) -> bool {
let text = format!("{e:#}");
text.contains("no snapshots found under") || text.contains("no snapshot manifest found")
}
#[cfg(test)]
mod tests {
use super::*;
use object_store::memory::InMemory;
use object_store::ObjectStoreExt;
fn s(names: &[(&str, SubjectState)]) -> Vec<(String, SubjectState)> {
names.iter().map(|(n, st)| (n.to_string(), *st)).collect()
}
const PRESENT: SubjectState = SubjectState::Populated { bytes: 4096 };
#[test]
fn an_empty_volume_hydrates_and_a_full_one_does_not() {
assert_eq!(
assess(&s(&[("a.db", SubjectState::Absent), ("b.db", SubjectState::Absent)])).unwrap(),
Readiness::Hydrate
);
assert_eq!(
assess(&s(&[("a.db", PRESENT), ("b.db", PRESENT)])).unwrap(),
Readiness::AlreadyPopulated
);
}
#[test]
fn a_partly_populated_volume_is_refused_rather_than_topped_up() {
let r = assess(&s(&[
("accounts.db", PRESENT),
("passkeys.db", SubjectState::Absent),
("sessions.db", PRESENT),
]))
.unwrap();
assert_eq!(
r,
Readiness::TornVolume {
populated: vec!["accounts.db".into(), "sessions.db".into()],
absent: vec!["passkeys.db".into()],
}
);
}
#[test]
fn a_zero_length_file_counts_as_absent() {
let dir = tempdir("zero");
std::fs::write(dir.join("a.db"), b"").unwrap();
let states = inspect_volume(&dir, &["a.db".to_string()]).unwrap();
assert_eq!(states[0].1, SubjectState::Absent);
assert_eq!(assess(&states).unwrap(), Readiness::Hydrate);
}
#[test]
fn inspect_volume_reports_sizes_for_what_is_there() {
let dir = tempdir("sizes");
std::fs::write(dir.join("a.db"), b"0123456789").unwrap();
let states = inspect_volume(&dir, &["a.db".into(), "b.db".into()]).unwrap();
assert_eq!(states[0].1, SubjectState::Populated { bytes: 10 });
assert_eq!(states[1].1, SubjectState::Absent);
}
#[test]
fn a_traversing_subject_is_refused_at_the_path_join_too() {
let dir = tempdir("traverse");
for bad in ["/etc/passwd", "../../etc/passwd", "./a.db"] {
let err = inspect_volume(&dir, &[bad.to_string()]).unwrap_err();
assert!(
format!("{err:#}").contains(bad),
"refusal should name the subject: {err:#}"
);
}
}
#[test]
fn no_subjects_at_all_is_an_error_not_a_silent_no_op() {
let err = assess(&[]).unwrap_err();
assert!(format!("{err:#}").contains("at least one database"), "{err:#}");
}
fn tempdir(tag: &str) -> PathBuf {
let p = std::env::temp_dir().join(format!(
"turso-hydrate-test-{tag}-{}-{:?}",
std::process::id(),
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_nanos()
));
std::fs::create_dir_all(&p).unwrap();
p
}
fn db_image(tag: &[u8]) -> Vec<u8> {
let mut v = b"SQLite format 3\0".to_vec();
v.extend_from_slice(tag);
v
}
async fn seed_snapshot(store: &Arc<dyn ObjectStore>, prefix: &str, tag: &[u8]) -> String {
let target = BackupTarget {
store: store.clone(),
prefix: prefix.to_string(),
};
crate::snapshot::upload_base_snapshot(&target, &db_image(tag))
.await
.unwrap()
}
fn request<'a>(
store: &Arc<dyn ObjectStore>,
volume: &'a Path,
subjects: &'a [String],
owner: &'a str,
) -> HydrateRequest<'a> {
HydrateRequest {
store: store.clone(),
store_prefix: "wl/acct",
volume_root: volume,
subjects,
tier: Tier::Snapshot,
owner,
}
}
#[tokio::test]
async fn an_empty_volume_is_filled_from_the_store_under_a_fresh_epoch() {
let store: Arc<dyn ObjectStore> = Arc::new(InMemory::new());
seed_snapshot(&store, "wl/acct/accounts.db", b"ACCOUNTS").await;
seed_snapshot(&store, "wl/acct/sessions.db", b"SESSIONS").await;
let dir = tempdir("fill");
let subjects = vec!["accounts.db".to_string(), "sessions.db".to_string()];
let out = hydrate(request(&store, &dir, &subjects, "us-west-003"))
.await
.unwrap();
let HydrateOutcome::Hydrated { epoch, subjects: r, .. } = out else {
panic!("expected Hydrated, got {out:?}");
};
assert_eq!(epoch, claim::FIRST_EPOCH);
assert_eq!(r.len(), 2);
assert_eq!(
std::fs::read(dir.join("accounts.db")).unwrap(),
db_image(b"ACCOUNTS")
);
assert_eq!(
std::fs::read(dir.join("sessions.db")).unwrap(),
db_image(b"SESSIONS")
);
assert_eq!(r[0].bytes, db_image(b"ACCOUNTS").len() as u64);
let held = claim::read_claim(&BackupTarget {
store: store.clone(),
prefix: "wl/acct".into(),
})
.await
.unwrap()
.unwrap();
assert_eq!(held.owner, "us-west-003");
}
#[tokio::test]
async fn an_ordinary_restart_neither_restores_nor_takes_a_claim() {
let store: Arc<dyn ObjectStore> = Arc::new(InMemory::new());
seed_snapshot(&store, "wl/acct/accounts.db", b"OLD-BACKUP").await;
let dir = tempdir("restart");
std::fs::write(dir.join("accounts.db"), b"LIVE-DATA-WRITTEN-SINCE").unwrap();
let subjects = vec!["accounts.db".to_string()];
let out = hydrate(request(&store, &dir, &subjects, "us-east-001"))
.await
.unwrap();
assert_eq!(out, HydrateOutcome::AlreadyPopulated);
assert_eq!(
std::fs::read(dir.join("accounts.db")).unwrap(),
b"LIVE-DATA-WRITTEN-SINCE"
);
assert_eq!(
claim::read_claim(&BackupTarget {
store: store.clone(),
prefix: "wl/acct".into()
})
.await
.unwrap(),
None,
"a restart must not mint an epoch"
);
}
#[tokio::test]
async fn a_torn_volume_is_refused_before_the_claim_is_touched() {
let store: Arc<dyn ObjectStore> = Arc::new(InMemory::new());
seed_snapshot(&store, "wl/acct/accounts.db", b"A").await;
seed_snapshot(&store, "wl/acct/sessions.db", b"S").await;
let dir = tempdir("torn");
std::fs::write(dir.join("accounts.db"), b"LIVE").unwrap();
let subjects = vec!["accounts.db".to_string(), "sessions.db".to_string()];
let out = hydrate(request(&store, &dir, &subjects, "n1")).await.unwrap();
let HydrateOutcome::Refused(refusal) = out else {
panic!("expected a refusal, got {out:?}");
};
assert!(refusal.headline().contains("sessions.db"), "{}", refusal.headline());
assert!(
!dir.join("sessions.db").exists(),
"a refused hydrate must not have written anything"
);
}
#[tokio::test]
async fn losing_the_claim_race_refuses_instead_of_hydrating() {
let store = racing_store(RivalOn::ClaimPut, 9, "us-east-001");
seed_snapshot(&store, "wl/acct/accounts.db", b"A").await;
let dir = tempdir("raced");
let subjects = vec!["accounts.db".to_string()];
let out = hydrate(request(&store, &dir, &subjects, "us-west-003"))
.await
.unwrap();
let HydrateOutcome::Refused(HydrateRefusal::ClaimLost { current }) = out else {
panic!("expected ClaimLost, got {out:?}");
};
assert_eq!(current.owner, "us-east-001");
assert_eq!(current.epoch, 9);
assert!(
!dir.join("accounts.db").exists(),
"a node that lost the claim must not have read a byte of state"
);
}
#[tokio::test]
async fn a_sequential_takeover_still_hydrates() {
let store: Arc<dyn ObjectStore> = Arc::new(InMemory::new());
seed_snapshot(&store, "wl/acct/accounts.db", b"A").await;
let dir = tempdir("failover");
let subjects = vec!["accounts.db".to_string()];
claim::acquire(
&BackupTarget {
store: store.clone(),
prefix: "wl/acct".into(),
},
"us-east-001",
)
.await
.unwrap();
let out = hydrate(request(&store, &dir, &subjects, "us-west-003"))
.await
.unwrap();
let HydrateOutcome::Hydrated { epoch, .. } = out else {
panic!("a sequential takeover must hydrate, got {out:?}");
};
assert_eq!(epoch, claim::FIRST_EPOCH + 1);
}
#[tokio::test]
async fn an_empty_prefix_is_nothing_in_the_store_not_an_error() {
let store: Arc<dyn ObjectStore> = Arc::new(InMemory::new());
let dir = tempdir("virgin");
let subjects = vec!["accounts.db".to_string()];
let out = hydrate(request(&store, &dir, &subjects, "n1")).await.unwrap();
let HydrateOutcome::NothingInTheStore { epoch, subjects: missing } = out else {
panic!("expected NothingInTheStore, got {out:?}");
};
assert_eq!(epoch, claim::FIRST_EPOCH);
assert_eq!(missing, vec!["accounts.db".to_string()]);
assert!(!dir.join("accounts.db").exists());
}
#[tokio::test]
async fn a_prefix_holding_only_some_subjects_is_refused() {
let store: Arc<dyn ObjectStore> = Arc::new(InMemory::new());
seed_snapshot(&store, "wl/acct/accounts.db", b"A").await;
let dir = tempdir("half");
let subjects = vec!["accounts.db".to_string(), "sessions.db".to_string()];
let out = hydrate(request(&store, &dir, &subjects, "n1")).await.unwrap();
let HydrateOutcome::Refused(HydrateRefusal::TornVolume { populated, absent }) = out else {
panic!("expected TornVolume, got {out:?}");
};
assert_eq!(populated, vec!["accounts.db".to_string()]);
assert_eq!(absent, vec!["sessions.db".to_string()]);
}
#[tokio::test]
async fn a_takeover_during_the_restore_is_caught_after_the_bytes_land() {
let store = racing_store(RivalOn::SnapshotRead, 9, "us-east-001");
seed_snapshot(&store, "wl/acct/accounts.db", b"ACCOUNTS").await;
let dir = tempdir("midrestore");
let subjects = vec!["accounts.db".to_string()];
let out = hydrate(request(&store, &dir, &subjects, "us-west-003"))
.await
.unwrap();
let HydrateOutcome::Refused(HydrateRefusal::FencedMidRestore { lost, restored }) = out
else {
panic!("expected FencedMidRestore, got {out:?}");
};
let ClaimLost::Superseded { current, ours } = &lost else {
panic!("expected Superseded, got {lost:?}");
};
assert_eq!(*ours, claim::FIRST_EPOCH);
assert_eq!(current.epoch, 9);
assert_eq!(current.owner, "us-east-001");
assert_eq!(restored.len(), 1, "the restore itself completed");
assert!(
dir.join("accounts.db").exists(),
"the restored bytes must be left alone; the refusal is what stops the workload"
);
assert!(
HydrateRefusal::FencedMidRestore { lost, restored }
.headline()
.contains("do not start")
);
}
#[derive(Clone, Copy, PartialEq, Eq)]
enum RivalOn {
ClaimPut,
SnapshotRead,
}
#[tokio::test]
async fn a_sink_that_ignores_conditional_puts_is_refused_before_any_claim() {
let store: Arc<dyn ObjectStore> = Arc::new(UnconditionalStore {
inner: Arc::new(InMemory::new()),
});
seed_snapshot(&store, "wl/acct/accounts.db", b"A").await;
let dir = tempdir("degraded");
let subjects = vec!["accounts.db".to_string()];
let out = hydrate(request(&store, &dir, &subjects, "n1")).await.unwrap();
let HydrateOutcome::Refused(refusal) = out else {
panic!("expected a refusal, got {out:?}");
};
assert!(
matches!(refusal, HydrateRefusal::SinkNotFenced { .. }),
"{refusal:?}"
);
assert!(refusal.headline().contains("conditional writes"), "{}", refusal.headline());
assert!(
!dir.join("accounts.db").exists(),
"nothing may be restored under a fence that is not real"
);
}
struct UnconditionalStore {
inner: Arc<dyn ObjectStore>,
}
impl std::fmt::Display for UnconditionalStore {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(f, "UnconditionalStore({})", self.inner)
}
}
impl std::fmt::Debug for UnconditionalStore {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(f, "UnconditionalStore({:?})", self.inner)
}
}
#[async_trait::async_trait]
impl ObjectStore for UnconditionalStore {
async fn put_opts(
&self,
location: &object_store::path::Path,
payload: object_store::PutPayload,
_opts: object_store::PutOptions,
) -> object_store::Result<object_store::PutResult> {
self.inner
.put_opts(location, payload, object_store::PutOptions::default())
.await
}
async fn put_multipart_opts(
&self,
location: &object_store::path::Path,
opts: object_store::PutMultipartOptions,
) -> object_store::Result<Box<dyn object_store::MultipartUpload>> {
self.inner.put_multipart_opts(location, opts).await
}
async fn get_opts(
&self,
location: &object_store::path::Path,
options: object_store::GetOptions,
) -> object_store::Result<object_store::GetResult> {
self.inner.get_opts(location, options).await
}
fn delete_stream(
&self,
locations: futures_util::stream::BoxStream<
'static,
object_store::Result<object_store::path::Path>,
>,
) -> futures_util::stream::BoxStream<'static, object_store::Result<object_store::path::Path>>
{
self.inner.delete_stream(locations)
}
fn list(
&self,
prefix: Option<&object_store::path::Path>,
) -> futures_util::stream::BoxStream<'static, object_store::Result<object_store::ObjectMeta>>
{
self.inner.list(prefix)
}
async fn list_with_delimiter(
&self,
prefix: Option<&object_store::path::Path>,
) -> object_store::Result<object_store::ListResult> {
self.inner.list_with_delimiter(prefix).await
}
async fn copy_opts(
&self,
from: &object_store::path::Path,
to: &object_store::path::Path,
options: object_store::CopyOptions,
) -> object_store::Result<()> {
self.inner.copy_opts(from, to, options).await
}
}
struct RacingStore {
inner: Arc<dyn ObjectStore>,
claim_key: object_store::path::Path,
when: RivalOn,
rival_epoch: u64,
rival_owner: String,
fired: std::sync::atomic::AtomicBool,
}
impl RacingStore {
async fn fire(&self) -> object_store::Result<()> {
if self
.fired
.swap(true, std::sync::atomic::Ordering::SeqCst)
{
return Ok(());
}
let body = format!("{} 1 {}\n", self.rival_epoch, self.rival_owner);
self.inner
.put(&self.claim_key, body.into_bytes().into())
.await?;
Ok(())
}
}
impl std::fmt::Display for RacingStore {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(f, "RacingStore({})", self.inner)
}
}
impl std::fmt::Debug for RacingStore {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(f, "RacingStore({:?})", self.inner)
}
}
#[async_trait::async_trait]
impl ObjectStore for RacingStore {
async fn put_opts(
&self,
location: &object_store::path::Path,
payload: object_store::PutPayload,
opts: object_store::PutOptions,
) -> object_store::Result<object_store::PutResult> {
if self.when == RivalOn::ClaimPut && location == &self.claim_key {
self.fire().await?;
}
self.inner.put_opts(location, payload, opts).await
}
async fn put_multipart_opts(
&self,
location: &object_store::path::Path,
opts: object_store::PutMultipartOptions,
) -> object_store::Result<Box<dyn object_store::MultipartUpload>> {
self.inner.put_multipart_opts(location, opts).await
}
async fn get_opts(
&self,
location: &object_store::path::Path,
options: object_store::GetOptions,
) -> object_store::Result<object_store::GetResult> {
self.inner.get_opts(location, options).await
}
fn delete_stream(
&self,
locations: futures_util::stream::BoxStream<
'static,
object_store::Result<object_store::path::Path>,
>,
) -> futures_util::stream::BoxStream<'static, object_store::Result<object_store::path::Path>>
{
self.inner.delete_stream(locations)
}
fn list(
&self,
prefix: Option<&object_store::path::Path>,
) -> futures_util::stream::BoxStream<'static, object_store::Result<object_store::ObjectMeta>>
{
self.inner.list(prefix)
}
async fn list_with_delimiter(
&self,
prefix: Option<&object_store::path::Path>,
) -> object_store::Result<object_store::ListResult> {
if self.when == RivalOn::SnapshotRead
&& prefix.is_some_and(|p| p.as_ref().ends_with("snapshots"))
{
self.fire().await?;
}
self.inner.list_with_delimiter(prefix).await
}
async fn copy_opts(
&self,
from: &object_store::path::Path,
to: &object_store::path::Path,
options: object_store::CopyOptions,
) -> object_store::Result<()> {
self.inner.copy_opts(from, to, options).await
}
}
fn racing_store(when: RivalOn, rival_epoch: u64, rival_owner: &str) -> Arc<dyn ObjectStore> {
Arc::new(RacingStore {
inner: Arc::new(InMemory::new()),
claim_key: object_store::path::Path::from("wl/acct/latest.owner-claim"),
when,
rival_epoch,
rival_owner: rival_owner.to_string(),
fired: std::sync::atomic::AtomicBool::new(false),
})
}
#[tokio::test]
async fn a_nested_subject_gets_its_directory_created() {
let store: Arc<dyn ObjectStore> = Arc::new(InMemory::new());
seed_snapshot(&store, "wl/acct/db/accounts.db", b"NESTED").await;
let dir = tempdir("nested");
let subjects = vec!["db/accounts.db".to_string()];
let out = hydrate(request(&store, &dir, &subjects, "n1")).await.unwrap();
assert!(matches!(out, HydrateOutcome::Hydrated { .. }), "{out:?}");
assert_eq!(
std::fs::read(dir.join("db/accounts.db")).unwrap(),
db_image(b"NESTED")
);
}
}