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