Skip to main content

laminar_connectors/mongodb/source/
lifecycle.rs

1//! Source contract, lifecycle, polling, and shutdown.
2
3use std::collections::BTreeMap;
4use std::sync::Arc;
5
6use arrow_schema::SchemaRef;
7use async_trait::async_trait;
8use tokio::sync::Notify;
9
10use crate::checkpoint::SourceCheckpoint;
11use crate::config::{ConnectorConfig, ConnectorState};
12use crate::connector::{
13    ConnectorTaskTracker, SourceBatch, SourceConnector, SourceContract, SourcePosition, SourceStart,
14};
15use crate::error::ConnectorError;
16
17use super::{
18    mongodb_stream_identity, parse_mongodb_checkpoint, reap_mongo_reader, MongoCheckpointPosition,
19    MongoDbCdcSource, MongoDbSourceConfig, MongoResumePosition, ParsedMongoCheckpoint,
20    COLLECTION_UUID_METADATA, DEPLOYMENT_IDENTITY_METADATA, MONGODB_CHECKPOINT_CONNECTOR,
21    MONGODB_CHECKPOINT_VERSION, READER_SHUTDOWN_TIMEOUT, RESUME_TOKEN_OFFSET,
22    START_AFTER_TOKEN_OFFSET, STREAM_IDENTITY_METADATA,
23};
24
25#[async_trait]
26impl SourceConnector for MongoDbCdcSource {
27    fn terminal_task_tracker(&self) -> Option<ConnectorTaskTracker> {
28        Some(self.task_tracker.clone())
29    }
30
31    fn recovery_identity_options(
32        &self,
33        config: &ConnectorConfig,
34    ) -> Result<Option<BTreeMap<String, String>>, ConnectorError> {
35        let mut parsed = if config.properties().is_empty() {
36            self.config.clone()
37        } else {
38            MongoDbSourceConfig::from_config(config)?
39        };
40        parsed.normalize_pipeline()?;
41        parsed.validate()?;
42        let pipeline = super::super::config::canonical_pipeline_json(&parsed.pipeline);
43
44        Ok(Some(BTreeMap::from([
45            ("collection".into(), parsed.collection),
46            ("database".into(), parsed.database),
47            (
48                "full.document.mode".into(),
49                parsed.full_document_mode.to_string(),
50            ),
51            ("pipeline".into(), pipeline),
52            ("wire.protocol".into(), "change-stream-expanded-v1".into()),
53        ])))
54    }
55
56    async fn start(&mut self, request: SourceStart) -> Result<(), ConnectorError> {
57        if self.state != ConnectorState::Created {
58            return Err(ConnectorError::InvalidState {
59                expected: ConnectorState::Created.to_string(),
60                actual: self.state.to_string(),
61            });
62        }
63        let (config, position, _) = request.into_parts();
64        let parsed_config = if config.properties().is_empty() {
65            let mut config = self.config.clone();
66            config.normalize_pipeline()?;
67            config.validate()?;
68            config
69        } else {
70            MongoDbSourceConfig::from_config(&config)?
71        };
72        let (
73            checkpoint_resume_token,
74            checkpoint_requires_start_after,
75            initial_resume_position,
76            expected_collection_uuid,
77            expected_deployment_identity,
78        ) = match position {
79            SourcePosition::Initial => (None, false, None, None, None),
80            SourcePosition::Resume {
81                attempt,
82                checkpoint,
83            } => {
84                let ParsedMongoCheckpoint {
85                    position,
86                    collection_uuid,
87                    deployment_identity,
88                } = parse_mongodb_checkpoint(&checkpoint, &parsed_config).map_err(|error| {
89                    ConnectorError::ConfigurationError(format!(
90                        "invalid MongoDB CDC checkpoint {attempt:?}: {error}"
91                    ))
92                })?;
93                match position {
94                    MongoCheckpointPosition::ResumeAfter(token) => {
95                        let driver_token = serde_json::from_str(&token).map_err(|error| {
96                            ConnectorError::ConfigurationError(format!(
97                                "invalid MongoDB CDC resume token in checkpoint {attempt:?}: \
98                                 {error}"
99                            ))
100                        })?;
101                        (
102                            Some(token),
103                            false,
104                            Some(MongoResumePosition::ResumeAfter(driver_token)),
105                            Some(collection_uuid),
106                            Some(deployment_identity),
107                        )
108                    }
109                    MongoCheckpointPosition::StartAfter(token) => {
110                        let driver_token = serde_json::from_str(&token).map_err(|error| {
111                            ConnectorError::ConfigurationError(format!(
112                                "invalid MongoDB CDC start-after token in checkpoint {attempt:?}: \
113                                 {error}"
114                            ))
115                        })?;
116                        (
117                            Some(token),
118                            true,
119                            Some(MongoResumePosition::StartAfter(driver_token)),
120                            Some(collection_uuid),
121                            Some(deployment_identity),
122                        )
123                    }
124                }
125            }
126        };
127
128        self.start_change_stream_reader(
129            parsed_config,
130            checkpoint_resume_token,
131            checkpoint_requires_start_after,
132            initial_resume_position,
133            expected_collection_uuid,
134            expected_deployment_identity,
135        )
136        .await?;
137
138        self.state = ConnectorState::Running;
139        tracing::info!(
140            database = %self.config.database,
141            collection = %self.config.collection,
142            full_document_mode = ?self.config.full_document_mode,
143            "MongoDB CDC source opened"
144        );
145
146        Ok(())
147    }
148
149    async fn poll_batch(
150        &mut self,
151        max_records: usize,
152    ) -> Result<Option<SourceBatch>, ConnectorError> {
153        self.drain_channel(max_records.saturating_sub(self.event_buffer.len()));
154        if let Some(batch) = self.drain_to_batch(max_records)? {
155            return Ok(Some(batch));
156        }
157        self.check_reader_error()?;
158        Ok(None)
159    }
160
161    fn schema(&self) -> SchemaRef {
162        Arc::clone(&self.schema)
163    }
164
165    fn checkpoint(&self) -> SourceCheckpoint {
166        let mut checkpoint = SourceCheckpoint::new();
167        let Some(collection_uuid) = self.collection_uuid else {
168            // A configured namespace is not a physical replay identity until admission has read
169            // the server-assigned collection UUID.
170            return checkpoint;
171        };
172        let Some(deployment_identity) = self.deployment_identity.as_ref() else {
173            return checkpoint;
174        };
175        if let Some(token) = self.checkpoint_resume_token.as_ref() {
176            checkpoint.set_offset(
177                if self.checkpoint_requires_start_after {
178                    START_AFTER_TOKEN_OFFSET
179                } else {
180                    RESUME_TOKEN_OFFSET
181                },
182                token,
183            );
184        } else {
185            // Before a fresh source has opened, it has no lossless replay position.
186            return checkpoint;
187        }
188        checkpoint.set_metadata("connector", MONGODB_CHECKPOINT_CONNECTOR);
189        checkpoint.set_metadata("version", MONGODB_CHECKPOINT_VERSION);
190        checkpoint.set_metadata("database", &self.config.database);
191        checkpoint.set_metadata("collection", &self.config.collection);
192        checkpoint.set_metadata(
193            COLLECTION_UUID_METADATA,
194            collection_uuid.hyphenated().to_string(),
195        );
196        checkpoint.set_metadata(DEPLOYMENT_IDENTITY_METADATA, deployment_identity.encode());
197        checkpoint.set_metadata(
198            STREAM_IDENTITY_METADATA,
199            mongodb_stream_identity(&self.config),
200        );
201        checkpoint
202    }
203
204    async fn close(&mut self) -> Result<(), ConnectorError> {
205        #[cfg(feature = "mongodb-cdc")]
206        let mut reader_join_error = None;
207        #[cfg(feature = "mongodb-cdc")]
208        {
209            if let Some(tx) = self.reader_shutdown.as_ref() {
210                tx.send_replace(true);
211            }
212            let mut detach_reader = false;
213            if let Some(handle) = self.reader_handle.as_mut() {
214                match tokio::time::timeout(READER_SHUTDOWN_TIMEOUT, &mut *handle).await {
215                    Ok(Ok(())) => {}
216                    Ok(Err(error)) if error.is_cancelled() => {}
217                    Ok(Err(error)) => reader_join_error = Some(error.to_string()),
218                    Err(_) => {
219                        tracing::warn!(
220                            "MongoDB CDC reader exceeded its close deadline; its tracked reaper retains shutdown ownership"
221                        );
222                        detach_reader = true;
223                    }
224                }
225            }
226            if detach_reader {
227                let handle = self
228                    .reader_handle
229                    .take()
230                    .expect("reader handle was present while awaiting it");
231                reap_mongo_reader(handle, &self.task_owner);
232            } else {
233                self.reader_handle = None;
234            }
235            self.reader_shutdown = None;
236            self.event_rx = None;
237            self.reader_error = None;
238        }
239
240        self.event_buffer.clear();
241        self.state = ConnectorState::Closed;
242        tracing::info!("MongoDB CDC source closed");
243        #[cfg(feature = "mongodb-cdc")]
244        if let Some(error) = reader_join_error {
245            return Err(ConnectorError::ReadError(format!(
246                "MongoDB CDC reader task failed during close: {error}"
247            )));
248        }
249        Ok(())
250    }
251
252    fn data_ready_notify(&self) -> Option<Arc<Notify>> {
253        Some(Arc::clone(&self.data_ready))
254    }
255
256    fn contract(&self, config: &ConnectorConfig) -> Result<SourceContract, ConnectorError> {
257        if config.properties().is_empty() {
258            self.config.validate()?;
259        } else {
260            MongoDbSourceConfig::from_config(config)?;
261        }
262        Err(ConnectorError::ConfigurationError(
263            "MongoDB CDC emits a raw JSON change envelope; canonical primary-keyed row/delete records are required"
264                .into(),
265        ))
266    }
267}