Skip to main content

laminar_connectors/mongodb/
mod.rs

1//! `MongoDB` CDC source and sink connectors.
2
3pub mod change_event;
4pub mod config;
5pub mod lookup;
6pub mod metrics;
7pub mod sink;
8pub mod source;
9pub mod timeseries;
10pub mod write_model;
11
12// Re-export primary types at module level.
13pub use config::{FullDocumentMode, MongoDbSinkConfig, MongoDbSourceConfig};
14pub use sink::MongoDbSink;
15pub use source::{mongodb_cdc_envelope_schema, MongoDbCdcSource};
16pub use timeseries::{CollectionKind, TimeSeriesConfig, TimeSeriesGranularity};
17pub use write_model::WriteMode;
18
19const MONGODB_LOOKUP_PROPERTIES: &[&str] = &[
20    "connection.uri",
21    "database",
22    "collection",
23    "laminar.source.name",
24    "_arrow_schema",
25    "_primary_key_columns",
26];
27
28use std::sync::Arc;
29
30use crate::config::{ConfigKeySpec, ConnectorInfo};
31use crate::registry::ConnectorRegistry;
32
33/// Registers the `MongoDB` CDC source connector with the given registry.
34///
35/// # Errors
36///
37/// Returns an error if the connector name is already registered or the registry is frozen.
38pub fn register_mongodb_cdc_source(
39    registry: &ConnectorRegistry,
40) -> Result<(), crate::error::ConnectorError> {
41    let info = ConnectorInfo {
42        name: "mongodb-cdc".to_string(),
43        display_name: "MongoDB CDC Source".to_string(),
44        version: env!("CARGO_PKG_VERSION").to_string(),
45        is_source: true,
46        is_sink: false,
47        config_keys: mongodb_cdc_config_keys(),
48    };
49
50    registry.register_source(
51        "mongodb-cdc",
52        info,
53        Arc::new(|registry: Option<&Arc<prometheus::Registry>>| {
54            Ok(Box::new(MongoDbCdcSource::new(
55                MongoDbSourceConfig::default(),
56                registry.map(Arc::as_ref),
57            )))
58        }),
59    )?;
60
61    // On-demand (partial cache mode) lookup source: find({ pk: { $in: [...] } }).
62    registry.register_lookup_source(
63        "mongodb",
64        ConnectorInfo {
65            name: "mongodb".to_string(),
66            display_name: "MongoDB Lookup Source".to_string(),
67            version: env!("CARGO_PKG_VERSION").to_string(),
68            is_source: true,
69            is_sink: false,
70            config_keys: mongodb_lookup_config_keys(),
71        },
72        Arc::new(MongoLookupFactory),
73    )
74}
75
76struct MongoLookupFactory;
77
78#[async_trait::async_trait]
79impl crate::registry::LookupSourceFactory for MongoLookupFactory {
80    async fn build(
81        &self,
82        config: crate::config::ConnectorConfig,
83        declared_schema: Option<arrow_schema::SchemaRef>,
84    ) -> Result<Arc<dyn laminar_core::lookup::source::LookupSourceDyn>, crate::error::ConnectorError>
85    {
86        use crate::mongodb::lookup::{MongoLookupSource, MongoLookupSourceConfig};
87
88        let schema = declared_schema.ok_or_else(|| {
89            crate::error::ConnectorError::ConfigurationError(
90                "mongodb lookup source requires a declared table schema".into(),
91            )
92        })?;
93
94        let pk_columns: Vec<String> = config
95            .get("_primary_key_columns")
96            .unwrap_or("")
97            .split(',')
98            .map(|s| s.trim().to_string())
99            .filter(|s| !s.is_empty())
100            .collect();
101        if pk_columns.is_empty() {
102            return Err(crate::error::ConnectorError::ConfigurationError(
103                "mongodb lookup source requires primary key columns".into(),
104            ));
105        }
106
107        config.reject_unknown_properties(MONGODB_LOOKUP_PROPERTIES, "MongoDB lookup")?;
108        let lookup_config = MongoLookupSourceConfig {
109            connection_uri: config.require("connection.uri")?.to_string(),
110            database: config.require("database")?.to_string(),
111            collection: config.require("collection")?.to_string(),
112            primary_key_columns: pk_columns,
113            schema,
114        };
115
116        let source = MongoLookupSource::open(lookup_config).await?;
117        Ok(Arc::new(source) as Arc<dyn laminar_core::lookup::source::LookupSourceDyn>)
118    }
119}
120
121/// Registers the `MongoDB` sink connector with the given registry.
122///
123/// # Errors
124///
125/// Returns an error if the connector name is already registered or the registry is frozen.
126pub fn register_mongodb_sink(
127    registry: &ConnectorRegistry,
128) -> Result<(), crate::error::ConnectorError> {
129    let info = ConnectorInfo {
130        name: "mongodb-sink".to_string(),
131        display_name: "MongoDB Sink".to_string(),
132        version: env!("CARGO_PKG_VERSION").to_string(),
133        is_source: false,
134        is_sink: true,
135        config_keys: mongodb_sink_config_keys(),
136    };
137
138    registry.register_sink(
139        "mongodb-sink",
140        info,
141        Arc::new(|config, registry: Option<&Arc<prometheus::Registry>>| {
142            MongoDbSink::from_connector_config(config, registry.map(Arc::as_ref))
143                .map(|sink| Box::new(sink) as Box<dyn crate::connector::SinkConnector>)
144        }),
145    )
146}
147
148fn mongodb_cdc_config_keys() -> Vec<ConfigKeySpec> {
149    vec![
150        ConfigKeySpec::required("connection.uri", "MongoDB connection URI"),
151        ConfigKeySpec::required("database", "Database name"),
152        ConfigKeySpec::required("collection", "Fixed collection name"),
153        ConfigKeySpec::optional(
154            "full.document.mode",
155            "Deterministic full document mode (delta or required post-image)",
156            "delta",
157        ),
158        ConfigKeySpec::optional(
159            "pipeline",
160            "JSON array of up to 64 $match stages (maximum 256 KiB)",
161            "[]",
162        ),
163        ConfigKeySpec::optional(
164            "max.buffered.bytes",
165            "Max retained decoded bytes before backpressure (1 MiB to 4 GiB)",
166            config::DEFAULT_MAX_BUFFERED_BYTES.to_string(),
167        ),
168    ]
169}
170
171fn mongodb_sink_config_keys() -> Vec<ConfigKeySpec> {
172    vec![
173        ConfigKeySpec::required("connection.uri", "MongoDB connection URI"),
174        ConfigKeySpec::required("database", "Target database name"),
175        ConfigKeySpec::required("collection", "Target collection name"),
176        ConfigKeySpec::optional("flush.interval.ms", "Max time between flushes (ms)", "250"),
177        ConfigKeySpec::optional(
178            "write.mode",
179            "Write operation mode (insert, upsert, cdc_replay)",
180            "insert",
181        ),
182        ConfigKeySpec::optional(
183            "write.mode.key_fields",
184            "Comma-separated key fields to match documents in upsert mode",
185            "",
186        ),
187        ConfigKeySpec::optional(
188            "timeseries.time_field",
189            "The field in each document containing the date",
190            "",
191        ),
192        ConfigKeySpec::optional(
193            "timeseries.meta_field",
194            "An optional field labeling the data source",
195            "",
196        ),
197        ConfigKeySpec::optional(
198            "timeseries.granularity",
199            "Bucketing granularity (seconds, minutes, hours, custom)",
200            "seconds",
201        ),
202        ConfigKeySpec::optional(
203            "timeseries.bucket_max_span_seconds",
204            "Max span of a single bucket in seconds (custom granularity)",
205            "",
206        ),
207        ConfigKeySpec::optional(
208            "timeseries.bucket_rounding_seconds",
209            "Rounding boundary in seconds (custom granularity)",
210            "",
211        ),
212        ConfigKeySpec::optional(
213            "timeseries.expire_after_seconds",
214            "TTL in seconds (automatically delete documents after this span)",
215            "",
216        ),
217        ConfigKeySpec::optional(
218            "sink.write.timeout.ms",
219            "Complete MongoDB sink write deadline in milliseconds",
220            "30000",
221        ),
222    ]
223}
224
225fn mongodb_lookup_config_keys() -> Vec<ConfigKeySpec> {
226    vec![
227        ConfigKeySpec::required("connection.uri", "MongoDB connection URI"),
228        ConfigKeySpec::required("database", "Database name"),
229        ConfigKeySpec::required("collection", "Collection name"),
230    ]
231}
232
233#[cfg(test)]
234mod tests {
235    use super::*;
236    use arrow_schema::{DataType, Field, Schema};
237
238    fn sink_factory_config() -> crate::config::ConnectorConfig {
239        let schema = Schema::new(vec![Field::new("id", DataType::Int64, false)]);
240        let mut config = crate::config::ConnectorConfig::new("mongodb-sink");
241        config.set("connection.uri", "mongodb://localhost:27017");
242        config.set("database", "db");
243        config.set("collection", "out");
244        config.set(
245            "_arrow_schema",
246            crate::config::encode_arrow_schema_ipc(&schema),
247        );
248        config
249    }
250
251    #[test]
252    fn test_register_mongodb_cdc_source() {
253        let registry = ConnectorRegistry::new();
254        register_mongodb_cdc_source(&registry).unwrap();
255
256        let info = registry.source_info("mongodb-cdc");
257        assert!(info.is_some());
258        let info = info.unwrap();
259        assert_eq!(info.name, "mongodb-cdc");
260        assert!(info.is_source);
261        assert!(!info.is_sink);
262        assert!(!info.config_keys.is_empty());
263    }
264
265    #[test]
266    fn test_register_mongodb_sink() {
267        let registry = ConnectorRegistry::new();
268        register_mongodb_sink(&registry).unwrap();
269
270        let info = registry.sink_info("mongodb-sink");
271        assert!(info.is_some());
272        let info = info.unwrap();
273        assert_eq!(info.name, "mongodb-sink");
274        assert!(info.is_sink);
275        assert!(!info.is_source);
276        assert!(!info.config_keys.is_empty());
277    }
278
279    #[test]
280    fn test_cdc_config_keys() {
281        let keys = mongodb_cdc_config_keys();
282        let required: Vec<&str> = keys
283            .iter()
284            .filter(|k| k.required)
285            .map(|k| k.key.as_str())
286            .collect();
287        assert!(required.contains(&"connection.uri"));
288        assert!(required.contains(&"database"));
289        assert!(required.contains(&"collection"));
290        let byte_budget = keys
291            .iter()
292            .find(|key| key.key == "max.buffered.bytes")
293            .expect("MongoDB CDC byte budget must be discoverable");
294        assert_eq!(
295            byte_budget
296                .default
297                .as_deref()
298                .and_then(|value| value.parse::<usize>().ok()),
299            Some(config::DEFAULT_MAX_BUFFERED_BYTES)
300        );
301        let pipeline = keys
302            .iter()
303            .find(|key| key.key == "pipeline")
304            .expect("MongoDB CDC pipeline must be discoverable");
305        assert_eq!(pipeline.default.as_deref(), Some("[]"));
306        for removed in [
307            "batch.size",
308            "max.buffered.events",
309            "max.await.time.ms",
310            "resume.token.store",
311            "split.large.events",
312            "max.poll.records",
313        ] {
314            assert!(keys.iter().all(|key| key.key != removed));
315        }
316    }
317
318    #[test]
319    fn lookup_config_keys_are_minimal() {
320        let keys = mongodb_lookup_config_keys();
321        assert_eq!(keys.len(), 3);
322        assert!(keys.iter().all(|key| key.required));
323        assert!(keys.iter().all(|key| !matches!(
324            key.key.as_str(),
325            "full.document.mode" | "max.buffered.bytes"
326        )));
327    }
328
329    #[test]
330    fn test_sink_config_keys() {
331        let keys = mongodb_sink_config_keys();
332        let required: Vec<&str> = keys
333            .iter()
334            .filter(|k| k.required)
335            .map(|k| k.key.as_str())
336            .collect();
337        assert!(required.contains(&"connection.uri"));
338        assert!(required.contains(&"database"));
339        assert!(required.contains(&"collection"));
340        for removed in [
341            "batch.size",
342            "ordered",
343            "write.mode.upsert_on_missing",
344            "write_concern.journal",
345            "write_concern.timeout_ms",
346        ] {
347            assert!(keys.iter().all(|key| key.key != removed));
348        }
349        let write_timeout = keys
350            .iter()
351            .find(|key| key.key == "sink.write.timeout.ms")
352            .expect("write timeout must be discoverable");
353        assert_eq!(write_timeout.default.as_deref(), Some("30000"));
354    }
355
356    #[test]
357    fn test_factory_creates_source() {
358        let registry = ConnectorRegistry::new();
359        register_mongodb_cdc_source(&registry).unwrap();
360
361        let config = crate::config::ConnectorConfig::new("mongodb-cdc");
362        let source = registry.create_source(&config, None);
363        assert!(source.is_ok());
364    }
365
366    #[test]
367    fn test_factory_creates_sink() {
368        let registry = ConnectorRegistry::new();
369        register_mongodb_sink(&registry).unwrap();
370
371        let sink = registry.create_sink(&sink_factory_config(), None).unwrap();
372        assert_eq!(sink.schema().field(0).name(), "id");
373    }
374
375    #[test]
376    fn sink_factory_rejects_missing_and_malformed_schema() {
377        let registry = ConnectorRegistry::new();
378        register_mongodb_sink(&registry).unwrap();
379
380        let mut missing = sink_factory_config();
381        let mut properties = missing.properties().clone();
382        properties.remove("_arrow_schema");
383        missing = crate::config::ConnectorConfig::with_properties("mongodb-sink", properties);
384        let missing_error = registry
385            .create_sink(&missing, None)
386            .err()
387            .expect("missing schema must fail")
388            .to_string();
389        assert!(missing_error.contains("_arrow_schema"), "{missing_error}");
390
391        let mut malformed = sink_factory_config();
392        malformed.set("_arrow_schema", "not-arrow-ipc");
393        let malformed_error = registry
394            .create_sink(&malformed, None)
395            .err()
396            .expect("malformed schema must fail")
397            .to_string();
398        assert!(
399            malformed_error.contains("invalid") && malformed_error.contains("_arrow_schema"),
400            "{malformed_error}"
401        );
402    }
403}