use std::time::Duration;
use zenkey_fleet::{Answer, declare_repeating};
async fn peer_pair(port: u16) -> (zenoh::Session, zenoh::Session) {
let listen = zenkey_fleet::session::open(&[], &[format!("tcp/127.0.0.1:{port}")], false)
.await
.expect("listener session");
let connect = zenkey_fleet::session::open(&[format!("tcp/127.0.0.1:{port}")], &[], false)
.await
.expect("connector session");
(listen, connect)
}
const HOST_A: &str = "v1/h-aaaaaaaaaaaa/@rpc/sysinfo/introspect";
const HOST_B: &str = "v1/h-bbbbbbbbbbbb/@rpc/sysinfo/introspect";
const SELECTOR: &str = "v1/*/@rpc/sysinfo/introspect";
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn a_complete_queryable_does_not_collapse_the_declared_fleet() {
let (a, b) = peer_pair(7481).await;
let _qa = a
.declare_queryable(HOST_A)
.complete(true)
.callback(|query| {
let q = query.clone();
tokio::spawn(async move {
q.reply(HOST_A, "from-a").await.unwrap();
});
})
.await
.expect("queryable a");
let _qb = a
.declare_queryable(HOST_B)
.callback(|query| {
let q = query.clone();
tokio::spawn(async move {
q.reply(HOST_B, "from-b").await.unwrap();
});
})
.await
.expect("queryable b");
let repeating = declare_repeating(&b, "", SELECTOR, Duration::from_secs(5))
.await
.expect("declare");
let answers = tokio::time::timeout(Duration::from_secs(5), async {
loop {
let answers = repeating.fetch().await.expect("fetch");
if answers.len() >= 2 {
break answers;
}
tokio::time::sleep(Duration::from_millis(20)).await;
}
})
.await
.expect("both queryables should answer within 5s");
let mut origins: Vec<&str> = answers.iter().map(|a| a.origin.as_str()).collect();
origins.sort_unstable();
assert_eq!(
origins,
vec!["h-aaaaaaaaaaaa", "h-bbbbbbbbbbbb"],
"target All + reply-key attribution must survive a complete queryable"
);
repeating.undeclare().await.expect("undeclare");
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn parameters_ride_per_get_not_in_the_declared_key() {
let (a, b) = peer_pair(7482).await;
let _q = a
.declare_queryable(HOST_A)
.callback(|query| {
let params = query.parameters().to_string();
let q = query.clone();
tokio::spawn(async move {
q.reply(HOST_A, params).await.unwrap();
});
})
.await
.expect("queryable");
let repeating = declare_repeating(&b, "", SELECTOR, Duration::from_secs(5))
.await
.expect("declare");
assert!(
!repeating.key().contains('?'),
"the declared keyexpr must never carry parameters"
);
let first = tokio::time::timeout(Duration::from_secs(5), async {
loop {
let answers = repeating
.fetch_with("round=1", None)
.await
.expect("fetch 1");
if !answers.is_empty() {
break answers;
}
tokio::time::sleep(Duration::from_millis(20)).await;
}
})
.await
.expect("a queryable should answer within 5s");
let second = repeating
.fetch_with("round=2", None)
.await
.expect("fetch 2");
for (answers, expected) in [(&first, "round=1"), (&second, "round=2")] {
assert_eq!(answers.len(), 1);
let Answer::Value(bytes) = &answers[0].answer else {
panic!("expected a value reply");
};
assert_eq!(
String::from_utf8_lossy(&bytes.to_bytes()),
expected,
"each get must carry its own parameters"
);
}
repeating.undeclare().await.expect("undeclare");
}