1use 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
17pub struct MeteredStorageHandler {
21 inner: Arc<dyn StorageHandler>,
22}
23
24impl MeteredStorageHandler {
25 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 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 #[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 #[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 #[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 #[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}