1use std::collections::{BTreeMap, VecDeque};
13use std::mem::size_of;
14use std::sync::Arc;
15
16use arrow_array::builder::{StringBuilder, UInt32Builder};
17use arrow_array::RecordBatch;
18use arrow_schema::{DataType, Field, Schema, SchemaRef};
19use async_trait::async_trait;
20#[cfg(feature = "mongodb-cdc")]
21use futures_util::TryStreamExt;
22use sha2::{Digest, Sha256};
23use tokio::sync::{Notify, OwnedSemaphorePermit, Semaphore};
24use uuid::Uuid;
25
26use crate::checkpoint::SourceCheckpoint;
27use crate::config::{ConnectorConfig, ConnectorState};
28use crate::connector::{
29 ConnectorTaskOwner, ConnectorTaskTracker, SourceBatch, SourceConnector, SourceConsistency,
30 SourceContract, SourcePosition, SourceStart, SourceTopology,
31};
32use crate::error::ConnectorError;
33
34use super::change_event::{MongoDbChangeEvent, OperationType};
35use super::config::MongoDbSourceConfig;
36use super::metrics::MongoDbCdcMetrics;
37
38const MAX_RESUME_TOKEN_BYTES: usize = 64 * 1024;
39const MONGODB_CHECKPOINT_CONNECTOR: &str = "mongodb-cdc";
40const MONGODB_CHECKPOINT_VERSION: &str = "4";
41const STREAM_IDENTITY_METADATA: &str = "stream_identity_sha256";
42const COLLECTION_UUID_METADATA: &str = "collection_uuid";
43const DEPLOYMENT_IDENTITY_METADATA: &str = "deployment_identity";
44const RESUME_TOKEN_OFFSET: &str = "resume_token";
45const START_AFTER_TOKEN_OFFSET: &str = "start_after_token";
46#[cfg(feature = "mongodb-cdc")]
47const MAX_MONGODB_WIRE_EVENT_BYTES: usize = 16 * 1024 * 1024;
48#[cfg(feature = "mongodb-cdc")]
49const CURSOR_MAX_AWAIT_TIME: std::time::Duration = std::time::Duration::from_secs(1);
50#[cfg(feature = "mongodb-cdc")]
51const READER_STARTUP_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(30);
52
53#[derive(Debug, Clone, PartialEq, Eq)]
54enum MongoCheckpointPosition {
55 ResumeAfter(String),
56 StartAfter(String),
57}
58
59#[derive(Debug, Clone, PartialEq, Eq)]
60enum MongoDeploymentIdentity {
61 ReplicaSet(String),
62 ShardedCluster(String),
63}
64
65impl MongoDeploymentIdentity {
66 fn encode(&self) -> String {
67 match self {
68 Self::ReplicaSet(id) => format!("replica-set:{id}"),
69 Self::ShardedCluster(id) => format!("sharded-cluster:{id}"),
70 }
71 }
72}
73
74#[derive(Debug, Clone, PartialEq, Eq)]
75struct ParsedMongoCheckpoint {
76 position: MongoCheckpointPosition,
77 collection_uuid: Uuid,
78 deployment_identity: MongoDeploymentIdentity,
79}
80
81fn mongodb_stream_identity(config: &MongoDbSourceConfig) -> String {
82 let mut digest = Sha256::new();
83 digest.update(b"laminardb-mongodb-change-stream-v4\0");
84 let full_document_mode = match config.full_document_mode {
85 super::config::FullDocumentMode::Delta => 0_u8,
86 super::config::FullDocumentMode::RequirePostImage => 1,
87 };
88 digest.update([full_document_mode]);
89 digest.update([1]); let pipeline = super::config::canonical_pipeline_json(&config.pipeline);
91 digest.update(
92 u64::try_from(pipeline.len())
93 .unwrap_or(u64::MAX)
94 .to_be_bytes(),
95 );
96 digest.update(pipeline.as_bytes());
97 format!("{:x}", digest.finalize())
98}
99
100fn canonical_resume_token(token: &str) -> Result<String, ConnectorError> {
101 if token.is_empty() || token.len() > MAX_RESUME_TOKEN_BYTES {
102 return Err(ConnectorError::ConfigurationError(format!(
103 "MongoDB CDC resume token size must be 1..={MAX_RESUME_TOKEN_BYTES} bytes"
104 )));
105 }
106 let value: serde_json::Value = serde_json::from_str(token).map_err(|error| {
107 ConnectorError::ConfigurationError(format!(
108 "MongoDB CDC resume token is not valid JSON: {error}"
109 ))
110 })?;
111 let serde_json::Value::Object(document) = &value else {
112 return Err(ConnectorError::ConfigurationError(
113 "MongoDB CDC resume token must be a JSON document".into(),
114 ));
115 };
116 if document.is_empty() {
117 return Err(ConnectorError::ConfigurationError(
118 "MongoDB CDC resume token document must not be empty".into(),
119 ));
120 }
121 let canonical = serde_json::to_string(&value).map_err(|error| {
122 ConnectorError::Internal(format!("serialize MongoDB CDC resume token: {error}"))
123 })?;
124 if canonical != token {
125 return Err(ConnectorError::ConfigurationError(
126 "MongoDB CDC resume token is not in canonical JSON form".into(),
127 ));
128 }
129 Ok(canonical)
130}
131
132fn parse_collection_uuid(encoded: &str) -> Result<Uuid, ConnectorError> {
133 let uuid = Uuid::parse_str(encoded).map_err(|error| {
134 ConnectorError::ConfigurationError(format!("invalid MongoDB CDC collection UUID: {error}"))
135 })?;
136 if uuid.hyphenated().to_string() != encoded {
137 return Err(ConnectorError::ConfigurationError(
138 "MongoDB CDC collection UUID is not in canonical lowercase hyphenated form".into(),
139 ));
140 }
141 Ok(uuid)
142}
143
144fn parse_deployment_identity(encoded: &str) -> Result<MongoDeploymentIdentity, ConnectorError> {
145 let (kind, id) = encoded.split_once(':').ok_or_else(|| {
146 ConnectorError::ConfigurationError(
147 "MongoDB CDC deployment identity must include its deployment type".into(),
148 )
149 })?;
150 if id.contains(':') {
151 return Err(ConnectorError::ConfigurationError(
152 "MongoDB CDC deployment identity has too many fields".into(),
153 ));
154 }
155 let object_id = mongodb::bson::oid::ObjectId::parse_str(id).map_err(|error| {
156 ConnectorError::ConfigurationError(format!(
157 "invalid MongoDB CDC deployment ObjectId: {error}"
158 ))
159 })?;
160 if object_id.to_hex() != id {
161 return Err(ConnectorError::ConfigurationError(
162 "MongoDB CDC deployment ObjectId is not canonical lowercase hex".into(),
163 ));
164 }
165 match kind {
166 "replica-set" => Ok(MongoDeploymentIdentity::ReplicaSet(id.to_string())),
167 "sharded-cluster" => Ok(MongoDeploymentIdentity::ShardedCluster(id.to_string())),
168 _ => Err(ConnectorError::ConfigurationError(format!(
169 "unknown MongoDB CDC deployment identity type '{kind}'"
170 ))),
171 }
172}
173
174fn parse_mongodb_checkpoint(
175 checkpoint: &SourceCheckpoint,
176 config: &MongoDbSourceConfig,
177) -> Result<ParsedMongoCheckpoint, ConnectorError> {
178 let expected_stream_identity = mongodb_stream_identity(config);
179 if checkpoint.get_metadata("connector") != Some(MONGODB_CHECKPOINT_CONNECTOR)
180 || checkpoint.get_metadata("version") != Some(MONGODB_CHECKPOINT_VERSION)
181 || checkpoint.get_metadata("database") != Some(config.database.as_str())
182 || checkpoint.get_metadata("collection") != Some(config.collection.as_str())
183 || checkpoint.get_metadata(STREAM_IDENTITY_METADATA)
184 != Some(expected_stream_identity.as_str())
185 {
186 return Err(ConnectorError::ConfigurationError(
187 "MongoDB CDC checkpoint identity or format does not match the configured source".into(),
188 ));
189 }
190 let collection_uuid = checkpoint
191 .get_metadata(COLLECTION_UUID_METADATA)
192 .ok_or_else(|| {
193 ConnectorError::ConfigurationError(
194 "MongoDB CDC checkpoint is missing its collection UUID".into(),
195 )
196 })
197 .and_then(parse_collection_uuid)?;
198 let deployment_identity = checkpoint
199 .get_metadata(DEPLOYMENT_IDENTITY_METADATA)
200 .ok_or_else(|| {
201 ConnectorError::ConfigurationError(
202 "MongoDB CDC checkpoint is missing its deployment identity".into(),
203 )
204 })
205 .and_then(parse_deployment_identity)?;
206 if checkpoint.metadata().len() != 7 {
207 return Err(ConnectorError::ConfigurationError(
208 "MongoDB CDC checkpoint contains unknown metadata fields".into(),
209 ));
210 }
211 if checkpoint.offsets().len() != 1 {
212 return Err(ConnectorError::ConfigurationError(
213 "MongoDB CDC checkpoint must contain exactly one resume token".into(),
214 ));
215 }
216 let position = if let Some(token) = checkpoint.get_offset(RESUME_TOKEN_OFFSET) {
217 canonical_resume_token(token).map(MongoCheckpointPosition::ResumeAfter)?
218 } else if let Some(token) = checkpoint.get_offset(START_AFTER_TOKEN_OFFSET) {
219 canonical_resume_token(token).map(MongoCheckpointPosition::StartAfter)?
220 } else {
221 return Err(ConnectorError::ConfigurationError(
222 "MongoDB CDC checkpoint contains an unknown position key".into(),
223 ));
224 };
225 Ok(ParsedMongoCheckpoint {
226 position,
227 collection_uuid,
228 deployment_identity,
229 })
230}
231
232enum BufferedMongoPayload {
233 Event(Box<MongoDbChangeEvent>),
234 HighWatermark {
235 token: String,
236 requires_start_after: bool,
237 },
238}
239
240struct BufferedMongoEvent {
241 payload: BufferedMongoPayload,
242 _byte_permit: OwnedSemaphorePermit,
243}
244
245impl BufferedMongoEvent {
246 fn new(event: MongoDbChangeEvent, byte_permit: OwnedSemaphorePermit) -> Self {
247 Self {
248 payload: BufferedMongoPayload::Event(Box::new(event)),
249 _byte_permit: byte_permit,
250 }
251 }
252
253 fn high_watermark(
254 token: String,
255 requires_start_after: bool,
256 byte_permit: OwnedSemaphorePermit,
257 ) -> Self {
258 Self {
259 payload: BufferedMongoPayload::HighWatermark {
260 token,
261 requires_start_after,
262 },
263 _byte_permit: byte_permit,
264 }
265 }
266
267 fn event(&self) -> Option<&MongoDbChangeEvent> {
268 match &self.payload {
269 BufferedMongoPayload::Event(event) => Some(event),
270 BufferedMongoPayload::HighWatermark { .. } => None,
271 }
272 }
273}
274
275fn checked_size_add(total: &mut usize, value: usize) -> Result<(), ConnectorError> {
276 *total = total.checked_add(value).ok_or_else(|| {
277 ConnectorError::ConfigurationError("MongoDB CDC event size overflow".into())
278 })?;
279 Ok(())
280}
281
282fn json_retained_bytes(value: &serde_json::Value) -> Result<usize, ConnectorError> {
283 let mut total = size_of::<serde_json::Value>();
284 match value {
285 serde_json::Value::String(value) => checked_size_add(&mut total, value.capacity())?,
286 serde_json::Value::Array(values) => {
287 checked_size_add(
288 &mut total,
289 values
290 .capacity()
291 .checked_mul(size_of::<serde_json::Value>())
292 .ok_or_else(|| {
293 ConnectorError::ConfigurationError(
294 "MongoDB CDC JSON array size overflow".into(),
295 )
296 })?,
297 )?;
298 for value in values {
299 checked_size_add(&mut total, json_retained_bytes(value)?)?;
300 }
301 }
302 serde_json::Value::Object(values) => {
303 for (key, value) in values {
306 checked_size_add(&mut total, size_of::<String>())?;
307 checked_size_add(&mut total, key.capacity())?;
308 checked_size_add(&mut total, json_retained_bytes(value)?)?;
309 }
310 }
311 serde_json::Value::Null | serde_json::Value::Bool(_) | serde_json::Value::Number(_) => {}
312 }
313 Ok(total)
314}
315
316fn mongo_event_retained_bytes(event: &MongoDbChangeEvent) -> Result<usize, ConnectorError> {
317 let mut total = size_of::<BufferedMongoEvent>();
318 if let OperationType::Other(value) = &event.operation_type {
319 checked_size_add(&mut total, value.capacity())?;
320 }
321 checked_size_add(&mut total, event.namespace.db.capacity())?;
322 checked_size_add(&mut total, event.namespace.coll.capacity())?;
323 checked_size_add(&mut total, event.document_key.capacity())?;
324 checked_size_add(
325 &mut total,
326 event.full_document.as_ref().map_or(0, String::capacity),
327 )?;
328 checked_size_add(&mut total, event.resume_token.capacity())?;
329
330 if let Some(update) = &event.update_description {
331 checked_size_add(
332 &mut total,
333 update
334 .updated_fields
335 .capacity()
336 .checked_mul(size_of::<(String, serde_json::Value)>())
337 .ok_or_else(|| {
338 ConnectorError::ConfigurationError(
339 "MongoDB CDC update field size overflow".into(),
340 )
341 })?,
342 )?;
343 for (key, value) in &update.updated_fields {
344 checked_size_add(&mut total, key.capacity())?;
345 checked_size_add(&mut total, json_retained_bytes(value)?)?;
346 }
347 checked_size_add(
348 &mut total,
349 update
350 .removed_fields
351 .capacity()
352 .checked_mul(size_of::<String>())
353 .ok_or_else(|| {
354 ConnectorError::ConfigurationError(
355 "MongoDB CDC removed field size overflow".into(),
356 )
357 })?,
358 )?;
359 for field in &update.removed_fields {
360 checked_size_add(&mut total, field.capacity())?;
361 }
362 checked_size_add(
363 &mut total,
364 update
365 .truncated_arrays
366 .capacity()
367 .checked_mul(size_of::<super::change_event::TruncatedArray>())
368 .ok_or_else(|| {
369 ConnectorError::ConfigurationError(
370 "MongoDB CDC truncated array size overflow".into(),
371 )
372 })?,
373 )?;
374 for array in &update.truncated_arrays {
375 checked_size_add(&mut total, array.field.capacity())?;
376 }
377 checked_size_add(
378 &mut total,
379 update
380 .disambiguated_paths
381 .capacity()
382 .checked_mul(size_of::<(String, serde_json::Value)>())
383 .ok_or_else(|| {
384 ConnectorError::ConfigurationError(
385 "MongoDB CDC disambiguated path size overflow".into(),
386 )
387 })?,
388 )?;
389 for (key, value) in &update.disambiguated_paths {
390 checked_size_add(&mut total, key.capacity())?;
391 checked_size_add(&mut total, json_retained_bytes(value)?)?;
392 }
393 }
394 Ok(total.max(1))
395}
396
397fn mongo_high_watermark_retained_bytes(token_capacity: usize) -> Result<usize, ConnectorError> {
398 let mut total = size_of::<BufferedMongoEvent>();
399 checked_size_add(&mut total, token_capacity)?;
400 Ok(total)
401}
402
403#[must_use]
417pub fn mongodb_cdc_envelope_schema() -> SchemaRef {
418 Arc::new(Schema::new(vec![
419 Field::new("_namespace", DataType::Utf8, false),
420 Field::new("_op", DataType::Utf8, false),
421 Field::new("_document_key", DataType::Utf8, false),
422 Field::new("_cluster_time_s", DataType::UInt32, false),
423 Field::new("_cluster_time_i", DataType::UInt32, false),
424 Field::new(
425 "_wall_time_ms",
426 DataType::Timestamp(arrow_schema::TimeUnit::Millisecond, None),
427 false,
428 ),
429 Field::new("_full_document", DataType::Utf8, true),
430 Field::new("_update_desc", DataType::Utf8, true),
431 Field::new("_resume_token", DataType::Utf8, false),
432 ]))
433}
434
435pub struct MongoDbCdcSource {
448 config: MongoDbSourceConfig,
450
451 state: ConnectorState,
453
454 schema: SchemaRef,
456
457 metrics: Arc<MongoDbCdcMetrics>,
459
460 event_buffer: VecDeque<BufferedMongoEvent>,
462
463 checkpoint_resume_token: Option<String>,
466
467 checkpoint_requires_start_after: bool,
469
470 collection_uuid: Option<Uuid>,
472
473 deployment_identity: Option<MongoDeploymentIdentity>,
475
476 byte_budget: Arc<Semaphore>,
478
479 data_ready: Arc<Notify>,
481
482 #[cfg(feature = "mongodb-cdc")]
484 reader_handle: Option<tokio::task::JoinHandle<()>>,
485
486 #[cfg(feature = "mongodb-cdc")]
488 event_rx: Option<ChangeStreamRx>,
489
490 #[cfg(feature = "mongodb-cdc")]
492 reader_shutdown: Option<tokio::sync::watch::Sender<bool>>,
493
494 #[cfg(feature = "mongodb-cdc")]
496 reader_error: Option<tokio::sync::watch::Receiver<Option<MongoReaderFailure>>>,
497
498 task_owner: ConnectorTaskOwner,
500 task_tracker: ConnectorTaskTracker,
501}
502
503impl Drop for MongoDbCdcSource {
504 fn drop(&mut self) {
505 #[cfg(feature = "mongodb-cdc")]
506 if let Some(shutdown) = self.reader_shutdown.take() {
507 shutdown.send_replace(true);
508 }
509 #[cfg(feature = "mongodb-cdc")]
510 if let Some(handle) = self.reader_handle.take() {
511 reap_mongo_reader(handle, &self.task_owner);
512 }
513 }
514}
515
516#[cfg(feature = "mongodb-cdc")]
517fn reap_mongo_reader(handle: tokio::task::JoinHandle<()>, task_owner: &ConnectorTaskOwner) {
518 let Some(reaper_guard) = task_owner.track() else {
519 tracing::warn!("MongoDB CDC task generation was sealed before reader reaping");
520 return;
521 };
522 let Ok(runtime) = tokio::runtime::Handle::try_current() else {
523 drop(reaper_guard);
526 return;
527 };
528 drop(runtime.spawn(async move {
529 let _reaper_guard = reaper_guard;
530 if let Err(error) = handle.await {
531 tracing::debug!(%error, "MongoDB CDC retired reader task reaped");
532 }
533 }));
534}
535
536#[cfg(feature = "mongodb-cdc")]
538type ChangeStreamTx = crossfire::MAsyncTx<crossfire::mpsc::Array<BufferedMongoEvent>>;
539#[cfg(feature = "mongodb-cdc")]
541type ChangeStreamRx = crossfire::AsyncRx<crossfire::mpsc::Array<BufferedMongoEvent>>;
542
543#[cfg(feature = "mongodb-cdc")]
544#[derive(Debug)]
545struct MongoReaderReady {
546 initial_resume_token: Option<String>,
547 collection_uuid: Uuid,
548 deployment_identity: MongoDeploymentIdentity,
549}
550
551#[cfg(feature = "mongodb-cdc")]
552#[derive(Clone, Debug)]
553enum MongoReaderFailure {
554 Configuration(String),
555 Connection(String),
556 Read(String),
557}
558
559#[cfg(feature = "mongodb-cdc")]
560impl MongoReaderFailure {
561 fn from_connector(error: &ConnectorError) -> Self {
562 match error {
563 ConnectorError::ConfigurationError(message) => Self::Configuration(message.clone()),
564 ConnectorError::ConnectionFailed(message) => Self::Connection(message.clone()),
565 ConnectorError::ReadError(message) => Self::Read(message.clone()),
566 error if error.is_transient() => Self::Read(error.to_string()),
567 error => Self::Configuration(error.to_string()),
568 }
569 }
570
571 fn into_connector(self) -> ConnectorError {
572 match self {
573 Self::Configuration(message) => ConnectorError::ConfigurationError(message),
574 Self::Connection(message) => ConnectorError::ConnectionFailed(message),
575 Self::Read(message) => ConnectorError::ReadError(message),
576 }
577 }
578}
579
580#[cfg(feature = "mongodb-cdc")]
581#[derive(Clone, Copy, Debug, PartialEq, Eq)]
582struct MongoCollectionObservation {
583 collection_uuid: Uuid,
584 post_images_enabled: bool,
585}
586
587#[cfg(feature = "mongodb-cdc")]
588#[derive(Clone, Debug, PartialEq, Eq)]
589struct MongoAdmissionObservation {
590 deployment_identity: MongoDeploymentIdentity,
591 collection: MongoCollectionObservation,
592}
593
594#[cfg(feature = "mongodb-cdc")]
595struct MongoReaderAdmissionGuard {
596 shutdown: Option<tokio::sync::watch::Sender<bool>>,
597}
598
599#[cfg(feature = "mongodb-cdc")]
600impl MongoReaderAdmissionGuard {
601 fn new(shutdown: tokio::sync::watch::Sender<bool>) -> Self {
602 Self {
603 shutdown: Some(shutdown),
604 }
605 }
606
607 fn disarm(&mut self) {
608 self.shutdown = None;
609 }
610}
611
612#[cfg(feature = "mongodb-cdc")]
613impl Drop for MongoReaderAdmissionGuard {
614 fn drop(&mut self) {
615 if let Some(shutdown) = self.shutdown.as_ref() {
616 shutdown.send_replace(true);
617 }
618 }
619}
620
621#[cfg(feature = "mongodb-cdc")]
622#[derive(Clone, Debug, PartialEq)]
623enum MongoResumePosition {
624 ResumeAfter(mongodb::change_stream::event::ResumeToken),
625 StartAfter(mongodb::change_stream::event::ResumeToken),
626}
627
628impl MongoDbCdcSource {
629 #[must_use]
631 pub fn new(config: MongoDbSourceConfig, registry: Option<&prometheus::Registry>) -> Self {
632 let byte_budget = Arc::new(Semaphore::new(config.max_buffered_bytes));
633 let (task_owner, task_tracker) = ConnectorTaskOwner::new();
634 Self {
635 byte_budget,
636 config,
637 state: ConnectorState::Created,
638 schema: mongodb_cdc_envelope_schema(),
639 metrics: Arc::new(MongoDbCdcMetrics::new(registry)),
640 event_buffer: VecDeque::new(),
641 checkpoint_resume_token: None,
642 checkpoint_requires_start_after: false,
643 collection_uuid: None,
644 deployment_identity: None,
645 data_ready: Arc::new(Notify::new()),
646 #[cfg(feature = "mongodb-cdc")]
647 reader_handle: None,
648 #[cfg(feature = "mongodb-cdc")]
649 event_rx: None,
650 #[cfg(feature = "mongodb-cdc")]
651 reader_shutdown: None,
652 #[cfg(feature = "mongodb-cdc")]
653 reader_error: None,
654 task_owner,
655 task_tracker,
656 }
657 }
658
659 #[cfg(test)]
660 fn buffered_events(&self) -> usize {
661 self.event_buffer
662 .iter()
663 .filter(|item| item.event().is_some())
664 .count()
665 }
666
667 #[cfg(test)]
669 fn enqueue_event(&mut self, event: MongoDbChangeEvent) -> Result<(), ConnectorError> {
670 let retained_bytes = mongo_event_retained_bytes(&event)?;
671 let byte_permits = u32::try_from(retained_bytes).map_err(|_| {
672 ConnectorError::ConfigurationError(format!(
673 "MongoDB CDC event exceeds the hard byte bound: event={retained_bytes}, limit={}",
674 self.config.max_buffered_bytes
675 ))
676 })?;
677 if retained_bytes > self.config.max_buffered_bytes {
678 return Err(ConnectorError::ConfigurationError(format!(
679 "MongoDB CDC event exceeds the hard byte bound: event={retained_bytes}, limit={}",
680 self.config.max_buffered_bytes
681 )));
682 }
683 let byte_permit = Arc::clone(&self.byte_budget)
684 .try_acquire_many_owned(byte_permits)
685 .map_err(|_| {
686 ConnectorError::ConfigurationError(format!(
687 "MongoDB CDC buffered bytes reached the hard bound: limit={}",
688 self.config.max_buffered_bytes
689 ))
690 })?;
691 self.metrics.record_event(event.operation_type.as_str());
692 self.event_buffer
693 .push_back(BufferedMongoEvent::new(event, byte_permit));
694 Ok(())
695 }
696
697 #[cfg(test)]
698 fn enqueue_high_watermark(&mut self, token: &str) -> Result<(), ConnectorError> {
699 let token = canonical_resume_token(token)?;
700 let retained_bytes = mongo_high_watermark_retained_bytes(token.capacity())?;
701 if retained_bytes > self.config.max_buffered_bytes {
702 return Err(ConnectorError::ConfigurationError(format!(
703 "MongoDB CDC high watermark exceeds the hard byte bound: item={retained_bytes}, \
704 limit={}",
705 self.config.max_buffered_bytes
706 )));
707 }
708 let permits = u32::try_from(retained_bytes).map_err(|_| {
709 ConnectorError::ConfigurationError("MongoDB CDC high watermark is too large".into())
710 })?;
711 let byte_permit = Arc::clone(&self.byte_budget)
712 .try_acquire_many_owned(permits)
713 .map_err(|_| {
714 ConnectorError::ConfigurationError(
715 "MongoDB CDC high watermark exceeded the byte budget".into(),
716 )
717 })?;
718 self.event_buffer
719 .push_back(BufferedMongoEvent::high_watermark(
720 token,
721 false,
722 byte_permit,
723 ));
724 Ok(())
725 }
726
727 fn drain_to_batch(
734 &mut self,
735 max_records: usize,
736 ) -> Result<Option<SourceBatch>, ConnectorError> {
737 if max_records == 0 || self.event_buffer.is_empty() {
738 return Ok(None);
739 }
740
741 let count = max_records.min(self.event_buffer.len());
742 let count = self
745 .event_buffer
746 .iter()
747 .take(count)
748 .position(|item| {
749 item.event()
750 .is_some_and(|event| event.operation_type == OperationType::Invalidate)
751 })
752 .map_or(count, |index| index + 1);
753 let items: Vec<BufferedMongoEvent> = self.event_buffer.drain(..count).collect();
754 let events: Vec<&MongoDbChangeEvent> =
755 items.iter().filter_map(BufferedMongoEvent::event).collect();
756 let (position_token, requires_start_after) = match items.last() {
757 Some(item) => {
758 let (token, start_after) = match &item.payload {
759 BufferedMongoPayload::Event(event) => (
760 event.resume_token.as_str(),
761 event.operation_type == OperationType::Invalidate,
762 ),
763 BufferedMongoPayload::HighWatermark {
764 token,
765 requires_start_after,
766 } => (token.as_str(), *requires_start_after),
767 };
768 match canonical_resume_token(token) {
769 Ok(token) => (token, start_after),
770 Err(error) => {
771 drop(events);
772 for item in items.into_iter().rev() {
773 self.event_buffer.push_front(item);
774 }
775 return Err(error);
776 }
777 }
778 }
779 None => return Ok(None),
780 };
781
782 if events.is_empty() {
783 self.checkpoint_resume_token = Some(position_token);
784 self.checkpoint_requires_start_after = requires_start_after;
785 return Ok(None);
786 }
787
788 let batch = match events_to_record_batch_refs(&events, &self.schema) {
789 Ok(batch) => batch,
790 Err(error) => {
791 drop(events);
792 for item in items.into_iter().rev() {
793 self.event_buffer.push_front(item);
794 }
795 return Err(error);
796 }
797 };
798 self.metrics.record_batch();
799 self.checkpoint_resume_token = Some(position_token);
800 self.checkpoint_requires_start_after = requires_start_after;
801
802 Ok(Some(SourceBatch::new(batch)))
803 }
804}
805
806#[cfg(test)]
808fn events_to_record_batch(
809 events: &[MongoDbChangeEvent],
810 schema: &SchemaRef,
811) -> Result<RecordBatch, ConnectorError> {
812 let events: Vec<&MongoDbChangeEvent> = events.iter().collect();
813 events_to_record_batch_refs(&events, schema)
814}
815
816fn events_to_record_batch_refs(
817 events: &[&MongoDbChangeEvent],
818 schema: &SchemaRef,
819) -> Result<RecordBatch, ConnectorError> {
820 let len = events.len();
821
822 let mut ns_builder = StringBuilder::with_capacity(len, len * 32);
823 let mut op_builder = StringBuilder::with_capacity(len, len * 4);
824 let mut dk_builder = StringBuilder::with_capacity(len, len * 64);
825 let mut cts_builder = UInt32Builder::with_capacity(len);
826 let mut ct_inc_builder = UInt32Builder::with_capacity(len);
827 let mut wt_builder = arrow_array::builder::TimestampMillisecondBuilder::with_capacity(len);
828 let mut fd_builder = StringBuilder::with_capacity(len, len * 128);
829 let mut ud_builder = StringBuilder::with_capacity(len, len * 64);
830 let mut rt_builder = StringBuilder::with_capacity(len, len * 64);
831
832 for event in events {
833 ns_builder.append_value(event.namespace.full_name());
834 op_builder.append_value(event.operation_type.as_str());
835 dk_builder.append_value(&event.document_key);
836 cts_builder.append_value(event.cluster_time_secs);
837 ct_inc_builder.append_value(event.cluster_time_inc);
838 wt_builder.append_value(event.wall_time_ms);
839
840 match &event.full_document {
841 Some(doc) => fd_builder.append_value(doc),
842 None => fd_builder.append_null(),
843 }
844
845 match &event.update_description {
846 Some(desc) => {
847 let json = serde_json::to_string(desc)
848 .map_err(|e| ConnectorError::Internal(format!("serialize update_desc: {e}")))?;
849 ud_builder.append_value(&json);
850 }
851 None => ud_builder.append_null(),
852 }
853
854 rt_builder.append_value(&event.resume_token);
855 }
856
857 RecordBatch::try_new(
858 Arc::clone(schema),
859 vec![
860 Arc::new(ns_builder.finish()),
861 Arc::new(op_builder.finish()),
862 Arc::new(dk_builder.finish()),
863 Arc::new(cts_builder.finish()),
864 Arc::new(ct_inc_builder.finish()),
865 Arc::new(wt_builder.finish()),
866 Arc::new(fd_builder.finish()),
867 Arc::new(ud_builder.finish()),
868 Arc::new(rt_builder.finish()),
869 ],
870 )
871 .map_err(|e| ConnectorError::Internal(format!("arrow batch: {e}")))
872}
873
874#[async_trait]
875impl SourceConnector for MongoDbCdcSource {
876 fn terminal_task_tracker(&self) -> Option<ConnectorTaskTracker> {
877 Some(self.task_tracker.clone())
878 }
879
880 fn recovery_identity_options(
881 &self,
882 config: &ConnectorConfig,
883 ) -> Result<Option<BTreeMap<String, String>>, ConnectorError> {
884 let mut parsed = if config.properties().is_empty() {
885 self.config.clone()
886 } else {
887 MongoDbSourceConfig::from_config(config)?
888 };
889 parsed.normalize_pipeline()?;
890 parsed.validate()?;
891 let pipeline = super::config::canonical_pipeline_json(&parsed.pipeline);
892
893 Ok(Some(BTreeMap::from([
894 ("collection".into(), parsed.collection),
895 ("database".into(), parsed.database),
896 (
897 "full.document.mode".into(),
898 parsed.full_document_mode.to_string(),
899 ),
900 ("pipeline".into(), pipeline),
901 ("wire.protocol".into(), "change-stream-expanded-v1".into()),
902 ])))
903 }
904
905 async fn start(&mut self, request: SourceStart) -> Result<(), ConnectorError> {
906 if self.state != ConnectorState::Created {
907 return Err(ConnectorError::InvalidState {
908 expected: ConnectorState::Created.to_string(),
909 actual: self.state.to_string(),
910 });
911 }
912 let (config, position, _) = request.into_parts();
913 let parsed_config = if config.properties().is_empty() {
914 let mut config = self.config.clone();
915 config.normalize_pipeline()?;
916 config.validate()?;
917 config
918 } else {
919 MongoDbSourceConfig::from_config(&config)?
920 };
921 let (
922 checkpoint_resume_token,
923 checkpoint_requires_start_after,
924 initial_resume_position,
925 expected_collection_uuid,
926 expected_deployment_identity,
927 ) = match position {
928 SourcePosition::Initial => (None, false, None, None, None),
929 SourcePosition::Resume {
930 attempt,
931 checkpoint,
932 } => {
933 let ParsedMongoCheckpoint {
934 position,
935 collection_uuid,
936 deployment_identity,
937 } = parse_mongodb_checkpoint(&checkpoint, &parsed_config).map_err(|error| {
938 ConnectorError::ConfigurationError(format!(
939 "invalid MongoDB CDC checkpoint {attempt:?}: {error}"
940 ))
941 })?;
942 match position {
943 MongoCheckpointPosition::ResumeAfter(token) => {
944 let driver_token = serde_json::from_str(&token).map_err(|error| {
945 ConnectorError::ConfigurationError(format!(
946 "invalid MongoDB CDC resume token in checkpoint {attempt:?}: \
947 {error}"
948 ))
949 })?;
950 (
951 Some(token),
952 false,
953 Some(MongoResumePosition::ResumeAfter(driver_token)),
954 Some(collection_uuid),
955 Some(deployment_identity),
956 )
957 }
958 MongoCheckpointPosition::StartAfter(token) => {
959 let driver_token = serde_json::from_str(&token).map_err(|error| {
960 ConnectorError::ConfigurationError(format!(
961 "invalid MongoDB CDC start-after token in checkpoint {attempt:?}: \
962 {error}"
963 ))
964 })?;
965 (
966 Some(token),
967 true,
968 Some(MongoResumePosition::StartAfter(driver_token)),
969 Some(collection_uuid),
970 Some(deployment_identity),
971 )
972 }
973 }
974 }
975 };
976
977 self.start_change_stream_reader(
978 parsed_config,
979 checkpoint_resume_token,
980 checkpoint_requires_start_after,
981 initial_resume_position,
982 expected_collection_uuid,
983 expected_deployment_identity,
984 )
985 .await?;
986
987 self.state = ConnectorState::Running;
988 tracing::info!(
989 database = %self.config.database,
990 collection = %self.config.collection,
991 full_document_mode = ?self.config.full_document_mode,
992 "MongoDB CDC source opened"
993 );
994
995 Ok(())
996 }
997
998 async fn poll_batch(
999 &mut self,
1000 max_records: usize,
1001 ) -> Result<Option<SourceBatch>, ConnectorError> {
1002 self.drain_channel(max_records.saturating_sub(self.event_buffer.len()));
1003 if let Some(batch) = self.drain_to_batch(max_records)? {
1004 return Ok(Some(batch));
1005 }
1006 self.check_reader_error()?;
1007 Ok(None)
1008 }
1009
1010 fn schema(&self) -> SchemaRef {
1011 Arc::clone(&self.schema)
1012 }
1013
1014 fn checkpoint(&self) -> SourceCheckpoint {
1015 let mut checkpoint = SourceCheckpoint::new();
1016 let Some(collection_uuid) = self.collection_uuid else {
1017 return checkpoint;
1020 };
1021 let Some(deployment_identity) = self.deployment_identity.as_ref() else {
1022 return checkpoint;
1023 };
1024 if let Some(token) = self.checkpoint_resume_token.as_ref() {
1025 checkpoint.set_offset(
1026 if self.checkpoint_requires_start_after {
1027 START_AFTER_TOKEN_OFFSET
1028 } else {
1029 RESUME_TOKEN_OFFSET
1030 },
1031 token,
1032 );
1033 } else {
1034 return checkpoint;
1036 }
1037 checkpoint.set_metadata("connector", MONGODB_CHECKPOINT_CONNECTOR);
1038 checkpoint.set_metadata("version", MONGODB_CHECKPOINT_VERSION);
1039 checkpoint.set_metadata("database", &self.config.database);
1040 checkpoint.set_metadata("collection", &self.config.collection);
1041 checkpoint.set_metadata(
1042 COLLECTION_UUID_METADATA,
1043 collection_uuid.hyphenated().to_string(),
1044 );
1045 checkpoint.set_metadata(DEPLOYMENT_IDENTITY_METADATA, deployment_identity.encode());
1046 checkpoint.set_metadata(
1047 STREAM_IDENTITY_METADATA,
1048 mongodb_stream_identity(&self.config),
1049 );
1050 checkpoint
1051 }
1052
1053 async fn close(&mut self) -> Result<(), ConnectorError> {
1054 #[cfg(feature = "mongodb-cdc")]
1055 let mut reader_join_error = None;
1056 #[cfg(feature = "mongodb-cdc")]
1057 {
1058 if let Some(tx) = self.reader_shutdown.as_ref() {
1059 tx.send_replace(true);
1060 }
1061 let mut detach_reader = false;
1062 if let Some(handle) = self.reader_handle.as_mut() {
1063 match tokio::time::timeout(READER_SHUTDOWN_TIMEOUT, &mut *handle).await {
1064 Ok(Ok(())) => {}
1065 Ok(Err(error)) if error.is_cancelled() => {}
1066 Ok(Err(error)) => reader_join_error = Some(error.to_string()),
1067 Err(_) => {
1068 tracing::warn!(
1069 "MongoDB CDC reader exceeded its close deadline; its tracked reaper retains shutdown ownership"
1070 );
1071 detach_reader = true;
1072 }
1073 }
1074 }
1075 if detach_reader {
1076 let handle = self
1077 .reader_handle
1078 .take()
1079 .expect("reader handle was present while awaiting it");
1080 reap_mongo_reader(handle, &self.task_owner);
1081 } else {
1082 self.reader_handle = None;
1083 }
1084 self.reader_shutdown = None;
1085 self.event_rx = None;
1086 self.reader_error = None;
1087 }
1088
1089 self.event_buffer.clear();
1090 self.state = ConnectorState::Closed;
1091 tracing::info!("MongoDB CDC source closed");
1092 #[cfg(feature = "mongodb-cdc")]
1093 if let Some(error) = reader_join_error {
1094 return Err(ConnectorError::ReadError(format!(
1095 "MongoDB CDC reader task failed during close: {error}"
1096 )));
1097 }
1098 Ok(())
1099 }
1100
1101 fn data_ready_notify(&self) -> Option<Arc<Notify>> {
1102 Some(Arc::clone(&self.data_ready))
1103 }
1104
1105 fn contract(&self, config: &ConnectorConfig) -> Result<SourceContract, ConnectorError> {
1106 if config.properties().is_empty() {
1107 self.config.validate()?;
1108 } else {
1109 MongoDbSourceConfig::from_config(config)?;
1110 }
1111 Ok(SourceContract::new(
1112 SourceConsistency::Replayable,
1113 SourceTopology::Singleton,
1114 ))
1115 }
1116}
1117
1118#[cfg(feature = "mongodb-cdc")]
1121fn clamp_source_startup_timeout(configured: Option<std::time::Duration>) -> std::time::Duration {
1122 configured
1123 .filter(|timeout| !timeout.is_zero())
1124 .map_or(READER_STARTUP_TIMEOUT, |timeout| {
1125 timeout.min(READER_STARTUP_TIMEOUT)
1126 })
1127}
1128
1129#[cfg(feature = "mongodb-cdc")]
1130async fn source_client_options(
1131 connection_uri: &str,
1132) -> Result<mongodb::options::ClientOptions, ConnectorError> {
1133 let mut options = mongodb::options::ClientOptions::parse(connection_uri)
1134 .await
1135 .map_err(|error| ConnectorError::ConfigurationError(format!("parse URI: {error}")))?;
1136 super::sink::harden_mongodb_tls(&mut options)?;
1137 options.connect_timeout = Some(clamp_source_startup_timeout(options.connect_timeout));
1138 options.server_selection_timeout = Some(clamp_source_startup_timeout(
1139 options.server_selection_timeout,
1140 ));
1141
1142 if let Some(pool) = options.max_pool_size {
1143 if pool <= 1 {
1144 tracing::warn!(
1145 max_pool_size = pool,
1146 "max_pool_size is very small; mongos may exhaust per-shard cursors"
1147 );
1148 }
1149 }
1150 Ok(options)
1151}
1152
1153#[cfg(feature = "mongodb-cdc")]
1154async fn source_database(
1155 connection_uri: &str,
1156 database: &str,
1157) -> Result<mongodb::Database, ConnectorError> {
1158 let options = source_client_options(connection_uri).await?;
1159 let client = mongodb::Client::with_options(options)
1160 .map_err(|error| ConnectorError::ConfigurationError(format!("create client: {error}")))?;
1161 Ok(client.database(database))
1162}
1163
1164#[cfg(feature = "mongodb-cdc")]
1165async fn await_mongo_reader_ready(
1166 ready_rx: tokio::sync::oneshot::Receiver<Result<MongoReaderReady, MongoReaderFailure>>,
1167 shutdown_tx: &tokio::sync::watch::Sender<bool>,
1168 handle: &mut tokio::task::JoinHandle<()>,
1169) -> Result<MongoReaderReady, ConnectorError> {
1170 let (error, include_join_error) =
1171 match tokio::time::timeout(READER_STARTUP_TIMEOUT, ready_rx).await {
1172 Ok(Ok(Ok(ready))) => return Ok(ready),
1173 Ok(Ok(Err(error))) => (error.into_connector(), false),
1174 Ok(Err(_)) => (
1175 ConnectorError::ReadError(
1176 "MongoDB CDC reader exited before opening the change stream".into(),
1177 ),
1178 true,
1179 ),
1180 Err(_) => (
1181 ConnectorError::ReadError(format!(
1182 "MongoDB CDC did not open a change stream within the {READER_STARTUP_TIMEOUT:?} startup deadline"
1183 )),
1184 false,
1185 ),
1186 };
1187
1188 shutdown_tx.send_replace(true);
1189 let Ok(join_result) = tokio::time::timeout(READER_SHUTDOWN_TIMEOUT, &mut *handle).await else {
1190 tracing::warn!(
1191 "MongoDB CDC admission reader exceeded its shutdown deadline; the retired generation remains tracked until it exits"
1192 );
1193 return Err(error);
1194 };
1195 let error = if include_join_error {
1196 match join_result {
1197 Err(join_error) => ConnectorError::ReadError(format!("{error}: {join_error}")),
1198 _ => error,
1199 }
1200 } else {
1201 error
1202 };
1203 Err(error)
1204}
1205
1206#[cfg(feature = "mongodb-cdc")]
1207impl MongoDbCdcSource {
1208 async fn start_change_stream_reader(
1210 &mut self,
1211 config: MongoDbSourceConfig,
1212 checkpoint_resume_token: Option<String>,
1213 checkpoint_requires_start_after: bool,
1214 initial_resume_position: Option<MongoResumePosition>,
1215 expected_collection_uuid: Option<Uuid>,
1216 expected_deployment_identity: Option<MongoDeploymentIdentity>,
1217 ) -> Result<(), ConnectorError> {
1218 if self.reader_handle.is_some() {
1219 return Err(ConnectorError::InvalidState {
1220 expected: ConnectorState::Created.to_string(),
1221 actual: "reader already started".into(),
1222 });
1223 }
1224 if !self.event_buffer.is_empty() {
1225 return Err(ConnectorError::ConfigurationError(
1226 "MongoDB CDC cannot start with pre-buffered test events".into(),
1227 ));
1228 }
1229 let max_buffered_bytes = config.max_buffered_bytes;
1230 let byte_budget = Arc::new(Semaphore::new(max_buffered_bytes));
1231
1232 let channel_capacity = config.reader_channel_capacity();
1233 let (tx, rx) = crossfire::mpsc::bounded_async::<BufferedMongoEvent>(channel_capacity);
1234 let (shutdown_tx, shutdown_rx) = tokio::sync::watch::channel(false);
1235 let (error_tx, error_rx) = tokio::sync::watch::channel(None);
1236 let (ready_tx, ready_rx) =
1237 tokio::sync::oneshot::channel::<Result<MongoReaderReady, MongoReaderFailure>>();
1238 let reader_config = config.clone();
1239 let data_ready = Arc::clone(&self.data_ready);
1240 let terminal_ready = Arc::clone(&self.data_ready);
1241 let metrics = Arc::clone(&self.metrics);
1242 let task_byte_budget = Arc::clone(&byte_budget);
1243
1244 let reader_guard = self.task_owner.track().ok_or_else(|| {
1245 ConnectorError::Internal("MongoDB CDC connector generation is already retired".into())
1246 })?;
1247
1248 let mut handle = tokio::spawn(async move {
1249 let _reader_guard = reader_guard;
1250 let result = async {
1251 let db =
1252 match source_database(&reader_config.connection_uri, &reader_config.database)
1253 .await
1254 {
1255 Ok(database) => database,
1256 Err(error) => {
1257 let admission_error = MongoReaderFailure::from_connector(&error);
1258 let _ = ready_tx.send(Err(admission_error));
1259 return Err(error);
1260 }
1261 };
1262 run_change_stream_reader(
1263 db,
1264 reader_config,
1265 tx,
1266 shutdown_rx,
1267 data_ready,
1268 metrics,
1269 task_byte_budget,
1270 max_buffered_bytes,
1271 initial_resume_position,
1272 expected_collection_uuid,
1273 expected_deployment_identity,
1274 ready_tx,
1275 )
1276 .await
1277 }
1278 .await;
1279 if let Err(e) = result {
1280 tracing::error!(error = %e, "change stream reader task failed");
1281 error_tx.send_replace(Some(MongoReaderFailure::from_connector(&e)));
1282 terminal_ready.notify_one();
1283 }
1284 });
1285 let mut admission_guard = MongoReaderAdmissionGuard::new(shutdown_tx.clone());
1286
1287 let ready = await_mongo_reader_ready(ready_rx, &shutdown_tx, &mut handle).await?;
1288
1289 admission_guard.disarm();
1290 self.config = config;
1291 self.checkpoint_resume_token = ready.initial_resume_token.or(checkpoint_resume_token);
1292 self.checkpoint_requires_start_after = checkpoint_requires_start_after;
1293 self.collection_uuid = Some(ready.collection_uuid);
1294 self.deployment_identity = Some(ready.deployment_identity);
1295 self.byte_budget = byte_budget;
1296 self.reader_handle = Some(handle);
1297 self.event_rx = Some(rx);
1298 self.reader_shutdown = Some(shutdown_tx);
1299 self.reader_error = Some(error_rx);
1300 Ok(())
1301 }
1302
1303 fn drain_channel(&mut self, max_events: usize) {
1305 let max_events = max_events.min(self.config.reader_channel_capacity());
1306 for _ in 0..max_events {
1307 let item = {
1308 let Some(rx) = self.event_rx.as_mut() else {
1309 break;
1310 };
1311 let Ok(item) = rx.try_recv() else {
1312 break;
1313 };
1314 item
1315 };
1316 if let Some(event) = item.event() {
1317 self.metrics.record_event(event.operation_type.as_str());
1318 }
1319 self.event_buffer.push_back(item);
1320 }
1321 }
1322
1323 fn check_reader_error(&mut self) -> Result<(), ConnectorError> {
1324 let error = self
1325 .reader_error
1326 .as_mut()
1327 .and_then(|receiver| receiver.borrow_and_update().clone());
1328 if let Some(error) = error {
1329 self.metrics.record_error();
1330 return Err(error.into_connector());
1331 }
1332 Ok(())
1333 }
1334}
1335
1336#[cfg(feature = "mongodb-cdc")]
1337fn mongodb_identity_command_is_permanent(code: i32, code_name: &str) -> bool {
1338 matches!(
1339 code,
1340 13 | 18 | 20 | 26 | 59 | 72 | 76 | 115 | 123 | 323 | 8000
1341 ) || matches!(
1342 code_name,
1343 "Unauthorized"
1344 | "AuthenticationFailed"
1345 | "IllegalOperation"
1346 | "NamespaceNotFound"
1347 | "CommandNotFound"
1348 | "InvalidOptions"
1349 | "NoReplicationEnabled"
1350 | "CommandNotSupported"
1351 | "NotAReplicaSet"
1352 | "APIStrictError"
1353 | "AtlasError"
1354 )
1355}
1356
1357#[cfg(feature = "mongodb-cdc")]
1358fn mongodb_identity_probe_is_permanent(error: &mongodb::error::Error) -> bool {
1359 match error.kind.as_ref() {
1360 mongodb::error::ErrorKind::Authentication { .. } => true,
1361 mongodb::error::ErrorKind::Command(command) => {
1362 mongodb_identity_command_is_permanent(command.code, &command.code_name)
1363 }
1364 _ => false,
1365 }
1366}
1367
1368#[cfg(feature = "mongodb-cdc")]
1369async fn observe_mongodb_deployment(
1370 db: &mongodb::Database,
1371) -> Result<MongoDeploymentIdentity, ConnectorError> {
1372 let hello = db
1373 .run_command(mongodb::bson::doc! { "hello": 1 })
1374 .await
1375 .map_err(|error| {
1376 if mongodb_identity_probe_is_permanent(&error) {
1377 return ConnectorError::ConfigurationError(format!(
1378 "MongoDB CDC cannot inspect deployment topology; verify credentials and \
1379 deployment command support: {error}"
1380 ));
1381 }
1382 ConnectorError::ConnectionFailed(format!(
1383 "inspect MongoDB deployment topology with hello: {error}"
1384 ))
1385 })?;
1386
1387 if hello.get_str("msg").ok() == Some("isdbgrid") {
1388 let version = db
1389 .client()
1390 .database("config")
1391 .collection::<mongodb::bson::Document>("version")
1392 .find_one(mongodb::bson::doc! { "_id": 1 })
1393 .projection(mongodb::bson::doc! { "clusterId": 1 })
1394 .await
1395 .map_err(|error| {
1396 if mongodb_identity_probe_is_permanent(&error) {
1397 ConnectorError::ConfigurationError(format!(
1398 "MongoDB CDC requires read access to config.version {{_id: 1}}.clusterId \
1399 to bind checkpoints to the sharded cluster identity: {error}"
1400 ))
1401 } else {
1402 ConnectorError::ConnectionFailed(format!(
1403 "read MongoDB sharded cluster identity from config.version: {error}"
1404 ))
1405 }
1406 })?
1407 .ok_or_else(|| {
1408 ConnectorError::ConfigurationError(
1409 "MongoDB config.version {_id: 1} is missing; cannot bind CDC checkpoints to \
1410 this sharded cluster"
1411 .into(),
1412 )
1413 })?;
1414 let cluster_id = version.get_object_id("clusterId").map_err(|error| {
1415 ConnectorError::ConfigurationError(format!(
1416 "MongoDB config.version.clusterId is missing or not an ObjectId: {error}"
1417 ))
1418 })?;
1419 return Ok(MongoDeploymentIdentity::ShardedCluster(cluster_id.to_hex()));
1420 }
1421
1422 if hello.get_str("setName").is_ok() {
1423 let response = db
1424 .client()
1425 .database("admin")
1426 .run_command(mongodb::bson::doc! { "replSetGetConfig": 1 })
1427 .await
1428 .map_err(|error| {
1429 if mongodb_identity_probe_is_permanent(&error) {
1430 ConnectorError::ConfigurationError(format!(
1431 "MongoDB CDC requires replSetGetConfig access to bind checkpoints to the \
1432 replica-set identity; Atlas M0 and Flex tiers do not support this \
1433 command: {error}"
1434 ))
1435 } else {
1436 ConnectorError::ConnectionFailed(format!(
1437 "read MongoDB replica-set identity with replSetGetConfig: {error}"
1438 ))
1439 }
1440 })?;
1441 let replica_set_id = response
1442 .get_document("config")
1443 .and_then(|config| config.get_document("settings"))
1444 .and_then(|settings| settings.get_object_id("replicaSetId"))
1445 .map_err(|error| {
1446 ConnectorError::ConfigurationError(format!(
1447 "MongoDB replSetGetConfig omitted settings.replicaSetId: {error}"
1448 ))
1449 })?;
1450 return Ok(MongoDeploymentIdentity::ReplicaSet(replica_set_id.to_hex()));
1451 }
1452
1453 Err(ConnectorError::ConfigurationError(
1454 "MongoDB CDC requires a replica set or sharded cluster; hello reported neither topology"
1455 .into(),
1456 ))
1457}
1458
1459#[cfg(feature = "mongodb-cdc")]
1460async fn observe_mongodb_admission(
1461 db: &mongodb::Database,
1462 database: &str,
1463 collection: &str,
1464) -> Result<MongoAdmissionObservation, ConnectorError> {
1465 let (deployment_identity, collection) = tokio::try_join!(
1466 observe_mongodb_deployment(db),
1467 observe_mongodb_collection(db, database, collection),
1468 )?;
1469 Ok(MongoAdmissionObservation {
1470 deployment_identity,
1471 collection,
1472 })
1473}
1474
1475#[cfg(feature = "mongodb-cdc")]
1477async fn observe_mongodb_collection(
1478 db: &mongodb::Database,
1479 database: &str,
1480 collection: &str,
1481) -> Result<MongoCollectionObservation, ConnectorError> {
1482 let mut cursor = db
1483 .list_collections()
1484 .filter(mongodb::bson::doc! { "name": collection })
1485 .batch_size(1)
1486 .await
1487 .map_err(|error| {
1488 if mongodb_identity_probe_is_permanent(&error) {
1489 ConnectorError::ConfigurationError(format!(
1490 "MongoDB CDC requires database-scoped listCollections access to bind \
1491 {database}.{collection} to its collection UUID: {error}"
1492 ))
1493 } else {
1494 ConnectorError::ConnectionFailed(format!(
1495 "inspect MongoDB collection {database}.{collection}: {error}"
1496 ))
1497 }
1498 })?;
1499 let spec = cursor
1500 .try_next()
1501 .await
1502 .map_err(|error| {
1503 ConnectorError::ConnectionFailed(format!(
1504 "read MongoDB collection identity for {database}.{collection}: {error}"
1505 ))
1506 })?
1507 .ok_or_else(|| {
1508 ConnectorError::ConfigurationError(format!(
1509 "MongoDB CDC collection {database}.{collection} does not exist; create the fixed \
1510 collection before starting the source"
1511 ))
1512 })?;
1513
1514 match spec.collection_type {
1515 mongodb::results::CollectionType::Collection => {}
1516 mongodb::results::CollectionType::Timeseries => {
1517 return Err(ConnectorError::ConfigurationError(format!(
1518 "time series collection {database}.{collection} does not support change streams"
1519 )));
1520 }
1521 mongodb::results::CollectionType::View => {
1522 return Err(ConnectorError::ConfigurationError(format!(
1523 "MongoDB CDC source {database}.{collection} must be a collection, not a view"
1524 )));
1525 }
1526 _ => {
1527 return Err(ConnectorError::ConfigurationError(format!(
1528 "MongoDB CDC source {database}.{collection} has an unsupported collection type"
1529 )));
1530 }
1531 }
1532
1533 let post_images_enabled = spec
1534 .options
1535 .change_stream_pre_and_post_images
1536 .as_ref()
1537 .is_some_and(|options| options.enabled);
1538 let binary = spec.info.uuid.ok_or_else(|| {
1539 ConnectorError::ConfigurationError(format!(
1540 "MongoDB collection {database}.{collection} did not expose an immutable UUID"
1541 ))
1542 })?;
1543 if binary.subtype != mongodb::bson::spec::BinarySubtype::Uuid || binary.bytes.len() != 16 {
1544 return Err(ConnectorError::ConfigurationError(format!(
1545 "MongoDB collection {database}.{collection} returned a non-standard collection UUID"
1546 )));
1547 }
1548 let collection_uuid = Uuid::from_slice(&binary.bytes).map_err(|error| {
1549 ConnectorError::ConfigurationError(format!(
1550 "invalid UUID for MongoDB collection {database}.{collection}: {error}"
1551 ))
1552 })?;
1553 Ok(MongoCollectionObservation {
1554 collection_uuid,
1555 post_images_enabled,
1556 })
1557}
1558
1559#[cfg(feature = "mongodb-cdc")]
1561const MAX_FAILURES: u32 = 10;
1562
1563#[cfg(feature = "mongodb-cdc")]
1564const READER_SHUTDOWN_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(5);
1565
1566#[cfg(feature = "mongodb-cdc")]
1567enum ChangeStreamRead {
1568 Stop,
1569 Reconnect,
1570}
1571
1572#[cfg(feature = "mongodb-cdc")]
1573fn reader_stopping(shutdown_rx: &tokio::sync::watch::Receiver<bool>) -> bool {
1574 *shutdown_rx.borrow() || shutdown_rx.has_changed().is_err()
1575}
1576
1577#[cfg(feature = "mongodb-cdc")]
1578fn change_stream_options(
1579 config: &MongoDbSourceConfig,
1580 position: Option<&MongoResumePosition>,
1581) -> mongodb::options::ChangeStreamOptions {
1582 let mut options = mongodb::options::ChangeStreamOptions::default();
1583 options.full_document = match config.full_document_mode {
1584 super::config::FullDocumentMode::Delta => None,
1585 super::config::FullDocumentMode::RequirePostImage => {
1586 Some(mongodb::options::FullDocumentType::Required)
1587 }
1588 };
1589 options.max_await_time = Some(CURSOR_MAX_AWAIT_TIME);
1590 options.batch_size = Some(config.cursor_batch_size());
1591 options.show_expanded_events = Some(true);
1592 match position {
1593 Some(MongoResumePosition::ResumeAfter(token)) => options.resume_after = Some(token.clone()),
1594 Some(MongoResumePosition::StartAfter(token)) => options.start_after = Some(token.clone()),
1595 None => {}
1596 }
1597 options
1598}
1599
1600#[cfg(feature = "mongodb-cdc")]
1601fn bootstrap_change_stream_options(
1602 config: &MongoDbSourceConfig,
1603) -> mongodb::options::ChangeStreamOptions {
1604 let mut options = change_stream_options(config, None);
1605 options.batch_size = Some(0);
1608 options
1609}
1610
1611#[cfg(feature = "mongodb-cdc")]
1612async fn forward_change_stream(
1613 cursor: &mut mongodb::change_stream::ChangeStream<
1614 mongodb::change_stream::event::ChangeStreamEvent<mongodb::bson::Document>,
1615 >,
1616 shutdown_rx: &mut tokio::sync::watch::Receiver<bool>,
1617 resume_position: &mut Option<MongoResumePosition>,
1618 tx: &ChangeStreamTx,
1619 data_ready: &Notify,
1620 consecutive_failures: &mut u32,
1621 byte_budget: &Arc<Semaphore>,
1622 max_buffered_bytes: usize,
1623 metrics: &MongoDbCdcMetrics,
1624) -> Result<ChangeStreamRead, ConnectorError> {
1625 loop {
1626 if reader_stopping(shutdown_rx) {
1627 tracing::info!("change stream reader shutting down");
1628 return Ok(ChangeStreamRead::Stop);
1629 }
1630
1631 let next = cursor.next_if_any().await;
1634 if reader_stopping(shutdown_rx) {
1635 tracing::info!("change stream reader shutting down after completed getMore");
1636 return Ok(ChangeStreamRead::Stop);
1637 }
1638
1639 match next {
1640 Ok(Some(event)) => {
1641 *consecutive_failures = 0;
1642 let event_token = event.id.clone();
1643 let wire_bytes = mongodb::bson::to_vec(&event)
1644 .map_err(|error| {
1645 ConnectorError::ReadError(format!(
1646 "serialize change stream event for byte accounting: {error}"
1647 ))
1648 })?
1649 .len();
1650 if wire_bytes > MAX_MONGODB_WIRE_EVENT_BYTES {
1651 return Err(ConnectorError::ReadError(format!(
1652 "MongoDB CDC wire event exceeds the supported unsplit BSON bound: \
1653 event={wire_bytes}, limit={MAX_MONGODB_WIRE_EVENT_BYTES}"
1654 )));
1655 }
1656 metrics.record_bytes(u64::try_from(wire_bytes).unwrap_or(u64::MAX));
1657 let change_event = parse_change_stream_event(&event)?;
1658 let invalidated = change_event.operation_type == OperationType::Invalidate;
1659 let Some(change_event) = acquire_mongo_event_ownership(
1660 change_event,
1661 byte_budget,
1662 max_buffered_bytes,
1663 shutdown_rx,
1664 )
1665 .await?
1666 else {
1667 return Ok(ChangeStreamRead::Stop);
1668 };
1669 if !send_event_or_shutdown(tx, change_event, shutdown_rx).await {
1670 return Ok(ChangeStreamRead::Stop);
1671 }
1672 *resume_position = Some(if invalidated {
1673 MongoResumePosition::StartAfter(event_token)
1674 } else {
1675 MongoResumePosition::ResumeAfter(cursor.resume_token().unwrap_or(event_token))
1676 });
1677 data_ready.notify_one();
1678 if invalidated {
1679 return Ok(ChangeStreamRead::Reconnect);
1680 }
1681 }
1682 Ok(None) => {
1683 let cursor_alive = cursor.is_alive();
1684 if !matches!(
1685 resume_position.as_ref(),
1686 Some(MongoResumePosition::StartAfter(_))
1687 ) {
1688 if let Some(token) = cursor.resume_token() {
1689 let requires_start_after = !cursor_alive;
1690 let changed = match resume_position.as_ref() {
1691 Some(MongoResumePosition::ResumeAfter(current)) => {
1692 requires_start_after || current != &token
1693 }
1694 Some(MongoResumePosition::StartAfter(_)) | None => true,
1695 };
1696 if changed {
1697 let encoded = serde_json::to_string(&token).map_err(|error| {
1698 ConnectorError::ReadError(format!(
1699 "serialize MongoDB post-batch resume token: {error}"
1700 ))
1701 })?;
1702 let encoded = canonical_resume_token(&encoded).map_err(|error| {
1703 ConnectorError::ReadError(format!(
1704 "invalid MongoDB post-batch resume token: {error}"
1705 ))
1706 })?;
1707 let Some(marker) = acquire_mongo_high_watermark_ownership(
1708 encoded,
1709 requires_start_after,
1710 byte_budget,
1711 max_buffered_bytes,
1712 shutdown_rx,
1713 )
1714 .await?
1715 else {
1716 return Ok(ChangeStreamRead::Stop);
1717 };
1718 if !send_event_or_shutdown(tx, marker, shutdown_rx).await {
1719 return Ok(ChangeStreamRead::Stop);
1720 }
1721 data_ready.notify_one();
1722 }
1723 *resume_position = Some(if requires_start_after {
1724 MongoResumePosition::StartAfter(token)
1725 } else {
1726 MongoResumePosition::ResumeAfter(token)
1727 });
1728 }
1729 }
1730 if !cursor_alive {
1731 tracing::info!("change stream cursor exhausted");
1732 return Ok(ChangeStreamRead::Reconnect);
1733 }
1734 *consecutive_failures = 0;
1735 }
1736 Err(error) => {
1737 tracing::error!(%error, "change stream error");
1738 return Ok(ChangeStreamRead::Reconnect);
1739 }
1740 }
1741 }
1742}
1743
1744#[cfg(feature = "mongodb-cdc")]
1745async fn send_event_or_shutdown(
1746 tx: &ChangeStreamTx,
1747 event: BufferedMongoEvent,
1748 shutdown_rx: &mut tokio::sync::watch::Receiver<bool>,
1749) -> bool {
1750 if reader_stopping(shutdown_rx) {
1751 return false;
1752 }
1753
1754 tokio::select! {
1755 biased;
1756 _ = shutdown_rx.changed() => false,
1757 result = tx.send(event) => {
1758 if result.is_err() {
1759 tracing::warn!("source channel closed, stopping reader");
1760 }
1761 result.is_ok()
1762 }
1763 }
1764}
1765
1766#[cfg(feature = "mongodb-cdc")]
1767async fn acquire_mongo_event_ownership(
1768 event: MongoDbChangeEvent,
1769 byte_budget: &Arc<Semaphore>,
1770 max_buffered_bytes: usize,
1771 shutdown_rx: &mut tokio::sync::watch::Receiver<bool>,
1772) -> Result<Option<BufferedMongoEvent>, ConnectorError> {
1773 let retained_bytes = mongo_event_retained_bytes(&event)?;
1774 let Some(byte_permit) =
1775 acquire_mongo_byte_permit(retained_bytes, byte_budget, max_buffered_bytes, shutdown_rx)
1776 .await?
1777 else {
1778 return Ok(None);
1779 };
1780
1781 Ok(Some(BufferedMongoEvent::new(event, byte_permit)))
1782}
1783
1784#[cfg(feature = "mongodb-cdc")]
1785async fn acquire_mongo_high_watermark_ownership(
1786 token: String,
1787 requires_start_after: bool,
1788 byte_budget: &Arc<Semaphore>,
1789 max_buffered_bytes: usize,
1790 shutdown_rx: &mut tokio::sync::watch::Receiver<bool>,
1791) -> Result<Option<BufferedMongoEvent>, ConnectorError> {
1792 let retained_bytes = mongo_high_watermark_retained_bytes(token.capacity())?;
1793 let Some(byte_permit) =
1794 acquire_mongo_byte_permit(retained_bytes, byte_budget, max_buffered_bytes, shutdown_rx)
1795 .await?
1796 else {
1797 return Ok(None);
1798 };
1799 Ok(Some(BufferedMongoEvent::high_watermark(
1800 token,
1801 requires_start_after,
1802 byte_permit,
1803 )))
1804}
1805
1806#[cfg(feature = "mongodb-cdc")]
1807async fn acquire_mongo_byte_permit(
1808 retained_bytes: usize,
1809 byte_budget: &Arc<Semaphore>,
1810 max_buffered_bytes: usize,
1811 shutdown_rx: &mut tokio::sync::watch::Receiver<bool>,
1812) -> Result<Option<OwnedSemaphorePermit>, ConnectorError> {
1813 if reader_stopping(shutdown_rx) {
1814 return Ok(None);
1815 }
1816 if retained_bytes > max_buffered_bytes {
1817 return Err(ConnectorError::ReadError(format!(
1818 "MongoDB CDC decoded item exceeds the hard byte bound: item={retained_bytes}, \
1819 limit={max_buffered_bytes}"
1820 )));
1821 }
1822 let permits = u32::try_from(retained_bytes).map_err(|_| {
1823 ConnectorError::ReadError(format!(
1824 "MongoDB CDC decoded item exceeds the hard byte bound: item={retained_bytes}, \
1825 limit={max_buffered_bytes}"
1826 ))
1827 })?;
1828 let byte_permit = tokio::select! {
1829 biased;
1830 _ = shutdown_rx.changed() => return Ok(None),
1831 permit = Arc::clone(byte_budget).acquire_many_owned(permits) => permit.map_err(|_| {
1832 ConnectorError::ReadError("MongoDB CDC byte budget closed".into())
1833 })?,
1834 };
1835 Ok(Some(byte_permit))
1836}
1837
1838#[cfg(feature = "mongodb-cdc")]
1839async fn retry_interrupted(
1840 shutdown_rx: &mut tokio::sync::watch::Receiver<bool>,
1841 delay: std::time::Duration,
1842) -> bool {
1843 tokio::select! {
1844 changed = shutdown_rx.changed() => changed.is_err() || *shutdown_rx.borrow(),
1845 () = tokio::time::sleep(delay) => false,
1846 }
1847}
1848
1849#[cfg(feature = "mongodb-cdc")]
1850fn parse_change_stream_pipeline(
1851 pipeline: &[serde_json::Value],
1852) -> Result<Vec<mongodb::bson::Document>, ConnectorError> {
1853 pipeline
1854 .iter()
1855 .enumerate()
1856 .map(|(index, value)| {
1857 mongodb::bson::to_document(value).map_err(|error| {
1858 ConnectorError::ConfigurationError(format!(
1859 "pipeline stage {index} cannot be represented as BSON: {error}"
1860 ))
1861 })
1862 })
1863 .collect()
1864}
1865
1866#[cfg(feature = "mongodb-cdc")]
1867fn verify_mongodb_collection_uuid(
1868 expected: Uuid,
1869 observed: Uuid,
1870 database: &str,
1871 collection: &str,
1872) -> Result<(), ConnectorError> {
1873 if expected == observed {
1874 return Ok(());
1875 }
1876 Err(ConnectorError::ConfigurationError(format!(
1877 "MongoDB CDC collection identity changed for {database}.{collection}: \
1878 checkpoint/bound UUID={expected}, observed UUID={observed}"
1879 )))
1880}
1881
1882#[cfg(feature = "mongodb-cdc")]
1883fn verify_mongodb_collection(
1884 config: &MongoDbSourceConfig,
1885 expected_uuid: Uuid,
1886 observation: &MongoCollectionObservation,
1887) -> Result<(), ConnectorError> {
1888 verify_mongodb_collection_uuid(
1889 expected_uuid,
1890 observation.collection_uuid,
1891 &config.database,
1892 &config.collection,
1893 )?;
1894 if config.full_document_mode == super::config::FullDocumentMode::RequirePostImage
1895 && !observation.post_images_enabled
1896 {
1897 return Err(ConnectorError::ConfigurationError(format!(
1898 "MongoDB CDC full.document.mode=required needs changeStreamPreAndPostImages enabled \
1899 on {}.{} before the source starts",
1900 config.database, config.collection
1901 )));
1902 }
1903 Ok(())
1904}
1905
1906#[cfg(feature = "mongodb-cdc")]
1907fn verify_mongodb_deployment_identity(
1908 expected: &MongoDeploymentIdentity,
1909 observed: &MongoDeploymentIdentity,
1910) -> Result<(), ConnectorError> {
1911 if expected == observed {
1912 return Ok(());
1913 }
1914 Err(ConnectorError::ConfigurationError(format!(
1915 "MongoDB CDC deployment identity changed: checkpoint/bound identity={}, observed \
1916 identity={}",
1917 expected.encode(),
1918 observed.encode()
1919 )))
1920}
1921
1922#[cfg(feature = "mongodb-cdc")]
1923fn verify_mongodb_admission(
1924 config: &MongoDbSourceConfig,
1925 expected_deployment: &MongoDeploymentIdentity,
1926 expected_uuid: Uuid,
1927 observation: &MongoAdmissionObservation,
1928) -> Result<(), ConnectorError> {
1929 verify_mongodb_deployment_identity(expected_deployment, &observation.deployment_identity)?;
1930 verify_mongodb_collection(config, expected_uuid, &observation.collection)
1931}
1932
1933#[cfg(feature = "mongodb-cdc")]
1934fn fresh_stream_anchor(
1935 cursor: &mongodb::change_stream::ChangeStream<
1936 mongodb::change_stream::event::ChangeStreamEvent<mongodb::bson::Document>,
1937 >,
1938) -> Result<(mongodb::change_stream::event::ResumeToken, String), ConnectorError> {
1939 let token = cursor.resume_token().ok_or_else(|| {
1943 ConnectorError::ReadError(
1944 "fresh MongoDB change stream omitted its initial postBatchResumeToken".into(),
1945 )
1946 })?;
1947 let encoded = serde_json::to_string(&token).map_err(|error| {
1948 ConnectorError::ReadError(format!(
1949 "serialize initial MongoDB post-batch resume token: {error}"
1950 ))
1951 })?;
1952 let encoded = canonical_resume_token(&encoded).map_err(|error| {
1953 ConnectorError::ReadError(format!(
1954 "invalid initial MongoDB post-batch resume token: {error}"
1955 ))
1956 })?;
1957 Ok((token, encoded))
1958}
1959
1960#[cfg(feature = "mongodb-cdc")]
1961fn report_mongo_reader_admission_error(
1962 ready_tx: &mut Option<
1963 tokio::sync::oneshot::Sender<Result<MongoReaderReady, MongoReaderFailure>>,
1964 >,
1965 error: &ConnectorError,
1966) {
1967 if let Some(ready_tx) = ready_tx.take() {
1968 let _ = ready_tx.send(Err(MongoReaderFailure::from_connector(error)));
1969 }
1970}
1971
1972#[cfg(feature = "mongodb-cdc")]
1978async fn run_change_stream_reader(
1979 db: mongodb::Database,
1980 config: MongoDbSourceConfig,
1981 tx: ChangeStreamTx,
1982 shutdown_rx: tokio::sync::watch::Receiver<bool>,
1983 data_ready: Arc<Notify>,
1984 metrics: Arc<MongoDbCdcMetrics>,
1985 byte_budget: Arc<Semaphore>,
1986 max_buffered_bytes: usize,
1987 initial_resume_position: Option<MongoResumePosition>,
1988 expected_collection_uuid: Option<Uuid>,
1989 expected_deployment_identity: Option<MongoDeploymentIdentity>,
1990 ready_tx: tokio::sync::oneshot::Sender<Result<MongoReaderReady, MongoReaderFailure>>,
1991) -> Result<(), ConnectorError> {
1992 let client = db.client().clone();
1993 let result = run_change_stream_reader_loop(
1994 db,
1995 config,
1996 tx,
1997 shutdown_rx,
1998 data_ready,
1999 metrics,
2000 byte_budget,
2001 max_buffered_bytes,
2002 initial_resume_position,
2003 expected_collection_uuid,
2004 expected_deployment_identity,
2005 ready_tx,
2006 )
2007 .await;
2008
2009 client.shutdown().await;
2012 result
2013}
2014
2015#[cfg(feature = "mongodb-cdc")]
2016async fn run_change_stream_reader_loop(
2017 db: mongodb::Database,
2018 config: MongoDbSourceConfig,
2019 tx: ChangeStreamTx,
2020 mut shutdown_rx: tokio::sync::watch::Receiver<bool>,
2021 data_ready: Arc<Notify>,
2022 metrics: Arc<MongoDbCdcMetrics>,
2023 byte_budget: Arc<Semaphore>,
2024 max_buffered_bytes: usize,
2025 initial_resume_position: Option<MongoResumePosition>,
2026 expected_collection_uuid: Option<Uuid>,
2027 expected_deployment_identity: Option<MongoDeploymentIdentity>,
2028 ready_tx: tokio::sync::oneshot::Sender<Result<MongoReaderReady, MongoReaderFailure>>,
2029) -> Result<(), ConnectorError> {
2030 let mut resume_position = initial_resume_position;
2031 let fresh_start = resume_position.is_none();
2032 let mut initial_resume_token = None;
2033 let mut ready_tx = Some(ready_tx);
2034 let pipeline = match parse_change_stream_pipeline(&config.pipeline) {
2035 Ok(pipeline) => pipeline,
2036 Err(error) => {
2037 report_mongo_reader_admission_error(&mut ready_tx, &error);
2038 return Err(error);
2039 }
2040 };
2041
2042 let current_db = db;
2043 let mut consecutive_failures: u32 = 0;
2044 let initial_observation = loop {
2045 match observe_mongodb_admission(¤t_db, &config.database, &config.collection).await {
2046 Ok(observation) => break observation,
2047 Err(error) if !error.is_transient() => {
2048 report_mongo_reader_admission_error(&mut ready_tx, &error);
2049 return Err(error);
2050 }
2051 Err(error) => {
2052 consecutive_failures += 1;
2053 if consecutive_failures >= MAX_FAILURES {
2054 report_mongo_reader_admission_error(&mut ready_tx, &error);
2055 return Err(error);
2056 }
2057 let backoff = crate::retry::Backoff::broker_reconnect().delay(consecutive_failures);
2058 tracing::warn!(
2059 attempt = consecutive_failures,
2060 ?backoff,
2061 error = %error,
2062 "failed to inspect MongoDB deployment or collection identity, retrying"
2063 );
2064 metrics.record_reconnect();
2065 if retry_interrupted(&mut shutdown_rx, backoff).await {
2066 return Ok(());
2067 }
2068 }
2069 }
2070 };
2071 consecutive_failures = 0;
2072
2073 let collection_uuid =
2074 expected_collection_uuid.unwrap_or(initial_observation.collection.collection_uuid);
2075 let deployment_identity = expected_deployment_identity
2076 .unwrap_or_else(|| initial_observation.deployment_identity.clone());
2077 if let Err(error) = verify_mongodb_admission(
2078 &config,
2079 &deployment_identity,
2080 collection_uuid,
2081 &initial_observation,
2082 ) {
2083 report_mongo_reader_admission_error(&mut ready_tx, &error);
2084 return Err(error);
2085 }
2086
2087 let mut verify_before_open = false;
2088
2089 'reconnect: loop {
2090 if verify_before_open {
2091 match observe_mongodb_admission(¤t_db, &config.database, &config.collection).await
2092 {
2093 Ok(observation) => {
2094 if let Err(error) = verify_mongodb_admission(
2095 &config,
2096 &deployment_identity,
2097 collection_uuid,
2098 &observation,
2099 ) {
2100 report_mongo_reader_admission_error(&mut ready_tx, &error);
2101 return Err(error);
2102 }
2103 verify_before_open = false;
2104 }
2105 Err(error) if !error.is_transient() => {
2106 report_mongo_reader_admission_error(&mut ready_tx, &error);
2107 return Err(error);
2108 }
2109 Err(error) => {
2110 consecutive_failures += 1;
2111 if consecutive_failures >= MAX_FAILURES {
2112 report_mongo_reader_admission_error(&mut ready_tx, &error);
2113 return Err(error);
2114 }
2115 let backoff =
2116 crate::retry::Backoff::broker_reconnect().delay(consecutive_failures);
2117 tracing::warn!(
2118 attempt = consecutive_failures,
2119 ?backoff,
2120 error = %error,
2121 "failed to verify MongoDB deployment or collection identity before reconnect"
2122 );
2123 metrics.record_reconnect();
2124 if retry_interrupted(&mut shutdown_rx, backoff).await {
2125 break 'reconnect;
2126 }
2127 continue 'reconnect;
2128 }
2129 }
2130 }
2131
2132 let bootstrap = fresh_start && ready_tx.is_some() && resume_position.is_none();
2133 let options = if bootstrap {
2134 bootstrap_change_stream_options(&config)
2135 } else {
2136 change_stream_options(&config, resume_position.as_ref())
2137 };
2138
2139 let cursor_result = current_db
2141 .collection::<mongodb::bson::Document>(&config.collection)
2142 .watch()
2143 .pipeline(pipeline.clone())
2144 .with_options(options)
2145 .await;
2146
2147 let mut cursor = match cursor_result {
2148 Ok(c) => c,
2149 Err(e) => {
2150 consecutive_failures += 1;
2151 if consecutive_failures >= MAX_FAILURES {
2152 let msg =
2153 format!("change stream open failed after {MAX_FAILURES} attempts: {e}");
2154 tracing::error!(%msg);
2155 let error = ConnectorError::ReadError(msg);
2156 report_mongo_reader_admission_error(&mut ready_tx, &error);
2157 return Err(error);
2158 }
2159 let backoff = crate::retry::Backoff::broker_reconnect().delay(consecutive_failures);
2160 tracing::warn!(
2161 attempt = consecutive_failures,
2162 ?backoff,
2163 error = %e,
2164 "failed to open change stream, retrying"
2165 );
2166 metrics.record_reconnect();
2167 if retry_interrupted(&mut shutdown_rx, backoff).await {
2168 break 'reconnect;
2169 }
2170 verify_before_open = true;
2171 continue 'reconnect;
2172 }
2173 };
2174
2175 match observe_mongodb_admission(¤t_db, &config.database, &config.collection).await {
2176 Ok(observation) => {
2177 if let Err(error) = verify_mongodb_admission(
2178 &config,
2179 &deployment_identity,
2180 collection_uuid,
2181 &observation,
2182 ) {
2183 report_mongo_reader_admission_error(&mut ready_tx, &error);
2184 return Err(error);
2185 }
2186 }
2187 Err(error) if !error.is_transient() => {
2188 report_mongo_reader_admission_error(&mut ready_tx, &error);
2189 return Err(error);
2190 }
2191 Err(error) => {
2192 consecutive_failures += 1;
2193 if consecutive_failures >= MAX_FAILURES {
2194 report_mongo_reader_admission_error(&mut ready_tx, &error);
2195 return Err(error);
2196 }
2197 let backoff = crate::retry::Backoff::broker_reconnect().delay(consecutive_failures);
2198 tracing::warn!(
2199 attempt = consecutive_failures,
2200 ?backoff,
2201 error = %error,
2202 "failed to verify MongoDB deployment or collection identity after opening change stream"
2203 );
2204 metrics.record_reconnect();
2205 if retry_interrupted(&mut shutdown_rx, backoff).await {
2206 break 'reconnect;
2207 }
2208 verify_before_open = true;
2209 continue 'reconnect;
2210 }
2211 }
2212 consecutive_failures = 0;
2213
2214 if bootstrap {
2215 match fresh_stream_anchor(&cursor) {
2216 Ok((token, encoded)) => {
2217 resume_position = Some(MongoResumePosition::ResumeAfter(token));
2218 initial_resume_token = Some(encoded);
2219 }
2220 Err(error) => {
2221 report_mongo_reader_admission_error(&mut ready_tx, &error);
2222 return Err(error);
2223 }
2224 }
2225 drop(cursor);
2226 continue 'reconnect;
2227 }
2228
2229 if let Some(ready_tx) = ready_tx.take() {
2230 let _ = ready_tx.send(Ok(MongoReaderReady {
2231 initial_resume_token: initial_resume_token.take(),
2232 collection_uuid,
2233 deployment_identity: deployment_identity.clone(),
2234 }));
2235 }
2236
2237 tracing::info!(
2238 database = %config.database,
2239 collection = %config.collection,
2240 resumed = resume_position.is_some(),
2241 "change stream reader started"
2242 );
2243
2244 if matches!(
2245 forward_change_stream(
2246 &mut cursor,
2247 &mut shutdown_rx,
2248 &mut resume_position,
2249 &tx,
2250 &data_ready,
2251 &mut consecutive_failures,
2252 &byte_budget,
2253 max_buffered_bytes,
2254 &metrics,
2255 )
2256 .await?,
2257 ChangeStreamRead::Stop
2258 ) {
2259 break 'reconnect;
2260 }
2261
2262 consecutive_failures += 1;
2264 if consecutive_failures >= MAX_FAILURES {
2265 let msg = format!("change stream failed after {MAX_FAILURES} consecutive failures");
2266 tracing::error!(%msg);
2267 return Err(ConnectorError::ReadError(msg));
2268 }
2269
2270 let backoff = crate::retry::Backoff::broker_reconnect().delay(consecutive_failures);
2271 tracing::warn!(
2272 resume_position = ?resume_position,
2273 attempt = consecutive_failures,
2274 ?backoff,
2275 "reconnecting change stream"
2276 );
2277 metrics.record_reconnect();
2278
2279 if retry_interrupted(&mut shutdown_rx, backoff).await {
2280 break 'reconnect;
2281 }
2282
2283 verify_before_open = true;
2286 }
2287
2288 if let Some(ready_tx) = ready_tx.take() {
2289 let _ = ready_tx.send(Err(MongoReaderFailure::Read(
2290 "change stream reader was shut down before the cursor opened".into(),
2291 )));
2292 }
2293 Ok(())
2294}
2295
2296#[cfg(feature = "mongodb-cdc")]
2298fn parse_change_stream_event(
2299 event: &mongodb::change_stream::event::ChangeStreamEvent<mongodb::bson::Document>,
2300) -> Result<MongoDbChangeEvent, ConnectorError> {
2301 use super::change_event::{Namespace, UpdateDescription};
2302 use mongodb::change_stream::event::OperationType as MongoOpType;
2303
2304 let operation_type = match &event.operation_type {
2305 MongoOpType::Insert => OperationType::Insert,
2306 MongoOpType::Update => OperationType::Update,
2307 MongoOpType::Replace => OperationType::Replace,
2308 MongoOpType::Delete => OperationType::Delete,
2309 MongoOpType::Drop => OperationType::Drop,
2310 MongoOpType::Rename => OperationType::Rename,
2311 MongoOpType::Invalidate => OperationType::Invalidate,
2312 MongoOpType::DropDatabase => OperationType::DropDatabase,
2313 MongoOpType::Other(value) => OperationType::Other(value.clone()),
2314 other => {
2315 return Err(ConnectorError::ReadError(format!(
2316 "unsupported MongoDB operation type: {other:?}"
2317 )));
2318 }
2319 };
2320
2321 let namespace = event.ns.as_ref().map_or_else(
2322 || Namespace {
2323 db: String::new(),
2324 coll: String::new(),
2325 },
2326 |ns| Namespace {
2327 db: ns.db.clone(),
2328 coll: ns.coll.clone().unwrap_or_default(),
2329 },
2330 );
2331
2332 let document_key = event.document_key.as_ref().map_or_else(
2333 || Ok(String::new()),
2334 |document| {
2335 serde_json::to_string(document)
2336 .map_err(|error| ConnectorError::ReadError(format!("document key: {error}")))
2337 },
2338 )?;
2339
2340 let full_document = event
2341 .full_document
2342 .as_ref()
2343 .map(|document| {
2344 serde_json::to_string(document)
2345 .map_err(|error| ConnectorError::ReadError(format!("full document: {error}")))
2346 })
2347 .transpose()?;
2348
2349 let update_description = event
2350 .update_description
2351 .as_ref()
2352 .map(|ud| -> Result<UpdateDescription, ConnectorError> {
2353 let updated_fields = ud
2354 .updated_fields
2355 .iter()
2356 .map(|(key, value)| {
2357 serde_json::to_value(value)
2358 .map(|value| (key.clone(), value))
2359 .map_err(|error| {
2360 ConnectorError::ReadError(format!("updated field '{key}': {error}"))
2361 })
2362 })
2363 .collect::<Result<_, _>>()?;
2364
2365 let removed_fields = ud.removed_fields.clone();
2366
2367 let truncated_arrays = ud
2368 .truncated_arrays
2369 .as_deref()
2370 .unwrap_or_default()
2371 .iter()
2372 .map(|t| {
2373 let new_size = u32::try_from(t.new_size).map_err(|_| {
2374 ConnectorError::ReadError(format!(
2375 "truncated array '{}' has negative newSize {}",
2376 t.field, t.new_size
2377 ))
2378 })?;
2379 Ok(super::change_event::TruncatedArray {
2380 field: t.field.clone(),
2381 new_size,
2382 })
2383 })
2384 .collect::<Result<Vec<_>, ConnectorError>>()?;
2385
2386 let disambiguated_paths = ud
2387 .disambiguated_paths
2388 .as_ref()
2389 .map(|paths| {
2390 paths
2391 .iter()
2392 .map(|(key, value)| {
2393 serde_json::to_value(value)
2394 .map(|value| (key.clone(), value))
2395 .map_err(|error| {
2396 ConnectorError::ReadError(format!(
2397 "disambiguated path '{key}': {error}"
2398 ))
2399 })
2400 })
2401 .collect::<Result<_, _>>()
2402 })
2403 .transpose()?
2404 .unwrap_or_default();
2405
2406 Ok(UpdateDescription {
2407 updated_fields,
2408 removed_fields,
2409 truncated_arrays,
2410 disambiguated_paths,
2411 })
2412 })
2413 .transpose()?;
2414
2415 let (cluster_time_secs, cluster_time_inc) = event
2416 .cluster_time
2417 .map_or((0, 0), |ts| (ts.time, ts.increment));
2418
2419 let wall_time_ms = event
2420 .wall_time
2421 .map_or(0, mongodb::bson::DateTime::timestamp_millis);
2422
2423 let resume_token = serde_json::to_string(&event.id)
2425 .map_err(|error| ConnectorError::ReadError(format!("resume token: {error}")))?;
2426
2427 Ok(MongoDbChangeEvent {
2428 operation_type,
2429 namespace,
2430 document_key,
2431 full_document,
2432 update_description,
2433 cluster_time_secs,
2434 cluster_time_inc,
2435 resume_token,
2436 wall_time_ms,
2437 })
2438}
2439
2440#[cfg(test)]
2441mod tests;