Skip to main content

objectstore_service/change_stream/
cost_tracker.rs

1//! The change stream implementation that reports through an [`InventoryTracker`].
2
3use std::fmt;
4use std::time::{Duration, SystemTime};
5
6use objectstore_inventory_tracker::{BoxError, InventoryTracker, Producer};
7
8use crate::change_stream::{
9    ChangeStream, CostTrackerStreamConfig, SCOPE_ORGANIZATION, SCOPE_PROJECT, scope_id,
10};
11use crate::id::ObjectId;
12
13/// Reports through an [`InventoryTracker`], which hashes each [`ObjectId`] both to
14/// anonymize it and to decide whether it is sampled. See [`objectstore_inventory_tracker`]
15/// for the record format.
16///
17/// Logs, counts, and swallows errors returned by the [`InventoryTracker`].
18pub struct CostTrackerStream<P: Producer> {
19    tracker: InventoryTracker<P>,
20}
21
22impl<P: Producer> CostTrackerStream<P> {
23    /// Reports changes through `producer`, as described by `config`.
24    pub fn new(producer: P, config: &CostTrackerStreamConfig) -> Self {
25        Self {
26            tracker: InventoryTracker::new(
27                producer,
28                &config.shared_resource_id,
29                config.sample_rate,
30            ),
31        }
32    }
33
34    /// Counts and logs a failed report, then returns.
35    ///
36    /// Note: success here doesn't necessarily mean that a message was emitted. The
37    /// tracker's sampling policy may have filtered an object out and chosen to emit
38    /// nothing. Or, a message may have been enqueued only to fail later. Those are
39    /// counted by the transport's delivery failure callback instead.
40    fn swallow(&self, op: &'static str, result: Result<(), P::Error>)
41    where
42        P::Error: Into<BoxError>,
43    {
44        if let Err(error) = result {
45            // Boxed because `SharedProducer` reports a `Box<dyn Error>`, which std does
46            // not implement `Error` for.
47            let error: BoxError = error.into();
48            // Records are dropped rather than retried
49            objectstore_metrics::count!(
50                "cost_tracker.dropped" += 1,
51                shared_resource_id = self.tracker.shared_resource_id().to_owned(),
52                op = op,
53            );
54            objectstore_log::warn!(!!&*error, op, "failed to publish change stream record");
55        }
56    }
57}
58
59impl<P: Producer> fmt::Debug for CostTrackerStream<P> {
60    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
61        f.debug_struct("CostTrackerStream")
62            .field("shared_resource_id", &self.tracker.shared_resource_id())
63            .field("sample_rate", &self.tracker.sample_rate())
64            .finish()
65    }
66}
67
68#[async_trait::async_trait]
69impl<P> ChangeStream for CostTrackerStream<P>
70where
71    P: Producer + Clone + Send + Sync + 'static,
72    P::Error: Into<BoxError> + Send + 'static,
73{
74    fn write(&self, id: &ObjectId, size: u64, expires_at: Option<SystemTime>) {
75        let result = self.tracker.write(
76            &id.as_storage_path().to_string(),
77            id.usecase(),
78            size,
79            SystemTime::now(),
80            expires_at,
81            scope_id(id, SCOPE_ORGANIZATION),
82            scope_id(id, SCOPE_PROJECT),
83        );
84        self.swallow("write", result);
85    }
86
87    fn update(&self, id: &ObjectId, expires_at: Option<SystemTime>) {
88        let result = self.tracker.update(
89            &id.as_storage_path().to_string(),
90            id.usecase(),
91            SystemTime::now(),
92            expires_at,
93            scope_id(id, SCOPE_ORGANIZATION),
94            scope_id(id, SCOPE_PROJECT),
95        );
96        self.swallow("update", result);
97    }
98
99    fn delete(&self, id: &ObjectId) {
100        let result = self.tracker.delete(
101            &id.as_storage_path().to_string(),
102            id.usecase(),
103            SystemTime::now(),
104        );
105        self.swallow("delete", result);
106    }
107
108    async fn join(&self, timeout: Duration) {
109        self.swallow("join", self.tracker.join(timeout).await);
110    }
111}
112
113#[cfg(test)]
114mod tests {
115    use objectstore_inventory_tracker::OpType;
116    use objectstore_inventory_tracker::test_utils::DummyProducer;
117
118    use super::*;
119
120    fn object_id(path: &str) -> ObjectId {
121        ObjectId::from_storage_path(path).expect("valid storage path")
122    }
123
124    fn stream(sample_rate: f64) -> (DummyProducer, CostTrackerStream<DummyProducer>) {
125        let producer = DummyProducer::default();
126        let stream = CostTrackerStream::new(
127            producer.clone(),
128            &CostTrackerStreamConfig {
129                shared_resource_id: "bigtable_objectstore".into(),
130                sample_rate,
131            },
132        );
133        (producer, stream)
134    }
135
136    #[test]
137    fn usecase_and_scopes_are_extracted_from_the_id() {
138        let (producer, stream) = stream(1.0);
139        let id = object_id("attachments/org.17/project.42/objects/abc");
140
141        stream.write(&id, 4096, None);
142
143        let record = &producer.records()[0];
144        assert_eq!(record.shared_resource_id, "bigtable_objectstore");
145        assert_eq!(record.app_feature, "attachments");
146        assert_eq!(record.organization_id, Some(17));
147        assert_eq!(record.project_id, Some(42));
148        assert_eq!(record.size, Some(4096));
149        assert_eq!(record.op_type, OpType::Write);
150    }
151
152    #[test]
153    fn the_storage_path_is_not_emitted() {
154        let (producer, stream) = stream(1.0);
155        let id = object_id("attachments/org.17/project.42/objects/abc");
156
157        stream.write(&id, 4096, None);
158
159        let record_id = &producer.records()[0].record_id;
160        assert_ne!(record_id, &id.as_storage_path().to_string());
161        assert!(!record_id.contains("attachments"));
162    }
163
164    #[test]
165    fn missing_or_unparseable_scopes_are_reported_as_absent() {
166        let (producer, stream) = stream(1.0);
167
168        for path in [
169            "attachments/objects/abc",
170            "attachments/organization.17/objects/abc",
171            "attachments/org.not-a-number/project.42/objects/abc",
172        ] {
173            stream.write(&object_id(path), 1, None);
174        }
175
176        let records = producer.records();
177        assert_eq!(
178            records.len(),
179            3,
180            "every object is reported regardless of scopes"
181        );
182        assert_eq!(records[0].organization_id, None, "no scopes at all");
183        assert_eq!(
184            records[1].organization_id, None,
185            "wrong scope key is not org"
186        );
187        assert_eq!(
188            records[2].organization_id, None,
189            "unparseable org is absent"
190        );
191        assert_eq!(
192            records[2].project_id,
193            Some(42),
194            "but project still resolves"
195        );
196        for record in &records {
197            assert_eq!(record.app_feature, "attachments");
198        }
199    }
200
201    #[test]
202    fn every_operation_on_an_object_reports_the_same_record() {
203        let (producer, stream) = stream(1.0);
204        let id = object_id("attachments/org.1/project.2/objects/abc");
205
206        stream.write(&id, 10, None);
207        stream.update(&id, Some(SystemTime::now()));
208        stream.delete(&id);
209
210        let records = producer.records();
211        assert_eq!(records.len(), 3);
212        assert_eq!(records[0].record_id, records[1].record_id);
213        assert_eq!(records[1].record_id, records[2].record_id);
214    }
215
216    #[test]
217    fn distinct_revisions_are_distinct_records() {
218        let (producer, stream) = stream(1.0);
219
220        stream.write(
221            &object_id("attachments/org.1/project.2/objects/abc/0199aaaa"),
222            1,
223            None,
224        );
225        stream.write(
226            &object_id("attachments/org.1/project.2/objects/abc/0199bbbb"),
227            1,
228            None,
229        );
230
231        let records = producer.records();
232        assert_ne!(records[0].record_id, records[1].record_id);
233    }
234
235    #[test]
236    fn update_omits_size_and_delete_omits_everything_optional() {
237        let (producer, stream) = stream(1.0);
238        let id = object_id("attachments/org.1/project.2/objects/abc");
239
240        stream.update(&id, Some(SystemTime::now()));
241        stream.delete(&id);
242
243        let records = producer.records();
244        assert_eq!(records[0].op_type, OpType::Update);
245        assert_eq!(records[0].size, None);
246        assert!(records[0].expiration_time.is_some());
247        assert_eq!(records[1].op_type, OpType::Delete);
248        assert_eq!(records[1].size, None);
249        assert_eq!(records[1].expiration_time, None);
250    }
251
252    #[test]
253    fn a_listener_sampled_at_zero_reports_nothing() {
254        let (producer, stream) = stream(0.0);
255
256        for i in 0..100 {
257            let id = object_id(&format!("attachments/org.1/project.2/objects/{i}"));
258            stream.write(&id, 1, None);
259            stream.update(&id, None);
260            stream.delete(&id);
261        }
262
263        assert!(producer.records().is_empty());
264    }
265}