objectstore_service/backend/
counting.rs1use std::num::NonZeroU64;
18use std::sync::Arc;
19
20use objectstore_types::metadata::Metadata;
21use objectstore_types::range::ByteRange;
22use objectstore_types::resumable::UploadProgress;
23use objectstore_types::time::Timestamp;
24
25use crate::backend::common::{
26 Backend, DeleteResponse, ExpiryUpdate, GetResponse, MetadataResponse, MultipartUploadBackend,
27 PutResponse, SetExpiryResponse,
28};
29use crate::error::Result;
30use crate::id::ObjectId;
31use crate::multipart::{
32 AbortMultipartResponse, CompleteMultipartResponse, CompletedPart, InitiateMultipartResponse,
33 ListPartsResponse, PartNumber, UploadId, UploadPartResponse,
34};
35use crate::resumable::{BackendToken, Session};
36use crate::stream::ClientStream;
37
38fn count(usecase: &str) {
43 objectstore_metrics::count!("cogs.usage" += 1, app_feature = usecase.to_owned());
44}
45
46#[derive(Debug)]
55pub struct CountingBackend {
56 inner: Arc<dyn Backend>,
57}
58
59impl CountingBackend {
60 pub fn new(inner: Box<dyn Backend>) -> Self {
63 let inner: Arc<dyn Backend> = Arc::from(inner);
64 Self { inner }
65 }
66}
67
68#[async_trait::async_trait]
69impl Backend for CountingBackend {
70 fn name(&self) -> &'static str {
71 self.inner.name()
72 }
73
74 fn upload_granularity(&self) -> u64 {
75 self.inner.upload_granularity()
76 }
77
78 async fn put_object(
79 &self,
80 id: &ObjectId,
81 metadata: &Metadata,
82 stream: ClientStream,
83 access_time: Timestamp,
84 ) -> Result<PutResponse> {
85 count(&id.context.usecase);
86 self.inner
87 .put_object(id, metadata, stream, access_time)
88 .await
89 }
90
91 async fn get_object(
92 &self,
93 id: &ObjectId,
94 access_time: Timestamp,
95 range: Option<ByteRange>,
96 ) -> Result<GetResponse> {
97 count(&id.context.usecase);
98 self.inner.get_object(id, access_time, range).await
99 }
100
101 async fn get_metadata(
102 &self,
103 id: &ObjectId,
104 access_time: Timestamp,
105 ) -> Result<MetadataResponse> {
106 count(&id.context.usecase);
107 self.inner.get_metadata(id, access_time).await
108 }
109
110 async fn set_expiry(
111 &self,
112 id: &ObjectId,
113 target: ExpiryUpdate,
114 access_time: Timestamp,
115 ) -> Result<SetExpiryResponse> {
116 count(&id.context.usecase);
117 self.inner.set_expiry(id, target, access_time).await
118 }
119
120 async fn delete_object(&self, id: &ObjectId, access_time: Timestamp) -> Result<DeleteResponse> {
121 count(&id.context.usecase);
122 self.inner.delete_object(id, access_time).await
123 }
124
125 async fn join(&self) {
126 self.inner.join().await;
127 }
128
129 fn as_multipart_upload_backend(&self) -> Result<&dyn MultipartUploadBackend> {
130 self.inner.as_multipart_upload_backend()?;
131 Ok(self)
132 }
133
134 async fn create_upload_session(
135 &self,
136 id: &ObjectId,
137 metadata: &Metadata,
138 upload_length: NonZeroU64,
139 ) -> Result<Option<BackendToken>> {
140 count(&id.context.usecase);
141 self.inner
142 .create_upload_session(id, metadata, upload_length)
143 .await
144 }
145
146 async fn put_chunk(
147 &self,
148 session: &Session,
149 offset: u64,
150 content_length: u64,
151 stream: ClientStream,
152 ) -> Result<UploadProgress> {
153 count(&session.object_id.context.usecase);
154 self.inner
155 .put_chunk(session, offset, content_length, stream)
156 .await
157 }
158
159 async fn upload_offset(&self, session: &Session) -> Result<UploadProgress> {
160 count(&session.object_id.context.usecase);
161 self.inner.upload_offset(session).await
162 }
163
164 async fn cancel_upload(&self, session: &Session) -> Result<()> {
165 count(&session.object_id.context.usecase);
166 self.inner.cancel_upload(session).await
167 }
168}
169
170#[async_trait::async_trait]
171impl MultipartUploadBackend for CountingBackend {
172 async fn initiate_multipart(
173 &self,
174 id: &ObjectId,
175 metadata: &Metadata,
176 ) -> Result<InitiateMultipartResponse> {
177 count(&id.context.usecase);
178 self.inner
179 .as_multipart_upload_backend()?
180 .initiate_multipart(id, metadata)
181 .await
182 }
183
184 async fn upload_part(
185 &self,
186 id: &ObjectId,
187 upload_id: &UploadId,
188 part_number: PartNumber,
189 content_length: u64,
190 content_md5: Option<&str>,
191 body: ClientStream,
192 ) -> Result<UploadPartResponse> {
193 count(&id.context.usecase);
194 self.inner
195 .as_multipart_upload_backend()?
196 .upload_part(
197 id,
198 upload_id,
199 part_number,
200 content_length,
201 content_md5,
202 body,
203 )
204 .await
205 }
206
207 async fn list_parts(
208 &self,
209 id: &ObjectId,
210 upload_id: &UploadId,
211 max_parts: Option<u32>,
212 part_number_marker: Option<PartNumber>,
213 ) -> Result<ListPartsResponse> {
214 count(&id.context.usecase);
215 self.inner
216 .as_multipart_upload_backend()?
217 .list_parts(id, upload_id, max_parts, part_number_marker)
218 .await
219 }
220
221 async fn abort_multipart(
222 &self,
223 id: &ObjectId,
224 upload_id: &UploadId,
225 ) -> Result<AbortMultipartResponse> {
226 count(&id.context.usecase);
227 self.inner
228 .as_multipart_upload_backend()?
229 .abort_multipart(id, upload_id)
230 .await
231 }
232
233 async fn complete_multipart(
234 &self,
235 id: &ObjectId,
236 upload_id: &UploadId,
237 parts: Vec<CompletedPart>,
238 access_time: Timestamp,
239 ) -> Result<CompleteMultipartResponse> {
240 count(&id.context.usecase);
241 self.inner
242 .as_multipart_upload_backend()?
243 .complete_multipart(id, upload_id, parts, access_time)
244 .await
245 }
246}
247
248#[cfg(test)]
249mod tests {
250 use objectstore_types::scope::{Scope, Scopes};
251
252 use super::*;
253 use crate::backend::in_memory::InMemoryBackend;
254 use crate::id::ObjectContext;
255 use crate::stream;
256
257 fn object_id(usecase: &str) -> ObjectId {
258 ObjectId::new(
259 ObjectContext {
260 usecase: usecase.into(),
261 scopes: Scopes::from_iter([Scope::create("org", "1").unwrap()]),
262 },
263 "key".into(),
264 )
265 }
266
267 fn capture(f: impl std::future::Future<Output = ()>) -> Vec<String> {
272 objectstore_metrics::with_capturing_test_client(|| {
273 tokio::runtime::Builder::new_current_thread()
274 .enable_all()
275 .build()
276 .unwrap()
277 .block_on(f);
278 })
279 }
280
281 #[test]
282 fn counts_each_core_operation_once() {
283 let captured = capture(async {
284 let backend = CountingBackend::new(Box::new(InMemoryBackend::new("in-memory")));
285 let id = object_id("attachments");
286
287 backend
288 .put_object(
289 &id,
290 &Metadata::default(),
291 stream::single("hi"),
292 Timestamp::now(),
293 )
294 .await
295 .unwrap();
296 backend
297 .get_object(&id, Timestamp::now(), None)
298 .await
299 .unwrap();
300 backend.get_metadata(&id, Timestamp::now()).await.unwrap();
301 backend.delete_object(&id, Timestamp::now()).await.unwrap();
302 });
303
304 let cogs = captured
305 .iter()
306 .filter(|m| m.starts_with("cogs.usage"))
307 .count();
308 assert_eq!(
309 cogs, 4,
310 "expected one count per operation, captured: {captured:?}"
311 );
312 assert!(
313 captured
314 .iter()
315 .all(|m| !m.starts_with("cogs.usage")
316 || m == "cogs.usage:+1|c|#app_feature:attachments"),
317 "captured: {captured:?}"
318 );
319 }
320
321 #[test]
322 fn counts_missing_reads_on_dispatch() {
323 let captured = capture(async {
324 let backend = CountingBackend::new(Box::new(InMemoryBackend::new("in-memory")));
325 let result = backend
327 .get_object(&object_id("attachments"), Timestamp::now(), None)
328 .await
329 .unwrap();
330 assert!(result.is_none());
331 });
332
333 assert_eq!(
334 captured
335 .iter()
336 .filter(|m| m.starts_with("cogs.usage"))
337 .count(),
338 1,
339 "captured: {captured:?}"
340 );
341 }
342
343 #[test]
344 fn new_usecase_is_app_feature() {
345 let captured = capture(async {
346 let backend = CountingBackend::new(Box::new(InMemoryBackend::new("in-memory")));
347 backend
348 .get_object(&object_id("new_usecase"), Timestamp::now(), None)
349 .await
350 .unwrap();
351 });
352
353 assert!(
354 captured
355 .iter()
356 .any(|m| m == "cogs.usage:+1|c|#app_feature:new_usecase"),
357 "captured: {captured:?}"
358 );
359 }
360
361 #[test]
362 fn counts_each_multipart_operation() {
363 let captured = capture(async {
364 let backend: Arc<dyn Backend> = Arc::new(CountingBackend::new(Box::new(
365 InMemoryBackend::new("in-memory"),
366 )));
367 let multipart = backend.as_multipart_upload_backend().unwrap();
368 let id = object_id("attachments");
369
370 let upload_id = multipart
371 .initiate_multipart(&id, &Metadata::default())
372 .await
373 .unwrap();
374 multipart
375 .upload_part(
376 &id,
377 &upload_id,
378 PartNumber::new(1).unwrap(),
379 2,
380 None,
381 stream::single("hi"),
382 )
383 .await
384 .unwrap();
385 multipart
386 .list_parts(&id, &upload_id, None, None)
387 .await
388 .unwrap();
389 multipart
390 .complete_multipart(&id, &upload_id, vec![], Timestamp::now())
391 .await
392 .unwrap();
393 let _ = multipart.abort_multipart(&id, &upload_id).await;
396 });
397
398 assert_eq!(
399 captured
400 .iter()
401 .filter(|m| m == &"cogs.usage:+1|c|#app_feature:attachments")
402 .count(),
403 5,
404 "expected one count per multipart operation, captured: {captured:?}"
405 );
406 }
407}