use cbor2::{from_reader, to_writer};
use futures::{StreamExt, TryStreamExt, stream::BoxStream};
use moka::{future::Cache, ops::compute::Op};
use object_store::{path::Path, *};
use rand::RngExt;
use serde::{Serialize, de::DeserializeOwned};
use std::{collections::HashMap, sync::Arc};
use crate::map_arc_error;
type MetadataValidator<M> = dyn Fn(&Path, &M) -> Result<()> + Send + Sync;
pub(crate) struct ListingMetaPolicy<M> {
reject_corrupt: bool,
validator: Option<Arc<MetadataValidator<M>>>,
}
impl<M> Clone for ListingMetaPolicy<M> {
fn clone(&self) -> Self {
Self {
reject_corrupt: self.reject_corrupt,
validator: self.validator.clone(),
}
}
}
impl<M> ListingMetaPolicy<M> {
pub(crate) fn unchecked() -> Self {
Self {
reject_corrupt: false,
validator: None,
}
}
pub(crate) fn verified(
reject_corrupt: bool,
validator: impl Fn(&Path, &M) -> Result<()> + Send + Sync + 'static,
) -> Self {
Self {
reject_corrupt,
validator: Some(Arc::new(validator)),
}
}
}
pub(crate) trait SidecarMeta: Serialize + DeserializeOwned + Send + Sync + 'static {
const STORE_NAME: &'static str;
fn e_tag(&self) -> Option<&str>;
fn size(&self) -> u64;
fn generation(&self) -> Option<&str>;
}
pub(crate) fn new_generation() -> String {
let ms = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_millis() as u64)
.unwrap_or(0);
let salt: u32 = rand::rng().random();
format!("{ms:016x}-{salt:08x}")
}
fn generation_timestamp_ms(generation: &str) -> Option<u64> {
let (ts, salt) = generation.split_once('-')?;
if ts.len() != 16 || salt.len() != 8 {
return None;
}
u64::from_str_radix(ts, 16).ok()
}
enum PayloadRef {
Generation(String),
Legacy,
Unknown,
}
pub(crate) struct SidecarStore<T: ObjectStore, M: SidecarMeta> {
pub(crate) store: T,
data_prefix: Path,
gen_prefix: Path,
meta_prefix: Path,
meta_cache: Cache<Path, Arc<M>>,
}
impl<T: ObjectStore, M: SidecarMeta> SidecarStore<T, M> {
pub(crate) fn new(store: T, meta_cache: Cache<Path, Arc<M>>) -> Self {
SidecarStore {
store,
data_prefix: Path::from("data"),
gen_prefix: Path::from("gen"),
meta_prefix: Path::from("meta"),
meta_cache,
}
}
pub(crate) fn meta_path(&self, location: &Path) -> Path {
self.meta_prefix.parts().chain(location.parts()).collect()
}
pub(crate) fn generation_path(&self, location: &Path, generation: &str) -> Path {
self.gen_prefix
.parts()
.chain(location.parts())
.chain(Path::from(generation).parts())
.collect()
}
pub(crate) fn legacy_path(&self, location: &Path) -> Path {
self.data_prefix.parts().chain(location.parts()).collect()
}
pub(crate) fn payload_path(&self, location: &Path, generation: Option<&str>) -> Path {
match generation {
Some(generation) => self.generation_path(location, generation),
None => self.legacy_path(location),
}
}
fn strip_meta_prefix(&self, path: Path) -> Path {
if let Some(suffix) = path.prefix_match(&self.meta_prefix) {
return suffix.collect();
}
path
}
fn split_generation(&self, path: &Path) -> Option<(Path, String)> {
let mut parts: Vec<_> = path.prefix_match(&self.gen_prefix)?.collect();
if parts.len() < 2 {
return None;
}
let generation = parts.pop()?.as_ref().to_string();
Some((parts.into_iter().collect(), generation))
}
async fn fetch_meta_bytes(&self, location: &Path) -> Result<bytes::Bytes> {
let meta_path = self.meta_path(location);
let data = self.store.get(&meta_path).await.map_err(|err| match err {
Error::NotFound { source, .. } => Error::NotFound {
path: location.to_string(),
source,
},
err => err,
})?;
data.bytes().await
}
fn decode_meta(&self, location: &Path, data: &[u8]) -> Result<M> {
from_reader(data).map_err(|err| Error::Generic {
store: M::STORE_NAME,
source: format!("Failed to deserialize Metadata for path {location}: {err:?}").into(),
})
}
async fn load_meta(&self, location: &Path) -> Result<M> {
let data = self.fetch_meta_bytes(location).await?;
self.decode_meta(location, &data)
}
pub(crate) async fn get_meta(&self, location: &Path) -> Result<Arc<M>> {
let meta = self
.meta_cache
.try_get_with(location.clone(), async {
let meta = self.load_meta(location).await?;
Ok(Arc::new(meta))
})
.await
.map_err(|err| map_arc_error(M::STORE_NAME, err))?;
Ok(meta)
}
pub(crate) async fn refresh_meta(&self, location: &Path) -> Result<Arc<M>> {
let rt = self
.meta_cache
.entry(location.clone())
.and_try_compute_with(|_| async {
let meta = self.load_meta(location).await?;
Ok::<_, Error>(Op::Put(Arc::new(meta)))
})
.await?;
Ok(rt.unwrap().value().clone())
}
pub(crate) async fn update_meta_with<F>(
&self,
location: &Path,
create: bool,
f: F,
) -> Result<Arc<M>>
where
F: AsyncFnOnce(Option<&M>) -> Result<M>,
{
let already_exists = || Error::AlreadyExists {
path: location.to_string(),
source: "object already exists".into(),
};
let mut replaced: Option<Path> = None;
let replaced_out = &mut replaced;
let mut f = Some(f);
let rt = self
.meta_cache
.entry(location.clone())
.and_try_compute_with(|_entry| async move {
let f = f.take().expect("update_meta_with closure invoked twice");
let mut meta_mode = PutMode::Overwrite;
let val = match self.fetch_meta_bytes(location).await {
Ok(data) => match self.decode_meta(location, &data) {
Ok(cur) => {
if create {
return Err(already_exists());
}
*replaced_out = Some(self.payload_path(location, cur.generation()));
f(Some(&cur)).await?
}
Err(err) => {
log::warn!(
"{}: replacing corrupted metadata for {location}: {err}",
M::STORE_NAME
);
f(None).await?
}
},
Err(Error::NotFound { .. }) => {
if create {
meta_mode = PutMode::Create;
}
f(None).await?
}
Err(err) => return Err(err),
};
let meta_path = self.meta_path(location);
let mut data = Vec::new();
to_writer(&val, &mut data).map_err(|err| Error::Generic {
store: M::STORE_NAME,
source: format!("Failed to serialize Metadata for path {location}: {err:?}")
.into(),
})?;
self.store
.put_opts(
&meta_path,
data.into(),
PutOptions {
mode: meta_mode,
..Default::default()
},
)
.await
.map_err(|err| match err {
Error::AlreadyExists { source, .. } => Error::AlreadyExists {
path: location.to_string(),
source,
},
err => err,
})?;
Ok::<_, Error>(Op::Put(Arc::new(val)))
})
.await?;
let rt = rt.unwrap().value().clone();
if let Some(old) = replaced
&& old != self.payload_path(location, rt.generation())
{
self.best_effort_delete(&old).await;
}
Ok(rt)
}
async fn best_effort_delete(&self, path: &Path) {
match self.store.delete(path).await {
Ok(()) | Err(Error::NotFound { .. }) => {}
Err(err) => log::warn!(
"{}: failed to delete replaced payload {path}: {err}",
M::STORE_NAME
),
}
}
pub(crate) async fn delete_object(&self, location: &Path) -> Result<()> {
let mut payload: Option<Path> = None;
let payload_out = &mut payload;
self.meta_cache
.entry(location.clone())
.and_try_compute_with(|_entry| async move {
match self.fetch_meta_bytes(location).await {
Ok(data) => match self.decode_meta(location, &data) {
Ok(cur) => {
*payload_out = Some(self.payload_path(location, cur.generation()));
}
Err(err) => {
log::warn!(
"{}: deleting object with corrupted metadata at {location}: {err}",
M::STORE_NAME
);
}
},
Err(Error::NotFound { source, .. }) => {
return Err(Error::NotFound {
path: location.to_string(),
source,
});
}
Err(err) => return Err(err),
}
match self.store.delete(&self.meta_path(location)).await {
Ok(()) | Err(Error::NotFound { .. }) => {}
Err(err) => return Err(err),
}
Ok::<_, Error>(Op::Remove)
})
.await?;
if let Some(path) = payload {
self.best_effort_delete(&path).await;
}
Ok(())
}
pub(crate) fn delete_stream(
self: Arc<Self>,
locations: BoxStream<'static, Result<Path>>,
) -> BoxStream<'static, Result<Path>> {
let inner = self;
locations
.map(move |location| {
let inner = inner.clone();
async move {
let location = location?;
inner.delete_object(&location).await?;
Ok(location)
}
})
.buffered(10)
.boxed()
}
pub(crate) fn list(
self: Arc<Self>,
prefix: Option<&Path>,
policy: ListingMetaPolicy<M>,
) -> BoxStream<'static, Result<ObjectMeta>> {
let prefix = self.meta_path(prefix.unwrap_or(&Path::default()));
let stream = self.store.list(Some(&prefix));
self.decorate_listing(stream, policy)
}
pub(crate) fn list_with_offset(
self: Arc<Self>,
prefix: Option<&Path>,
offset: &Path,
policy: ListingMetaPolicy<M>,
) -> BoxStream<'static, Result<ObjectMeta>> {
let offset = self.meta_path(offset);
let prefix = self.meta_path(prefix.unwrap_or(&Path::default()));
let stream = self.store.list_with_offset(Some(&prefix), &offset);
self.decorate_listing(stream, policy)
}
fn decorate_listing(
self: Arc<Self>,
stream: BoxStream<'static, Result<ObjectMeta>>,
policy: ListingMetaPolicy<M>,
) -> BoxStream<'static, Result<ObjectMeta>> {
let inner = self;
stream
.map_ok(move |obj| {
let store = inner.clone();
let policy = policy.clone();
async move { store.listing_entry(obj, &policy).await }
})
.try_buffered(8) .try_filter_map(|entry| async move { Ok(entry) })
.boxed()
}
async fn listing_entry(
&self,
obj: ObjectMeta,
policy: &ListingMetaPolicy<M>,
) -> Result<Option<ObjectMeta>> {
let location = self.strip_meta_prefix(obj.location);
let meta: Arc<M> = if let Some(meta) = self.meta_cache.get(&location).await {
meta
} else {
match self.fetch_meta_bytes(&location).await {
Ok(data) => match self.decode_meta(&location, &data) {
Ok(meta) => Arc::new(meta),
Err(err) => {
if policy.reject_corrupt {
return Err(err);
}
log::warn!(
"{}: skipping object with corrupted metadata in listing: {location}: {err}",
M::STORE_NAME
);
return Ok(None);
}
},
Err(Error::NotFound { .. }) => return Ok(None),
Err(err) => return Err(err),
}
};
if let Some(validator) = &policy.validator {
validator(&location, &meta)?;
}
Ok(Some(ObjectMeta {
location,
last_modified: obj.last_modified,
size: meta.size(),
e_tag: meta.e_tag().map(String::from),
version: None,
}))
}
pub(crate) async fn list_with_delimiter(
&self,
prefix: Option<&Path>,
policy: ListingMetaPolicy<M>,
) -> Result<ListResult> {
let prefix = self.meta_path(prefix.unwrap_or(&Path::default()));
let rt = self.store.list_with_delimiter(Some(&prefix)).await?;
let common_prefixes = rt
.common_prefixes
.into_iter()
.map(|p| self.strip_meta_prefix(p))
.collect::<Vec<_>>();
let mut indexed =
futures::stream::iter(rt.objects.into_iter().enumerate().map(move |(idx, obj)| {
let policy = policy.clone();
async move {
let entry = self.listing_entry(obj, &policy).await?;
Ok::<_, Error>((idx, entry))
}
}))
.buffer_unordered(8)
.try_collect::<Vec<_>>()
.await?;
indexed.sort_by_key(|(idx, _)| *idx);
let objects = indexed.into_iter().filter_map(|(_, entry)| entry).collect();
Ok(ListResult {
common_prefixes,
objects,
extensions: Extensions::default(),
})
}
pub(crate) async fn copy_payload<F>(
&self,
from: &Path,
to: &Path,
verify: F,
) -> Result<(Arc<M>, String)>
where
F: Fn(&Path, &M) -> Result<()>,
{
let mut retried = false;
loop {
let src = self.get_meta(from).await?;
verify(from, &src)?;
let src_path = self.payload_path(from, src.generation());
let generation = new_generation();
let dst_path = self.generation_path(to, &generation);
match self.store.copy(&src_path, &dst_path).await {
Ok(()) => return Ok((src, generation)),
Err(Error::NotFound { source, .. }) => {
if !retried && src.generation().is_some() {
retried = true;
self.refresh_meta(from).await?;
continue;
}
return Err(Error::NotFound {
path: from.to_string(),
source,
});
}
Err(err) => return Err(err),
}
}
}
pub(crate) async fn check_self_rename(
&self,
location: &Path,
options: &RenameOptions,
) -> Result<()> {
self.get_meta(location).await?;
match options.target_mode {
RenameTargetMode::Overwrite => Ok(()),
RenameTargetMode::Create => Err(Error::AlreadyExists {
path: location.to_string(),
source: "rename target already exists".into(),
}),
}
}
pub(crate) async fn collect_garbage(&self) -> Result<usize> {
let floor_ms = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_millis() as u64)
.unwrap_or(0);
let mut referenced: HashMap<Path, PayloadRef> = HashMap::new();
let mut metas = self.store.list(Some(&self.meta_prefix));
while let Some(obj) = metas.try_next().await? {
let location = self.strip_meta_prefix(obj.location);
let state = match self.fetch_meta_bytes(&location).await {
Ok(data) => match self.decode_meta(&location, &data) {
Ok(meta) => meta
.generation()
.map(|g| PayloadRef::Generation(g.to_string()))
.unwrap_or(PayloadRef::Legacy),
Err(_) => PayloadRef::Unknown,
},
Err(Error::NotFound { .. }) => continue, Err(err) => return Err(err),
};
referenced.insert(location, state);
}
let mut candidates: Vec<(Path, Path, Option<String>)> = Vec::new();
let mut gens = self.store.list(Some(&self.gen_prefix));
while let Some(obj) = gens.try_next().await? {
let parsed = self
.split_generation(&obj.location)
.and_then(|(loc, g)| generation_timestamp_ms(&g).map(|ts| (loc, g, ts)));
let Some((location, generation, ts)) = parsed else {
log::warn!(
"{}: skipping unrecognized object under generation prefix: {}",
M::STORE_NAME,
obj.location
);
continue;
};
if ts >= floor_ms {
continue;
}
match referenced.get(&location) {
Some(PayloadRef::Generation(g)) if *g == generation => continue,
Some(PayloadRef::Unknown) => continue,
_ => candidates.push((obj.location, location, Some(generation))),
}
}
let mut legacy = self.store.list(Some(&self.data_prefix));
while let Some(obj) = legacy.try_next().await? {
let location = match obj.location.prefix_match(&self.data_prefix) {
Some(suffix) => suffix.collect::<Path>(),
None => continue,
};
match referenced.get(&location) {
Some(PayloadRef::Legacy) | Some(PayloadRef::Unknown) => continue,
_ => candidates.push((obj.location, location, None)),
}
}
let mut deleted = 0usize;
for (full_path, location, generation) in candidates {
if self.is_referenced(&location, generation.as_deref()).await? {
continue;
}
match self.store.delete(&full_path).await {
Ok(()) => deleted += 1,
Err(Error::NotFound { .. }) => {}
Err(err) => return Err(err),
}
}
Ok(deleted)
}
async fn is_referenced(&self, location: &Path, generation: Option<&str>) -> Result<bool> {
match self.fetch_meta_bytes(location).await {
Ok(data) => match self.decode_meta(location, &data) {
Ok(meta) => Ok(meta.generation() == generation),
Err(_) => Ok(true),
},
Err(Error::NotFound { .. }) => Ok(false),
Err(err) => Err(err),
}
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn generation_ids_are_unique_and_carry_timestamps() {
let a = new_generation();
let b = new_generation();
assert_ne!(a, b);
let ts = generation_timestamp_ms(&a).unwrap();
let now = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_millis() as u64;
assert!(
ts <= now && ts + 60_000 > now,
"timestamp {ts} vs now {now}"
);
assert_eq!(generation_timestamp_ms("not-a-generation"), None);
assert_eq!(generation_timestamp_ms("0123"), None);
assert_eq!(generation_timestamp_ms("zzzzzzzzzzzzzzzz-00000000"), None);
}
}