1pub 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
12pub 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
33pub 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 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
121pub 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(®istry).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(®istry).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(®istry).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(®istry).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(®istry).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}