micromegas_ingestion/
data_lake_connection.rs1use anyhow::{Context, Result};
2use micromegas_object_cache::CacheClientStore;
3use micromegas_object_cache::prefetch::{ObjectPrefetch, PrefetchItem, PrefixPrefetch};
4use micromegas_telemetry::blob_storage::BlobStorage;
5use micromegas_tracing::prelude::*;
6use object_store::ObjectStore;
7use sqlx::PgPool;
8use std::sync::Arc;
9use tokio::task::JoinHandle;
10
11#[derive(Debug, Clone)]
13pub struct DataLakeConnection {
14 pub db_pool: PgPool,
15 pub blob_storage: Arc<BlobStorage>,
16 prefetch: Option<Arc<dyn ObjectPrefetch>>,
20}
21
22impl DataLakeConnection {
23 pub fn new(db_pool: PgPool, blob_storage: Arc<BlobStorage>) -> Self {
24 Self {
25 db_pool,
26 blob_storage,
27 prefetch: None,
28 }
29 }
30
31 pub fn new_with_prefetch(
34 db_pool: PgPool,
35 blob_storage: Arc<BlobStorage>,
36 prefetch: Option<Arc<dyn ObjectPrefetch>>,
37 ) -> Self {
38 Self {
39 db_pool,
40 blob_storage,
41 prefetch,
42 }
43 }
44
45 pub fn warm_object(&self, key: &str, size: i64) -> Option<JoinHandle<()>> {
58 let prefetch = self.prefetch.as_ref()?.clone();
59 if size <= 0 {
60 return None; }
62 let key = key.to_string(); let item = PrefetchItem {
64 key: key.clone(),
65 size: size as u64,
66 ranges: None,
67 };
68 imetric!("object_warm_requested", "count", 1_u64);
69 Some(spawn_with_context(async move {
70 match prefetch.prefetch(vec![item]).await {
71 Ok(resp) => debug!(
72 "write-time warm enqueued accepted={} rejected={} dropped={}",
73 resp.accepted, resp.rejected, resp.dropped
74 ),
75 Err(e) => debug!("write-time warm failed for {key}: {e}"),
78 }
79 }))
80 }
81}
82
83pub(crate) fn make_cache(
87 direct: Arc<dyn ObjectStore>,
88) -> (Arc<dyn ObjectStore>, Option<Arc<dyn ObjectPrefetch>>) {
89 let cache_url = std::env::var("MICROMEGAS_OBJECT_CACHE_URL").ok();
90 let api_key = std::env::var("MICROMEGAS_OBJECT_CACHE_API_KEY").ok();
91 match cache_url {
92 Some(url) if api_key.is_some() => {
93 let client = Arc::new(CacheClientStore::new(url, api_key, direct));
94 (
95 client.clone() as Arc<dyn ObjectStore>,
96 Some(client as Arc<dyn ObjectPrefetch>),
97 )
98 }
99 Some(url) => {
100 warn!(
102 "MICROMEGAS_OBJECT_CACHE_URL is set ({url}) but MICROMEGAS_OBJECT_CACHE_API_KEY is missing: the object cache is disabled and requests will go directly to the store"
103 );
104 (direct, None)
105 }
106 None => (direct, None),
107 }
108}
109
110pub async fn connect_to_data_lake(
112 db_uri: &str,
113 object_store_url: &str,
114) -> Result<DataLakeConnection> {
115 info!("connecting to blob storage");
116 let (raw_store, root) = BlobStorage::parse_url_opts(object_store_url)
117 .with_context(|| "connecting to blob storage")?;
118 let (layered, prefetch_client) = make_cache(raw_store);
119 let blob_storage = Arc::new(BlobStorage::new(layered, root.clone()));
120 let prefetch =
121 prefetch_client.map(|p| Arc::new(PrefixPrefetch::new(p, root)) as Arc<dyn ObjectPrefetch>);
122 let pool = sqlx::postgres::PgPoolOptions::new()
123 .connect(db_uri)
124 .await
125 .with_context(|| String::from("Connecting to telemetry database"))?;
126 Ok(DataLakeConnection::new_with_prefetch(
127 pool,
128 blob_storage,
129 prefetch,
130 ))
131}