Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

1 change: 1 addition & 0 deletions dt-tests/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -37,6 +37,7 @@ url = { workspace = true }
rust_decimal = { workspace = true }
uuid = { workspace = true }
openssl = { workspace = true }
zookeeper-client = { workspace = true }

[dev-dependencies]
rand = "0.9.2"
Expand Down
16 changes: 16 additions & 0 deletions dt-tests/docker-compose.integration.yml
Original file line number Diff line number Diff line change
Expand Up @@ -1404,6 +1404,22 @@ services:
stdin_open: true
tty: true

zk-src:
image: zookeeper:3.9
container_name: zk-src-it
ports:
- "2181:2181"
environment:
ZOO_MY_ID: 1

zk-dst:
image: zookeeper:3.9
container_name: zk-dst-it
ports:
- "2182:2181"
environment:
ZOO_MY_ID: 2

networks:
default:
name: ape-dts-integration-network
4 changes: 4 additions & 0 deletions dt-tests/tests/.env
Original file line number Diff line number Diff line change
Expand Up @@ -158,3 +158,7 @@ clickhouse_url=http://admin:123456@127.0.0.1:8123

# tidb
tidb_sinker_url=mysql://root:@127.0.0.1:4000?ssl-mode=disabled

# zookeeper
zk_src_url=127.0.0.1:2181
zk_dst_url=127.0.0.1:2182
1 change: 1 addition & 0 deletions dt-tests/tests/integration_test.rs
Original file line number Diff line number Diff line change
Expand Up @@ -20,3 +20,4 @@ mod pg_to_starrocks;
mod redis_to_redis;
mod test_config_util;
mod test_runner;
mod zk_to_zk;
2 changes: 2 additions & 0 deletions dt-tests/tests/test_runner/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -21,3 +21,5 @@ pub mod redis_statistic_runner;
pub mod redis_test_runner;
pub mod redis_test_util;
pub mod test_base;
pub mod zk_cycle_test_runner;
pub mod zk_test_runner;
10 changes: 9 additions & 1 deletion dt-tests/tests/test_runner/test_base.rs
Original file line number Diff line number Diff line change
Expand Up @@ -13,7 +13,7 @@ use super::{
rdb_redis_test_runner::RdbRedisTestRunner, rdb_sql_test_runner::RdbSqlTestRunner,
rdb_starrocks_test_runner::RdbStarRocksTestRunner, rdb_struct_test_runner::RdbStructTestRunner,
rdb_test_runner::RdbTestRunner, redis_statistic_runner::RedisStatisticTestRunner,
redis_test_runner::RedisTestRunner,
redis_test_runner::RedisTestRunner, zk_test_runner::ZkTestRunner,
};

pub struct TestBase {}
Expand Down Expand Up @@ -431,4 +431,12 @@ impl TestBase {
runner.dcl_check_sql_execution().await.unwrap();
runner.close().await.unwrap();
}

pub async fn run_zk_cdc_test(test_dir: &str, start_millis: u64, parse_millis: u64) {
let runner = ZkTestRunner::new(test_dir).await.unwrap();
runner
.run_cdc_test(start_millis, parse_millis)
.await
.unwrap();
}
}
194 changes: 194 additions & 0 deletions dt-tests/tests/test_runner/zk_cycle_test_runner.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,194 @@
use dt_common::utils::time_util::TimeUtil;
use tokio::task::JoinHandle;
use zookeeper_client as zk;

use crate::test_config_util::TestConfigUtil;

use super::zk_test_runner::ZkTestRunner;

const SHADOW_PREFIX: &str = "/__ape_dts_shadow";
const MARKER_PATH: &str = "/__ape_dts_marker";

pub struct ZkCycleTestRunner {}

