use simple_someip::e2e::{E2ECheckStatus, E2EKey, E2EProfile, Profile4Config};
use simple_someip::protocol::{Header, Message, MessageId, sd};
use simple_someip::server::ServerConfig;
use simple_someip::{
Client, ClientUpdate, ClientUpdates, PayloadWireFormat, RawPayload, Server, VecSdHeader,
};
use std::net::{Ipv4Addr, SocketAddrV4};
fn empty_sd_header() -> VecSdHeader {
VecSdHeader {
flags: sd::Flags::new_sd(sd::RebootFlag::RecentlyRebooted),
entries: vec![],
options: vec![],
}
}
type TestClient = Client<RawPayload>;
const SERVER_IP: Ipv4Addr = Ipv4Addr::new(127, 0, 0, 2);
async fn create_server(service_id: u16, instance_id: u16) -> (Server, u16) {
let config = ServerConfig::new(SERVER_IP, 0, service_id, instance_id);
let mut server: Server = Server::new(config).await.expect("Server::new failed");
let port = match server.unicast_local_addr().expect("local_addr failed") {
std::net::SocketAddr::V4(a) => a.port(),
_ => panic!("expected IPv4"),
};
server.set_local_port(port);
(server, port)
}
async fn wait_for_subscribers(
publisher: &simple_someip::server::EventPublisher,
service_id: u16,
instance_id: u16,
event_group_id: u16,
) -> bool {
for _ in 0..20 {
if publisher
.has_subscribers(service_id, instance_id, event_group_id)
.await
{
return true;
}
tokio::time::sleep(std::time::Duration::from_millis(50)).await;
}
false
}
async fn recv_unicast(updates: &mut ClientUpdates<RawPayload>) -> ClientUpdate<RawPayload> {
let deadline = tokio::time::Instant::now() + std::time::Duration::from_secs(5);
loop {
match tokio::time::timeout_at(deadline, updates.recv()).await {
Ok(Some(update @ ClientUpdate::Unicast { .. })) => return update,
Ok(Some(_)) => continue,
Ok(None) => panic!("update channel closed before the Unicast event"),
Err(_) => panic!("timed out waiting for the Unicast event"),
}
}
}
#[tokio::test]
async fn test_client_server_subscribe_and_receive_event() {
let (mut server, server_port) = create_server(0x5B, 1).await;
let publisher = server.publisher();
let server_handle = tokio::spawn(async move { server.run().await });
let (client, mut updates) = TestClient::new(Ipv4Addr::LOCALHOST);
let server_addr = SocketAddrV4::new(SERVER_IP, server_port);
client.add_endpoint(0x5B, 1, server_addr, 0).await.unwrap();
client.subscribe(0x5B, 1, 1, 3, 0x01, 0).await.unwrap();
assert!(
wait_for_subscribers(&publisher, 0x5B, 1, 0x01).await,
"server should have registered the subscriber"
);
let _ = tokio::time::timeout(std::time::Duration::from_millis(250), updates.recv()).await;
let event_msg = Message::<RawPayload>::new_sd(0x0001, &empty_sd_header());
let sent = publisher
.publish_event(0x5B, 1, 0x01, &event_msg)
.await
.expect("publish_event failed");
assert_eq!(sent, 1);
recv_unicast(&mut updates).await;
client.unbind_discovery().await.unwrap();
client.shut_down();
server_handle.abort();
}
#[tokio::test]
async fn test_client_send_sd_auto_binds_discovery() {
let (mut server, server_port) = create_server(0x5B, 1).await;
let server_handle = tokio::spawn(async move { server.run().await });
let (client, _updates) = TestClient::new(Ipv4Addr::LOCALHOST);
let sd_header = VecSdHeader {
flags: sd::Flags::new_sd(sd::RebootFlag::RecentlyRebooted),
entries: vec![sd::Entry::SubscribeEventGroup(sd::EventGroupEntry::new(
0x5B, 1, 1, 3, 0x01,
))],
options: vec![sd::Options::IpV4Endpoint {
ip: Ipv4Addr::LOCALHOST,
protocol: sd::TransportProtocol::Udp,
port: 12345,
}],
};
let target = SocketAddrV4::new(SERVER_IP, server_port);
client
.send_sd_message(target, sd_header)
.await
.expect("send_sd_message should auto-bind discovery and succeed");
client.shut_down();
server_handle.abort();
}
#[tokio::test]
async fn test_client_bind_unbind_lifecycle_with_server() {
let (mut server, server_port) = create_server(0x5B, 1).await;
let server_handle = tokio::spawn(async move { server.run().await });
let (client, _updates) = TestClient::new(Ipv4Addr::LOCALHOST);
client.bind_discovery().await.unwrap();
let server_addr = SocketAddrV4::new(SERVER_IP, server_port);
client.add_endpoint(0x5B, 1, server_addr, 0).await.unwrap();
client.subscribe(0x5B, 1, 1, 3, 0x01, 0).await.unwrap();
client.unbind_discovery().await.unwrap();
client.bind_discovery().await.unwrap();
client.set_interface(Ipv4Addr::LOCALHOST).await.unwrap();
client.shut_down();
server_handle.abort();
}
#[tokio::test]
async fn test_add_endpoint_and_send_to_service() {
let (mut server, server_port) = create_server(0x5B, 1).await;
let publisher = server.publisher();
let server_handle = tokio::spawn(async move { server.run().await });
let (client, mut updates) = TestClient::new(Ipv4Addr::LOCALHOST);
client.bind_discovery().await.unwrap();
let server_addr = SocketAddrV4::new(SERVER_IP, server_port);
client.add_endpoint(0x5B, 1, server_addr, 0).await.unwrap();
client.subscribe(0x5B, 1, 1, 3, 0x01, 0).await.unwrap();
assert!(
wait_for_subscribers(&publisher, 0x5B, 1, 0x01).await,
"server should have registered the subscriber"
);
let _ = tokio::time::timeout(std::time::Duration::from_millis(250), updates.recv()).await;
let event_msg = Message::<RawPayload>::new_sd(0x0001, &empty_sd_header());
let sent = publisher
.publish_event(0x5B, 1, 0x01, &event_msg)
.await
.expect("publish_event failed");
assert_eq!(sent, 1);
recv_unicast(&mut updates).await;
client.remove_endpoint(0x5B, 1).await.unwrap();
let msg = Message::<RawPayload>::new_sd(0x0001, &empty_sd_header());
let result = client.send_to_service(0x5B, 1, msg).await;
assert!(
matches!(result, Err(simple_someip::client::Error::ServiceNotFound)),
"expected ServiceNotFound after remove, got {result:?}"
);
let _: fn() -> Option<simple_someip::PendingResponse<RawPayload>> = || None;
client.shut_down();
server_handle.abort();
}
#[tokio::test]
async fn test_subscribe_auto_binds_discovery() {
let (mut server, server_port) = create_server(0x5B, 1).await;
let publisher = server.publisher();
let server_handle = tokio::spawn(async move { server.run().await });
let (client, mut updates) = TestClient::new(Ipv4Addr::LOCALHOST);
let server_addr = SocketAddrV4::new(SERVER_IP, server_port);
client.add_endpoint(0x5B, 1, server_addr, 0).await.unwrap();
client.subscribe(0x5B, 1, 1, 3, 0x01, 0).await.unwrap();
assert!(
wait_for_subscribers(&publisher, 0x5B, 1, 0x01).await,
"server should have registered the subscriber"
);
let _ = tokio::time::timeout(std::time::Duration::from_millis(250), updates.recv()).await;
let event_msg = Message::<RawPayload>::new_sd(0x0001, &empty_sd_header());
let sent = publisher
.publish_event(0x5B, 1, 0x01, &event_msg)
.await
.expect("publish_event failed");
assert_eq!(sent, 1);
recv_unicast(&mut updates).await;
client.shut_down();
server_handle.abort();
}
#[tokio::test]
async fn test_client_request_resolves_via_unicast_reply() {
let (mut server, server_port) = create_server(0x5B, 1).await;
let publisher = server.publisher();
let server_handle = tokio::spawn(async move { server.run().await });
let (client, mut updates) = TestClient::new(Ipv4Addr::LOCALHOST);
let server_addr = SocketAddrV4::new(SERVER_IP, server_port);
client.add_endpoint(0x5B, 1, server_addr, 0).await.unwrap();
client.subscribe(0x5B, 1, 1, 3, 0x01, 0).await.unwrap();
assert!(
wait_for_subscribers(&publisher, 0x5B, 1, 0x01).await,
"server should have registered the subscriber"
);
let _ = tokio::time::timeout(std::time::Duration::from_millis(250), updates.recv()).await;
let msg = Message::<RawPayload>::new_sd(0x0001, &empty_sd_header());
let pending = client
.send_to_service(0x5B, 1, msg)
.await
.expect("send_to_service failed");
let event_msg = Message::<RawPayload>::new_sd(0x0001, &empty_sd_header());
publisher
.publish_event(0x5B, 1, 0x01, &event_msg)
.await
.expect("publish_event failed");
let update = tokio::time::timeout(std::time::Duration::from_secs(2), updates.recv())
.await
.expect("timeout waiting for unicast update");
assert!(update.is_some(), "expected an update");
drop(pending);
client.shut_down();
server_handle.abort();
}
#[tokio::test]
async fn test_e2e_protect_on_publish_and_check_on_receive() {
let (mut server, server_port) = create_server(0x5B, 1).await;
let publisher = server.publisher();
let key = E2EKey {
service_id: 0x5B,
method_or_event_id: 0x0001,
};
let profile = E2EProfile::Profile4(Profile4Config::new(0x12345678, 15));
server.register_e2e(key, profile.clone());
let server_handle = tokio::spawn(async move { server.run().await });
let (client, mut updates) = TestClient::new(Ipv4Addr::LOCALHOST);
client.register_e2e(key, profile);
let server_addr = SocketAddrV4::new(SERVER_IP, server_port);
client.add_endpoint(0x5B, 1, server_addr, 0).await.unwrap();
client.subscribe(0x5B, 1, 1, 3, 0x01, 0).await.unwrap();
assert!(
wait_for_subscribers(&publisher, 0x5B, 1, 0x01).await,
"server should have registered the subscriber"
);
let _ = tokio::time::timeout(std::time::Duration::from_millis(250), updates.recv()).await;
let payload_bytes = [0xAA, 0xBB];
let msg_id = MessageId::new_from_service_and_method(0x5B, 0x0001);
let raw_payload = RawPayload::from_payload_bytes(msg_id, &payload_bytes).unwrap();
let header = Header::new_event(0x5B, 0x0001, 0, 0x01, 0x01, payload_bytes.len());
let event_msg = Message::new(header, raw_payload);
let sent = publisher
.publish_event(0x5B, 1, 0x01, &event_msg)
.await
.expect("publish_event failed");
assert_eq!(sent, 1);
match recv_unicast(&mut updates).await {
ClientUpdate::Unicast { e2e_status, .. } => {
assert!(
e2e_status.is_some(),
"expected e2e_status to be populated when E2E is configured"
);
assert_eq!(
e2e_status.unwrap(),
E2ECheckStatus::Ok,
"E2E check should pass for correctly protected message"
);
}
other => unreachable!("recv_unicast only returns Unicast, got {other:?}"),
}
client.shut_down();
server_handle.abort();
}
#[tokio::test]
async fn test_multiple_subscribers_receive_events() {
let (mut server, server_port) = create_server(0x5B, 1).await;
let publisher = server.publisher();
let server_handle = tokio::spawn(async move { server.run().await });
let server_addr = SocketAddrV4::new(SERVER_IP, server_port);
let (client1, mut updates1) = TestClient::new(Ipv4Addr::LOCALHOST);
client1.add_endpoint(0x5B, 1, server_addr, 0).await.unwrap();
client1.subscribe(0x5B, 1, 1, 3, 0x01, 0).await.unwrap();
let (client2, mut updates2) = TestClient::new(Ipv4Addr::LOCALHOST);
client2.add_endpoint(0x5B, 1, server_addr, 0).await.unwrap();
client2.subscribe(0x5B, 1, 1, 3, 0x01, 0).await.unwrap();
for _ in 0..40 {
if publisher.subscriber_count(0x5B, 1, 0x01).await >= 2 {
break;
}
tokio::time::sleep(std::time::Duration::from_millis(50)).await;
}
assert!(
publisher.subscriber_count(0x5B, 1, 0x01).await >= 2,
"expected at least 2 subscribers"
);
let _ = tokio::time::timeout(std::time::Duration::from_millis(250), updates1.recv()).await;
let _ = tokio::time::timeout(std::time::Duration::from_millis(250), updates2.recv()).await;
let event_msg = Message::<RawPayload>::new_sd(0x0001, &empty_sd_header());
let sent = publisher
.publish_event(0x5B, 1, 0x01, &event_msg)
.await
.expect("publish_event failed");
assert!(sent >= 2, "expected sent >= 2, got {sent}");
recv_unicast(&mut updates1).await;
recv_unicast(&mut updates2).await;
client1.shut_down();
client2.shut_down();
server_handle.abort();
}
#[tokio::test]
async fn test_updates_drain_after_shutdown() {
let (client, mut updates) = TestClient::new(Ipv4Addr::LOCALHOST);
client.shut_down();
let result = tokio::time::timeout(std::time::Duration::from_secs(2), updates.recv())
.await
.expect("timeout waiting for None");
assert!(result.is_none(), "expected None after shutdown");
}
#[tokio::test]
async fn test_cloned_client_works() {
let (mut server, server_port) = create_server(0x5B, 1).await;
let server_handle = tokio::spawn(async move { server.run().await });
let (client, _updates) = TestClient::new(Ipv4Addr::LOCALHOST);
let client2 = client.clone();
let server_addr = SocketAddrV4::new(SERVER_IP, server_port);
client.add_endpoint(0x5B, 1, server_addr, 0).await.unwrap();
client2.subscribe(0x5B, 1, 1, 3, 0x01, 0).await.unwrap();
client.shut_down();
server_handle.abort();
}
#[tokio::test]
async fn test_subscribe_specific_port_reuse() {
let (mut server, server_port) = create_server(0x5B, 1).await;
let server_handle = tokio::spawn(async move { server.run().await });
let (client, _updates) = TestClient::new(Ipv4Addr::LOCALHOST);
let server_addr = SocketAddrV4::new(SERVER_IP, server_port);
client.add_endpoint(0x5B, 1, server_addr, 0).await.unwrap();
let specific_port = 44444;
client
.subscribe(0x5B, 1, 1, 3, 0x01, specific_port)
.await
.unwrap();
client
.subscribe(0x5B, 1, 1, 3, 0x02, specific_port)
.await
.unwrap();
client.shut_down();
server_handle.abort();
}