use std::collections::BTreeMap;
use chrono::{DateTime, Utc};
use futures::StreamExt;
use object_store::{ObjectStoreExt, PutMode, PutOptions};
use slatedb::seq_tracker::FindOption;
use super::MirrorError;
use super::copier::Copier;
use super::layout::{self, ManifestObjects, object_path};
use crate::config::{CopierKind, ResolvedMirrorTarget};
use crate::services::DatabaseHandle;
#[derive(Clone, Debug)]
pub enum RestorePoint {
Latest,
Manifest(u64),
Time(DateTime<Utc>),
}
impl RestorePoint {
pub fn parse(at: &str) -> Result<Self, String> {
if let Ok(id) = at.parse::<u64>() {
return Ok(RestorePoint::Manifest(id));
}
DateTime::parse_from_rfc3339(at)
.map(|t| RestorePoint::Time(t.with_timezone(&Utc)))
.map_err(|_| format!("--at {at:?} is neither a manifest id nor an RFC 3339 timestamp"))
}
}
#[derive(Clone, Copy, Debug, Default)]
pub struct RestoreOutcome {
pub manifest_id: u64,
pub manifests_committed: u64,
pub copied_objects: u64,
pub copied_bytes: u64,
}
pub async fn restore(
backup: &DatabaseHandle,
dest: &DatabaseHandle,
point: RestorePoint,
) -> Result<RestoreOutcome, MirrorError> {
let mut existing = dest.store.list(Some(&dest.path));
if existing.next().await.transpose()?.is_some() {
return Err(MirrorError::DestinationNotEmpty {
url: dest.url.clone(),
});
}
let manifests = layout::list_manifests(backup).await?;
if manifests.is_empty() {
return Err(MirrorError::NotADatabase {
url: backup.url.clone(),
});
}
let chosen = resolve_point(backup, &manifests, &point).await?;
let manifest = backup
.admin
.read_manifest(Some(chosen))
.await?
.ok_or_else(|| MirrorError::NoRestorePoint {
at: format!("{point:?}"),
reason: format!("manifest {chosen} vanished from the backup"),
})?;
let now = Utc::now();
let mut members: BTreeMap<u64, ManifestObjects> = BTreeMap::new();
members.insert(chosen, layout::manifest_objects(&manifest));
for cp in manifest.checkpoints() {
if !layout::checkpoint_live(cp, now) || cp.manifest_id == chosen {
continue;
}
match backup.admin.read_manifest(Some(cp.manifest_id)).await? {
Some(pinned) => {
members.insert(cp.manifest_id, layout::manifest_objects(&pinned));
}
None => {
return Err(MirrorError::NoRestorePoint {
at: format!("{point:?}"),
reason: format!(
"manifest {chosen} is closure support, not a restore point: its live \
checkpoint pins manifest {}, which the backup no longer has",
cp.manifest_id
),
});
}
}
}
let mut objects = ManifestObjects::default();
for member in members.values() {
objects.extend(member);
}
let settings = ResolvedMirrorTarget {
copier: CopierKind::Builtin,
..ResolvedMirrorTarget::default()
};
let copier = Copier::new(&settings, None, backup, dest);
let copied = copier.copy(&objects.rel_names()).await?;
let mut committed = 0;
for &id in members.keys() {
let rel = layout::manifest_rel(id);
let bytes = backup
.store
.get(&object_path(backup, &rel))
.await?
.bytes()
.await?;
dest.store
.put_opts(
&object_path(dest, &rel),
bytes.into(),
PutOptions::from(PutMode::Create),
)
.await?;
committed += 1;
}
Ok(RestoreOutcome {
manifest_id: chosen,
manifests_committed: committed,
copied_objects: copied.objects,
copied_bytes: copied.bytes,
})
}
async fn resolve_point(
backup: &DatabaseHandle,
manifests: &[(u64, object_store::ObjectMeta)],
point: &RestorePoint,
) -> Result<u64, MirrorError> {
let latest = manifests.last().expect("nonempty").0;
match point {
RestorePoint::Latest => Ok(latest),
RestorePoint::Manifest(id) => {
if manifests.iter().any(|(m, _)| m == id) {
Ok(*id)
} else {
Err(MirrorError::NoRestorePoint {
at: id.to_string(),
reason: format!("the backup has no manifest {id}"),
})
}
}
RestorePoint::Time(ts) => {
let head = backup
.admin
.read_manifest(Some(latest))
.await?
.ok_or_else(|| MirrorError::NotADatabase {
url: backup.url.clone(),
})?;
for (id, _) in manifests.iter().rev() {
let Some(m) = backup.admin.read_manifest(Some(*id)).await? else {
continue;
};
let Some(m_ts) = head
.sequence_tracker()
.find_ts(m.last_l0_seq(), FindOption::RoundDown)
else {
continue;
};
if m_ts <= *ts {
return Ok(*id);
}
}
Err(MirrorError::NoRestorePoint {
at: ts.to_rfc3339(),
reason: "the timestamp predates the backup's tracked history".to_string(),
})
}
}
}