Skip to main content

delta_kernel/metrics/
metered_storage.rs

1//! [`MeteredStorageHandler`] wraps any [`StorageHandler`] so it emits the kernel's
2//! standard `"storage"` tracing spans. Usually reached via [`MeteredDeltaEngine`];
3//! construct directly when wrapping a `StorageHandler` outside an `Engine`.
4//!
5//! [`MeteredDeltaEngine`]: crate::metrics::MeteredDeltaEngine
6
7use std::sync::Arc;
8use std::time::Instant;
9
10use bytes::Bytes;
11use url::Url;
12
13use crate::metrics::events::{StorageCopyCompleted, StorageListCompleted, StorageReadCompleted};
14use crate::metrics::{emit_storage_span, MetricsIterator};
15use crate::{CancellationTokenRef, DeltaResult, FileMeta, FileSlice, StorageHandler};
16
17/// Decorator over an engine-provided `Arc<dyn StorageHandler>` that emits the kernel's
18/// standard `"storage"` spans on operations that produce metrics. `put`, `head`, and `delete`
19/// are pass-through and emit nothing.
20pub struct MeteredStorageHandler {
21    inner: Arc<dyn StorageHandler>,
22}
23
24impl MeteredStorageHandler {
25    /// Wrap `inner`. Debug-asserts that `inner` is not already a
26    /// [`MeteredStorageHandler`] so spans are emitted exactly once.
27    pub fn new(inner: Arc<dyn StorageHandler>) -> Self {
28        debug_assert!(
29            !inner.any_ref().is::<MeteredStorageHandler>(),
30            "MeteredStorageHandler wraps another MeteredStorageHandler; \
31             remove the outer wrap to avoid double-counting metrics",
32        );
33        Self { inner }
34    }
35}
36
37impl std::fmt::Debug for MeteredStorageHandler {
38    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
39        f.debug_struct("MeteredStorageHandler")
40            .finish_non_exhaustive()
41    }
42}
43
44impl StorageHandler for MeteredStorageHandler {
45    fn list_from(
46        &self,
47        path: &Url,
48    ) -> DeltaResult<Box<dyn Iterator<Item = DeltaResult<FileMeta>>>> {
49        let start = Instant::now();
50        let inner = self.inner.list_from(path)?;
51        Ok(Box::new(MetricsIterator::<_, FileMeta>::new(
52            inner,
53            StorageListCompleted::NAME,
54            start,
55        )))
56    }
57
58    // Forward the token by identity so a caller can still downcast it to recover their own.
59    fn list_from_with_cancellation(
60        &self,
61        path: &Url,
62        cancellation_token: Option<CancellationTokenRef>,
63    ) -> DeltaResult<Box<dyn Iterator<Item = DeltaResult<FileMeta>>>> {
64        let start = Instant::now();
65        let inner = self
66            .inner
67            .list_from_with_cancellation(path, cancellation_token)?;
68        Ok(Box::new(MetricsIterator::<_, FileMeta>::new(
69            inner,
70            StorageListCompleted::NAME,
71            start,
72        )))
73    }
74
75    fn read_files(
76        &self,
77        files: Vec<FileSlice>,
78    ) -> DeltaResult<Box<dyn Iterator<Item = DeltaResult<Bytes>>>> {
79        let start = Instant::now();
80        let inner = self.inner.read_files(files)?;
81        Ok(Box::new(MetricsIterator::<_, Bytes>::new(
82            inner,
83            StorageReadCompleted::NAME,
84            start,
85        )))
86    }
87
88    fn read_files_with_cancellation(
89        &self,
90        files: Vec<FileSlice>,
91        cancellation_token: Option<CancellationTokenRef>,
92    ) -> DeltaResult<Box<dyn Iterator<Item = DeltaResult<Bytes>>>> {
93        let start = Instant::now();
94        let inner = self
95            .inner
96            .read_files_with_cancellation(files, cancellation_token)?;
97        Ok(Box::new(MetricsIterator::<_, Bytes>::new(
98            inner,
99            StorageReadCompleted::NAME,
100            start,
101        )))
102    }
103
104    fn copy_atomic(&self, src: &Url, dest: &Url) -> DeltaResult<()> {
105        let start = Instant::now();
106        let result = self.inner.copy_atomic(src, dest);
107        emit_storage_span(StorageCopyCompleted::NAME, start.elapsed(), 0, 0);
108        result
109    }
110
111    fn put(&self, path: &Url, data: Bytes, overwrite: bool) -> DeltaResult<()> {
112        self.inner.put(path, data, overwrite)
113    }
114
115    fn head(&self, path: &Url) -> DeltaResult<FileMeta> {
116        self.inner.head(path)
117    }
118
119    fn delete(&self, path: &Url) -> DeltaResult<()> {
120        self.inner.delete(path)
121    }
122}
123
124#[cfg(test)]
125mod tests {
126    use std::sync::Arc;
127
128    use super::*;
129    use crate::metrics::MetricEvent;
130    use crate::unit_test_utils::{install_thread_local_metrics_reporter, CapturingReporter};
131
132    /// Storage handler that returns N preconfigured FileMeta / byte slices.
133    #[derive(Debug)]
134    struct StubStorageHandler {
135        list_results: Vec<FileMeta>,
136        read_results: Vec<Bytes>,
137    }
138
139    impl StorageHandler for StubStorageHandler {
140        fn list_from(
141            &self,
142            _path: &Url,
143        ) -> DeltaResult<Box<dyn Iterator<Item = DeltaResult<FileMeta>>>> {
144            let results: Vec<_> = self.list_results.iter().cloned().map(Ok).collect();
145            Ok(Box::new(results.into_iter()))
146        }
147
148        fn read_files(
149            &self,
150            _files: Vec<FileSlice>,
151        ) -> DeltaResult<Box<dyn Iterator<Item = DeltaResult<Bytes>>>> {
152            let results: Vec<_> = self.read_results.iter().cloned().map(Ok).collect();
153            Ok(Box::new(results.into_iter()))
154        }
155
156        fn copy_atomic(&self, _src: &Url, _dest: &Url) -> DeltaResult<()> {
157            Ok(())
158        }
159
160        fn put(&self, _path: &Url, _data: Bytes, _overwrite: bool) -> DeltaResult<()> {
161            Ok(())
162        }
163
164        fn head(&self, _path: &Url) -> DeltaResult<FileMeta> {
165            unreachable!("not exercised in these tests")
166        }
167
168        fn delete(&self, _path: &Url) -> DeltaResult<()> {
169            Ok(())
170        }
171    }
172
173    fn fake_url() -> Url {
174        Url::parse("memory:///_delta_log/").unwrap()
175    }
176
177    fn fake_file_meta(name: &str) -> FileMeta {
178        FileMeta {
179            location: Url::parse(&format!("memory:///_delta_log/{name}")).unwrap(),
180            last_modified: 0,
181            size: 0,
182        }
183    }
184
185    fn install_capture() -> (Arc<CapturingReporter>, tracing::subscriber::DefaultGuard) {
186        let reporter = Arc::new(CapturingReporter::default());
187        let guard = install_thread_local_metrics_reporter(reporter.clone());
188        (reporter, guard)
189    }
190
191    #[test]
192    fn list_from_emits_storage_list_completed() {
193        let (reporter, _guard) = install_capture();
194        let inner: Arc<dyn StorageHandler> = Arc::new(StubStorageHandler {
195            list_results: vec![
196                fake_file_meta("00000000000000000000.json"),
197                fake_file_meta("00000000000000000001.json"),
198            ],
199            read_results: vec![],
200        });
201        let storage = MeteredStorageHandler::new(inner);
202
203        let iter = storage.list_from(&fake_url()).unwrap();
204        let _: Vec<_> = iter.collect();
205
206        let events = reporter.events();
207        let listed = events
208            .iter()
209            .find(|e| matches!(e, MetricEvent::StorageListCompleted(_)))
210            .expect("expected StorageListCompleted event");
211        let MetricEvent::StorageListCompleted(e) = listed else {
212            unreachable!();
213        };
214        assert_eq!(e.num_files, 2);
215    }
216
217    #[test]
218    fn read_files_emits_storage_read_completed() {
219        let (reporter, _guard) = install_capture();
220        let inner: Arc<dyn StorageHandler> = Arc::new(StubStorageHandler {
221            list_results: vec![],
222            read_results: vec![Bytes::from(vec![0u8; 32]), Bytes::from(vec![0u8; 8])],
223        });
224        let storage = MeteredStorageHandler::new(inner);
225
226        let iter = storage.read_files(vec![]).unwrap();
227        let _: Vec<_> = iter.collect();
228
229        let events = reporter.events();
230        let read = events
231            .iter()
232            .find(|e| matches!(e, MetricEvent::StorageReadCompleted(_)))
233            .expect("expected StorageReadCompleted event");
234        let MetricEvent::StorageReadCompleted(e) = read else {
235            unreachable!();
236        };
237        assert_eq!(e.num_files, 2);
238        assert_eq!(e.bytes_read, 40);
239    }
240
241    #[test]
242    fn copy_atomic_emits_storage_copy_completed() {
243        let (reporter, _guard) = install_capture();
244        let inner: Arc<dyn StorageHandler> = Arc::new(StubStorageHandler {
245            list_results: vec![],
246            read_results: vec![],
247        });
248        let storage = MeteredStorageHandler::new(inner);
249
250        storage.copy_atomic(&fake_url(), &fake_url()).unwrap();
251
252        let events = reporter.events();
253        assert!(events
254            .iter()
255            .any(|e| matches!(e, MetricEvent::StorageCopyCompleted(_))));
256    }
257
258    #[test]
259    #[should_panic(expected = "wraps another MeteredStorageHandler")]
260    fn new_panics_on_double_wrap() {
261        let inner: Arc<dyn StorageHandler> = Arc::new(StubStorageHandler {
262            list_results: vec![],
263            read_results: vec![],
264        });
265        let once: Arc<dyn StorageHandler> = Arc::new(MeteredStorageHandler::new(inner));
266        let _twice = MeteredStorageHandler::new(once);
267    }
268
269    /// Records the token it is handed, so a test can check what survived the metered wrapper.
270    #[derive(Default)]
271    struct TokenCapturingStorageHandler {
272        seen: std::sync::Mutex<Option<CancellationTokenRef>>,
273    }
274
275    impl StorageHandler for TokenCapturingStorageHandler {
276        fn list_from(
277            &self,
278            _path: &Url,
279        ) -> DeltaResult<Box<dyn Iterator<Item = DeltaResult<FileMeta>>>> {
280            Ok(Box::new(std::iter::empty()))
281        }
282
283        fn list_from_with_cancellation(
284            &self,
285            _path: &Url,
286            cancellation_token: Option<CancellationTokenRef>,
287        ) -> DeltaResult<Box<dyn Iterator<Item = DeltaResult<FileMeta>>>> {
288            *self.seen.lock().unwrap() = cancellation_token;
289            Ok(Box::new(std::iter::empty()))
290        }
291
292        fn read_files(
293            &self,
294            _files: Vec<FileSlice>,
295        ) -> DeltaResult<Box<dyn Iterator<Item = DeltaResult<Bytes>>>> {
296            Ok(Box::new(std::iter::empty()))
297        }
298
299        fn read_files_with_cancellation(
300            &self,
301            _files: Vec<FileSlice>,
302            cancellation_token: Option<CancellationTokenRef>,
303        ) -> DeltaResult<Box<dyn Iterator<Item = DeltaResult<Bytes>>>> {
304            *self.seen.lock().unwrap() = cancellation_token;
305            Ok(Box::new(std::iter::empty()))
306        }
307
308        fn copy_atomic(&self, _src: &Url, _dest: &Url) -> DeltaResult<()> {
309            Ok(())
310        }
311
312        fn put(&self, _path: &Url, _data: Bytes, _overwrite: bool) -> DeltaResult<()> {
313            Ok(())
314        }
315
316        fn head(&self, _path: &Url) -> DeltaResult<FileMeta> {
317            unreachable!("not exercised in these tests")
318        }
319
320        fn delete(&self, _path: &Url) -> DeltaResult<()> {
321            Ok(())
322        }
323    }
324
325    // The wrapper forwards to the inner cancellation-aware methods and passes the token through by
326    // identity, so a caller can still downcast it to recover what it supplied.
327    #[rstest::rstest]
328    #[case::list(true)]
329    #[case::read(false)]
330    fn forwards_cancellation_token_by_identity(#[case] list: bool) {
331        let stub = Arc::new(TokenCapturingStorageHandler::default());
332        let storage = MeteredStorageHandler::new(stub.clone());
333        let token: CancellationTokenRef =
334            Arc::new(crate::unit_test_utils::TestCancellationToken::default());
335
336        if list {
337            let iter = storage
338                .list_from_with_cancellation(&fake_url(), Some(token.clone()))
339                .unwrap();
340            let _: Vec<_> = iter.collect();
341        } else {
342            let iter = storage
343                .read_files_with_cancellation(vec![], Some(token.clone()))
344                .unwrap();
345            let _: Vec<_> = iter.collect();
346        }
347
348        let seen = stub
349            .seen
350            .lock()
351            .unwrap()
352            .clone()
353            .expect("inner handler should have received the token");
354        assert!(
355            Arc::ptr_eq(&token, &seen),
356            "the metered wrapper must not wrap or replace the token"
357        );
358    }
359
360    // Metrics are still emitted when the cancellation-aware variants are used.
361    #[test]
362    fn cancellation_variants_still_emit_metrics() {
363        let (reporter, _guard) = install_capture();
364        let inner: Arc<dyn StorageHandler> = Arc::new(StubStorageHandler {
365            list_results: vec![fake_file_meta("00000000000000000000.json")],
366            read_results: vec![Bytes::from(vec![0u8; 4])],
367        });
368        let storage = MeteredStorageHandler::new(inner);
369
370        let _: Vec<_> = storage
371            .list_from_with_cancellation(&fake_url(), None)
372            .unwrap()
373            .collect();
374        let _: Vec<_> = storage
375            .read_files_with_cancellation(vec![], None)
376            .unwrap()
377            .collect();
378
379        let events = reporter.events();
380        assert!(events
381            .iter()
382            .any(|e| matches!(e, MetricEvent::StorageListCompleted(_))));
383        assert!(events
384            .iter()
385            .any(|e| matches!(e, MetricEvent::StorageReadCompleted(_))));
386    }
387}