use faucet_conformance::{
assert_bookmark_roundtrip, assert_bounded_memory, assert_config_schema_valid_value,
assert_errors_not_panics,
};
use faucet_source_mongodb_cdc::{MongoCdcSource, MongoCdcSourceConfig};
use mongodb::Client;
use mongodb::bson::{Document, doc};
use serde_json::json;
use std::time::Duration;
use testcontainers::ContainerAsync;
use testcontainers::runners::AsyncRunner;
use testcontainers_modules::mongo::Mongo;
const DB: &str = "app";
const COLL: &str = "events";
const BATCH: usize = 250;
const TOTAL: usize = 5000;
#[test]
fn conformance_config_schema_valid() {
let schema = serde_json::to_value(schemars::schema_for!(MongoCdcSourceConfig)).unwrap();
assert_config_schema_valid_value(&schema, "faucet-source-mongodb-cdc");
}
async fn start_repl_set() -> (ContainerAsync<Mongo>, String) {
let container = Mongo::repl_set()
.start()
.await
.expect("mongo replica-set container start");
let port = container
.get_host_port_ipv4(27017)
.await
.expect("mongo port");
let uri = format!("mongodb://127.0.0.1:{port}/?directConnection=true");
(container, uri)
}
fn config(uri: &str) -> MongoCdcSourceConfig {
serde_json::from_value(json!({
"connection_uri": uri,
"scope": { "type": "collection", "database": DB, "collection": COLL },
"start_from": { "type": "now" },
"idle_timeout": 20,
"max_await_time_ms": 500,
"batch_size": BATCH,
}))
.expect("config")
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn conformance_bounded_memory() {
let (_container, uri) = start_repl_set().await;
{
let client = Client::with_uri_str(&uri).await.expect("seed client");
client
.database(DB)
.collection::<Document>(COLL)
.insert_one(doc! { "_id": 0, "seed": true })
.await
.expect("seed insert");
}
let source = MongoCdcSource::new(config(&uri)).await.expect("source");
let writer_uri = uri.clone();
let writer = tokio::spawn(async move {
tokio::time::sleep(Duration::from_secs(5)).await;
let client = Client::with_uri_str(&writer_uri)
.await
.expect("writer client");
let coll = client.database(DB).collection::<Document>(COLL);
for i in 1..=TOTAL {
coll.insert_one(doc! { "_id": i as i64, "name": "n" })
.await
.expect("insert");
}
});
assert_bounded_memory(&source, BATCH, TOTAL).await;
writer.await.expect("writer task");
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn conformance_bookmark_roundtrip() {
const N: usize = 300;
let (_container, uri) = start_repl_set().await;
{
let client = Client::with_uri_str(&uri).await.expect("seed client");
client
.database(DB)
.collection::<Document>(COLL)
.insert_one(doc! { "_id": 0, "seed": true })
.await
.expect("seed insert");
}
let source = MongoCdcSource::new(config(&uri)).await.expect("source");
let writer_uri = uri.clone();
let writer = tokio::spawn(async move {
tokio::time::sleep(Duration::from_secs(5)).await;
let client = Client::with_uri_str(&writer_uri)
.await
.expect("writer client");
let coll = client.database(DB).collection::<Document>(COLL);
for i in 1..=N {
coll.insert_one(doc! { "_id": i as i64, "name": "n" })
.await
.expect("insert");
}
});
assert_bookmark_roundtrip(&source).await;
writer.await.expect("writer task");
}
#[tokio::test(flavor = "multi_thread")]
async fn conformance_errors_not_panics() {
let (_container, uri) = start_repl_set().await;
{
let client = Client::with_uri_str(&uri).await.expect("seed client");
client
.database(DB)
.collection::<Document>(COLL)
.insert_one(doc! { "_id": 0, "seed": true })
.await
.expect("seed insert");
}
let config: MongoCdcSourceConfig = serde_json::from_value(json!({
"connection_uri": uri,
"scope": { "type": "collection", "database": DB, "collection": COLL },
"start_from": { "type": "now" },
"idle_timeout": 20,
"max_await_time_ms": 500,
"batch_size": 250,
"aggregation_pipeline": [ { "$notARealStage": {} } ],
}))
.expect("config");
let source = MongoCdcSource::new(config)
.await
.expect("new + hello succeed against a live replica set");
assert_errors_not_panics(&source).await;
}