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, 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) {
let config = ServerConfig::new(service_id, instance_id)
.with_interface(SERVER_IP)
.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(service_id, 1, server_addr, 0)
.await
.unwrap();
client
.subscribe(service_id, 1, 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(service_id, 1, server_addr, 0)
.await
.unwrap();
client
.subscribe(service_id, 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 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(service_id, 1, server_addr, 0)
.await
.unwrap();
client
.subscribe(service_id, 1, 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(service_id, 1).await.unwrap();
let msg = Message::<RawPayload>::new_sd(0x0001, &empty_sd_header());
let result = client.send_to_service(service_id, 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, 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(service_id, 1, server_addr, 0)
.await
.unwrap();
client
.subscribe(service_id, 1, 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(service_id, 1, server_addr, 0)
.await
.unwrap();
client
.subscribe(service_id, 1, 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(service_id, 1, 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(service_id, 1, server_addr, 0)
.await
.unwrap();
client
.subscribe(service_id, 1, 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(service_id, 1, server_addr, 0)
.await
.unwrap();
client1
.subscribe(service_id, 1, 1, 3, 0x01, 0)
.await
.unwrap();
let (client2, mut updates2, run_fut2) = TestClient::new(Ipv4Addr::LOCALHOST);
tokio::spawn(run_fut2);
client2
.add_endpoint(service_id, 1, server_addr, 0)
.await
.unwrap();
client2
.subscribe(service_id, 1, 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(service_id, 1, server_addr, 0)
.await
.unwrap();
client2
.subscribe(service_id, 1, 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(service_id, 1, server_addr, 0)
.await
.unwrap();
let specific_port = 44444;
client
.subscribe(service_id, 1, 1, 3, 0x01, specific_port)
.await
.unwrap();
client
.subscribe(service_id, 1, 1, 3, 0x02, specific_port)
.await
.unwrap();
client.shut_down();
server_handle.abort();
}