laminar_connectors/mongodb/source/
lifecycle.rs1use 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 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 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}