use std::sync::Arc;
use futures_util::stream::BoxStream;
use taquba::object_store::{
Error as ObjectStoreError, ObjectMeta, ObjectStore, ObjectStoreExt, path::Path,
};
use crate::error::Result;
#[derive(Clone)]
pub(crate) struct ObjectPrefix {
store: Arc<dyn ObjectStore>,
prefix: String,
}
impl ObjectPrefix {
pub(crate) fn new(store: Arc<dyn ObjectStore>, prefix: impl Into<String>) -> Self {
Self {
store,
prefix: prefix.into(),
}
}
pub(crate) fn prefix(&self) -> &str {
&self.prefix
}
pub(crate) fn path(&self, suffix: &str) -> Path {
Path::from(format!("{}/{suffix}", self.prefix))
}
pub(crate) async fn get(&self, path: &Path) -> Result<Option<Vec<u8>>> {
match self.store.get(path).await {
Ok(result) => Ok(Some(result.bytes().await?.to_vec())),
Err(ObjectStoreError::NotFound { .. }) => Ok(None),
Err(err) => Err(err.into()),
}
}
pub(crate) async fn put(&self, path: &Path, value: &[u8]) -> Result<()> {
self.store.put(path, value.to_vec().into()).await?;
Ok(())
}
pub(crate) async fn delete(&self, path: &Path) -> Result<bool> {
match self.store.delete(path).await {
Ok(()) => Ok(true),
Err(ObjectStoreError::NotFound { .. }) => Ok(false),
Err(err) => Err(err.into()),
}
}
pub(crate) fn list(
&self,
prefix: &Path,
) -> BoxStream<'static, std::result::Result<ObjectMeta, ObjectStoreError>> {
self.store.list(Some(prefix))
}
}