use std::collections::{BTreeMap, BTreeSet};
use chrono::{DateTime, Utc};
use futures::TryStreamExt;
use object_store::path::Path as StorePath;
use object_store::{ObjectMeta, ObjectStoreExt};
use slatedb::{Checkpoint, VersionedManifest};
use super::MirrorError;
use crate::services::DatabaseHandle;
pub const MANIFEST_DIR: &str = "manifest";
pub const WAL_DIR: &str = "wal";
pub const COMPACTED_DIR: &str = "compacted";
pub fn manifest_rel(id: u64) -> String {
format!("{MANIFEST_DIR}/{id:020}.manifest")
}
pub fn wal_rel(id: u64) -> String {
format!("{WAL_DIR}/{id:020}.sst")
}
pub fn compacted_rel(ulid: &str) -> String {
format!("{COMPACTED_DIR}/{ulid}.sst")
}
pub fn object_path(db: &DatabaseHandle, rel: &str) -> StorePath {
StorePath::from(format!("{}/{}", db.path, rel))
}
const HEAD_PARALLELISM: usize = 16;
pub(crate) async fn head_misses<T, F>(
db: &DatabaseHandle,
items: impl IntoIterator<Item = T>,
rel: F,
) -> Result<Vec<T>, MirrorError>
where
F: Fn(&T) -> String,
{
use futures::StreamExt;
let misses: Vec<Option<T>> = futures::stream::iter(items)
.map(|item| {
let path = object_path(db, &rel(&item));
async move {
match db.store.head(&path).await {
Ok(_) => Ok(None),
Err(object_store::Error::NotFound { .. }) => Ok(Some(item)),
Err(e) => Err(MirrorError::from(e)),
}
}
})
.buffer_unordered(HEAD_PARALLELISM)
.try_collect()
.await?;
Ok(misses.into_iter().flatten().collect())
}
pub fn parse_manifest_name(name: &str) -> Option<u64> {
name.strip_suffix(".manifest")?.parse().ok()
}
pub fn parse_wal_name(name: &str) -> Option<u64> {
name.strip_suffix(".sst")?.parse().ok()
}
pub fn parse_compacted_name(name: &str) -> Option<String> {
let ulid = name.strip_suffix(".sst")?;
(ulid.len() == 26 && ulid.chars().all(|c| c.is_ascii_alphanumeric())).then(|| ulid.to_string())
}
async fn list_dir(db: &DatabaseHandle, dir: &str) -> Result<Vec<ObjectMeta>, MirrorError> {
let prefix = StorePath::from(format!("{}/{dir}", db.path));
match db.store.list(Some(&prefix)).try_collect().await {
Ok(metas) => Ok(metas),
Err(object_store::Error::NotFound { .. }) => Ok(Vec::new()),
Err(e) => Err(e.into()),
}
}
pub async fn list_manifests(db: &DatabaseHandle) -> Result<Vec<(u64, ObjectMeta)>, MirrorError> {
let mut out: Vec<(u64, ObjectMeta)> = list_dir(db, MANIFEST_DIR)
.await?
.into_iter()
.filter_map(|meta| {
let id = meta.location.filename().and_then(parse_manifest_name)?;
Some((id, meta))
})
.collect();
out.sort_by_key(|(id, _)| *id);
Ok(out)
}
pub async fn max_manifest_id(db: &DatabaseHandle) -> Result<Option<u64>, MirrorError> {
Ok(list_manifests(db).await?.last().map(|(id, _)| *id))
}
pub async fn list_wals(db: &DatabaseHandle) -> Result<Vec<(u64, ObjectMeta)>, MirrorError> {
let mut out: Vec<(u64, ObjectMeta)> = list_dir(db, WAL_DIR)
.await?
.into_iter()
.filter_map(|meta| {
let id = meta.location.filename().and_then(parse_wal_name)?;
Some((id, meta))
})
.collect();
out.sort_by_key(|(id, _)| *id);
Ok(out)
}
pub async fn list_compacted(db: &DatabaseHandle) -> Result<Vec<(String, ObjectMeta)>, MirrorError> {
Ok(list_dir(db, COMPACTED_DIR)
.await?
.into_iter()
.filter_map(|meta| {
let ulid = meta.location.filename().and_then(parse_compacted_name)?;
Some((ulid, meta))
})
.collect())
}
#[derive(Clone, Debug, Default, PartialEq)]
pub struct ManifestObjects {
pub compacted: BTreeSet<String>,
pub wal: BTreeSet<u64>,
}
impl ManifestObjects {
pub fn difference(&self, other: &ManifestObjects) -> ManifestObjects {
ManifestObjects {
compacted: self
.compacted
.difference(&other.compacted)
.cloned()
.collect(),
wal: self.wal.difference(&other.wal).copied().collect(),
}
}
pub fn extend(&mut self, other: &ManifestObjects) {
self.compacted.extend(other.compacted.iter().cloned());
self.wal.extend(other.wal.iter().copied());
}
pub fn is_empty(&self) -> bool {
self.compacted.is_empty() && self.wal.is_empty()
}
pub fn len(&self) -> usize {
self.compacted.len() + self.wal.len()
}
pub fn rel_names(&self) -> Vec<String> {
let mut names: Vec<String> = self.wal.iter().map(|&id| wal_rel(id)).collect();
names.extend(self.compacted.iter().map(|u| compacted_rel(u)));
names
}
}
pub fn manifest_objects(m: &VersionedManifest) -> ManifestObjects {
use slatedb::manifest::{SsTableId, SsTableView};
let mut objects = ManifestObjects::default();
let mut add = |view: &SsTableView| match &view.sst.id {
SsTableId::Compacted(ulid) => {
objects.compacted.insert(ulid.to_string());
}
SsTableId::Wal(id) => {
objects.wal.insert(*id);
}
};
for (l0, runs) in std::iter::once((m.l0(), m.compacted()))
.chain(m.segments().iter().map(|s| (s.l0(), s.compacted())))
{
l0.iter().for_each(&mut add);
runs.iter()
.flat_map(|run| run.sst_views.iter())
.for_each(&mut add);
}
for wal_id in (m.replay_after_wal_id() + 1)..m.next_wal_sst_id() {
objects.wal.insert(wal_id);
}
objects
}
pub fn checkpoint_live(cp: &Checkpoint, now: DateTime<Utc>) -> bool {
cp.expire_time.is_none_or(|t| t > now)
}
pub fn check_source(url: &str, m: &VersionedManifest) -> Result<(), MirrorError> {
if !m.external_dbs().is_empty() {
return Err(MirrorError::ExcludedSource {
url: url.to_string(),
reason: "it is a clone (external_dbs is set)".to_string(),
});
}
if m.wal_object_store_uri().is_some() {
return Err(MirrorError::ExcludedSource {
url: url.to_string(),
reason: "it uses a separate WAL object store".to_string(),
});
}
Ok(())
}
#[derive(Default)]
pub struct ManifestCache {
by_id: BTreeMap<u64, VersionedManifest>,
}
impl ManifestCache {
pub async fn read(
&mut self,
db: &DatabaseHandle,
id: u64,
) -> Result<Option<&VersionedManifest>, MirrorError> {
match self.by_id.entry(id) {
std::collections::btree_map::Entry::Occupied(entry) => Ok(Some(entry.into_mut())),
std::collections::btree_map::Entry::Vacant(entry) => {
match db.admin.read_manifest(Some(id)).await? {
Some(manifest) => Ok(Some(entry.insert(manifest))),
None => Ok(None),
}
}
}
}
}
pub struct Closure {
pub head: u64,
pub manifests: BTreeMap<u64, ManifestObjects>,
}
impl Closure {
pub fn objects(&self) -> ManifestObjects {
let mut all = ManifestObjects::default();
for objects in self.manifests.values() {
all.extend(objects);
}
all
}
}
pub async fn closure(
db: &DatabaseHandle,
head: &VersionedManifest,
now: DateTime<Utc>,
cache: &mut ManifestCache,
) -> Result<Option<Closure>, MirrorError> {
let mut manifests = BTreeMap::new();
manifests.insert(head.id(), manifest_objects(head));
for cp in head.checkpoints() {
if !checkpoint_live(cp, now) || cp.manifest_id == head.id() {
continue;
}
match cache.read(db, cp.manifest_id).await? {
Some(pinned) => {
manifests.insert(cp.manifest_id, manifest_objects(pinned));
}
None => return Ok(None),
}
}
Ok(Some(Closure {
head: head.id(),
manifests,
}))
}