Skip to main content

laminar_connectors/mongodb/
source.rs

1//! `MongoDB` CDC source connector implementation.
2//!
3//! Implements [`SourceConnector`] for streaming change events from `MongoDB`
4//! change streams into `LaminarDB` as Arrow `RecordBatch`es.
5//!
6//! # Cancellation Safety
7//!
8//! Connector lifecycle futures never directly poll the `MongoDB` driver. Driver
9//! I/O lives in an owned reader task; cancellation aborts that task so no
10//! connection or cursor outlives its connector.
11
12use 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]); // showExpandedEvents is always enabled.
90    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            // serde_json::Map does not expose its allocation capacity. Charging each live
304            // entry plus its recursively owned values is stable across map backends.
305            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/// Returns the Arrow schema for `MongoDB` CDC envelope records.
404///
405/// | Column              | Type   | Nullable | Description                        |
406/// |---------------------|--------|----------|------------------------------------|
407/// | `_namespace`        | Utf8   | no       | `database.collection`              |
408/// | `_op`               | Utf8   | no       | Operation code (I/U/R/D/DROP/...)  |
409/// | `_document_key`     | Utf8   | no       | Document key JSON                  |
410/// | `_cluster_time_s`   | UInt32 | no       | Cluster time seconds               |
411/// | `_cluster_time_i`   | UInt32 | no       | Cluster time increment             |
412/// | `_wall_time_ms`     | Timestamp(ms) | no | Wall clock timestamp             |
413/// | `_full_document`    | Utf8   | yes      | Full document JSON                 |
414/// | `_update_desc`      | Utf8   | yes      | Update description JSON            |
415/// | `_resume_token`     | Utf8   | no       | Opaque resume token JSON           |
416#[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
435/// `MongoDB` CDC source connector.
436///
437/// Streams change events from a `MongoDB` change stream using the
438/// `SourceConnector` trait. Events are buffered internally and
439/// converted to Arrow `RecordBatch`es on `poll_batch`.
440///
441/// # Sharded Cluster Note
442///
443/// On sharded clusters, `mongos` opens per-shard cursors and merges
444/// results transparently. Ensure `max_pool_size` is at least as large
445/// as the expected number of concurrent change streams to avoid
446/// connection starvation.
447pub struct MongoDbCdcSource {
448    /// Connector configuration.
449    config: MongoDbSourceConfig,
450
451    /// Current lifecycle state.
452    state: ConnectorState,
453
454    /// Output schema (CDC envelope).
455    schema: SchemaRef,
456
457    /// Lock-free metrics.
458    metrics: Arc<MongoDbCdcMetrics>,
459
460    /// Buffered change events awaiting `poll_batch`.
461    event_buffer: VecDeque<BufferedMongoEvent>,
462
463    /// Latest ordered event or post-batch token consumed by `poll_batch`.
464    /// The reader's newer cursor token is deliberately not shared with this field.
465    checkpoint_resume_token: Option<String>,
466
467    /// Invalidation tokens must be restored with `startAfter`, never `resumeAfter`.
468    checkpoint_requires_start_after: bool,
469
470    /// Physical identity of the fixed collection admitted by `listCollections`.
471    collection_uuid: Option<Uuid>,
472
473    /// Immutable server-issued identity of the replica set or sharded cluster.
474    deployment_identity: Option<MongoDeploymentIdentity>,
475
476    /// Shared ownership limits span the reader channel and poll buffer.
477    byte_budget: Arc<Semaphore>,
478
479    /// Notification handle signalled when data arrives from the stream.
480    data_ready: Arc<Notify>,
481
482    /// Background change stream reader task handle (feature-gated).
483    #[cfg(feature = "mongodb-cdc")]
484    reader_handle: Option<tokio::task::JoinHandle<()>>,
485
486    /// Channel receiver for change events from the background task.
487    #[cfg(feature = "mongodb-cdc")]
488    event_rx: Option<ChangeStreamRx>,
489
490    /// Shutdown signal for the background reader task.
491    #[cfg(feature = "mongodb-cdc")]
492    reader_shutdown: Option<tokio::sync::watch::Sender<bool>>,
493
494    /// Terminal reader failure, independent of the bounded event queue.
495    #[cfg(feature = "mongodb-cdc")]
496    reader_error: Option<tokio::sync::watch::Receiver<Option<MongoReaderFailure>>>,
497
498    /// Admission authority and terminal observer for this connector generation.
499    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        // The reader owns a separate task guard. Runtime destruction drops the
524        // reader future and resolves that proof without a timer or join guess.
525        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/// Cloneable async sender for the change stream reader → `poll_batch` queue.
537#[cfg(feature = "mongodb-cdc")]
538type ChangeStreamTx = crossfire::MAsyncTx<crossfire::mpsc::Array<BufferedMongoEvent>>;
539/// Single-consumer async receiver for the change stream reader → `poll_batch` queue.
540#[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    /// Creates a new `MongoDB` CDC source with the given configuration.
630    #[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    /// Enqueues a change event for focused source tests without bypassing production bounds.
668    #[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    /// Drains up to `max_records` events from the buffer and converts
728    /// them to an Arrow `RecordBatch`.
729    ///
730    /// # Errors
731    ///
732    /// Returns `ConnectorError` if Arrow batch construction fails.
733    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        // An invalidate token changes the legal resume option. End the batch exactly there even
743        // when the background reader has already reopened with startAfter and queued later data.
744        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/// Converts a batch of [`MongoDbChangeEvent`]s to an Arrow `RecordBatch`.
807#[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            // A configured namespace is not a physical replay identity until admission has read
1018            // the server-assigned collection UUID.
1019            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            // Before a fresh source has opened, it has no lossless replay position.
1035            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// ── Feature-gated I/O (real MongoDB driver) ──
1119
1120#[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    /// Starts the background change stream reader task.
1209    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    /// Drains events from the background reader channel into the buffer.
1304    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/// Read the immutable identity and post-image capability for one fixed collection.
1476#[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/// Maximum consecutive failures before the reader gives up.
1560#[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    // MongoDB guarantees an empty firstBatch for batchSize=0, so its PBRT is an exact opening
1606    // cut and cannot skip concurrently buffered events.
1607    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        // Poll getMore to completion during normal operation; maxAwaitTime keeps cooperative
1632        // shutdown prompt. The connector aborts and joins the owned task at its hard deadline.
1633        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    // The bootstrap aggregate uses batchSize=0, so MongoDB returns an empty firstBatch and its
1940    // exact postBatchResumeToken. Refuse an inclusive timestamp fallback: it can replay the final
1941    // write that preceded admission.
1942    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/// Background task that reads from the `MongoDB` change stream and sends
1973/// events to the source via a channel.
1974///
1975/// Uses a `'reconnect` / `'recv` double-loop pattern (mirroring the
1976/// Postgres CDC source) with exponential backoff capped at 30 seconds.
1977#[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    // The loop owns every database, collection, and cursor handle. Once it
2010    // returns, shutdown can drain the driver's own async cleanup tasks.
2011    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(&current_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(&current_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        // Open the change stream cursor.
2140        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(&current_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        // Exited recv loop due to error or cursor exhaustion — attempt reconnect.
2263        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        // The MongoDB client owns topology monitoring and reconnects its pool.
2284        // Reusing it avoids spawning untracked driver generations on each retry.
2285        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/// Parses a `ChangeStreamEvent<Document>` into a [`MongoDbChangeEvent`].
2297#[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    // Serialize the ResumeToken via serde (it implements Serialize).
2424    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;