Skip to main content

objectstore_service/backend/
counting.rs

1//! Defines [`CountingBackend`], a decorator for the [`Backend`] trait that emits Cost of Goods Sold
2//! (COGS) compute usage metrics to the `cogs.usage` counter.
3//!
4//! [`CountingBackend`] is meant to wrap the outer-most [`Backend`] implementation owned by
5//! [`StorageService`] so that every tracked backend operation, whether a single-object operation
6//! called by [`StorageService`] or a batched operation streamed by [`StreamExecutor`], is counted
7//! once. Notably, any operation that fails before it gets to [`StorageService`] (e.g. an auth or
8//! rate limit failure at a higher layer) is not counted.
9//!
10//! For COGS purposes we use operation count as a proxy for compute cost under the assumption that
11//! each request we serve has a basically flat CPU cost. Large payloads take longer, but they can be
12//! streamed in the background while other requests are served so they don't really cost more.
13//!
14//! [`StorageService`]: crate::service::StorageService
15//! [`StreamExecutor`]: crate::streaming::StreamExecutor
16
17use 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
38/// Increments `cogs.usage` by one operation for the given `usecase`.
39///
40/// Under the hood, the `usecase` is used as the `app_feature`. This allows to identify distinct
41/// products and map them in the for the COGs pipeline.
42fn count(usecase: &str) {
43    objectstore_metrics::count!("cogs.usage" += 1, app_feature = usecase.to_owned());
44}
45
46/// A [`Backend`] decorator that counts each operation performed for COGS. Also implements
47/// [`MultipartUploadBackend`]. See the [module documentation](self) for how it should be used.
48///
49/// [`CountingBackend`]'s implementation clashes with how the [`MultipartUploadBackend`] trait is
50/// connected to the [`Backend`] trait. The workaround is to give `CountingBackend` (up to) two
51/// `Arc`s that point to the inner backend:
52/// - `inner: Arc<dyn Backend>`
53/// - `inner_multipart: Option<Arc<dyn MultipartUploadBackend>>` if `inner` supports it
54#[derive(Debug)]
55pub struct CountingBackend {
56    inner: Arc<dyn Backend>,
57}
58
59impl CountingBackend {
60    /// Creates a [`CountingBackend`] that wraps `inner` and increments `cogs.usage`
61    /// before delegating operations to it.
62    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    /// Runs `f` on a current-thread runtime while capturing emitted metrics.
268    ///
269    /// The capturing client is thread-local, so the futures must run on the same
270    /// thread that installs it.
271    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            // Nothing stored: the read returns `None` but is still billed.
326            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            // The upload was completed above, so aborting it is a no-op; counting
394            // happens on dispatch regardless, which is what this asserts.
395            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}