Skip to main content

objectstore_service/change_stream/
factory.rs

1//! Constructs the [`ChangeStream`] implementation(s) for a backend based on the available
2//! service-wide sink config and per-backend stream config.
3
4use 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/// Where every backend's change stream records are carried for cost tracking.
19///
20/// Service-wide: a transport owns connections and a send queue worth sharing.
21#[cfg(feature = "storage-cogs")]
22#[derive(Debug, Clone, Deserialize, Serialize)]
23#[serde(tag = "type", rename_all = "lowercase")]
24pub enum CostTrackerConfig {
25    /// Reports onto a Kafka topic.
26    Kafka(objectstore_inventory_tracker::kafka::KafkaConfig),
27}
28
29/// Builds the [`ChangeStream`] impl(s) a backend reports to.
30///
31/// Without a usable transport every backend gets a [`NoopStream`].
32#[derive(Clone, Default)]
33pub struct ChangeStreamFactory {
34    #[cfg(feature = "storage-cogs")]
35    producer: Option<SharedProducer>,
36}
37
38impl ChangeStreamFactory {
39    /// Builds the transport described by `config`.
40    ///
41    /// Fails open: an unusable transport is logged, not fatal.
42    #[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    /// Builds the stream `config` asks for, or a [`NoopStream`] if it cannot be built.
51    #[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    /// Reporting is not compiled in, so every backend reports nothing.
68    #[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/// Creates the shared Kafka producer, or logs why there will be no reporting.
84#[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    // Delivery is asynchronous, so an accepted record can still fail to arrive. Without
92    // this those show up only as a shortfall in the downstream data.
93    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/// A [`ChangeStreamFactory`] that reports into the returned producer.
111#[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}