impl ZkCycleTestRunner {
pub async fn run_cycle_cdc_test(test_dir: &str, start_millis: u64, parse_millis: u64) {
let sub_paths = TestConfigUtil::get_absolute_sub_dir(test_dir);
assert!(
!sub_paths.is_empty(),
"cycle test requires at least one sub-directory in {}",
test_dir
);

let mut handlers: Vec<JoinHandle<()>> = vec![];
let mut runners: Vec<(String, ZkTestRunner)> = vec![];

for sub_path in &sub_paths {
let runner = ZkTestRunner::new(format!("{}/{}", test_dir, sub_path.1).as_str())
.await
.unwrap();
runners.push((sub_path.1.clone(), runner));
}

let first_runner = &runners[0].1;
let src = zk::Client::connect(&first_runner.src_url).await.unwrap();
let dst = zk::Client::connect(&first_runner.dst_url).await.unwrap();

let options = zk::CreateMode::Persistent.with_acls(zk::Acls::anyone_all());
ZkCycleTestRunner::ensure_path(&src, "/app", &options).await;
ZkCycleTestRunner::ensure_path(&dst, "/app", &options).await;
ZkCycleTestRunner::delete_children(&src, "/app").await;
ZkCycleTestRunner::delete_children(&dst, "/app").await;

ZkCycleTestRunner::delete_children(&src, SHADOW_PREFIX).await;
ZkCycleTestRunner::delete_node(&src, SHADOW_PREFIX).await;
ZkCycleTestRunner::delete_children(&dst, SHADOW_PREFIX).await;
ZkCycleTestRunner::delete_node(&dst, SHADOW_PREFIX).await;
ZkCycleTestRunner::delete_node(&src, MARKER_PATH).await;
ZkCycleTestRunner::delete_node(&dst, MARKER_PATH).await;

for (_, runner) in &runners {
handlers.push(runner.base.spawn_task().await.unwrap());
TimeUtil::sleep_millis(start_millis).await;
}

src.create("/app/from-node1", b"hello-from-1", &options)
.await
.unwrap();
dst.create("/app/from-node2", b"hello-from-2", &options)
.await
.unwrap();

TimeUtil::sleep_millis(parse_millis).await;

// data sync assertions
let (src_data_1, _) = src.get_data("/app/from-node1").await.unwrap();
let (dst_data_1, _) = dst.get_data("/app/from-node1").await.unwrap();
assert_eq!(src_data_1, dst_data_1, "from-node1 data mismatch");

let (src_data_2, _) = src.get_data("/app/from-node2").await.unwrap();
let (dst_data_2, _) = dst.get_data("/app/from-node2").await.unwrap();
assert_eq!(src_data_2, dst_data_2, "from-node2 data mismatch");

// shadow metadata assertions — verify shadow znodes exist on both sides
Self::assert_shadow_exists(&dst, "/app/from-node1").await;
Self::assert_shadow_exists(&src, "/app/from-node2").await;

// marker assertions — verify data markers written on both sides
Self::assert_marker_exists(&src).await;
Self::assert_marker_exists(&dst).await;

// anti-loop stability — sleep again and verify data hasn't changed (no bounce-back loop)
let snapshot_src_1 = src.get_data("/app/from-node1").await.unwrap();
let snapshot_dst_2 = dst.get_data("/app/from-node2").await.unwrap();
TimeUtil::sleep_millis(parse_millis / 2).await;
let check_src_1 = src.get_data("/app/from-node1").await.unwrap();
let check_dst_2 = dst.get_data("/app/from-node2").await.unwrap();
assert_eq!(
snapshot_src_1.1.version, check_src_1.1.version,
"anti-loop: /app/from-node1 on src was modified again (version changed {} -> {})",
snapshot_src_1.1.version, check_src_1.1.version
);
assert_eq!(
snapshot_dst_2.1.version, check_dst_2.1.version,
"anti-loop: /app/from-node2 on dst was modified again (version changed {} -> {})",
snapshot_dst_2.1.version, check_dst_2.1.version
);

for handler in handlers {
handler.abort();
while !handler.is_finished() {
TimeUtil::sleep_millis(1).await;
}
}
}

async fn assert_shadow_exists(client: &zk::Client, data_path: &str) {
let shadow_path = format!("{}{}", SHADOW_PREFIX, data_path);
let (data, _) = client
.get_data(&shadow_path)
.await
.unwrap_or_else(|e| panic!("shadow znode {} should exist: {}", shadow_path, e));
let json: serde_json::Value = serde_json::from_slice(&data)
.unwrap_or_else(|e| panic!("shadow {} invalid JSON: {}", shadow_path, e));
assert!(
json.get("source_id")
.and_then(|v| v.as_str())
.map_or(false, |v| !v.is_empty()),
"shadow {} missing or empty source_id",
shadow_path
);
assert!(
json.get("source_order_millis")
.and_then(|v| v.as_i64())
.map_or(false, |v| v > 0),
"shadow {} missing or zero source_order_millis",
shadow_path
);
assert!(
json.get("version").and_then(|v| v.as_i64()).is_some(),
"shadow {} missing or invalid version",
shadow_path
);
assert_eq!(
json.get("deleted").and_then(|v| v.as_bool()),
Some(false),
"shadow {} should have deleted=false",
shadow_path
);
}

async fn assert_marker_exists(client: &zk::Client) {
let (data, _) = client
.get_data(MARKER_PATH)
.await
.unwrap_or_else(|e| panic!("marker {} should exist: {}", MARKER_PATH, e));
let json: serde_json::Value = serde_json::from_slice(&data)
.unwrap_or_else(|e| panic!("marker {} invalid JSON: {}", MARKER_PATH, e));
assert!(
json.get("source_id")
.and_then(|v| v.as_str())
.map_or(false, |v| !v.is_empty()),
"marker {} missing or empty source_id",
MARKER_PATH
);
}

async fn delete_node(client: &zk::Client, path: &str) {
match client.delete(path, None).await {
Ok(()) => {}
Err(ref e) if matches!(e, zk::Error::NoNode) => {}
Err(e) => panic!("delete_node {} failed: {}", path, e),
}
}

async fn ensure_path(client: &zk::Client, path: &str, options: &zk::CreateOptions<'_>) {
match client.create(path, &[], options).await {
Ok(_) => {}
Err(ref e) if matches!(e, zk::Error::NodeExists) => {}
Err(e) => panic!("ensure_path {} failed: {}", path, e),
}
}

fn delete_children<'a>(
client: &'a zk::Client,
path: &'a str,
) -> std::pin::Pin<Box<dyn std::future::Future<Output = ()> + 'a>> {
Box::pin(async move {
let children = match client.get_children(path).await {
Ok((children, _)) => children,
Err(ref e) if matches!(e, zk::Error::NoNode) => return,
Err(e) => panic!("get_children {} failed: {}", path, e),
};
for child in children {
let child_path = format!("{}/{}", path, child);
Self::delete_children(client, &child_path).await;
match client.delete(&child_path, None).await {
Ok(()) => {}
Err(ref e) if matches!(e, zk::Error::NoNode) => {}
Err(e) => panic!("delete {} failed: {}", child_path, e),
}
}
})
}
}
Loading
Loading