objectstore_service/change_stream/
factory.rs1use std::fmt;
5use std::sync::Arc;
6
7#[cfg(feature = "storage-cogs")]
8use objectstore_inventory_tracker::SharedProducer;
9#[cfg(all(test, feature = "storage-cogs"))]
10use objectstore_inventory_tracker::test_utils;
11#[cfg(feature = "storage-cogs")]
12use serde::{Deserialize, Serialize};
13
14#[cfg(feature = "storage-cogs")]
15use super::CostTrackerStream;
16use super::{ChangeStream, CostTrackerStreamConfig, NoopStream};
17
18#[cfg(feature = "storage-cogs")]
22#[derive(Debug, Clone, Deserialize, Serialize)]
23#[serde(tag = "type", rename_all = "lowercase")]
24pub enum CostTrackerConfig {
25 Kafka(objectstore_inventory_tracker::kafka::KafkaConfig),
27}
28
29#[derive(Clone, Default)]
33pub struct ChangeStreamFactory {
34 #[cfg(feature = "storage-cogs")]
35 producer: Option<SharedProducer>,
36}
37
38impl ChangeStreamFactory {
39 #[cfg(feature = "storage-cogs")]
43 pub fn new(config: &CostTrackerConfig) -> Self {
44 let CostTrackerConfig::Kafka(kafka) = config;
45 Self {
46 producer: build_kafka_producer(kafka),
47 }
48 }
49
50 #[cfg(feature = "storage-cogs")]
52 pub fn build(&self, config: Option<&CostTrackerStreamConfig>) -> Arc<dyn ChangeStream> {
53 match (config, self.producer.clone()) {
54 (Some(config), Some(producer)) => Arc::new(CostTrackerStream::new(producer, config)),
55 (None, None) => Arc::new(NoopStream),
56 (c, p) => {
57 objectstore_log::warn!(
58 stream_configured = c.is_some(),
59 producer_configured = p.is_some(),
60 "incomplete change stream configuration, returning NoopStream",
61 );
62 Arc::new(NoopStream)
63 }
64 }
65 }
66
67 #[cfg(not(feature = "storage-cogs"))]
69 pub fn build(&self, _config: Option<&CostTrackerStreamConfig>) -> Arc<dyn ChangeStream> {
70 Arc::new(NoopStream)
71 }
72}
73
74impl fmt::Debug for ChangeStreamFactory {
75 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
76 let mut f = f.debug_struct("ChangeStreamFactory");
77 #[cfg(feature = "storage-cogs")]
78 f.field("producer", &self.producer.is_some());
79 f.finish()
80 }
81}
82
83#[cfg(feature = "storage-cogs")]
85fn build_kafka_producer(
86 config: &objectstore_inventory_tracker::kafka::KafkaConfig,
87) -> Option<SharedProducer> {
88 use objectstore_inventory_tracker::Producer as _;
89 use objectstore_inventory_tracker::kafka::KafkaProducer;
90
91 let on_delivery_failure = Box::new(|_: &_| {
94 objectstore_metrics::count!("cost_tracker.undelivered" += 1);
95 });
96
97 match KafkaProducer::try_new(config.clone(), Some(on_delivery_failure)) {
98 Ok(producer) => Some(producer.shared()),
99 Err(error) => {
100 objectstore_log::error!(
101 !!&error,
102 "failed to create the change stream kafka producer; \
103 backends with a change stream will report nothing"
104 );
105 None
106 }
107 }
108}
109
110#[cfg(all(test, feature = "storage-cogs"))]
112pub(crate) fn dummy_factory() -> (ChangeStreamFactory, test_utils::DummyProducer) {
113 use objectstore_inventory_tracker::Producer as _;
114
115 let producer = test_utils::DummyProducer::default();
116 let factory = ChangeStreamFactory {
117 producer: Some(producer.clone().shared()),
118 };
119
120 (factory, producer)
121}
122
123#[cfg(all(test, feature = "storage-cogs"))]
124mod tests {
125 use super::*;
126
127 fn config() -> CostTrackerStreamConfig {
128 CostTrackerStreamConfig {
129 shared_resource_id: "bigtable_objectstore".into(),
130 sample_rate: 1.0,
131 }
132 }
133
134 fn reports(stream: &Arc<dyn ChangeStream>) -> bool {
135 !format!("{stream:?}").contains("NoopStream")
136 }
137
138 #[test]
139 fn a_backend_without_a_change_stream_config_reports_nothing() {
140 let (factory, _producer) = dummy_factory();
141
142 assert!(!reports(&factory.build(None)));
143 }
144
145 #[test]
146 fn a_configured_backend_without_a_transport_reports_nothing() {
147 let factory = ChangeStreamFactory::default();
148
149 assert!(!reports(&factory.build(Some(&config()))));
150 }
151
152 #[test]
153 fn a_configured_backend_reports_through_the_transport() {
154 let (factory, producer) = dummy_factory();
155 let stream = factory.build(Some(&config()));
156
157 assert!(reports(&stream));
158
159 stream.delete(&crate::id::ObjectId::from_storage_path("attachments/objects/abc").unwrap());
160
161 let records = producer.records();
162 assert_eq!(records.len(), 1);
163 assert_eq!(records[0].shared_resource_id, "bigtable_objectstore");
164 }
165
166 #[test]
167 fn an_unusable_transport_disables_reporting_instead_of_failing() {
168 let factory = ChangeStreamFactory::new(&CostTrackerConfig::Kafka(
169 objectstore_inventory_tracker::kafka::KafkaConfig {
170 topic: "shared-resources-inventory".into(),
171 bootstrap_servers: vec!["127.0.0.1:9092".into()],
172 override_params: [("not.a.real.property".to_owned(), "1".to_owned())].into(),
173 },
174 ));
175
176 assert!(!reports(&factory.build(Some(&config()))));
177 }
178}