objectstore_service/change_stream/
cost_tracker.rs1use 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
13pub struct CostTrackerStream<P: Producer> {
19 tracker: InventoryTracker<P>,
20}
21
22impl<P: Producer> CostTrackerStream<P> {
23 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 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 let error: BoxError = error.into();
48 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}