mirror of
https://github.com/logos-messaging/logos-delivery-rust-bindings.git
synced 2026-07-30 06:53:29 +00:00
Bumps vendor to pick up channelExists, cherry-picked onto the nim-ffi
0.2.0 lineage: the chore/channel-exists branch forked before that
migration, so its pre-0.2.0 {.ffi.} shape would not compile here.
The typed model carries a real bool, so channel_exists_async returns
Result<bool, _> rather than the "true"/"false" string the original
would have produced.
Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
275 lines
10 KiB
Rust
275 lines
10 KiB
Rust
use base64::Engine;
|
|
use regex::Regex;
|
|
use serial_test::serial;
|
|
use std::sync::{Arc, Mutex};
|
|
use std::time::Duration;
|
|
use waku_bindings::{
|
|
ChannelMessageReceivedPayload, ChannelMessageSentPayload, ChannelSendRequest,
|
|
LogosDeliveryCtx, WakuNodeConfig,
|
|
};
|
|
|
|
const TEST_CHANNEL_ID: &str = "test-channel";
|
|
const TEST_CONTENT_TOPIC: &str = "/test/1/channels/proto";
|
|
const TEST_SENDER_ID: &str = "test-sender";
|
|
const NOOP_ENCRYPTION: &str = "noop";
|
|
const OTHER_SENDER_ID: &str = "other-sender";
|
|
const CHANNEL_PAYLOAD: &[u8] = b"Hi from a reliable channel!";
|
|
const TIMEOUT: Duration = Duration::from_secs(30);
|
|
const SHARDS_IN_NETWORK: usize = 8;
|
|
/// Cargo runs test binaries in parallel and serial_test only serialises within
|
|
/// one, so each binary needs its own persistency root. Every node in a binary
|
|
/// must share it: the singleton refuses to be re-targeted.
|
|
const STORAGE_PATH: &str = "./data-channels-test";
|
|
|
|
fn new_node(tcp_port: usize) -> LogosDeliveryCtx {
|
|
// A channel's shard is derived from its content topic, so the node needs
|
|
// autosharding on, and has to sit on every shard a topic might hash to.
|
|
let config = serde_json::to_string(&WakuNodeConfig {
|
|
tcp_port: Some(tcp_port),
|
|
num_shards_in_network: Some(SHARDS_IN_NETWORK as u16),
|
|
shards: (0..SHARDS_IN_NETWORK).collect(),
|
|
local_storage_path: Some(STORAGE_PATH.to_string()),
|
|
..Default::default()
|
|
})
|
|
.expect("config should serialise");
|
|
|
|
LogosDeliveryCtx::create(config, TIMEOUT).expect("node should instantiate")
|
|
}
|
|
|
|
fn channel_message(payload: &[u8]) -> ChannelSendRequest {
|
|
ChannelSendRequest {
|
|
payload: base64::engine::general_purpose::STANDARD.encode(payload),
|
|
ephemeral: false,
|
|
}
|
|
}
|
|
|
|
/// Dials `to` from `from`, rewriting the advertised address to loopback so NAT
|
|
/// or a firewall cannot interfere.
|
|
async fn connect(from: &LogosDeliveryCtx, to: &LogosDeliveryCtx) -> Result<(), String> {
|
|
let addresses = to.waku_listen_addresses_async().await?;
|
|
let address = addresses
|
|
.split(',')
|
|
.next()
|
|
.ok_or("peer should report a listen address")?;
|
|
|
|
let re = Regex::new(r"\b(?:\d{1,3}\.){3}\d{1,3}\b").unwrap();
|
|
let address = re.replace_all(address, "127.0.0.1").to_string();
|
|
|
|
from.waku_connect_async(address, 10_000).await.map(|_| ())
|
|
}
|
|
|
|
#[tokio::test]
|
|
#[serial]
|
|
async fn channel_create_send_close() {
|
|
let node = new_node(60070);
|
|
|
|
// Channel traffic is reported through events, so capture them to prove the
|
|
// typed listener registers and delivers ChannelMessageSentPayload.
|
|
let sent: Arc<Mutex<Vec<ChannelMessageSentPayload>>> = Arc::new(Mutex::new(Vec::new()));
|
|
let sent_cloned = sent.clone();
|
|
node.add_on_channel_message_sent_listener(move |event| {
|
|
sent_cloned.lock().unwrap().push(event.clone());
|
|
});
|
|
|
|
node.start_node_async().await.expect("node should start");
|
|
|
|
assert!(
|
|
!node
|
|
.channel_exists_async(TEST_CHANNEL_ID.to_string())
|
|
.await
|
|
.expect("exists should answer for an unknown channel"),
|
|
"channel should not exist before it is created"
|
|
);
|
|
|
|
let channel_id = node
|
|
.channel_create_async(
|
|
TEST_CHANNEL_ID.to_string(),
|
|
TEST_CONTENT_TOPIC.to_string(),
|
|
TEST_SENDER_ID.to_string(),
|
|
NOOP_ENCRYPTION.to_string(),
|
|
)
|
|
.await
|
|
.expect("channel should be created");
|
|
assert_eq!(channel_id, TEST_CHANNEL_ID);
|
|
|
|
assert!(
|
|
node.channel_exists_async(TEST_CHANNEL_ID.to_string())
|
|
.await
|
|
.expect("exists should answer after create"),
|
|
"channel should exist once created"
|
|
);
|
|
|
|
let request_id = node
|
|
.channel_send_async(TEST_CHANNEL_ID.to_string(), channel_message(CHANNEL_PAYLOAD))
|
|
.await
|
|
.expect("channel send should succeed");
|
|
assert!(!request_id.is_empty(), "send should return a request id");
|
|
|
|
// Give the send a moment to finalise and emit its outcome event.
|
|
tokio::time::sleep(Duration::from_secs(2)).await;
|
|
|
|
node.channel_close_async(TEST_CHANNEL_ID.to_string())
|
|
.await
|
|
.expect("channel should close");
|
|
|
|
// Closing drops the channel from the manager even though its SDS state is
|
|
// persisted, so this is a liveness check rather than "was it ever created".
|
|
assert!(
|
|
!node
|
|
.channel_exists_async(TEST_CHANNEL_ID.to_string())
|
|
.await
|
|
.expect("exists should answer after close"),
|
|
"channel should not exist once closed"
|
|
);
|
|
|
|
// Scoped so the guard is dropped before the await below. Nothing is asserted
|
|
// about arrival: a lone node has no peer to confirm delivery with, so
|
|
// onChannelMessageSent legitimately never fires. Delivery is covered by
|
|
// channel_message_reaches_peer.
|
|
{
|
|
let sent = sent.lock().unwrap();
|
|
println!("channel sent events observed: {sent:?}");
|
|
}
|
|
|
|
node.stop_node_async().await.expect("node should stop");
|
|
// The context is torn down by LogosDeliveryCtx's Drop impl.
|
|
}
|
|
|
|
/// Two nodes sharing a channel: a channel message should cross the wire and
|
|
/// surface through the typed onChannelMessageReceived listener.
|
|
///
|
|
/// Ignored: this cannot work in-process, and the channels API is not at fault.
|
|
/// Persistency is a process-wide singleton keyed on its root dir, and it refuses
|
|
/// to be re-targeted, so both nodes share one SDS persistence job whose rows are
|
|
/// keyed by channel id alone. The receiver's ensureChannel therefore loads the
|
|
/// *sender's* history before the duplicate check, and SDS correctly rejects the
|
|
/// message as a replay -- silently, since a duplicate is not an error.
|
|
///
|
|
/// Everything either side of that works: autosharding, per-channel encryption,
|
|
/// the send, connectivity, the channel's ingress listener, and the typed event
|
|
/// path (proven by default_echo over the same route). Verified by instrumenting
|
|
/// the vendor: the listener fires, both filters pass, and SDS reports
|
|
/// `duplicate`.
|
|
///
|
|
/// Un-ignore by driving the two nodes as separate processes, or once SDS
|
|
/// persistence is keyed by participant as well as channel.
|
|
#[tokio::test]
|
|
#[serial]
|
|
#[ignore = "needs two processes: in-process nodes share SDS persistence; see the doc comment"]
|
|
async fn channel_message_reaches_peer() -> Result<(), String> {
|
|
// SDS state is persisted per channel id, so a fresh id keeps a previous
|
|
// run's causal history from parking this run's message.
|
|
let channel_id = format!(
|
|
"test-channel-{}",
|
|
std::time::SystemTime::now()
|
|
.duration_since(std::time::UNIX_EPOCH)
|
|
.unwrap()
|
|
.as_nanos()
|
|
);
|
|
let sender = new_node(60080);
|
|
let receiver = new_node(60090);
|
|
|
|
let received: Arc<Mutex<Vec<ChannelMessageReceivedPayload>>> = Arc::new(Mutex::new(Vec::new()));
|
|
let received_cloned = received.clone();
|
|
receiver.add_on_channel_message_received_listener(move |event| {
|
|
received_cloned.lock().unwrap().push(event.clone());
|
|
});
|
|
|
|
// The raw message is observed too: it proves delivery reaches the messaging
|
|
// layer, which is what localises the failure to channel ingress.
|
|
let raw: Arc<Mutex<Vec<usize>>> = Arc::new(Mutex::new(Vec::new()));
|
|
let raw_cloned = raw.clone();
|
|
receiver.add_on_message_received_listener(move |e| {
|
|
let len = base64::engine::general_purpose::STANDARD
|
|
.decode(&e.message.payload)
|
|
.map(|p| p.len())
|
|
.unwrap_or(0);
|
|
raw_cloned.lock().unwrap().push(len);
|
|
});
|
|
|
|
// A failed send reports itself here; without this the test could only time
|
|
// out and say nothing about why.
|
|
let errors: Arc<Mutex<Vec<String>>> = Arc::new(Mutex::new(Vec::new()));
|
|
let errors_cloned = errors.clone();
|
|
sender.add_on_channel_message_error_listener(move |event| {
|
|
errors_cloned.lock().unwrap().push(event.error.clone());
|
|
});
|
|
|
|
let msg_errors: Arc<Mutex<Vec<String>>> = Arc::new(Mutex::new(Vec::new()));
|
|
let msg_errors_cloned = msg_errors.clone();
|
|
sender.add_on_message_error_listener(move |event| {
|
|
msg_errors_cloned.lock().unwrap().push(event.error.clone());
|
|
});
|
|
|
|
sender.start_node_async().await?;
|
|
receiver.start_node_async().await?;
|
|
|
|
// Both ends open the same channel on the same content topic, under distinct
|
|
// participant ids.
|
|
sender
|
|
.channel_create_async(
|
|
channel_id.clone(),
|
|
TEST_CONTENT_TOPIC.to_string(),
|
|
TEST_SENDER_ID.to_string(),
|
|
NOOP_ENCRYPTION.to_string(),
|
|
)
|
|
.await?;
|
|
receiver
|
|
.channel_create_async(
|
|
channel_id.clone(),
|
|
TEST_CONTENT_TOPIC.to_string(),
|
|
OTHER_SENDER_ID.to_string(),
|
|
NOOP_ENCRYPTION.to_string(),
|
|
)
|
|
.await?;
|
|
|
|
// Creating a channel does not subscribe: ingress arrives through the
|
|
// messaging layer, so both ends have to subscribe to the content topic.
|
|
sender.subscribe_async(TEST_CONTENT_TOPIC.to_string()).await?;
|
|
receiver
|
|
.subscribe_async(TEST_CONTENT_TOPIC.to_string())
|
|
.await?;
|
|
|
|
connect(&receiver, &sender).await?;
|
|
|
|
// Wait for the mesh to form before publishing, or the message goes nowhere.
|
|
tokio::time::sleep(Duration::from_secs(5)).await;
|
|
|
|
sender
|
|
.channel_send_async(channel_id.clone(), channel_message(CHANNEL_PAYLOAD))
|
|
.await?;
|
|
|
|
let mut delivered = None;
|
|
for _ in 0..100 {
|
|
if let Some(event) = received.lock().unwrap().first().cloned() {
|
|
delivered = Some(event);
|
|
break;
|
|
}
|
|
tokio::time::sleep(Duration::from_millis(100)).await;
|
|
}
|
|
|
|
let event = delivered.ok_or_else(|| {
|
|
let errors = errors.lock().unwrap();
|
|
let msg_errors = msg_errors.lock().unwrap();
|
|
let raw = raw.lock().unwrap();
|
|
format!(
|
|
"never reached the peer; raw payload sizes seen by receiver: {raw:?}; \
|
|
channel errors: {errors:?}; messaging errors: {msg_errors:?}"
|
|
)
|
|
})?;
|
|
assert_eq!(event.channel_id, channel_id);
|
|
assert_eq!(event.sender_id, TEST_SENDER_ID);
|
|
let payload = base64::engine::general_purpose::STANDARD
|
|
.decode(&event.payload)
|
|
.map_err(|e| e.to_string())?;
|
|
assert_eq!(payload, CHANNEL_PAYLOAD);
|
|
|
|
sender.channel_close_async(channel_id.clone()).await?;
|
|
receiver.channel_close_async(channel_id).await?;
|
|
|
|
sender.stop_node_async().await?;
|
|
receiver.stop_node_async().await?;
|
|
|
|
Ok(())
|
|
}
|