micromegas_telemetry/
blob_storage.rs1use anyhow::Result;
2use futures::StreamExt;
3use futures::stream;
4use object_store::prefix::PrefixStore;
5use object_store::{ObjectStore, ObjectStoreExt, path::Path};
6use std::sync::Arc;
7
8pub fn parse_object_store_url(uri: &str) -> Result<(Arc<dyn ObjectStore>, Path)> {
12 parse_object_store_url_parsed(&url::Url::parse(uri)?)
13}
14
15pub fn parse_object_store_url_parsed(url: &url::Url) -> Result<(Arc<dyn ObjectStore>, Path)> {
18 let (store, prefix) =
19 object_store::parse_url_opts(url, std::env::vars().map(|(k, v)| (k.to_lowercase(), v)))?;
20 Ok((Arc::new(store), prefix))
21}
22
23#[derive(Debug)]
28pub struct BlobStorage {
29 blob_store: Arc<dyn ObjectStore>,
30}
31
32impl BlobStorage {
33 pub fn new(blob_store: Arc<dyn ObjectStore>, blob_store_root: Path) -> Self {
35 Self {
36 blob_store: Arc::new(PrefixStore::new(blob_store, blob_store_root)),
37 }
38 }
39
40 pub fn connect(object_store_url: &str) -> Result<Self> {
42 Self::connect_with_layer(object_store_url, |s| s)
43 }
44
45 pub fn parse_url_opts(object_store_url: &str) -> Result<(Arc<dyn ObjectStore>, Path)> {
52 parse_object_store_url(object_store_url)
53 }
54
55 pub fn connect_with_layer(
59 object_store_url: &str,
60 layer: impl FnOnce(Arc<dyn ObjectStore>) -> Arc<dyn ObjectStore>,
61 ) -> Result<Self> {
62 let (blob_store, blob_store_root) = Self::parse_url_opts(object_store_url)?;
63 let layered = layer(blob_store);
64 Ok(Self {
65 blob_store: Arc::new(PrefixStore::new(layered, blob_store_root)),
66 })
67 }
68
69 pub fn inner(&self) -> Arc<dyn ObjectStore> {
71 self.blob_store.clone()
72 }
73
74 pub async fn put(&self, obj_path: &str, buffer: bytes::Bytes) -> Result<()> {
76 self.blob_store
77 .put(&Path::from(obj_path), buffer.into())
78 .await?;
79 Ok(())
80 }
81
82 pub async fn read_blob(&self, obj_path: &str) -> Result<bytes::Bytes> {
84 let get_result = self.blob_store.get(&Path::from(obj_path)).await?;
85 Ok(get_result.bytes().await?)
86 }
87
88 pub async fn delete(&self, obj_path: &str) -> Result<()> {
90 self.blob_store.delete(&Path::from(obj_path)).await?;
91 Ok(())
92 }
93
94 pub async fn probe(&self) -> anyhow::Result<()> {
97 match self.blob_store.list(None).next().await {
98 Some(Ok(_)) | None => Ok(()),
99 Some(Err(e)) => Err(e.into()),
100 }
101 }
102
103 pub async fn delete_batch(&self, objects: &[String]) -> Result<()> {
105 let paths: Vec<_> = objects
106 .iter()
107 .map(|obj_path| Ok(Path::from(obj_path.as_str())))
108 .collect();
109 let path_stream = stream::iter(paths);
110 let mut stream = self.blob_store.delete_stream(Box::pin(path_stream));
111 while let Some(res) = stream.next().await {
112 if let Err(e) = res {
113 match e {
114 object_store::Error::NotFound { path: _, source: _ } => Ok(()),
115 ref _other_error => Err(e),
116 }?
117 }
118 }
119 Ok(())
120 }
121}