use simple_someip::ServiceEndpointKey;
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, TokioChannels,
VecSdHeader,
};
use std::net::{Ipv4Addr, SocketAddr, SocketAddrV4};
use std::sync::atomic::{AtomicU16, Ordering};
fn next_service_id() -> u16 {
static NEXT: AtomicU16 = AtomicU16::new(0x5B);
NEXT.fetch_add(1, Ordering::Relaxed)
}
fn empty_sd_header() -> VecSdHeader {
VecSdHeader {
flags: sd::Flags::new_sd(sd::RebootFlag::RecentlyRebooted),
entries: vec![],
options: vec![],
}
}
type TestClient = Client<
RawPayload,
std::sync::Arc<std::sync::Mutex<simple_someip::e2e::E2ERegistry>>,
std::sync::Arc<std::sync::RwLock<Ipv4Addr>>,
TokioChannels,
>;
type TestServer = Server<
simple_someip::TokioTransport,
simple_someip::TokioTimer,
std::sync::Arc<std::sync::Mutex<simple_someip::e2e::E2ERegistry>>,
std::sync::Arc<tokio::sync::RwLock<simple_someip::server::SubscriptionManager>>,
>;
type TestEventPublisher = simple_someip::server::EventPublisher<
std::sync::Arc<std::sync::Mutex<simple_someip::e2e::E2ERegistry>>,
std::sync::Arc<tokio::sync::RwLock<simple_someip::server::SubscriptionManager>>,
std::sync::Arc<simple_someip::TokioSocket>,
simple_someip::TokioSocket,
>;
const SERVER_IP: Ipv4Addr = Ipv4Addr::new(127, 0, 0, 2);
async fn create_server(service_id: u16, instance_id: u16) -> (TestServer, u16) {
create_server_on(SERVER_IP, service_id, instance_id).await
}
async fn create_server_on(
interface: Ipv4Addr,
service_id: u16,
instance_id: u16,
) -> (TestServer, u16) {
let config = ServerConfig::new(service_id, instance_id)
.with_interface(interface)
.with_local_port(0);
let (server, _handles, _run): (TestServer, _, _) =
TestServer::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, port)
}
async fn wait_for_subscribers(
publisher: &TestEventPublisher,
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, TokioChannels>,
) -> 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 service_id = next_service_id();
let (server, server_port) = create_server(service_id, 1).await;
let publisher = server.publisher();
let server_handle = tokio::spawn(async move { server.run().await });
let (client, mut updates, run_fut) = TestClient::new(Ipv4Addr::LOCALHOST);
let _run_handle = tokio::spawn(run_fut);
let server_addr = SocketAddrV4::new(SERVER_IP, server_port);
client
.add_endpoint(
ServiceEndpointKey::udp(service_id, SocketAddr::V4(server_addr)),
1,
0,
)
.await
.unwrap();
client
.subscribe(
ServiceEndpointKey::udp(service_id, SocketAddr::V4(server_addr)),
1,
3,
0x01,
0,
)
.await
.unwrap();
assert!(
wait_for_subscribers(&publisher, service_id, 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(service_id, 1, 0x01, &event_msg)
.await
.expect("publish_event failed");
assert_eq!(sent, 1);
let update = recv_unicast(&mut updates).await;
assert!(
matches!(update, ClientUpdate::Unicast { .. }),
"expected Unicast, got {update:?}"
);
client.unbind_discovery().await.unwrap();
client.shut_down();
server_handle.abort();
}
#[tokio::test]
async fn test_client_send_sd_auto_binds_discovery() {
let service_id = next_service_id();
let (server, server_port) = create_server(service_id, 1).await;
let server_handle = tokio::spawn(async move { server.run().await });
let (client, _updates, run_fut) = TestClient::new(Ipv4Addr::LOCALHOST);
let _run_handle = tokio::spawn(run_fut);
let sd_header = VecSdHeader {
flags: sd::Flags::new_sd(sd::RebootFlag::RecentlyRebooted),
entries: vec![sd::Entry::SubscribeEventGroup(sd::EventGroupEntry::new(
service_id, 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 service_id = next_service_id();
let (server, server_port) = create_server(service_id, 1).await;
let server_handle = tokio::spawn(async move { server.run().await });
let (client, _updates, run_fut) = TestClient::new(Ipv4Addr::LOCALHOST);
let _run_handle = tokio::spawn(run_fut);
client.bind_discovery().await.unwrap();
let server_addr = SocketAddrV4::new(SERVER_IP, server_port);
client
.add_endpoint(
ServiceEndpointKey::udp(service_id, SocketAddr::V4(server_addr)),
1,
0,
)
.await
.unwrap();
client
.subscribe(
ServiceEndpointKey::udp(service_id, SocketAddr::V4(server_addr)),
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 service_id = next_service_id();
let (server, server_port) = create_server(service_id, 1).await;
let publisher = server.publisher();
let server_handle = tokio::spawn(async move { server.run().await });
let (client, mut updates, run_fut) = TestClient::new(Ipv4Addr::LOCALHOST);
let _run_handle = tokio::spawn(run_fut);
client.bind_discovery().await.unwrap();
let server_addr = SocketAddrV4::new(SERVER_IP, server_port);
client
.add_endpoint(
ServiceEndpointKey::udp(service_id, SocketAddr::V4(server_addr)),
1,
0,
)
.await
.unwrap();
client
.subscribe(
ServiceEndpointKey::udp(service_id, SocketAddr::V4(server_addr)),
1,
3,
0x01,
0,
)
.await
.unwrap();
assert!(
wait_for_subscribers(&publisher, service_id, 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(service_id, 1, 0x01, &event_msg)
.await
.expect("publish_event failed");
assert_eq!(sent, 1);
let update = recv_unicast(&mut updates).await;
assert!(
matches!(update, ClientUpdate::Unicast { .. }),
"expected Unicast, got {update:?}"
);
client
.remove_endpoint(ServiceEndpointKey::udp(
service_id,
SocketAddr::V4(server_addr),
))
.await
.unwrap();
let msg = Message::<RawPayload>::new_sd(0x0001, &empty_sd_header());
let result = client
.send_to_service(
ServiceEndpointKey::udp(service_id, SocketAddr::V4(server_addr)),
msg,
)
.await;
assert!(
matches!(result, Err(simple_someip::client::Error::ServiceNotFound)),
"expected ServiceNotFound after remove, got {result:?}"
);
let _: fn() -> Option<simple_someip::PendingResponse<RawPayload, TokioChannels>> = || None;
client.shut_down();
server_handle.abort();
}
#[tokio::test]
async fn test_subscribe_auto_binds_discovery() {
let service_id = next_service_id();
let (server, server_port) = create_server(service_id, 1).await;
let publisher = server.publisher();
let server_handle = tokio::spawn(async move { server.run().await });
let (client, mut updates, run_fut) = TestClient::new(Ipv4Addr::LOCALHOST);
let _run_handle = tokio::spawn(run_fut);
let server_addr = SocketAddrV4::new(SERVER_IP, server_port);
client
.add_endpoint(
ServiceEndpointKey::udp(service_id, SocketAddr::V4(server_addr)),
1,
0,
)
.await
.unwrap();
client
.subscribe(
ServiceEndpointKey::udp(service_id, SocketAddr::V4(server_addr)),
1,
3,
0x01,
0,
)
.await
.unwrap();
assert!(
wait_for_subscribers(&publisher, service_id, 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(service_id, 1, 0x01, &event_msg)
.await
.expect("publish_event failed");
assert_eq!(sent, 1);
let update = recv_unicast(&mut updates).await;
assert!(
matches!(update, ClientUpdate::Unicast { .. }),
"expected Unicast, got {update:?}"
);
client.shut_down();
server_handle.abort();
}
#[tokio::test]
async fn test_client_request_resolves_via_unicast_reply() {
let service_id = next_service_id();
let (server, server_port) = create_server(service_id, 1).await;
let publisher = server.publisher();
let server_handle = tokio::spawn(async move { server.run().await });
let (client, mut updates, run_fut) = TestClient::new(Ipv4Addr::LOCALHOST);
let _run_handle = tokio::spawn(run_fut);
let server_addr = SocketAddrV4::new(SERVER_IP, server_port);
client
.add_endpoint(
ServiceEndpointKey::udp(service_id, SocketAddr::V4(server_addr)),
1,
0,
)
.await
.unwrap();
client
.subscribe(
ServiceEndpointKey::udp(service_id, SocketAddr::V4(server_addr)),
1,
3,
0x01,
0,
)
.await
.unwrap();
assert!(
wait_for_subscribers(&publisher, service_id, 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(
ServiceEndpointKey::udp(service_id, SocketAddr::V4(server_addr)),
msg,
)
.await
.expect("send_to_service failed");
let event_msg = Message::<RawPayload>::new_sd(0x0001, &empty_sd_header());
publisher
.publish_event(service_id, 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 service_id = next_service_id();
let (server, server_port) = create_server(service_id, 1).await;
let publisher = server.publisher();
let key = E2EKey {
service_id,
method_or_event_id: 0x0001,
};
let profile = E2EProfile::Profile4(Profile4Config::new(0x12345678, 15));
server
.register_e2e(key, profile.clone())
.expect("E2E registry has capacity for one entry");
let server_handle = tokio::spawn(async move { server.run().await });
let (client, mut updates, run_fut) = TestClient::new(Ipv4Addr::LOCALHOST);
let _run_handle = tokio::spawn(run_fut);
client
.register_e2e(key, profile)
.expect("E2E registry has capacity for one entry");
let server_addr = SocketAddrV4::new(SERVER_IP, server_port);
client
.add_endpoint(
ServiceEndpointKey::udp(service_id, SocketAddr::V4(server_addr)),
1,
0,
)
.await
.unwrap();
client
.subscribe(
ServiceEndpointKey::udp(service_id, SocketAddr::V4(server_addr)),
1,
3,
0x01,
0,
)
.await
.unwrap();
assert!(
wait_for_subscribers(&publisher, service_id, 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(service_id, 0x0001);
let raw_payload = RawPayload::from_payload_bytes(msg_id, &payload_bytes).unwrap();
let header = Header::new_event(service_id, 0x0001, 0, 0x01, 0x01, payload_bytes.len());
let event_msg = Message::new(header, raw_payload);
let sent = publisher
.publish_event(service_id, 1, 0x01, &event_msg)
.await
.expect("publish_event failed");
assert_eq!(sent, 1);
let update = recv_unicast(&mut updates).await;
match update {
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 service_id = next_service_id();
let (server, server_port) = create_server(service_id, 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, run_fut1) = TestClient::new(Ipv4Addr::LOCALHOST);
tokio::spawn(run_fut1);
client1
.add_endpoint(
ServiceEndpointKey::udp(service_id, SocketAddr::V4(server_addr)),
1,
0,
)
.await
.unwrap();
client1
.subscribe(
ServiceEndpointKey::udp(service_id, SocketAddr::V4(server_addr)),
1,
3,
0x01,
0,
)
.await
.unwrap();
let (client2, mut updates2, run_fut2) = TestClient::new(Ipv4Addr::LOCALHOST);
tokio::spawn(run_fut2);
client2
.add_endpoint(
ServiceEndpointKey::udp(service_id, SocketAddr::V4(server_addr)),
1,
0,
)
.await
.unwrap();
client2
.subscribe(
ServiceEndpointKey::udp(service_id, SocketAddr::V4(server_addr)),
1,
3,
0x01,
0,
)
.await
.unwrap();
let expected = ServerConfig::SUBSCRIBERS_PER_GROUP_CAP.min(2);
for _ in 0..40 {
if publisher.subscriber_count(service_id, 1, 0x01).await >= expected {
break;
}
tokio::time::sleep(std::time::Duration::from_millis(50)).await;
}
assert!(
publisher.subscriber_count(service_id, 1, 0x01).await >= expected,
"expected at least {expected} subscriber(s) (SUBSCRIBERS_PER_GROUP cap = {})",
ServerConfig::SUBSCRIBERS_PER_GROUP_CAP
);
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(service_id, 1, 0x01, &event_msg)
.await
.expect("publish_event failed");
assert!(sent >= expected, "expected sent >= {expected}, got {sent}");
let u1 = recv_unicast(&mut updates1).await;
assert!(
matches!(u1, ClientUpdate::Unicast { .. }),
"client1 expected Unicast, got {u1:?}"
);
if expected >= 2 {
let u2 = recv_unicast(&mut updates2).await;
assert!(
matches!(u2, ClientUpdate::Unicast { .. }),
"client2 expected Unicast, got {u2:?}"
);
} else {
eprintln!(
"SUBSCRIBERS_PER_GROUP cap = 1: skipping the second-subscriber fan-out \
assertion (rebuild with SIMPLE_SOMEIP_MAX_SUBS>=2 to exercise it)"
);
}
client1.shut_down();
client2.shut_down();
server_handle.abort();
}
#[tokio::test]
async fn test_updates_drain_after_shutdown() {
let (client, mut updates, run_fut) = TestClient::new(Ipv4Addr::LOCALHOST);
let _run_handle = tokio::spawn(run_fut);
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 service_id = next_service_id();
let (server, server_port) = create_server(service_id, 1).await;
let server_handle = tokio::spawn(async move { server.run().await });
let (client, _updates, run_fut) = TestClient::new(Ipv4Addr::LOCALHOST);
let _run_handle = tokio::spawn(run_fut);
let client2 = client.clone();
let server_addr = SocketAddrV4::new(SERVER_IP, server_port);
client
.add_endpoint(
ServiceEndpointKey::udp(service_id, SocketAddr::V4(server_addr)),
1,
0,
)
.await
.unwrap();
client2
.subscribe(
ServiceEndpointKey::udp(service_id, SocketAddr::V4(server_addr)),
1,
3,
0x01,
0,
)
.await
.unwrap();
client.shut_down();
server_handle.abort();
}
#[tokio::test]
async fn test_subscribe_specific_port_reuse() {
let service_id = next_service_id();
let (server, server_port) = create_server(service_id, 1).await;
let server_handle = tokio::spawn(async move { server.run().await });
let (client, _updates, run_fut) = TestClient::new(Ipv4Addr::LOCALHOST);
let _run_handle = tokio::spawn(run_fut);
let server_addr = SocketAddrV4::new(SERVER_IP, server_port);
client
.add_endpoint(
ServiceEndpointKey::udp(service_id, SocketAddr::V4(server_addr)),
1,
0,
)
.await
.unwrap();
let specific_port = 44444;
client
.subscribe(
ServiceEndpointKey::udp(service_id, SocketAddr::V4(server_addr)),
1,
3,
0x01,
specific_port,
)
.await
.unwrap();
client
.subscribe(
ServiceEndpointKey::udp(service_id, SocketAddr::V4(server_addr)),
1,
3,
0x02,
specific_port,
)
.await
.unwrap();
client.shut_down();
server_handle.abort();
}
#[tokio::test]
async fn test_two_devices_same_service_instance_addressed_independently() {
let service_id = next_service_id();
const DEVICE_A: Ipv4Addr = Ipv4Addr::new(127, 0, 0, 10);
const DEVICE_B: Ipv4Addr = Ipv4Addr::new(127, 0, 0, 11);
let (server_a, port_a) = create_server_on(DEVICE_A, service_id, 1).await;
let publisher_a = server_a.publisher();
let server_a_handle = tokio::spawn(async move { server_a.run().await });
let (server_b, port_b) = create_server_on(DEVICE_B, service_id, 1).await;
let publisher_b = server_b.publisher();
let server_b_handle = tokio::spawn(async move { server_b.run().await });
let addr_a = SocketAddrV4::new(DEVICE_A, port_a);
let addr_b = SocketAddrV4::new(DEVICE_B, port_b);
let (client, mut updates, run_fut) = TestClient::new(Ipv4Addr::LOCALHOST);
let _run_handle = tokio::spawn(run_fut);
client
.add_endpoint(
ServiceEndpointKey::udp(service_id, SocketAddr::V4(addr_a)),
1,
0,
)
.await
.unwrap();
client
.add_endpoint(
ServiceEndpointKey::udp(service_id, SocketAddr::V4(addr_b)),
1,
0,
)
.await
.unwrap();
client
.subscribe(
ServiceEndpointKey::udp(service_id, SocketAddr::V4(addr_a)),
1,
3,
0x01,
0,
)
.await
.unwrap();
client
.subscribe(
ServiceEndpointKey::udp(service_id, SocketAddr::V4(addr_b)),
1,
3,
0x02,
0,
)
.await
.unwrap();
assert!(
wait_for_subscribers(&publisher_a, service_id, 1, 0x01).await,
"server A should have registered the subscriber"
);
assert!(
wait_for_subscribers(&publisher_b, service_id, 1, 0x02).await,
"server B should have registered the subscriber"
);
let event_msg = Message::<RawPayload>::new_sd(0x0001, &empty_sd_header());
let sent_a = publisher_a
.publish_event(service_id, 1, 0x01, &event_msg)
.await
.expect("publish_event from device A failed");
assert_eq!(sent_a, 1);
match recv_unicast(&mut updates).await {
ClientUpdate::Unicast { source, .. } => {
assert_eq!(
source.ip(),
std::net::IpAddr::V4(DEVICE_A),
"event should have arrived from device A"
);
}
other => panic!("expected Unicast from device A, got {other:?}"),
}
let sent_b = publisher_b
.publish_event(service_id, 1, 0x02, &event_msg)
.await
.expect("publish_event from device B failed");
assert_eq!(sent_b, 1);
match recv_unicast(&mut updates).await {
ClientUpdate::Unicast { source, .. } => {
assert_eq!(
source.ip(),
std::net::IpAddr::V4(DEVICE_B),
"event should have arrived from device B"
);
}
other => panic!("expected Unicast from device B, got {other:?}"),
}
client
.remove_endpoint(ServiceEndpointKey::udp(service_id, SocketAddr::V4(addr_a)))
.await
.unwrap();
let msg = Message::<RawPayload>::new_sd(0x0001, &empty_sd_header());
let result_a = client
.send_to_service(
ServiceEndpointKey::udp(service_id, SocketAddr::V4(addr_a)),
msg,
)
.await;
assert!(
matches!(result_a, Err(simple_someip::client::Error::ServiceNotFound)),
"expected ServiceNotFound for removed device A, got {result_a:?}"
);
let msg = Message::<RawPayload>::new_sd(0x0001, &empty_sd_header());
let result_b = client
.send_to_service(
ServiceEndpointKey::udp(service_id, SocketAddr::V4(addr_b)),
msg,
)
.await;
assert!(
result_b.is_ok(),
"device B must remain reachable after removing device A's endpoint: {result_b:?}"
);
client.shut_down();
server_a_handle.abort();
server_b_handle.abort();
}