github-actions[bot] commented on code in PR #4225:
URL: https://github.com/apache/iggy/pull/4225#discussion_r4126826690


##########
core/connectors/sources/mongodb_source/src/lib.rs:
##########
@@ -0,0 +1,496 @@
+// 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 futures::stream::TryStreamExt;
+use iggy_common::{DateTime, Utc};
+use iggy_connector_sdk::{
+    ConnectorState, Error, ProducedMessage, ProducedMessages, Schema, Source,
+    source::SourceBatchResult, source_connector,
+};
+use mongodb::{
+    Client, Collection,
+    bson::{Bson, Document, doc},
+    options::ClientOptions,
+};
+use secrecy::{ExposeSecret, SecretString};
+use serde::{Deserialize, Serialize};
+use std::num::NonZeroU32;
+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,
+    // `_id` of the document at `last_poll_timestamp`, so documents sharing 
that
+    // millisecond but cut off by `batch_size` are picked up on the next poll.
+    #[serde(default)]
+    last_id: Option<String>,
+}
+
+#[derive(Debug, Clone, Serialize, Deserialize)]
+pub struct MongoDbSourceConfig {
+    #[serde(serialize_with = "iggy_common::serde_secret::serialize_secret")]

Review Comment:
   warning: The `connection_uri` field carries `serialize_secret`, which writes 
the plaintext URI, while the README says the connector redacts it when 
serialized. Use `serialize_redacted`, or drop `Serialize` from the config.



##########
core/connectors/sources/mongodb_source/src/lib.rs:
##########
@@ -0,0 +1,496 @@
+// 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 futures::stream::TryStreamExt;
+use iggy_common::{DateTime, Utc};
+use iggy_connector_sdk::{
+    ConnectorState, Error, ProducedMessage, ProducedMessages, Schema, Source,
+    source::SourceBatchResult, source_connector,
+};
+use mongodb::{
+    Client, Collection,
+    bson::{Bson, Document, doc},
+    options::ClientOptions,
+};
+use secrecy::{ExposeSecret, SecretString};
+use serde::{Deserialize, Serialize};
+use std::num::NonZeroU32;
+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,
+    // `_id` of the document at `last_poll_timestamp`, so documents sharing 
that
+    // millisecond but cut off by `batch_size` are picked up on the next poll.
+    #[serde(default)]
+    last_id: Option<String>,
+}
+
+#[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 batch_size: Option<NonZeroU32>,
+    pub polling_interval: Option<String>,
+}
+
+#[derive(Debug)]
+pub struct MongoDbSource {
+    id: u32,
+    config: MongoDbSourceConfig,
+    client: Option<Client>,
+    polling_interval: Duration,
+    state: Mutex<State>,
+    pending_state: Mutex<Option<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,
+                last_id: None,
+            })),
+            pending_state: Mutex::new(None),
+        }
+    }
+
+    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 check_collection(&self, client: &Client) -> Result<(), Error> {
+        let database = client.database(&self.config.database);
+        database
+            .run_command(doc! { "ping": 1 })
+            .await
+            .map_err(|e| Error::InitError(format!("Failed to ping MongoDB: 
{e}")))?;
+        let collections = database
+            .list_collection_names()
+            .filter(doc! { "name": &self.config.collection })
+            .await
+            .map_err(|e| Error::InitError(format!("Failed to list collections: 
{e}")))?;
+        if collections.is_empty() {
+            return Err(Error::InvalidConfigValue(format!(
+                "collection '{}' not found in database '{}'",
+                self.config.collection, self.config.database
+            )));
+        }
+        Ok(())
+    }
+
+    /// The cursor only advances on BSON `Date` values, so a string, number or
+    /// missing `timestamp_field` would leave the connector silently producing 
nothing.
+    async fn check_timestamp_field(&self, client: &Client) -> Result<(), 
Error> {
+        let Some(timestamp_field) = &self.config.timestamp_field else {
+            return Ok(());
+        };
+        let coll: Collection<Document> = client
+            .database(&self.config.database)
+            .collection(&self.config.collection);
+        let sample = coll
+            .find_one(self.config.query.clone().unwrap_or_default())
+            .await
+            .map_err(|e| Error::InitError(format!("Failed to read a sample 
document: {e}")))?;
+        match sample {
+            Some(doc) => validate_timestamp_field(&doc, timestamp_field),
+            None => Ok(()),
+        }
+    }
+
+    async fn search_documents(
+        &self,
+        client: &Client,
+    ) -> Result<(Vec<ProducedMessage>, State), Error> {
+        let state = self.state.lock().await.clone();
+        let batch_size = i64::from(self.config.batch_size.map_or(100, 
NonZeroU32::get));
+
+        let coll: Collection<Document> = client
+            .database(&self.config.database)
+            .collection(&self.config.collection);
+
+        let mut cursor = if let Some(timestamp_field) = 
&self.config.timestamp_field {
+            let cursor_filter = cursor_filter(timestamp_field, &state);
+            let filter = match &self.config.query {
+                Some(query) if !query.is_empty() => doc! { "$and": 
[query.clone(), cursor_filter] },
+                _ => cursor_filter,
+            };
+
+            coll.find(filter)
+                .limit(batch_size)
+                .sort(doc! { timestamp_field: 1, "_id": 1 })
+                .await
+                .map_err(|e| Error::Storage(format!("Failed to execute search: 
{e}")))?
+        } else {
+            coll.find(self.config.query.clone().unwrap_or_default())
+                .await
+                .map_err(|e| Error::Storage(format!("Failed to execute search: 
{e}")))?
+        };
+
+        let mut messages = Vec::new();
+        let mut latest_position = None;
+
+        while let Some(doc) = cursor
+            .try_next()
+            .await
+            .map_err(|e| Error::Storage(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(),
+                    )
+            {
+                // Results are sorted by (timestamp, _id), so the last one 
seen is the newest.
+                latest_position = Some((
+                    timestamp,
+                    doc.get("_id").and_then(|id| 
serde_json::to_string(id).ok()),
+                ));
+            }
+
+            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 (last_poll_timestamp, last_id) = match latest_position {
+            Some((timestamp, id)) => (Some(timestamp), id),
+            None => (state.last_poll_timestamp, state.last_id),
+        };
+        let candidate_state = State {
+            last_poll_timestamp,
+            total_documents_fetched: state.total_documents_fetched + 
messages.len(),
+            poll_count: state.poll_count + 1,
+            last_id,
+        };
+        Ok((messages, candidate_state))
+    }
+}
+
+#[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.check_collection(&client).await?;
+        self.check_timestamp_field(&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, candidate_state) = self.search_documents(client).await?;
+
+        let persisted_state = 
self.serialize_state(&candidate_state).ok_or_else(|| {
+            Error::Serialization("failed to serialize MongoDB source 
state".to_string())
+        })?;
+        *self.pending_state.lock().await = Some(candidate_state);
+
+        Ok(ProducedMessages {
+            schema: Schema::Json,
+            messages,
+            state: Some(persisted_state),

Review Comment:
   warning: `poll()` returns a state on every call because `poll_count` 
increments, so the runtime rewrites and fsyncs the state file on each idle 
poll. Return `state: None` when the batch holds no documents.



##########
core/connectors/sources/mongodb_source/README.md:
##########
@@ -0,0 +1,195 @@
+# MongoDB Source Connector with State Management
+
+This MongoDB source connector polls a MongoDB collection, produces each 
document as a JSON message into Apache Iggy, and persists its progress so 
ingestion resumes where it left off after a restart.
+
+## Features
+
+- **Incremental Data Processing**: Track the last processed timestamp 
(`timestamp_field`) to avoid reprocessing documents
+- **Timestamp-ordered Batches**: Incremental polls filter with `$gt`, sort 
ascending on the timestamp field, and cap each batch with `batch_size`
+- **Custom Filters**: Restrict the documents read with a MongoDB `query` filter
+- **Connection Pooling**: Configurable driver pool size via `max_pool_size`
+- **Persistent State Storage**: State is persisted by the connectors runtime 
(file or HTTP backend)
+- **State Recovery**: Resume processing from the last known timestamp after 
restart
+- **JSON Output**: Every document is emitted with the `json` schema
+
+## Configuration
+
+### Basic Configuration
+
+```toml
+type = "source"
+key = "mongodb"
+enabled = true
+version = 0
+name = "MongoDB source"
+path = "target/release/libiggy_connector_mongodb_source"
+
+[[streams]]
+stream = "mongodb_stream"
+topic = "documents"
+schema = "json"
+batch_length = 100
+linger_time = "5ms"
+
+[plugin_config]
+connection_uri = "mongodb://admin:admin123@localhost:27017"
+database = "test_source"
+collection = "test_messages"
+max_pool_size = 10
+polling_interval = "30s"
+batch_size = 100
+timestamp_field = "timestamp"
+query = { status = "active" }
+```
+
+| Field              | Required | Default         | Description                
                                                 |
+| ------------------ | -------- | --------------- | 
--------------------------------------------------------------------------- |
+| `connection_uri`   | yes      |                 | MongoDB connection string. 
Treated as a secret and redacted when serialized |
+| `database`         | yes      |                 | Database to read from      
                                                 |
+| `collection`       | yes      |                 | Collection to read from    
                                                 |
+| `max_pool_size`    | no       | driver default  | Maximum number of 
connections in the driver pool                            |
+| `query`            | no       | `{}`            | MongoDB filter document 
applied to every poll                               |
+| `timestamp_field`  | no       | none            | BSON `Date` field used for 
incremental polling                              |
+| `batch_size`       | no       | `100`           | Maximum documents per 
incremental poll. Must be greater than 0              |
+| `polling_interval` | no       | `"10s"`         | Delay before each poll 
(humantime format). Invalid values fall back to 10s  |
+
+### State Management Configuration
+
+The plugin has no state settings of its own. State is configured once for the 
whole connectors runtime, in the runtime `config.toml`:
+
+```toml
+[state]
+path = "local_state"  # used by storage = "file"
+storage = "file"      # "file" | "http"
+```
+
+## State Information
+
+The connector tracks the following state information:
+
+### Processing State
+
+- `last_poll_timestamp`: Highest `timestamp_field` value seen so far
+- `total_documents_fetched`: Total number of documents produced
+- `poll_count`: Number of polling cycles executed
+
+### Error Tracking
+
+The connector does not persist error counters. Errors are returned from 
`poll()` to the runtime, which logs them.
+
+### Performance Statistics
+
+The connector does not record performance statistics in its state. Use the 
runtime logs and metrics instead.
+
+## Storage Backends
+
+### File Storage (Default)
+
+State is written to `{path}/source_{key}.state`, where `key` is the connector 
`key`. The runtime writes to a temporary file and renames it, so a crash never 
leaves a half-written state file.
+
+```toml
+[state]
+storage = "file"
+path = "local_state"
+```
+
+### HTTP Storage
+
+State is stored at `{url}/source_{key}` with optimistic concurrency 
(ETag/If-Match) and idempotent retries.
+
+```toml
+[state]
+storage = "http"
+
+[state.http]
+url = "http://127.0.0.1:8080/connectors/state";
+timeout = "5s"
+```
+
+## Usage Examples
+
+### Basic Usage with State Management
+
+The connectors runtime normally drives this lifecycle. Calling it directly 
looks like this:
+
+```rust
+use iggy_connector_mongodb_source::{MongoDbSource, MongoDbSourceConfig};
+
+// `state` is the previously persisted ConnectorState, if any
+let mut connector = MongoDbSource::new(id, config, state);
+
+// Open connector (creates the MongoDB client)
+connector.open().await?;
+
+// Poll documents. The returned ProducedMessages carries the updated state
+let produced = connector.poll().await?;
+
+// Close connector
+connector.close().await?;
+```
+
+### Building
+
+```bash
+cargo build --release -p iggy_connector_mongodb_source
+```
+
+The resulting `target/release/libiggy_connector_mongodb_source` is what the 
connector `path` points to.
+
+## State File Format
+
+State is serialized with MessagePack, so the file is binary. Its content is 
equivalent to:
+
+```json
+{
+  "last_poll_timestamp": "2024-01-15T10:30:00Z",
+  "total_documents_fetched": 15000,
+  "poll_count": 150
+}
+```
+
+Messages themselves are the documents serialized with `serde_json`, so BSON 
types such as `ObjectId` and `Date` appear in extended JSON form, for example 
`{"_id": {"$oid": "65a4f0c2e1b2c3d4e5f60718"}}`.
+
+## Best Practices
+
+1. **Key Uniqueness**: Use a unique connector `key` per instance, since the 
state file is named after it
+2. **Index the Timestamp Field**: Create an index on `timestamp_field` so the 
sorted `$gt` query stays fast
+3. **Use BSON Dates**: Store `timestamp_field` as a BSON `Date`. Other types 
do not advance the state
+4. **Storage Location**: Point `[state].path` at persistent storage in 
production
+5. **Batch Tuning**: Balance `batch_size` and `polling_interval` against your 
write rate
+6. **Credentials**: Keep real credentials in `connection_uri` out of committed 
config files
+
+## Troubleshooting
+
+### Common Issues
+
+1. **Duplicate Messages on Every Poll**: Without `timestamp_field`, each poll 
reads every document that matches `query`. Set `timestamp_field` for 
incremental ingestion
+2. **Large First Poll**: With no saved timestamp, the first poll runs without 
`batch_size` or sort and reads the whole matching collection

Review Comment:
   nit: Troubleshooting item 2 says the first poll ignores `batch_size` and the 
sort, but `search_documents()` applies both when `timestamp_field` is set. Name 
the mode without `timestamp_field` instead. Also at line 188.



##########
core/connectors/sources/mongodb_source/src/lib.rs:
##########
@@ -0,0 +1,496 @@
+// 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 futures::stream::TryStreamExt;
+use iggy_common::{DateTime, Utc};
+use iggy_connector_sdk::{
+    ConnectorState, Error, ProducedMessage, ProducedMessages, Schema, Source,
+    source::SourceBatchResult, source_connector,
+};
+use mongodb::{
+    Client, Collection,
+    bson::{Bson, Document, doc},
+    options::ClientOptions,
+};
+use secrecy::{ExposeSecret, SecretString};
+use serde::{Deserialize, Serialize};
+use std::num::NonZeroU32;
+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,
+    // `_id` of the document at `last_poll_timestamp`, so documents sharing 
that
+    // millisecond but cut off by `batch_size` are picked up on the next poll.
+    #[serde(default)]
+    last_id: Option<String>,
+}
+
+#[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 batch_size: Option<NonZeroU32>,
+    pub polling_interval: Option<String>,
+}
+
+#[derive(Debug)]
+pub struct MongoDbSource {
+    id: u32,
+    config: MongoDbSourceConfig,
+    client: Option<Client>,
+    polling_interval: Duration,
+    state: Mutex<State>,
+    pending_state: Mutex<Option<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,
+                last_id: None,
+            })),
+            pending_state: Mutex::new(None),
+        }
+    }
+
+    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 check_collection(&self, client: &Client) -> Result<(), Error> {
+        let database = client.database(&self.config.database);
+        database
+            .run_command(doc! { "ping": 1 })
+            .await
+            .map_err(|e| Error::InitError(format!("Failed to ping MongoDB: 
{e}")))?;
+        let collections = database
+            .list_collection_names()
+            .filter(doc! { "name": &self.config.collection })
+            .await
+            .map_err(|e| Error::InitError(format!("Failed to list collections: 
{e}")))?;
+        if collections.is_empty() {
+            return Err(Error::InvalidConfigValue(format!(
+                "collection '{}' not found in database '{}'",
+                self.config.collection, self.config.database
+            )));
+        }
+        Ok(())
+    }
+
+    /// The cursor only advances on BSON `Date` values, so a string, number or
+    /// missing `timestamp_field` would leave the connector silently producing 
nothing.
+    async fn check_timestamp_field(&self, client: &Client) -> Result<(), 
Error> {
+        let Some(timestamp_field) = &self.config.timestamp_field else {
+            return Ok(());
+        };
+        let coll: Collection<Document> = client
+            .database(&self.config.database)
+            .collection(&self.config.collection);
+        let sample = coll
+            .find_one(self.config.query.clone().unwrap_or_default())

Review Comment:
   warning: `check_timestamp_field()` validates one arbitrary matching 
document, so `open()` fails on a collection whose sample lacks a BSON `Date`, 
even when dated documents exist. Query for one matching document with `​$type: 
"date"` and fail only when none exists.



##########
core/connectors/sources/mongodb_source/README.md:
##########
@@ -0,0 +1,195 @@
+# MongoDB Source Connector with State Management
+
+This MongoDB source connector polls a MongoDB collection, produces each 
document as a JSON message into Apache Iggy, and persists its progress so 
ingestion resumes where it left off after a restart.
+
+## Features
+
+- **Incremental Data Processing**: Track the last processed timestamp 
(`timestamp_field`) to avoid reprocessing documents
+- **Timestamp-ordered Batches**: Incremental polls filter with `$gt`, sort 
ascending on the timestamp field, and cap each batch with `batch_size`
+- **Custom Filters**: Restrict the documents read with a MongoDB `query` filter
+- **Connection Pooling**: Configurable driver pool size via `max_pool_size`
+- **Persistent State Storage**: State is persisted by the connectors runtime 
(file or HTTP backend)
+- **State Recovery**: Resume processing from the last known timestamp after 
restart
+- **JSON Output**: Every document is emitted with the `json` schema
+
+## Configuration
+
+### Basic Configuration
+
+```toml
+type = "source"
+key = "mongodb"
+enabled = true
+version = 0
+name = "MongoDB source"
+path = "target/release/libiggy_connector_mongodb_source"
+
+[[streams]]
+stream = "mongodb_stream"
+topic = "documents"
+schema = "json"
+batch_length = 100
+linger_time = "5ms"
+
+[plugin_config]
+connection_uri = "mongodb://admin:admin123@localhost:27017"
+database = "test_source"
+collection = "test_messages"
+max_pool_size = 10
+polling_interval = "30s"
+batch_size = 100
+timestamp_field = "timestamp"
+query = { status = "active" }
+```
+
+| Field              | Required | Default         | Description                
                                                 |
+| ------------------ | -------- | --------------- | 
--------------------------------------------------------------------------- |
+| `connection_uri`   | yes      |                 | MongoDB connection string. 
Treated as a secret and redacted when serialized |
+| `database`         | yes      |                 | Database to read from      
                                                 |
+| `collection`       | yes      |                 | Collection to read from    
                                                 |
+| `max_pool_size`    | no       | driver default  | Maximum number of 
connections in the driver pool                            |
+| `query`            | no       | `{}`            | MongoDB filter document 
applied to every poll                               |
+| `timestamp_field`  | no       | none            | BSON `Date` field used for 
incremental polling                              |
+| `batch_size`       | no       | `100`           | Maximum documents per 
incremental poll. Must be greater than 0              |
+| `polling_interval` | no       | `"10s"`         | Delay before each poll 
(humantime format). Invalid values fall back to 10s  |
+
+### State Management Configuration
+
+The plugin has no state settings of its own. State is configured once for the 
whole connectors runtime, in the runtime `config.toml`:
+
+```toml
+[state]
+path = "local_state"  # used by storage = "file"
+storage = "file"      # "file" | "http"
+```
+
+## State Information
+
+The connector tracks the following state information:
+
+### Processing State
+
+- `last_poll_timestamp`: Highest `timestamp_field` value seen so far
+- `total_documents_fetched`: Total number of documents produced
+- `poll_count`: Number of polling cycles executed
+
+### Error Tracking
+
+The connector does not persist error counters. Errors are returned from 
`poll()` to the runtime, which logs them.
+
+### Performance Statistics
+
+The connector does not record performance statistics in its state. Use the 
runtime logs and metrics instead.
+
+## Storage Backends
+
+### File Storage (Default)
+
+State is written to `{path}/source_{key}.state`, where `key` is the connector 
`key`. The runtime writes to a temporary file and renames it, so a crash never 
leaves a half-written state file.
+
+```toml
+[state]
+storage = "file"
+path = "local_state"
+```
+
+### HTTP Storage
+
+State is stored at `{url}/source_{key}` with optimistic concurrency 
(ETag/If-Match) and idempotent retries.
+
+```toml
+[state]
+storage = "http"
+
+[state.http]
+url = "http://127.0.0.1:8080/connectors/state";
+timeout = "5s"
+```
+
+## Usage Examples
+
+### Basic Usage with State Management
+
+The connectors runtime normally drives this lifecycle. Calling it directly 
looks like this:
+
+```rust
+use iggy_connector_mongodb_source::{MongoDbSource, MongoDbSourceConfig};
+
+// `state` is the previously persisted ConnectorState, if any
+let mut connector = MongoDbSource::new(id, config, state);
+
+// Open connector (creates the MongoDB client)
+connector.open().await?;
+
+// Poll documents. The returned ProducedMessages carries the updated state
+let produced = connector.poll().await?;
+
+// Close connector
+connector.close().await?;
+```
+
+### Building
+
+```bash
+cargo build --release -p iggy_connector_mongodb_source
+```
+
+The resulting `target/release/libiggy_connector_mongodb_source` is what the 
connector `path` points to.
+
+## State File Format
+
+State is serialized with MessagePack, so the file is binary. Its content is 
equivalent to:
+
+```json
+{
+  "last_poll_timestamp": "2024-01-15T10:30:00Z",
+  "total_documents_fetched": 15000,
+  "poll_count": 150
+}
+```
+
+Messages themselves are the documents serialized with `serde_json`, so BSON 
types such as `ObjectId` and `Date` appear in extended JSON form, for example 
`{"_id": {"$oid": "65a4f0c2e1b2c3d4e5f60718"}}`.
+
+## Best Practices
+
+1. **Key Uniqueness**: Use a unique connector `key` per instance, since the 
state file is named after it
+2. **Index the Timestamp Field**: Create an index on `timestamp_field` so the 
sorted `$gt` query stays fast
+3. **Use BSON Dates**: Store `timestamp_field` as a BSON `Date`. Other types 
do not advance the state
+4. **Storage Location**: Point `[state].path` at persistent storage in 
production
+5. **Batch Tuning**: Balance `batch_size` and `polling_interval` against your 
write rate
+6. **Credentials**: Keep real credentials in `connection_uri` out of committed 
config files
+
+## Troubleshooting
+
+### Common Issues
+
+1. **Duplicate Messages on Every Poll**: Without `timestamp_field`, each poll 
reads every document that matches `query`. Set `timestamp_field` for 
incremental ingestion
+2. **Large First Poll**: With no saved timestamp, the first poll runs without 
`batch_size` or sort and reads the whole matching collection
+3. **State Not Advancing**: `timestamp_field` values that are strings or 
numbers are ignored. Only BSON `Date` values update `last_poll_timestamp`
+4. **Skipped Documents**: The filter uses `$gt`, so documents sharing the last 
seen timestamp that fall past a `batch_size` boundary are not read on the next 
poll. Prefer unique, high-resolution timestamps

Review Comment:
   nit: Troubleshooting item 4 says same-timestamp documents past a 
`batch_size` boundary stay unread, but `cursor_filter()` re-reads them through 
the `_id` tie-break. Describe the tie-break, and add `last_id` to the state 
lists at lines 72-74 and 143-149.



##########
core/connectors/sources/mongodb_source/Cargo.toml:
##########
@@ -0,0 +1,50 @@
+# 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"
+description = "Iggy MongoDB source connector for polling collections into 
message streams"
+edition = "2024"
+license = "Apache-2.0"
+keywords = ["iggy", "messaging", "streaming", "mongodb", "source"]
+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"
+publish = false
+
+[package.metadata.cargo-machete]
+ignored = ["dashmap"]
+
+[lib]
+crate-type = ["cdylib", "lib"]
+
+[dependencies]
+async-trait = { workspace = true }
+dashmap = { workspace = true }
+futures.workspace = true
+humantime = { workspace = true }
+iggy_common = { workspace = true }
+iggy_connector_sdk = { workspace = true }
+mongodb.workspace = true
+secrecy = { workspace = true }
+serde = { workspace = true }
+serde_json = { workspace = true }
+tokio = { workspace = true }
+tracing = { workspace = true }

Review Comment:
   nit: The manifest omits `[lints]​ workspace = true`, so the workspace 
`warnings = "deny"` rule does not apply to this crate. Add the table, as every 
other source crate does.



##########
core/connectors/sources/mongodb_source/src/lib.rs:
##########
@@ -0,0 +1,496 @@
+// 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 futures::stream::TryStreamExt;
+use iggy_common::{DateTime, Utc};
+use iggy_connector_sdk::{
+    ConnectorState, Error, ProducedMessage, ProducedMessages, Schema, Source,
+    source::SourceBatchResult, source_connector,
+};
+use mongodb::{
+    Client, Collection,
+    bson::{Bson, Document, doc},
+    options::ClientOptions,
+};
+use secrecy::{ExposeSecret, SecretString};
+use serde::{Deserialize, Serialize};
+use std::num::NonZeroU32;
+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,
+    // `_id` of the document at `last_poll_timestamp`, so documents sharing 
that
+    // millisecond but cut off by `batch_size` are picked up on the next poll.
+    #[serde(default)]
+    last_id: Option<String>,
+}
+
+#[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 batch_size: Option<NonZeroU32>,
+    pub polling_interval: Option<String>,
+}
+
+#[derive(Debug)]
+pub struct MongoDbSource {
+    id: u32,
+    config: MongoDbSourceConfig,
+    client: Option<Client>,
+    polling_interval: Duration,
+    state: Mutex<State>,
+    pending_state: Mutex<Option<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,
+                last_id: None,
+            })),
+            pending_state: Mutex::new(None),
+        }
+    }
+
+    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 check_collection(&self, client: &Client) -> Result<(), Error> {
+        let database = client.database(&self.config.database);
+        database
+            .run_command(doc! { "ping": 1 })
+            .await
+            .map_err(|e| Error::InitError(format!("Failed to ping MongoDB: 
{e}")))?;
+        let collections = database
+            .list_collection_names()
+            .filter(doc! { "name": &self.config.collection })
+            .await
+            .map_err(|e| Error::InitError(format!("Failed to list collections: 
{e}")))?;
+        if collections.is_empty() {
+            return Err(Error::InvalidConfigValue(format!(
+                "collection '{}' not found in database '{}'",
+                self.config.collection, self.config.database
+            )));
+        }
+        Ok(())
+    }
+
+    /// The cursor only advances on BSON `Date` values, so a string, number or
+    /// missing `timestamp_field` would leave the connector silently producing 
nothing.
+    async fn check_timestamp_field(&self, client: &Client) -> Result<(), 
Error> {
+        let Some(timestamp_field) = &self.config.timestamp_field else {
+            return Ok(());
+        };
+        let coll: Collection<Document> = client
+            .database(&self.config.database)
+            .collection(&self.config.collection);
+        let sample = coll
+            .find_one(self.config.query.clone().unwrap_or_default())
+            .await
+            .map_err(|e| Error::InitError(format!("Failed to read a sample 
document: {e}")))?;
+        match sample {
+            Some(doc) => validate_timestamp_field(&doc, timestamp_field),
+            None => Ok(()),
+        }
+    }
+
+    async fn search_documents(
+        &self,
+        client: &Client,
+    ) -> Result<(Vec<ProducedMessage>, State), Error> {
+        let state = self.state.lock().await.clone();
+        let batch_size = i64::from(self.config.batch_size.map_or(100, 
NonZeroU32::get));
+
+        let coll: Collection<Document> = client
+            .database(&self.config.database)
+            .collection(&self.config.collection);
+
+        let mut cursor = if let Some(timestamp_field) = 
&self.config.timestamp_field {
+            let cursor_filter = cursor_filter(timestamp_field, &state);
+            let filter = match &self.config.query {
+                Some(query) if !query.is_empty() => doc! { "$and": 
[query.clone(), cursor_filter] },
+                _ => cursor_filter,
+            };
+
+            coll.find(filter)
+                .limit(batch_size)
+                .sort(doc! { timestamp_field: 1, "_id": 1 })
+                .await
+                .map_err(|e| Error::Storage(format!("Failed to execute search: 
{e}")))?
+        } else {
+            coll.find(self.config.query.clone().unwrap_or_default())
+                .await
+                .map_err(|e| Error::Storage(format!("Failed to execute search: 
{e}")))?
+        };
+
+        let mut messages = Vec::new();
+        let mut latest_position = None;
+
+        while let Some(doc) = cursor
+            .try_next()
+            .await
+            .map_err(|e| Error::Storage(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(),
+                    )
+            {
+                // Results are sorted by (timestamp, _id), so the last one 
seen is the newest.
+                latest_position = Some((
+                    timestamp,
+                    doc.get("_id").and_then(|id| 
serde_json::to_string(id).ok()),
+                ));
+            }
+
+            let payload = serde_json::to_vec(&doc).map_err(|e| {
+                Error::Serialization(format!("Failed to serialize document: 
{}", e))
+            })?;
+
+            let message = ProducedMessage {
+                id: None,

Review Comment:
   nit: Every produced message leaves `id` as `None`, so a consumer has no key 
to dedupe the documents this source replays after a NACK. Set `id` from the 
document `_id`.



##########
core/connectors/sources/mongodb_source/src/lib.rs:
##########
@@ -0,0 +1,496 @@
+// 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 futures::stream::TryStreamExt;
+use iggy_common::{DateTime, Utc};
+use iggy_connector_sdk::{
+    ConnectorState, Error, ProducedMessage, ProducedMessages, Schema, Source,
+    source::SourceBatchResult, source_connector,
+};
+use mongodb::{
+    Client, Collection,
+    bson::{Bson, Document, doc},
+    options::ClientOptions,
+};
+use secrecy::{ExposeSecret, SecretString};
+use serde::{Deserialize, Serialize};
+use std::num::NonZeroU32;
+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,
+    // `_id` of the document at `last_poll_timestamp`, so documents sharing 
that
+    // millisecond but cut off by `batch_size` are picked up on the next poll.
+    #[serde(default)]
+    last_id: Option<String>,
+}
+
+#[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 batch_size: Option<NonZeroU32>,
+    pub polling_interval: Option<String>,
+}
+
+#[derive(Debug)]
+pub struct MongoDbSource {
+    id: u32,
+    config: MongoDbSourceConfig,
+    client: Option<Client>,
+    polling_interval: Duration,
+    state: Mutex<State>,
+    pending_state: Mutex<Option<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,
+                last_id: None,
+            })),
+            pending_state: Mutex::new(None),
+        }
+    }
+
+    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 check_collection(&self, client: &Client) -> Result<(), Error> {
+        let database = client.database(&self.config.database);
+        database
+            .run_command(doc! { "ping": 1 })
+            .await
+            .map_err(|e| Error::InitError(format!("Failed to ping MongoDB: 
{e}")))?;
+        let collections = database
+            .list_collection_names()
+            .filter(doc! { "name": &self.config.collection })
+            .await
+            .map_err(|e| Error::InitError(format!("Failed to list collections: 
{e}")))?;
+        if collections.is_empty() {
+            return Err(Error::InvalidConfigValue(format!(
+                "collection '{}' not found in database '{}'",
+                self.config.collection, self.config.database
+            )));
+        }
+        Ok(())
+    }
+
+    /// The cursor only advances on BSON `Date` values, so a string, number or
+    /// missing `timestamp_field` would leave the connector silently producing 
nothing.
+    async fn check_timestamp_field(&self, client: &Client) -> Result<(), 
Error> {
+        let Some(timestamp_field) = &self.config.timestamp_field else {
+            return Ok(());
+        };
+        let coll: Collection<Document> = client
+            .database(&self.config.database)
+            .collection(&self.config.collection);
+        let sample = coll
+            .find_one(self.config.query.clone().unwrap_or_default())
+            .await
+            .map_err(|e| Error::InitError(format!("Failed to read a sample 
document: {e}")))?;
+        match sample {
+            Some(doc) => validate_timestamp_field(&doc, timestamp_field),
+            None => Ok(()),
+        }
+    }
+
+    async fn search_documents(
+        &self,
+        client: &Client,
+    ) -> Result<(Vec<ProducedMessage>, State), Error> {
+        let state = self.state.lock().await.clone();
+        let batch_size = i64::from(self.config.batch_size.map_or(100, 
NonZeroU32::get));
+
+        let coll: Collection<Document> = client
+            .database(&self.config.database)
+            .collection(&self.config.collection);
+
+        let mut cursor = if let Some(timestamp_field) = 
&self.config.timestamp_field {
+            let cursor_filter = cursor_filter(timestamp_field, &state);
+            let filter = match &self.config.query {
+                Some(query) if !query.is_empty() => doc! { "$and": 
[query.clone(), cursor_filter] },
+                _ => cursor_filter,
+            };
+
+            coll.find(filter)
+                .limit(batch_size)
+                .sort(doc! { timestamp_field: 1, "_id": 1 })
+                .await
+                .map_err(|e| Error::Storage(format!("Failed to execute search: 
{e}")))?
+        } else {
+            coll.find(self.config.query.clone().unwrap_or_default())
+                .await
+                .map_err(|e| Error::Storage(format!("Failed to execute search: 
{e}")))?
+        };
+
+        let mut messages = Vec::new();
+        let mut latest_position = None;
+
+        while let Some(doc) = cursor
+            .try_next()
+            .await
+            .map_err(|e| Error::Storage(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(),
+                    )
+            {
+                // Results are sorted by (timestamp, _id), so the last one 
seen is the newest.
+                latest_position = Some((
+                    timestamp,
+                    doc.get("_id").and_then(|id| 
serde_json::to_string(id).ok()),
+                ));
+            }
+
+            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 (last_poll_timestamp, last_id) = match latest_position {
+            Some((timestamp, id)) => (Some(timestamp), id),
+            None => (state.last_poll_timestamp, state.last_id),
+        };
+        let candidate_state = State {
+            last_poll_timestamp,
+            total_documents_fetched: state.total_documents_fetched + 
messages.len(),
+            poll_count: state.poll_count + 1,
+            last_id,
+        };
+        Ok((messages, candidate_state))
+    }
+}
+
+#[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.check_collection(&client).await?;
+        self.check_timestamp_field(&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, candidate_state) = self.search_documents(client).await?;
+
+        let persisted_state = 
self.serialize_state(&candidate_state).ok_or_else(|| {
+            Error::Serialization("failed to serialize MongoDB source 
state".to_string())
+        })?;
+        *self.pending_state.lock().await = Some(candidate_state);
+
+        Ok(ProducedMessages {
+            schema: Schema::Json,
+            messages,
+            state: Some(persisted_state),
+        })
+    }
+
+    async fn on_batch_result(&self, result: SourceBatchResult) -> Result<(), 
Error> {
+        let candidate_state = self.pending_state.lock().await.take();
+        if result == SourceBatchResult::Ack
+            && let Some(candidate_state) = candidate_state
+        {
+            *self.state.lock().await = candidate_state;
+        }
+        Ok(())
+    }
+
+    async fn close(&mut self) -> Result<(), Error> {
+        info!("Mongodb Connector with ID: {} is closing", self.id);
+
+        let state = self.state.lock().await;
+
+        info!(
+            "Mongodb source connector ID: {} closed. Total documents 
processed: {}",
+            self.id, state.total_documents_fetched
+        );
+        Ok(())
+    }
+}
+
+fn validate_timestamp_field(doc: &Document, timestamp_field: &str) -> 
Result<(), Error> {
+    match doc.get(timestamp_field) {
+        Some(Bson::DateTime(_)) => Ok(()),
+        Some(value) => Err(Error::InvalidConfigValue(format!(
+            "timestamp_field '{timestamp_field}' must hold a BSON Date, found 
{:?}",
+            value.element_type()
+        ))),
+        None => Err(Error::InvalidConfigValue(format!(
+            "timestamp_field '{timestamp_field}' is missing from documents in 
the collection"
+        ))),
+    }
+}
+
+fn cursor_filter(timestamp_field: &str, state: &State) -> Document {
+    let Some(last_timestamp) = state.last_poll_timestamp else {
+        return doc! { timestamp_field: { "$type": "date" } };
+    };
+    let last_timestamp = 
mongodb::bson::DateTime::from_millis(last_timestamp.timestamp_millis());
+    match state
+        .last_id
+        .as_deref()
+        .and_then(|id| serde_json::from_str::<Bson>(id).ok())
+    {
+        Some(last_id) => doc! {
+            "$or": [
+                { timestamp_field: { "$gt": last_timestamp } },
+                { timestamp_field: last_timestamp, "_id": { "$gt": last_id } },
+            ]
+        },
+        // State saved before `last_id` existed. `$gte` re-delivers the 
boundary
+        // documents once instead of risking skipping them.
+        None => doc! { timestamp_field: { "$gte": last_timestamp } },
+    }
+}
+
+#[cfg(test)]
+mod tests {

Review Comment:
   nit: The tests cover ACK/NACK and the `_id` filter, but two of the four 
canonical source state tests are missing: `given_no_state_should_start_fresh` 
and `given_invalid_state_should_start_fresh`. Add them, following 
`sources/random_source/src/lib.rs`.



##########
core/connectors/sources/mongodb_source/Cargo.toml:
##########
@@ -0,0 +1,50 @@
+# 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"
+description = "Iggy MongoDB source connector for polling collections into 
message streams"
+edition = "2024"
+license = "Apache-2.0"
+keywords = ["iggy", "messaging", "streaming", "mongodb", "source"]
+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"
+publish = false
+
+[package.metadata.cargo-machete]
+ignored = ["dashmap"]
+
+[lib]
+crate-type = ["cdylib", "lib"]
+
+[dependencies]
+async-trait = { workspace = true }
+dashmap = { workspace = true }

Review Comment:
   simplification: `dashmap` is declared and hidden from cargo-machete, but 
nothing in the crate uses it. Delete the dependency and the `ignored` entry.



##########
core/integration/tests/connectors/mongodb/mongodb_source.rs:
##########
@@ -0,0 +1,308 @@
+// 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::{POLL_ATTEMPTS, POLL_INTERVAL_MS, TEST_MESSAGE_COUNT};
+use crate::connectors::fixtures::MongoDbSourcePreCreatedFixture;
+use iggy_common::MessageClient;
+use iggy_common::{Consumer, Identifier, PollingStrategy};
+use integration::harness::seeds;
+use integration::iggy_harness;
+use std::time::Duration;
+use tokio::time::sleep;
+
+#[iggy_harness(
+    server(connectors_runtime(config_path = 
"tests/connectors/mongodb/source.toml")),
+    seed = seeds::connector_stream
+)]
+async fn mongodb_source_produces_messages_to_iggy(
+    harness: &TestHarness,
+    fixture: MongoDbSourcePreCreatedFixture,
+) {
+    let client = harness.root_client().await.unwrap();
+
+    fixture
+        .insert_documents(TEST_MESSAGE_COUNT)
+        .await
+        .expect("Failed to insert documents");
+
+    let doc_count = fixture
+        .get_document_count()
+        .await
+        .expect("Failed to get document count");
+    assert_eq!(
+        doc_count, TEST_MESSAGE_COUNT,
+        "Expected {TEST_MESSAGE_COUNT} documents in MongoDB"
+    );
+
+    let stream_id: Identifier = seeds::names::STREAM.try_into().unwrap();
+    let topic_id: Identifier = seeds::names::TOPIC.try_into().unwrap();
+    let consumer_id: Identifier = "test_consumer".try_into().unwrap();
+
+    let mut received: Vec<serde_json::Value> = Vec::new();
+    for _ in 0..POLL_ATTEMPTS {
+        if let Ok(polled) = client
+            .poll_messages(
+                &stream_id,
+                &topic_id,
+                None,
+                &Consumer::new(consumer_id.clone()),
+                &PollingStrategy::next(),
+                10,
+                true,
+            )
+            .await
+        {
+            for msg in polled.messages {
+                if let Ok(json) = serde_json::from_slice(&msg.payload) {
+                    received.push(json);
+                }
+            }
+            if received.len() >= TEST_MESSAGE_COUNT {
+                break;
+            }
+        }
+        sleep(Duration::from_millis(POLL_INTERVAL_MS)).await;
+    }
+
+    assert!(
+        received.len() >= TEST_MESSAGE_COUNT,
+        "Expected at least {TEST_MESSAGE_COUNT} messages, got {}",
+        received.len()
+    );
+
+    for (i, record) in received.iter().take(TEST_MESSAGE_COUNT).enumerate() {
+        let expected_id = (i + 1) as i64;
+        let expected_name = format!("doc_{}", i + 1);
+
+        assert_eq!(
+            record.get("id").and_then(|v| v.as_i64()),
+            Some(expected_id),
+            "ID mismatch at record {i}"
+        );
+        assert_eq!(
+            record.get("name").and_then(|v| v.as_str()),
+            Some(expected_name.as_str()),
+            "Name mismatch at record {i}"
+        );
+    }
+}
+
+#[iggy_harness(
+    server(connectors_runtime(config_path = 
"tests/connectors/mongodb/source.toml")),
+    seed = seeds::connector_stream
+)]
+async fn mongodb_source_handles_empty_collection(
+    harness: &TestHarness,
+    fixture: MongoDbSourcePreCreatedFixture,
+) {
+    let client = harness.root_client().await.unwrap();
+
+    let doc_count = fixture
+        .get_document_count()
+        .await
+        .expect("Failed to get document count");
+    assert_eq!(doc_count, 0, "Expected empty collection");
+
+    let stream_id: Identifier = seeds::names::STREAM.try_into().unwrap();
+    let topic_id: Identifier = seeds::names::TOPIC.try_into().unwrap();
+    let consumer_id: Identifier = "test_consumer".try_into().unwrap();
+
+    sleep(Duration::from_millis(100)).await;
+
+    let polled = client
+        .poll_messages(
+            &stream_id,
+            &topic_id,
+            None,
+            &Consumer::new(consumer_id),
+            &PollingStrategy::next(),
+            10,
+            false,
+        )
+        .await;
+
+    assert!(

Review Comment:
   nit: The empty-collection test asserts only that a poll against Iggy 
succeeds, which holds even when the source fails to start, so the test cannot 
fail. Assert that the topic stays empty after several poll intervals.



-- 
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]

Reply via email to