hubcio commented on code in PR #4193:
URL: https://github.com/apache/iggy/pull/4193#discussion_r4016266888
##########
core/integration/tests/server/specific.rs:
##########
@@ -141,3 +164,99 @@ async fn restart_offset_skip(harness: &mut TestHarness) {
async fn segment_rotation_scenario(harness: &TestHarness) {
segment_rotation_race_scenario::run(harness).await;
}
+
+async fn assert_tls_client_sharding_and_cleanup(
+ harness: &TestHarness,
+ clients: &[IggyClient],
+ transport: &str,
+) {
+ let mut client_ids = HashSet::with_capacity(clients.len());
+ for client in clients {
+ client_ids.insert(client.get_me().await.unwrap().client_id);
+ }
+ assert_eq!(client_ids.len(), clients.len(), "client IDs must be unique");
+
+ // The wire exposes only the sequence tail. Match it to the router's
+ // full ID and thread name to prove execution placement after handoff.
+ let install_marker = format!("installing delegated {transport} client fd");
+ tokio::time::timeout(TLS_CLEANUP_TIMEOUT, async {
+ loop {
+ let logs = harness.server().stdout_plain();
+ let mut counts = [0; TLS_TEST_SHARDS];
+ let mut installed = HashSet::new();
+ for line in logs.lines().filter(|line|
line.contains(&install_marker)) {
+ let fields: Vec<_> = line.split_whitespace().collect();
+ let full_id: u128 = fields
+ .iter()
+ .find_map(|field| field.strip_prefix("client_id="))
+ .expect("install log must carry client_id")
+ .parse()
+ .unwrap();
+ let wire_id = u32::try_from(full_id &
u128::from(u32::MAX)).unwrap();
+ if !client_ids.contains(&wire_id) {
+ continue;
+ }
+ let owner = usize::try_from(full_id >> 112).unwrap();
+ assert!(owner < TLS_TEST_SHARDS, "invalid owner: {line}");
+ assert!(
+ fields.contains(&format!("shard={owner}").as_str()),
+ "{line}"
+ );
+ assert!(
+ fields.contains(&format!("shard-{owner}").as_str()),
+ "{line}"
+ );
+ assert!(installed.insert(wire_id), "client installed twice:
{line}");
+ counts[owner] += 1;
+ }
+ if installed == client_ids {
+ assert_eq!(
Review Comment:
a starting-phase shift keeps eight consecutive allocations evenly split
across four shards. the local update reads the actual shard count and retains
the exact distribution check.
##########
core/integration/tests/server/specific.rs:
##########
@@ -141,3 +164,99 @@ async fn restart_offset_skip(harness: &mut TestHarness) {
async fn segment_rotation_scenario(harness: &TestHarness) {
segment_rotation_race_scenario::run(harness).await;
}
+
+async fn assert_tls_client_sharding_and_cleanup(
+ harness: &TestHarness,
+ clients: &[IggyClient],
+ transport: &str,
+) {
+ let mut client_ids = HashSet::with_capacity(clients.len());
+ for client in clients {
+ client_ids.insert(client.get_me().await.unwrap().client_id);
+ }
+ assert_eq!(client_ids.len(), clients.len(), "client IDs must be unique");
+
+ // The wire exposes only the sequence tail. Match it to the router's
+ // full ID and thread name to prove execution placement after handoff.
+ let install_marker = format!("installing delegated {transport} client fd");
+ tokio::time::timeout(TLS_CLEANUP_TIMEOUT, async {
+ loop {
+ let logs = harness.server().stdout_plain();
+ let mut counts = [0; TLS_TEST_SHARDS];
+ let mut installed = HashSet::new();
+ for line in logs.lines().filter(|line|
line.contains(&install_marker)) {
+ let fields: Vec<_> = line.split_whitespace().collect();
+ let full_id: u128 = fields
+ .iter()
+ .find_map(|field| field.strip_prefix("client_id="))
+ .expect("install log must carry client_id")
+ .parse()
+ .unwrap();
+ let wire_id = u32::try_from(full_id &
u128::from(u32::MAX)).unwrap();
+ if !client_ids.contains(&wire_id) {
+ continue;
+ }
+ let owner = usize::try_from(full_id >> 112).unwrap();
+ assert!(owner < TLS_TEST_SHARDS, "invalid owner: {line}");
+ assert!(
+ fields.contains(&format!("shard={owner}").as_str()),
+ "{line}"
+ );
+ assert!(
+ fields.contains(&format!("shard-{owner}").as_str()),
+ "{line}"
+ );
+ assert!(installed.insert(wire_id), "client installed twice:
{line}");
+ counts[owner] += 1;
+ }
+ if installed == client_ids {
+ assert_eq!(
+ counts,
+ [TLS_TEST_CLIENTS / TLS_TEST_SHARDS; TLS_TEST_SHARDS]
+ );
+ break;
+ }
+ tokio::time::sleep(TLS_POLL_INTERVAL).await;
+ }
+ })
+ .await
+ .expect("every encrypted client must be installed on its owning shard
thread");
+
+ for client in clients {
+ client.disconnect().await.unwrap();
+ }
+ // Polling opens a separate SDK connection. A fresh observer has none,
+ // so an exact count also detects leaked partition-client sessions.
+ let observer = harness.root_client().await.unwrap();
+ let observer_id = observer.get_me().await.unwrap().client_id;
+ tokio::time::timeout(TLS_CLEANUP_TIMEOUT, async {
+ loop {
+ let connected = observer.get_clients().await.unwrap();
+ if connected.len() == 1 {
Review Comment:
the local update keeps one observer per shard and requires the exact set of
observer IDs. a missing shard reply or leaked session can no longer satisfy the
cleanup check.
--
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]