hubcio commented on code in PR #4225: URL: https://github.com/apache/iggy/pull/4225#discussion_r4062028051
########## core/connectors/sources/mongodb_source/src/lib.rs: ########## @@ -0,0 +1,242 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +use async_trait::async_trait; +use core::mem::drop; +use futures::stream::TryStreamExt; +use iggy_common::{DateTime, Utc}; +use iggy_connector_sdk::{ + ConnectorState, Error, ProducedMessage, ProducedMessages, Schema, Source, source_connector, +}; +use mongodb::{Client, Collection, bson::Document, bson::doc, options::ClientOptions}; +use secrecy::{ExposeSecret, SecretString}; +use serde::{Deserialize, Serialize}; +use std::str::FromStr; +use std::time::Duration; +use tokio::sync::Mutex; +use tracing::info; + +source_connector!(MongodbSource); + +#[derive(Debug, Clone, Serialize, Deserialize)] +struct State { + last_poll_timestamp: Option<DateTime<Utc>>, + total_documents_fetched: usize, + poll_count: usize, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct MongodbSourceConfig { + #[serde(serialize_with = "iggy_common::serde_secret::serialize_secret")] + pub connection_uri: SecretString, + pub database: String, + pub collection: String, + pub max_pool_size: Option<u32>, + pub query: Option<Document>, + pub timestamp_field: Option<String>, + pub limit: Option<i64>, + pub polling_interval: Option<String>, +} + +#[derive(Debug)] +pub struct MongodbSource { + id: u32, + config: MongodbSourceConfig, + client: Option<Client>, + polling_interval: Duration, + state: Mutex<State>, +} + +const CONNECTOR_NAME: &str = "MongoDB source"; + +impl MongodbSource { + pub fn new(id: u32, config: MongodbSourceConfig, state: Option<ConnectorState>) -> Self { + let polling_interval = config + .polling_interval + .as_deref() + .unwrap_or("10s") + .parse::<humantime::Duration>() + .unwrap_or_else(|_| humantime::Duration::from_str("10s").unwrap()) + .into(); + + let restored_state = state + .and_then(|s| s.deserialize::<State>(CONNECTOR_NAME, id)) + .inspect(|s| { + info!( + "Restored state for {CONNECTOR_NAME} connector with ID: {id}. \ + Documents fetched: {}, poll count: {}", + s.total_documents_fetched, s.poll_count + ); + }); + + MongodbSource { + id, + config, + client: None, + polling_interval, + state: Mutex::new(restored_state.unwrap_or(State { + last_poll_timestamp: None, + total_documents_fetched: 0, + poll_count: 0, + })), + } + } + + fn serialize_state(&self, state: &State) -> Option<ConnectorState> { + ConnectorState::serialize(state, CONNECTOR_NAME, self.id) + } + + async fn create_client(&self) -> Result<Client, Error> { + let mut client_options: ClientOptions = + ClientOptions::parse(self.config.connection_uri.expose_secret()) + .await + .map_err(|e| Error::InitError(format!("Failed to parse connection URI: {e}")))?; + if let Some(pool_size) = self.config.max_pool_size { + client_options.max_pool_size = Some(pool_size); + } + let client = Client::with_options(client_options) + .map_err(|e| Error::InitError(format!("Failed to create client: {e}")))?; + Ok(client) + } + + async fn search_documents(&self, client: &Client) -> Result<Vec<ProducedMessage>, Error> { + let state = self.state.lock().await; + let limit = &self.config.limit.unwrap_or(100); + + let coll: Collection<Document> = client + .database(&self.config.database.to_string()) + .collection(&self.config.collection.to_string()); + + let mut cursor = if let Some(timestamp_field) = &self.config.timestamp_field + && let Some(last_timestamp) = state.last_poll_timestamp + { + let mut filter = self.config.query.clone().unwrap_or_default(); + filter.insert( + timestamp_field.clone(), + doc! { "$gt": mongodb::bson::DateTime::from_millis(last_timestamp.timestamp_millis()) }, + ); + + coll.find(filter.clone()) + .limit(*limit) + .sort(doc! { timestamp_field: 1 }) + .await + .map_err(|e| Error::InitError(format!("Failed to execute search: {e}")))? + } else { + let filter = &self.config.query.clone().unwrap_or_else(|| doc! {}); + coll.find(filter.clone()) + .await + .map_err(|e| Error::InitError(format!("Failed to execute search: {e}")))? + }; + + drop(state); + + let mut messages = Vec::new(); + let mut latest_timestamp = None; + + while let Some(doc) = cursor + .try_next() + .await + .map_err(|e| Error::InitError(format!("Failed to move cursor {e}")))? + { + if let Some(timestamp_field) = &self.config.timestamp_field + && let Some(timestamp_dt) = doc.get(timestamp_field).and_then(|v| v.as_datetime()) + && let Some(timestamp) = + iggy_common::DateTime::<iggy_common::Utc>::from_timestamp_millis( + timestamp_dt.timestamp_millis(), + ) + && latest_timestamp.is_none_or(|current| timestamp > current) + { + latest_timestamp = Some(timestamp); + } + + let payload = serde_json::to_vec(&doc).map_err(|e| { + Error::Serialization(format!("Failed to serialize document: {}", e)) + })?; + + let message = ProducedMessage { + id: None, + headers: None, + checksum: None, + timestamp: None, + origin_timestamp: None, + payload, + }; + messages.push(message); + } + let mut state = self.state.lock().await; Review Comment: critical: `search_documents` writes `state.last_poll_timestamp` inside `poll`, but `on_batch_result` keeps the default no-op, so a runtime `Nack` re-polls past the dropped batch. stage the cursor in `poll` and commit it on `Ack` only, as `random_source` does. ########## core/connectors/sources/mongodb_source/src/lib.rs: ########## @@ -0,0 +1,242 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +use async_trait::async_trait; +use core::mem::drop; +use futures::stream::TryStreamExt; +use iggy_common::{DateTime, Utc}; +use iggy_connector_sdk::{ + ConnectorState, Error, ProducedMessage, ProducedMessages, Schema, Source, source_connector, +}; +use mongodb::{Client, Collection, bson::Document, bson::doc, options::ClientOptions}; +use secrecy::{ExposeSecret, SecretString}; +use serde::{Deserialize, Serialize}; +use std::str::FromStr; +use std::time::Duration; +use tokio::sync::Mutex; +use tracing::info; + +source_connector!(MongodbSource); + +#[derive(Debug, Clone, Serialize, Deserialize)] +struct State { + last_poll_timestamp: Option<DateTime<Utc>>, + total_documents_fetched: usize, + poll_count: usize, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct MongodbSourceConfig { + #[serde(serialize_with = "iggy_common::serde_secret::serialize_secret")] + pub connection_uri: SecretString, + pub database: String, + pub collection: String, + pub max_pool_size: Option<u32>, + pub query: Option<Document>, + pub timestamp_field: Option<String>, + pub limit: Option<i64>, + pub polling_interval: Option<String>, +} + +#[derive(Debug)] +pub struct MongodbSource { + id: u32, + config: MongodbSourceConfig, + client: Option<Client>, + polling_interval: Duration, + state: Mutex<State>, +} + +const CONNECTOR_NAME: &str = "MongoDB source"; + +impl MongodbSource { + pub fn new(id: u32, config: MongodbSourceConfig, state: Option<ConnectorState>) -> Self { + let polling_interval = config + .polling_interval + .as_deref() + .unwrap_or("10s") + .parse::<humantime::Duration>() + .unwrap_or_else(|_| humantime::Duration::from_str("10s").unwrap()) + .into(); + + let restored_state = state + .and_then(|s| s.deserialize::<State>(CONNECTOR_NAME, id)) + .inspect(|s| { + info!( + "Restored state for {CONNECTOR_NAME} connector with ID: {id}. \ + Documents fetched: {}, poll count: {}", + s.total_documents_fetched, s.poll_count + ); + }); + + MongodbSource { + id, + config, + client: None, + polling_interval, + state: Mutex::new(restored_state.unwrap_or(State { + last_poll_timestamp: None, + total_documents_fetched: 0, + poll_count: 0, + })), + } + } + + fn serialize_state(&self, state: &State) -> Option<ConnectorState> { + ConnectorState::serialize(state, CONNECTOR_NAME, self.id) + } + + async fn create_client(&self) -> Result<Client, Error> { + let mut client_options: ClientOptions = + ClientOptions::parse(self.config.connection_uri.expose_secret()) + .await + .map_err(|e| Error::InitError(format!("Failed to parse connection URI: {e}")))?; + if let Some(pool_size) = self.config.max_pool_size { + client_options.max_pool_size = Some(pool_size); + } + let client = Client::with_options(client_options) + .map_err(|e| Error::InitError(format!("Failed to create client: {e}")))?; + Ok(client) + } + + async fn search_documents(&self, client: &Client) -> Result<Vec<ProducedMessage>, Error> { + let state = self.state.lock().await; + let limit = &self.config.limit.unwrap_or(100); + + let coll: Collection<Document> = client + .database(&self.config.database.to_string()) + .collection(&self.config.collection.to_string()); + + let mut cursor = if let Some(timestamp_field) = &self.config.timestamp_field + && let Some(last_timestamp) = state.last_poll_timestamp + { + let mut filter = self.config.query.clone().unwrap_or_default(); + filter.insert( + timestamp_field.clone(), + doc! { "$gt": mongodb::bson::DateTime::from_millis(last_timestamp.timestamp_millis()) }, + ); + + coll.find(filter.clone()) + .limit(*limit) + .sort(doc! { timestamp_field: 1 }) + .await + .map_err(|e| Error::InitError(format!("Failed to execute search: {e}")))? + } else { + let filter = &self.config.query.clone().unwrap_or_else(|| doc! {}); + coll.find(filter.clone()) + .await + .map_err(|e| Error::InitError(format!("Failed to execute search: {e}")))? + }; + + drop(state); + + let mut messages = Vec::new(); + let mut latest_timestamp = None; + + while let Some(doc) = cursor + .try_next() + .await + .map_err(|e| Error::InitError(format!("Failed to move cursor {e}")))? + { + if let Some(timestamp_field) = &self.config.timestamp_field Review Comment: critical: `as_datetime()` returns `None` for a string, an int, or a missing field, so `last_poll_timestamp` never advances and every poll takes the unbounded branch. check the field type in `open`, or parse strings and numbers. ########## core/connectors/sources/mongodb_source/src/lib.rs: ########## @@ -0,0 +1,242 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +use async_trait::async_trait; +use core::mem::drop; +use futures::stream::TryStreamExt; +use iggy_common::{DateTime, Utc}; +use iggy_connector_sdk::{ + ConnectorState, Error, ProducedMessage, ProducedMessages, Schema, Source, source_connector, +}; +use mongodb::{Client, Collection, bson::Document, bson::doc, options::ClientOptions}; +use secrecy::{ExposeSecret, SecretString}; +use serde::{Deserialize, Serialize}; +use std::str::FromStr; +use std::time::Duration; +use tokio::sync::Mutex; +use tracing::info; + +source_connector!(MongodbSource); + +#[derive(Debug, Clone, Serialize, Deserialize)] +struct State { + last_poll_timestamp: Option<DateTime<Utc>>, + total_documents_fetched: usize, + poll_count: usize, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct MongodbSourceConfig { + #[serde(serialize_with = "iggy_common::serde_secret::serialize_secret")] + pub connection_uri: SecretString, + pub database: String, + pub collection: String, + pub max_pool_size: Option<u32>, + pub query: Option<Document>, + pub timestamp_field: Option<String>, + pub limit: Option<i64>, + pub polling_interval: Option<String>, +} + +#[derive(Debug)] +pub struct MongodbSource { + id: u32, + config: MongodbSourceConfig, + client: Option<Client>, + polling_interval: Duration, + state: Mutex<State>, +} + +const CONNECTOR_NAME: &str = "MongoDB source"; + +impl MongodbSource { + pub fn new(id: u32, config: MongodbSourceConfig, state: Option<ConnectorState>) -> Self { + let polling_interval = config + .polling_interval + .as_deref() + .unwrap_or("10s") + .parse::<humantime::Duration>() + .unwrap_or_else(|_| humantime::Duration::from_str("10s").unwrap()) + .into(); + + let restored_state = state + .and_then(|s| s.deserialize::<State>(CONNECTOR_NAME, id)) + .inspect(|s| { + info!( + "Restored state for {CONNECTOR_NAME} connector with ID: {id}. \ + Documents fetched: {}, poll count: {}", + s.total_documents_fetched, s.poll_count + ); + }); + + MongodbSource { + id, + config, + client: None, + polling_interval, + state: Mutex::new(restored_state.unwrap_or(State { + last_poll_timestamp: None, + total_documents_fetched: 0, + poll_count: 0, + })), + } + } + + fn serialize_state(&self, state: &State) -> Option<ConnectorState> { + ConnectorState::serialize(state, CONNECTOR_NAME, self.id) + } + + async fn create_client(&self) -> Result<Client, Error> { + let mut client_options: ClientOptions = + ClientOptions::parse(self.config.connection_uri.expose_secret()) + .await + .map_err(|e| Error::InitError(format!("Failed to parse connection URI: {e}")))?; + if let Some(pool_size) = self.config.max_pool_size { + client_options.max_pool_size = Some(pool_size); + } + let client = Client::with_options(client_options) + .map_err(|e| Error::InitError(format!("Failed to create client: {e}")))?; + Ok(client) + } + + async fn search_documents(&self, client: &Client) -> Result<Vec<ProducedMessage>, Error> { + let state = self.state.lock().await; + let limit = &self.config.limit.unwrap_or(100); + + let coll: Collection<Document> = client + .database(&self.config.database.to_string()) + .collection(&self.config.collection.to_string()); + + let mut cursor = if let Some(timestamp_field) = &self.config.timestamp_field + && let Some(last_timestamp) = state.last_poll_timestamp + { + let mut filter = self.config.query.clone().unwrap_or_default(); + filter.insert( Review Comment: critical: the cursor stores the newest timestamp of a batch that `limit` truncated, so documents sharing that millisecond beyond the cutoff never come back. keep the last `_id` as a tiebreaker, or use `$gte` with a skip count. ########## Cargo.toml: ########## @@ -52,6 +52,7 @@ members = [ "core/connectors/sources/elasticsearch_source", "core/connectors/sources/http_source", "core/connectors/sources/influxdb_source", + "core/connectors/sources/mongodb_source", Review Comment: critical: this adds the workspace member and leaves `Cargo.lock` stale, so every `--locked` build fails at this head. commit the regenerated `Cargo.lock`. ########## core/connectors/sources/mongodb_source/Cargo.toml: ########## @@ -0,0 +1,49 @@ +# Licensed to the Apache Software Foundation (ASF) under one +# or more contributor license agreements. See the NOTICE file +# distributed with this work for additional information +# regarding copyright ownership. The ASF licenses this file +# to you under the Apache License, Version 2.0 (the +# "License"); you may not use this file except in compliance +# with the License. You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, +# software distributed under the License is distributed on an +# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +# KIND, either express or implied. See the License for the +# specific language governing permissions and limitations +# under the License. + +[package] +name = "iggy_connector_mongodb_source" +version = "0.1.0" +license = "Apache-2.0" +edition = "2024" +categories = ["command-line-utilities", "database", "network-programming"] +homepage = "https://iggy.apache.org" +documentation = "https://iggy.apache.org/docs" +repository = "https://github.com/apache/iggy" +readme = "../../README.md" + +[package.metadata.cargo-machete] +ignored = ["dashmap", "once_cell", "simd-json"] + +[lib] +crate-type = ["cdylib", "lib"] + +[dependencies] +async-trait = { workspace = true } +bson = "3.1.0" Review Comment: warning: the direct `bson = "3.1.0"` pulls a second bson major while `mongodb` re-exports `bson 2.15.0`, so the build compiles two incompatible bson crates. drop the direct dependency and reach bson through `mongodb::bson`. ########## core/connectors/sources/mongodb_source/src/lib.rs: ########## @@ -0,0 +1,242 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +use async_trait::async_trait; +use core::mem::drop; +use futures::stream::TryStreamExt; +use iggy_common::{DateTime, Utc}; +use iggy_connector_sdk::{ + ConnectorState, Error, ProducedMessage, ProducedMessages, Schema, Source, source_connector, +}; +use mongodb::{Client, Collection, bson::Document, bson::doc, options::ClientOptions}; +use secrecy::{ExposeSecret, SecretString}; +use serde::{Deserialize, Serialize}; +use std::str::FromStr; +use std::time::Duration; +use tokio::sync::Mutex; +use tracing::info; + +source_connector!(MongodbSource); + +#[derive(Debug, Clone, Serialize, Deserialize)] +struct State { + last_poll_timestamp: Option<DateTime<Utc>>, + total_documents_fetched: usize, + poll_count: usize, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct MongodbSourceConfig { + #[serde(serialize_with = "iggy_common::serde_secret::serialize_secret")] + pub connection_uri: SecretString, + pub database: String, + pub collection: String, + pub max_pool_size: Option<u32>, + pub query: Option<Document>, + pub timestamp_field: Option<String>, + pub limit: Option<i64>, + pub polling_interval: Option<String>, +} + +#[derive(Debug)] +pub struct MongodbSource { + id: u32, + config: MongodbSourceConfig, + client: Option<Client>, + polling_interval: Duration, + state: Mutex<State>, +} + +const CONNECTOR_NAME: &str = "MongoDB source"; + +impl MongodbSource { + pub fn new(id: u32, config: MongodbSourceConfig, state: Option<ConnectorState>) -> Self { + let polling_interval = config + .polling_interval + .as_deref() + .unwrap_or("10s") + .parse::<humantime::Duration>() + .unwrap_or_else(|_| humantime::Duration::from_str("10s").unwrap()) + .into(); + + let restored_state = state + .and_then(|s| s.deserialize::<State>(CONNECTOR_NAME, id)) + .inspect(|s| { + info!( + "Restored state for {CONNECTOR_NAME} connector with ID: {id}. \ + Documents fetched: {}, poll count: {}", + s.total_documents_fetched, s.poll_count + ); + }); + + MongodbSource { + id, + config, + client: None, + polling_interval, + state: Mutex::new(restored_state.unwrap_or(State { + last_poll_timestamp: None, + total_documents_fetched: 0, + poll_count: 0, + })), + } + } + + fn serialize_state(&self, state: &State) -> Option<ConnectorState> { + ConnectorState::serialize(state, CONNECTOR_NAME, self.id) + } + + async fn create_client(&self) -> Result<Client, Error> { + let mut client_options: ClientOptions = + ClientOptions::parse(self.config.connection_uri.expose_secret()) + .await + .map_err(|e| Error::InitError(format!("Failed to parse connection URI: {e}")))?; + if let Some(pool_size) = self.config.max_pool_size { + client_options.max_pool_size = Some(pool_size); + } + let client = Client::with_options(client_options) + .map_err(|e| Error::InitError(format!("Failed to create client: {e}")))?; + Ok(client) + } + + async fn search_documents(&self, client: &Client) -> Result<Vec<ProducedMessage>, Error> { + let state = self.state.lock().await; + let limit = &self.config.limit.unwrap_or(100); + + let coll: Collection<Document> = client + .database(&self.config.database.to_string()) + .collection(&self.config.collection.to_string()); + + let mut cursor = if let Some(timestamp_field) = &self.config.timestamp_field + && let Some(last_timestamp) = state.last_poll_timestamp + { + let mut filter = self.config.query.clone().unwrap_or_default(); + filter.insert( + timestamp_field.clone(), + doc! { "$gt": mongodb::bson::DateTime::from_millis(last_timestamp.timestamp_millis()) }, + ); + + coll.find(filter.clone()) + .limit(*limit) + .sort(doc! { timestamp_field: 1 }) + .await + .map_err(|e| Error::InitError(format!("Failed to execute search: {e}")))? + } else { + let filter = &self.config.query.clone().unwrap_or_else(|| doc! {}); + coll.find(filter.clone()) Review Comment: critical: the else branch runs `coll.find(filter)` with no `.limit()`, so a first poll on a large collection loads all of it into one batch. apply the same limit and sort on both branches. ########## core/connectors/sources/mongodb_source/src/lib.rs: ########## @@ -0,0 +1,242 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +use async_trait::async_trait; +use core::mem::drop; +use futures::stream::TryStreamExt; +use iggy_common::{DateTime, Utc}; +use iggy_connector_sdk::{ + ConnectorState, Error, ProducedMessage, ProducedMessages, Schema, Source, source_connector, +}; +use mongodb::{Client, Collection, bson::Document, bson::doc, options::ClientOptions}; +use secrecy::{ExposeSecret, SecretString}; +use serde::{Deserialize, Serialize}; +use std::str::FromStr; +use std::time::Duration; +use tokio::sync::Mutex; +use tracing::info; + +source_connector!(MongodbSource); + +#[derive(Debug, Clone, Serialize, Deserialize)] +struct State { + last_poll_timestamp: Option<DateTime<Utc>>, + total_documents_fetched: usize, + poll_count: usize, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct MongodbSourceConfig { + #[serde(serialize_with = "iggy_common::serde_secret::serialize_secret")] + pub connection_uri: SecretString, + pub database: String, + pub collection: String, + pub max_pool_size: Option<u32>, + pub query: Option<Document>, + pub timestamp_field: Option<String>, + pub limit: Option<i64>, + pub polling_interval: Option<String>, +} + +#[derive(Debug)] +pub struct MongodbSource { + id: u32, + config: MongodbSourceConfig, + client: Option<Client>, + polling_interval: Duration, + state: Mutex<State>, +} + +const CONNECTOR_NAME: &str = "MongoDB source"; + +impl MongodbSource { + pub fn new(id: u32, config: MongodbSourceConfig, state: Option<ConnectorState>) -> Self { + let polling_interval = config + .polling_interval + .as_deref() + .unwrap_or("10s") + .parse::<humantime::Duration>() + .unwrap_or_else(|_| humantime::Duration::from_str("10s").unwrap()) + .into(); + + let restored_state = state + .and_then(|s| s.deserialize::<State>(CONNECTOR_NAME, id)) + .inspect(|s| { + info!( + "Restored state for {CONNECTOR_NAME} connector with ID: {id}. \ + Documents fetched: {}, poll count: {}", + s.total_documents_fetched, s.poll_count + ); + }); + + MongodbSource { + id, + config, + client: None, + polling_interval, + state: Mutex::new(restored_state.unwrap_or(State { + last_poll_timestamp: None, + total_documents_fetched: 0, + poll_count: 0, + })), + } + } + + fn serialize_state(&self, state: &State) -> Option<ConnectorState> { + ConnectorState::serialize(state, CONNECTOR_NAME, self.id) + } + + async fn create_client(&self) -> Result<Client, Error> { + let mut client_options: ClientOptions = + ClientOptions::parse(self.config.connection_uri.expose_secret()) + .await + .map_err(|e| Error::InitError(format!("Failed to parse connection URI: {e}")))?; + if let Some(pool_size) = self.config.max_pool_size { + client_options.max_pool_size = Some(pool_size); + } + let client = Client::with_options(client_options) + .map_err(|e| Error::InitError(format!("Failed to create client: {e}")))?; + Ok(client) + } + + async fn search_documents(&self, client: &Client) -> Result<Vec<ProducedMessage>, Error> { + let state = self.state.lock().await; + let limit = &self.config.limit.unwrap_or(100); + + let coll: Collection<Document> = client + .database(&self.config.database.to_string()) + .collection(&self.config.collection.to_string()); + + let mut cursor = if let Some(timestamp_field) = &self.config.timestamp_field + && let Some(last_timestamp) = state.last_poll_timestamp + { + let mut filter = self.config.query.clone().unwrap_or_default(); + filter.insert( + timestamp_field.clone(), + doc! { "$gt": mongodb::bson::DateTime::from_millis(last_timestamp.timestamp_millis()) }, + ); + + coll.find(filter.clone()) + .limit(*limit) + .sort(doc! { timestamp_field: 1 }) + .await + .map_err(|e| Error::InitError(format!("Failed to execute search: {e}")))? + } else { + let filter = &self.config.query.clone().unwrap_or_else(|| doc! {}); + coll.find(filter.clone()) + .await + .map_err(|e| Error::InitError(format!("Failed to execute search: {e}")))? + }; + + drop(state); + + let mut messages = Vec::new(); + let mut latest_timestamp = None; + + while let Some(doc) = cursor + .try_next() + .await + .map_err(|e| Error::InitError(format!("Failed to move cursor {e}")))? + { + if let Some(timestamp_field) = &self.config.timestamp_field + && let Some(timestamp_dt) = doc.get(timestamp_field).and_then(|v| v.as_datetime()) + && let Some(timestamp) = + iggy_common::DateTime::<iggy_common::Utc>::from_timestamp_millis( + timestamp_dt.timestamp_millis(), + ) + && latest_timestamp.is_none_or(|current| timestamp > current) + { + latest_timestamp = Some(timestamp); + } + + let payload = serde_json::to_vec(&doc).map_err(|e| { + Error::Serialization(format!("Failed to serialize document: {}", e)) + })?; + + let message = ProducedMessage { + id: None, + headers: None, + checksum: None, + timestamp: None, + origin_timestamp: None, + payload, + }; + messages.push(message); + } + let mut state = self.state.lock().await; + state.total_documents_fetched += messages.len(); + state.poll_count += 1; + if let Some(timestamp) = latest_timestamp { + state.last_poll_timestamp = Some(timestamp); + } + Ok(messages) + } +} + +#[async_trait] +impl Source for MongodbSource { + async fn open(&mut self) -> Result<(), Error> { Review Comment: warning: `open` builds a lazy client and validates nothing, so a wrong database or collection name passes and every later `find` returns an empty batch. ping the server and confirm the database and collection in `open`. ########## core/connectors/sources/mongodb_source/src/lib.rs: ########## @@ -0,0 +1,242 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +use async_trait::async_trait; +use core::mem::drop; +use futures::stream::TryStreamExt; +use iggy_common::{DateTime, Utc}; +use iggy_connector_sdk::{ + ConnectorState, Error, ProducedMessage, ProducedMessages, Schema, Source, source_connector, +}; +use mongodb::{Client, Collection, bson::Document, bson::doc, options::ClientOptions}; +use secrecy::{ExposeSecret, SecretString}; +use serde::{Deserialize, Serialize}; +use std::str::FromStr; +use std::time::Duration; +use tokio::sync::Mutex; +use tracing::info; + +source_connector!(MongodbSource); + +#[derive(Debug, Clone, Serialize, Deserialize)] +struct State { + last_poll_timestamp: Option<DateTime<Utc>>, + total_documents_fetched: usize, + poll_count: usize, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct MongodbSourceConfig { + #[serde(serialize_with = "iggy_common::serde_secret::serialize_secret")] + pub connection_uri: SecretString, + pub database: String, + pub collection: String, + pub max_pool_size: Option<u32>, + pub query: Option<Document>, + pub timestamp_field: Option<String>, + pub limit: Option<i64>, + pub polling_interval: Option<String>, +} + +#[derive(Debug)] +pub struct MongodbSource { + id: u32, + config: MongodbSourceConfig, + client: Option<Client>, + polling_interval: Duration, + state: Mutex<State>, +} + +const CONNECTOR_NAME: &str = "MongoDB source"; + +impl MongodbSource { + pub fn new(id: u32, config: MongodbSourceConfig, state: Option<ConnectorState>) -> Self { + let polling_interval = config + .polling_interval + .as_deref() + .unwrap_or("10s") + .parse::<humantime::Duration>() + .unwrap_or_else(|_| humantime::Duration::from_str("10s").unwrap()) + .into(); + + let restored_state = state + .and_then(|s| s.deserialize::<State>(CONNECTOR_NAME, id)) + .inspect(|s| { + info!( + "Restored state for {CONNECTOR_NAME} connector with ID: {id}. \ + Documents fetched: {}, poll count: {}", + s.total_documents_fetched, s.poll_count + ); + }); + + MongodbSource { + id, + config, + client: None, + polling_interval, + state: Mutex::new(restored_state.unwrap_or(State { + last_poll_timestamp: None, + total_documents_fetched: 0, + poll_count: 0, + })), + } + } + + fn serialize_state(&self, state: &State) -> Option<ConnectorState> { + ConnectorState::serialize(state, CONNECTOR_NAME, self.id) + } + + async fn create_client(&self) -> Result<Client, Error> { + let mut client_options: ClientOptions = + ClientOptions::parse(self.config.connection_uri.expose_secret()) + .await + .map_err(|e| Error::InitError(format!("Failed to parse connection URI: {e}")))?; + if let Some(pool_size) = self.config.max_pool_size { + client_options.max_pool_size = Some(pool_size); + } + let client = Client::with_options(client_options) + .map_err(|e| Error::InitError(format!("Failed to create client: {e}")))?; + Ok(client) + } + + async fn search_documents(&self, client: &Client) -> Result<Vec<ProducedMessage>, Error> { + let state = self.state.lock().await; + let limit = &self.config.limit.unwrap_or(100); + + let coll: Collection<Document> = client + .database(&self.config.database.to_string()) + .collection(&self.config.collection.to_string()); + + let mut cursor = if let Some(timestamp_field) = &self.config.timestamp_field + && let Some(last_timestamp) = state.last_poll_timestamp + { + let mut filter = self.config.query.clone().unwrap_or_default(); + filter.insert( + timestamp_field.clone(), + doc! { "$gt": mongodb::bson::DateTime::from_millis(last_timestamp.timestamp_millis()) }, + ); + + coll.find(filter.clone()) + .limit(*limit) + .sort(doc! { timestamp_field: 1 }) + .await + .map_err(|e| Error::InitError(format!("Failed to execute search: {e}")))? Review Comment: nit: a failed search returns `Error::InitError`, which names an initialization fault for a runtime query. use `Error::Storage` here, as `elasticsearch_source` does on the same path. also at lines 142, 153. ########## core/connectors/sources/mongodb_source/Cargo.toml: ########## @@ -0,0 +1,49 @@ +# Licensed to the Apache Software Foundation (ASF) under one +# or more contributor license agreements. See the NOTICE file +# distributed with this work for additional information +# regarding copyright ownership. The ASF licenses this file +# to you under the Apache License, Version 2.0 (the +# "License"); you may not use this file except in compliance +# with the License. You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, +# software distributed under the License is distributed on an +# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +# KIND, either express or implied. See the License for the +# specific language governing permissions and limitations +# under the License. + +[package] +name = "iggy_connector_mongodb_source" Review Comment: nit: every other connector crate carries `publish = false`, `description` and `keywords`, and this one carries none of them. add them. ########## core/connectors/sources/mongodb_source/src/lib.rs: ########## @@ -0,0 +1,242 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +use async_trait::async_trait; +use core::mem::drop; +use futures::stream::TryStreamExt; +use iggy_common::{DateTime, Utc}; +use iggy_connector_sdk::{ + ConnectorState, Error, ProducedMessage, ProducedMessages, Schema, Source, source_connector, +}; +use mongodb::{Client, Collection, bson::Document, bson::doc, options::ClientOptions}; +use secrecy::{ExposeSecret, SecretString}; +use serde::{Deserialize, Serialize}; +use std::str::FromStr; +use std::time::Duration; +use tokio::sync::Mutex; +use tracing::info; + +source_connector!(MongodbSource); + +#[derive(Debug, Clone, Serialize, Deserialize)] +struct State { + last_poll_timestamp: Option<DateTime<Utc>>, + total_documents_fetched: usize, + poll_count: usize, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct MongodbSourceConfig { + #[serde(serialize_with = "iggy_common::serde_secret::serialize_secret")] + pub connection_uri: SecretString, + pub database: String, + pub collection: String, + pub max_pool_size: Option<u32>, + pub query: Option<Document>, + pub timestamp_field: Option<String>, + pub limit: Option<i64>, + pub polling_interval: Option<String>, +} + +#[derive(Debug)] +pub struct MongodbSource { + id: u32, + config: MongodbSourceConfig, + client: Option<Client>, + polling_interval: Duration, + state: Mutex<State>, +} + +const CONNECTOR_NAME: &str = "MongoDB source"; + +impl MongodbSource { + pub fn new(id: u32, config: MongodbSourceConfig, state: Option<ConnectorState>) -> Self { + let polling_interval = config + .polling_interval + .as_deref() + .unwrap_or("10s") + .parse::<humantime::Duration>() + .unwrap_or_else(|_| humantime::Duration::from_str("10s").unwrap()) + .into(); + + let restored_state = state + .and_then(|s| s.deserialize::<State>(CONNECTOR_NAME, id)) + .inspect(|s| { + info!( + "Restored state for {CONNECTOR_NAME} connector with ID: {id}. \ + Documents fetched: {}, poll count: {}", + s.total_documents_fetched, s.poll_count + ); + }); + + MongodbSource { + id, + config, + client: None, + polling_interval, + state: Mutex::new(restored_state.unwrap_or(State { + last_poll_timestamp: None, + total_documents_fetched: 0, + poll_count: 0, + })), + } + } + + fn serialize_state(&self, state: &State) -> Option<ConnectorState> { + ConnectorState::serialize(state, CONNECTOR_NAME, self.id) + } + + async fn create_client(&self) -> Result<Client, Error> { + let mut client_options: ClientOptions = + ClientOptions::parse(self.config.connection_uri.expose_secret()) + .await + .map_err(|e| Error::InitError(format!("Failed to parse connection URI: {e}")))?; + if let Some(pool_size) = self.config.max_pool_size { + client_options.max_pool_size = Some(pool_size); + } + let client = Client::with_options(client_options) + .map_err(|e| Error::InitError(format!("Failed to create client: {e}")))?; + Ok(client) + } + + async fn search_documents(&self, client: &Client) -> Result<Vec<ProducedMessage>, Error> { + let state = self.state.lock().await; + let limit = &self.config.limit.unwrap_or(100); + + let coll: Collection<Document> = client + .database(&self.config.database.to_string()) + .collection(&self.config.collection.to_string()); + + let mut cursor = if let Some(timestamp_field) = &self.config.timestamp_field + && let Some(last_timestamp) = state.last_poll_timestamp + { + let mut filter = self.config.query.clone().unwrap_or_default(); + filter.insert( + timestamp_field.clone(), + doc! { "$gt": mongodb::bson::DateTime::from_millis(last_timestamp.timestamp_millis()) }, + ); + + coll.find(filter.clone()) Review Comment: simplification: `filter` is cloned here and at 140 although `find` takes it by value, and lines 121-122 build `String` names that the driver takes as `&str`. pass the filter by value and borrow the names. ########## core/connectors/sources/mongodb_source/Cargo.toml: ########## @@ -0,0 +1,49 @@ +# Licensed to the Apache Software Foundation (ASF) under one +# or more contributor license agreements. See the NOTICE file +# distributed with this work for additional information +# regarding copyright ownership. The ASF licenses this file +# to you under the Apache License, Version 2.0 (the +# "License"); you may not use this file except in compliance +# with the License. You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, +# software distributed under the License is distributed on an +# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +# KIND, either express or implied. See the License for the +# specific language governing permissions and limitations +# under the License. + +[package] +name = "iggy_connector_mongodb_source" +version = "0.1.0" +license = "Apache-2.0" +edition = "2024" +categories = ["command-line-utilities", "database", "network-programming"] +homepage = "https://iggy.apache.org" +documentation = "https://iggy.apache.org/docs" +repository = "https://github.com/apache/iggy" +readme = "../../README.md" + +[package.metadata.cargo-machete] +ignored = ["dashmap", "once_cell", "simd-json"] Review Comment: simplification: `dashmap` and `simd-json` are declared but never named, and the machete ignore list hides them along with a stale `once_cell` entry. delete both dependencies and the ignore block. ########## core/connectors/sources/mongodb_source/src/lib.rs: ########## @@ -0,0 +1,242 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +use async_trait::async_trait; +use core::mem::drop; +use futures::stream::TryStreamExt; +use iggy_common::{DateTime, Utc}; +use iggy_connector_sdk::{ + ConnectorState, Error, ProducedMessage, ProducedMessages, Schema, Source, source_connector, +}; +use mongodb::{Client, Collection, bson::Document, bson::doc, options::ClientOptions}; +use secrecy::{ExposeSecret, SecretString}; +use serde::{Deserialize, Serialize}; +use std::str::FromStr; +use std::time::Duration; +use tokio::sync::Mutex; +use tracing::info; + +source_connector!(MongodbSource); + +#[derive(Debug, Clone, Serialize, Deserialize)] +struct State { + last_poll_timestamp: Option<DateTime<Utc>>, + total_documents_fetched: usize, + poll_count: usize, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct MongodbSourceConfig { + #[serde(serialize_with = "iggy_common::serde_secret::serialize_secret")] + pub connection_uri: SecretString, + pub database: String, + pub collection: String, + pub max_pool_size: Option<u32>, + pub query: Option<Document>, + pub timestamp_field: Option<String>, + pub limit: Option<i64>, + pub polling_interval: Option<String>, +} + +#[derive(Debug)] +pub struct MongodbSource { + id: u32, + config: MongodbSourceConfig, + client: Option<Client>, + polling_interval: Duration, + state: Mutex<State>, +} + +const CONNECTOR_NAME: &str = "MongoDB source"; + +impl MongodbSource { + pub fn new(id: u32, config: MongodbSourceConfig, state: Option<ConnectorState>) -> Self { + let polling_interval = config + .polling_interval + .as_deref() + .unwrap_or("10s") + .parse::<humantime::Duration>() + .unwrap_or_else(|_| humantime::Duration::from_str("10s").unwrap()) + .into(); + + let restored_state = state + .and_then(|s| s.deserialize::<State>(CONNECTOR_NAME, id)) + .inspect(|s| { + info!( + "Restored state for {CONNECTOR_NAME} connector with ID: {id}. \ + Documents fetched: {}, poll count: {}", + s.total_documents_fetched, s.poll_count + ); + }); + + MongodbSource { + id, + config, + client: None, + polling_interval, + state: Mutex::new(restored_state.unwrap_or(State { + last_poll_timestamp: None, + total_documents_fetched: 0, + poll_count: 0, + })), + } + } + + fn serialize_state(&self, state: &State) -> Option<ConnectorState> { + ConnectorState::serialize(state, CONNECTOR_NAME, self.id) + } + + async fn create_client(&self) -> Result<Client, Error> { + let mut client_options: ClientOptions = + ClientOptions::parse(self.config.connection_uri.expose_secret()) + .await + .map_err(|e| Error::InitError(format!("Failed to parse connection URI: {e}")))?; + if let Some(pool_size) = self.config.max_pool_size { + client_options.max_pool_size = Some(pool_size); + } + let client = Client::with_options(client_options) + .map_err(|e| Error::InitError(format!("Failed to create client: {e}")))?; + Ok(client) + } + + async fn search_documents(&self, client: &Client) -> Result<Vec<ProducedMessage>, Error> { + let state = self.state.lock().await; + let limit = &self.config.limit.unwrap_or(100); + + let coll: Collection<Document> = client + .database(&self.config.database.to_string()) + .collection(&self.config.collection.to_string()); + + let mut cursor = if let Some(timestamp_field) = &self.config.timestamp_field + && let Some(last_timestamp) = state.last_poll_timestamp + { + let mut filter = self.config.query.clone().unwrap_or_default(); + filter.insert( + timestamp_field.clone(), + doc! { "$gt": mongodb::bson::DateTime::from_millis(last_timestamp.timestamp_millis()) }, + ); + + coll.find(filter.clone()) + .limit(*limit) + .sort(doc! { timestamp_field: 1 }) + .await + .map_err(|e| Error::InitError(format!("Failed to execute search: {e}")))? + } else { + let filter = &self.config.query.clone().unwrap_or_else(|| doc! {}); + coll.find(filter.clone()) + .await + .map_err(|e| Error::InitError(format!("Failed to execute search: {e}")))? + }; + + drop(state); + + let mut messages = Vec::new(); + let mut latest_timestamp = None; + + while let Some(doc) = cursor + .try_next() + .await + .map_err(|e| Error::InitError(format!("Failed to move cursor {e}")))? + { + if let Some(timestamp_field) = &self.config.timestamp_field + && let Some(timestamp_dt) = doc.get(timestamp_field).and_then(|v| v.as_datetime()) + && let Some(timestamp) = + iggy_common::DateTime::<iggy_common::Utc>::from_timestamp_millis( + timestamp_dt.timestamp_millis(), + ) + && latest_timestamp.is_none_or(|current| timestamp > current) + { + latest_timestamp = Some(timestamp); + } + + let payload = serde_json::to_vec(&doc).map_err(|e| { + Error::Serialization(format!("Failed to serialize document: {}", e)) + })?; + + let message = ProducedMessage { + id: None, + headers: None, + checksum: None, + timestamp: None, + origin_timestamp: None, + payload, + }; + messages.push(message); + } + let mut state = self.state.lock().await; + state.total_documents_fetched += messages.len(); + state.poll_count += 1; + if let Some(timestamp) = latest_timestamp { + state.last_poll_timestamp = Some(timestamp); + } + Ok(messages) + } +} + +#[async_trait] +impl Source for MongodbSource { + async fn open(&mut self) -> Result<(), Error> { + info!( + "Opening Mongodb source connector with ID: {}, collection: {}", + self.id, self.config.collection + ); + + let client = self.create_client().await?; + self.client = Some(client); + + Ok(()) + } + + async fn poll(&self) -> Result<ProducedMessages, Error> { + let poll_interval = self.polling_interval; + tokio::time::sleep(poll_interval).await; + + let client = self + .client + .as_ref() + .ok_or_else(|| Error::Storage("Mongodb client not initialized".to_string()))?; + + let messages = match self.search_documents(client).await { Review Comment: simplification: this `match` passes the error straight through, so the `Err` arm adds nothing. replace it with `let messages = self.search_documents(client).await?;`. ########## core/integration/tests/connectors/fixtures/mongodb/source.rs: ########## @@ -0,0 +1,200 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +use super::container::{ + DEFAULT_SOURCE_COLLECTION, DEFAULT_SOURCE_DATABASE, DEFAULT_TEST_STREAM, DEFAULT_TEST_TOPIC, + ENV_SOURCE_COLLECTION, ENV_SOURCE_CONNECTION_URI, ENV_SOURCE_DATABASE, ENV_SOURCE_LIMIT, + ENV_SOURCE_PATH, ENV_SOURCE_POLLING_INTERVAL, ENV_SOURCE_STREAMS_0_SCHEMA, + ENV_SOURCE_STREAMS_0_STREAM, ENV_SOURCE_STREAMS_0_TOPIC, ENV_SOURCE_TIMESTAMP_FIELD, + MongoDbContainer, MongoDbOps, +}; +use async_trait::async_trait; +use integration::harness::{TestBinaryError, TestFixture}; +use mongodb::bson::{DateTime as BsonDateTime, Document, doc}; +use std::collections::HashMap; + +/// MongoDB source fixture for basic document polling. +pub struct MongodbSourceFixture { + container: MongoDbContainer, +} + +impl MongoDbOps for MongodbSourceFixture { + fn container(&self) -> &MongoDbContainer { + &self.container + } +} + +impl MongodbSourceFixture { + #[allow(dead_code)] Review Comment: simplification: `database_name()` and `collection_name()` have no callers, and the two `#[allow(dead_code)]` attributes hide that. delete both. ########## core/connectors/runtime/example_config/connectors/mongodb_source.toml: ########## @@ -0,0 +1,37 @@ +# Licensed to the Apache Software Foundation (ASF) under one +# or more contributor license agreements. See the NOTICE file +# distributed with this work for additional information +# regarding copyright ownership. The ASF licenses this file +# to you under the Apache License, Version 2.0 (the +# "License"); you may not use this file except in compliance +# with the License. You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, +# software distributed under the License is distributed on an +# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +# KIND, either express or implied. See the License for the +# specific language governing permissions and limitations +# under the License. + +type = "source" +key = "mongodb" +enabled = true +version = 0 +name = "MongoDB source" +path = "target/release/libiggy_connector_mongodb_source" + +[[streams]] +stream = "test_stream" +topic = "test_topic" +schema = "json" +batch_length = 100 + +[plugin_config] +connection_uri = "mongodb://admin:admin123@localhost:27017" +database = "test_source" +collection = "test_messages" +limit = 100 +polling_interval = "100ms" +timestamp_field = "timestamp" Review Comment: warning: the file ends without a trailing newline, so `taplo fmt --check` and the trailing-newline check both fail. add the newline. ########## core/connectors/sources/mongodb_source/src/lib.rs: ########## @@ -0,0 +1,242 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +use async_trait::async_trait; +use core::mem::drop; +use futures::stream::TryStreamExt; +use iggy_common::{DateTime, Utc}; +use iggy_connector_sdk::{ + ConnectorState, Error, ProducedMessage, ProducedMessages, Schema, Source, source_connector, +}; +use mongodb::{Client, Collection, bson::Document, bson::doc, options::ClientOptions}; +use secrecy::{ExposeSecret, SecretString}; +use serde::{Deserialize, Serialize}; +use std::str::FromStr; +use std::time::Duration; +use tokio::sync::Mutex; +use tracing::info; + +source_connector!(MongodbSource); + +#[derive(Debug, Clone, Serialize, Deserialize)] +struct State { + last_poll_timestamp: Option<DateTime<Utc>>, + total_documents_fetched: usize, + poll_count: usize, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct MongodbSourceConfig { + #[serde(serialize_with = "iggy_common::serde_secret::serialize_secret")] + pub connection_uri: SecretString, + pub database: String, + pub collection: String, + pub max_pool_size: Option<u32>, + pub query: Option<Document>, + pub timestamp_field: Option<String>, + pub limit: Option<i64>, + pub polling_interval: Option<String>, +} + +#[derive(Debug)] +pub struct MongodbSource { + id: u32, + config: MongodbSourceConfig, + client: Option<Client>, + polling_interval: Duration, + state: Mutex<State>, +} + +const CONNECTOR_NAME: &str = "MongoDB source"; + +impl MongodbSource { + pub fn new(id: u32, config: MongodbSourceConfig, state: Option<ConnectorState>) -> Self { + let polling_interval = config + .polling_interval + .as_deref() + .unwrap_or("10s") + .parse::<humantime::Duration>() + .unwrap_or_else(|_| humantime::Duration::from_str("10s").unwrap()) + .into(); + + let restored_state = state + .and_then(|s| s.deserialize::<State>(CONNECTOR_NAME, id)) + .inspect(|s| { + info!( + "Restored state for {CONNECTOR_NAME} connector with ID: {id}. \ + Documents fetched: {}, poll count: {}", + s.total_documents_fetched, s.poll_count + ); + }); + + MongodbSource { + id, + config, + client: None, + polling_interval, + state: Mutex::new(restored_state.unwrap_or(State { + last_poll_timestamp: None, + total_documents_fetched: 0, + poll_count: 0, + })), + } + } + + fn serialize_state(&self, state: &State) -> Option<ConnectorState> { + ConnectorState::serialize(state, CONNECTOR_NAME, self.id) + } + + async fn create_client(&self) -> Result<Client, Error> { + let mut client_options: ClientOptions = + ClientOptions::parse(self.config.connection_uri.expose_secret()) + .await + .map_err(|e| Error::InitError(format!("Failed to parse connection URI: {e}")))?; + if let Some(pool_size) = self.config.max_pool_size { + client_options.max_pool_size = Some(pool_size); + } + let client = Client::with_options(client_options) + .map_err(|e| Error::InitError(format!("Failed to create client: {e}")))?; + Ok(client) + } + + async fn search_documents(&self, client: &Client) -> Result<Vec<ProducedMessage>, Error> { + let state = self.state.lock().await; + let limit = &self.config.limit.unwrap_or(100); + + let coll: Collection<Document> = client + .database(&self.config.database.to_string()) + .collection(&self.config.collection.to_string()); + + let mut cursor = if let Some(timestamp_field) = &self.config.timestamp_field + && let Some(last_timestamp) = state.last_poll_timestamp + { + let mut filter = self.config.query.clone().unwrap_or_default(); + filter.insert( + timestamp_field.clone(), + doc! { "$gt": mongodb::bson::DateTime::from_millis(last_timestamp.timestamp_millis()) }, + ); + + coll.find(filter.clone()) + .limit(*limit) + .sort(doc! { timestamp_field: 1 }) + .await + .map_err(|e| Error::InitError(format!("Failed to execute search: {e}")))? + } else { + let filter = &self.config.query.clone().unwrap_or_else(|| doc! {}); + coll.find(filter.clone()) + .await + .map_err(|e| Error::InitError(format!("Failed to execute search: {e}")))? + }; + + drop(state); + + let mut messages = Vec::new(); + let mut latest_timestamp = None; + + while let Some(doc) = cursor + .try_next() + .await + .map_err(|e| Error::InitError(format!("Failed to move cursor {e}")))? + { + if let Some(timestamp_field) = &self.config.timestamp_field + && let Some(timestamp_dt) = doc.get(timestamp_field).and_then(|v| v.as_datetime()) + && let Some(timestamp) = + iggy_common::DateTime::<iggy_common::Utc>::from_timestamp_millis( + timestamp_dt.timestamp_millis(), + ) + && latest_timestamp.is_none_or(|current| timestamp > current) + { + latest_timestamp = Some(timestamp); + } + + let payload = serde_json::to_vec(&doc).map_err(|e| { + Error::Serialization(format!("Failed to serialize document: {}", e)) + })?; + + let message = ProducedMessage { + id: None, + headers: None, + checksum: None, + timestamp: None, + origin_timestamp: None, + payload, + }; + messages.push(message); + } + let mut state = self.state.lock().await; + state.total_documents_fetched += messages.len(); + state.poll_count += 1; + if let Some(timestamp) = latest_timestamp { + state.last_poll_timestamp = Some(timestamp); + } + Ok(messages) + } +} + +#[async_trait] +impl Source for MongodbSource { + async fn open(&mut self) -> Result<(), Error> { + info!( + "Opening Mongodb source connector with ID: {}, collection: {}", + self.id, self.config.collection + ); + + let client = self.create_client().await?; + self.client = Some(client); + + Ok(()) + } + + async fn poll(&self) -> Result<ProducedMessages, Error> { + let poll_interval = self.polling_interval; + tokio::time::sleep(poll_interval).await; + + let client = self + .client + .as_ref() + .ok_or_else(|| Error::Storage("Mongodb client not initialized".to_string()))?; + + let messages = match self.search_documents(client).await { + Ok(msgs) => msgs, + Err(e) => { + return Err(e); + } + }; + + let persisted_state = { + let state = self.state.lock().await; + self.serialize_state(&state) + }; + + Ok(ProducedMessages { + schema: Schema::Json, + messages, + state: persisted_state, + }) + } + async fn close(&mut self) -> Result<(), Error> { + info!("Mongodb Connector with ID: {} is closing", self.id); + + let state = self.state.lock().await; + + info!( + "PostgreSQL source connector ID: {} closed. Total documents processed: {}", Review Comment: nit: this close log prints "PostgreSQL source connector ID" from the MongoDB source. print `CONNECTOR_NAME` instead. ########## core/connectors/sources/mongodb_source/src/lib.rs: ########## @@ -0,0 +1,242 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +use async_trait::async_trait; +use core::mem::drop; +use futures::stream::TryStreamExt; +use iggy_common::{DateTime, Utc}; +use iggy_connector_sdk::{ + ConnectorState, Error, ProducedMessage, ProducedMessages, Schema, Source, source_connector, +}; +use mongodb::{Client, Collection, bson::Document, bson::doc, options::ClientOptions}; +use secrecy::{ExposeSecret, SecretString}; +use serde::{Deserialize, Serialize}; +use std::str::FromStr; +use std::time::Duration; +use tokio::sync::Mutex; +use tracing::info; + +source_connector!(MongodbSource); + +#[derive(Debug, Clone, Serialize, Deserialize)] +struct State { + last_poll_timestamp: Option<DateTime<Utc>>, + total_documents_fetched: usize, + poll_count: usize, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct MongodbSourceConfig { + #[serde(serialize_with = "iggy_common::serde_secret::serialize_secret")] + pub connection_uri: SecretString, + pub database: String, + pub collection: String, + pub max_pool_size: Option<u32>, + pub query: Option<Document>, + pub timestamp_field: Option<String>, + pub limit: Option<i64>, + pub polling_interval: Option<String>, +} + +#[derive(Debug)] +pub struct MongodbSource { Review Comment: nit: `MongodbSource` breaks the casing already used for every MongoDB identifier in the repo, such as `MongoDbSink` and `MongoDbOps`. rename to `MongoDbSource`, `MongoDbSourceConfig` and `MongoDbSourceFixture`. ########## core/connectors/sources/mongodb_source/src/lib.rs: ########## @@ -0,0 +1,242 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +use async_trait::async_trait; +use core::mem::drop; +use futures::stream::TryStreamExt; +use iggy_common::{DateTime, Utc}; +use iggy_connector_sdk::{ + ConnectorState, Error, ProducedMessage, ProducedMessages, Schema, Source, source_connector, +}; +use mongodb::{Client, Collection, bson::Document, bson::doc, options::ClientOptions}; +use secrecy::{ExposeSecret, SecretString}; +use serde::{Deserialize, Serialize}; +use std::str::FromStr; +use std::time::Duration; +use tokio::sync::Mutex; +use tracing::info; + +source_connector!(MongodbSource); + +#[derive(Debug, Clone, Serialize, Deserialize)] +struct State { + last_poll_timestamp: Option<DateTime<Utc>>, + total_documents_fetched: usize, + poll_count: usize, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct MongodbSourceConfig { + #[serde(serialize_with = "iggy_common::serde_secret::serialize_secret")] + pub connection_uri: SecretString, + pub database: String, + pub collection: String, + pub max_pool_size: Option<u32>, + pub query: Option<Document>, + pub timestamp_field: Option<String>, + pub limit: Option<i64>, Review Comment: nit: the other sources call this cap `batch_size`, and a `limit` of 0 reaches the driver as no limit at all. rename it and reject 0 at load. ########## core/connectors/sources/mongodb_source/src/lib.rs: ########## @@ -0,0 +1,242 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +use async_trait::async_trait; +use core::mem::drop; +use futures::stream::TryStreamExt; +use iggy_common::{DateTime, Utc}; +use iggy_connector_sdk::{ + ConnectorState, Error, ProducedMessage, ProducedMessages, Schema, Source, source_connector, +}; +use mongodb::{Client, Collection, bson::Document, bson::doc, options::ClientOptions}; +use secrecy::{ExposeSecret, SecretString}; +use serde::{Deserialize, Serialize}; +use std::str::FromStr; +use std::time::Duration; +use tokio::sync::Mutex; +use tracing::info; + +source_connector!(MongodbSource); + +#[derive(Debug, Clone, Serialize, Deserialize)] +struct State { + last_poll_timestamp: Option<DateTime<Utc>>, + total_documents_fetched: usize, + poll_count: usize, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct MongodbSourceConfig { + #[serde(serialize_with = "iggy_common::serde_secret::serialize_secret")] + pub connection_uri: SecretString, + pub database: String, + pub collection: String, + pub max_pool_size: Option<u32>, + pub query: Option<Document>, + pub timestamp_field: Option<String>, + pub limit: Option<i64>, + pub polling_interval: Option<String>, +} + +#[derive(Debug)] +pub struct MongodbSource { + id: u32, + config: MongodbSourceConfig, + client: Option<Client>, + polling_interval: Duration, + state: Mutex<State>, +} + +const CONNECTOR_NAME: &str = "MongoDB source"; + +impl MongodbSource { + pub fn new(id: u32, config: MongodbSourceConfig, state: Option<ConnectorState>) -> Self { + let polling_interval = config + .polling_interval + .as_deref() + .unwrap_or("10s") + .parse::<humantime::Duration>() + .unwrap_or_else(|_| humantime::Duration::from_str("10s").unwrap()) + .into(); + + let restored_state = state + .and_then(|s| s.deserialize::<State>(CONNECTOR_NAME, id)) + .inspect(|s| { + info!( + "Restored state for {CONNECTOR_NAME} connector with ID: {id}. \ + Documents fetched: {}, poll count: {}", + s.total_documents_fetched, s.poll_count + ); + }); + + MongodbSource { + id, + config, + client: None, + polling_interval, + state: Mutex::new(restored_state.unwrap_or(State { + last_poll_timestamp: None, + total_documents_fetched: 0, + poll_count: 0, + })), + } + } + + fn serialize_state(&self, state: &State) -> Option<ConnectorState> { + ConnectorState::serialize(state, CONNECTOR_NAME, self.id) + } + + async fn create_client(&self) -> Result<Client, Error> { + let mut client_options: ClientOptions = + ClientOptions::parse(self.config.connection_uri.expose_secret()) + .await + .map_err(|e| Error::InitError(format!("Failed to parse connection URI: {e}")))?; + if let Some(pool_size) = self.config.max_pool_size { + client_options.max_pool_size = Some(pool_size); + } + let client = Client::with_options(client_options) + .map_err(|e| Error::InitError(format!("Failed to create client: {e}")))?; + Ok(client) + } + + async fn search_documents(&self, client: &Client) -> Result<Vec<ProducedMessage>, Error> { + let state = self.state.lock().await; Review Comment: simplification: the state guard is held across `find` only because the second lock at line 180 forces the manual `drop`. read `last_poll_timestamp` into a local first, then the `drop` and its import both go away. also at lines 19, 145. -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected]
