use faucet_conformance::{
assert_bookmark_roundtrip, assert_bounded_memory, assert_config_schema_valid_value,
assert_errors_not_panics,
};
use faucet_source_mysql_cdc::{MysqlCdcSource, MysqlCdcSourceConfig};
use mysql_async::{Conn, Opts, prelude::Queryable};
use serde_json::json;
use std::time::Duration;
use testcontainers::{ContainerAsync, ImageExt, runners::AsyncRunner};
use testcontainers_modules::mysql::Mysql;
const BATCH: usize = 250;
const TOTAL: usize = 600;
#[test]
fn conformance_config_schema_valid() {
let schema = serde_json::to_value(schemars::schema_for!(MysqlCdcSourceConfig)).unwrap();
assert_config_schema_valid_value(&schema, "faucet-source-mysql-cdc");
}
async fn start_mysql_cdc() -> (ContainerAsync<Mysql>, String) {
let container = Mysql::default()
.with_tag("8.1")
.with_cmd([
"--server-id=1",
"--log-bin=mysql-bin",
"--binlog-format=ROW",
"--binlog-row-image=FULL",
"--binlog-row-metadata=FULL",
])
.start()
.await
.expect("mysql CDC container start");
let port = container
.get_host_port_ipv4(3306)
.await
.expect("mysql port");
let url = format!("mysql://root@127.0.0.1:{port}/test");
(container, url)
}
async fn connect(url: &str) -> Conn {
Conn::new(Opts::from_url(url).expect("parse URL"))
.await
.expect("connect")
}
fn build_config(url: &str) -> MysqlCdcSourceConfig {
serde_json::from_value(json!({
"connection_url": url,
"server_id": 5005,
"start_position": { "type": "current" },
"idle_timeout": 20,
"batch_size": BATCH,
}))
.expect("config")
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn conformance_bounded_memory() {
let (_container, url) = start_mysql_cdc().await;
{
let mut conn = connect(&url).await;
conn.query_drop("CREATE TABLE test.events (id INT PRIMARY KEY)")
.await
.expect("create table");
}
let source = MysqlCdcSource::new(build_config(&url))
.await
.expect("source new");
let writer_url = url.clone();
let writer = tokio::spawn(async move {
tokio::time::sleep(Duration::from_secs(5)).await;
let mut conn = connect(&writer_url).await;
for i in 0..TOTAL {
conn.query_drop(format!("INSERT INTO test.events (id) VALUES ({i})"))
.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, url) = start_mysql_cdc().await;
{
let mut conn = connect(&url).await;
conn.query_drop("CREATE TABLE test.events (id INT PRIMARY KEY)")
.await
.expect("create table");
}
let source = MysqlCdcSource::new(build_config(&url))
.await
.expect("source new");
let writer_url = url.clone();
let writer = tokio::spawn(async move {
tokio::time::sleep(Duration::from_secs(5)).await;
let mut conn = connect(&writer_url).await;
for i in 0..N {
conn.query_drop(format!("INSERT INTO test.events (id) VALUES ({i})"))
.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, url) = start_mysql_cdc().await;
{
let mut conn = connect(&url).await;
conn.query_drop("CREATE TABLE test.events (id INT PRIMARY KEY)")
.await
.expect("create table");
}
let config: MysqlCdcSourceConfig = serde_json::from_value(json!({
"connection_url": url,
"server_id": 6006,
"start_position": { "type": "file_pos", "file": "faucet-nonexistent-bin.999999", "pos": 4 },
"idle_timeout": 20,
"batch_size": 250,
}))
.expect("config");
let source = MysqlCdcSource::new(config)
.await
.expect("new + preflight succeed against a live server");
assert_errors_not_panics(&source).await;
}