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 objectstore_types::time::Timestamp;
12
13use crate::id::ObjectId;
14
15pub struct CostTrackerStream<P: Producer> {
21 tracker: InventoryTracker<P>,
22}
23
24impl<P: Producer> CostTrackerStream<P> {
25 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 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 let error: BoxError = error.into();
50 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}