mod bind_dispatch;
mod error;
mod inner;
mod service_registry;
mod session;
mod socket_manager;
pub use error::Error;
pub use inner::ControlMessage;
pub use service_registry::ServiceEndpointKey;
pub use socket_manager::{ReceivedMessage, SendMessage};
use crate::Timer;
#[cfg(feature = "client-tokio")]
use crate::e2e::E2ERegistry;
use crate::e2e::{E2ECheckStatus, E2EKey, E2EProfile};
use crate::log::info;
#[cfg(feature = "client-tokio")]
use crate::tokio_transport::{TokioChannels, TokioSpawner, TokioTimer};
use crate::transport::{
BoundedPooled, ChannelFactory, E2ERegistryHandle, InterfaceHandle, MpscSend, OneshotPooled,
OneshotRecv, Spawner, TransportFactory, TransportSocket, UnboundedPooled, UnboundedRecv,
};
use crate::{protocol, protocol::Message, traits::PayloadWireFormat};
use core::net::{Ipv4Addr, SocketAddr, SocketAddrV4};
use inner::Inner;
#[cfg(feature = "client-tokio")]
use std::sync::{Arc, Mutex, RwLock};
pub trait ClientChannelTypes<P: PayloadWireFormat + 'static>: ChannelFactory
where
Result<(), Error>: OneshotPooled<Self>,
Result<P, Error>: OneshotPooled<Self>,
Result<protocol::sd::RebootFlag, Error>: OneshotPooled<Self>,
ControlMessage<P, Self>: BoundedPooled<Self, 4>,
SendMessage<P, Self>: BoundedPooled<Self, 16>,
Result<ReceivedMessage<P>, Error>: BoundedPooled<Self, 16>,
ClientUpdate<P>: UnboundedPooled<Self>,
{
}
impl<P, C> ClientChannelTypes<P> for C
where
P: PayloadWireFormat + 'static,
C: ChannelFactory,
Result<(), Error>: OneshotPooled<C>,
Result<P, Error>: OneshotPooled<C>,
Result<protocol::sd::RebootFlag, Error>: OneshotPooled<C>,
ControlMessage<P, C>: BoundedPooled<C, 4>,
SendMessage<P, C>: BoundedPooled<C, 16>,
Result<ReceivedMessage<P>, Error>: BoundedPooled<C, 16>,
ClientUpdate<P>: UnboundedPooled<C>,
{
}
pub struct PendingResponse<P: Send + 'static, C: ChannelFactory> {
receiver: C::OneshotReceiver<Result<P, Error>>,
}
impl<P: Send + 'static, C: ChannelFactory> core::fmt::Debug for PendingResponse<P, C> {
fn fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
f.debug_struct("PendingResponse").finish_non_exhaustive()
}
}
impl<P: Send + 'static, C: ChannelFactory> PendingResponse<P, C> {
pub async fn response(self) -> Result<P, Error> {
self.receiver.recv().await.map_err(|_| Error::Shutdown)?
}
}
pub struct DiscoveryMessage<P: PayloadWireFormat> {
pub source: SocketAddr,
pub someip_header: protocol::Header,
pub sd_header: P::SdHeader,
}
impl<P: PayloadWireFormat> core::fmt::Debug for DiscoveryMessage<P> {
fn fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
f.debug_struct("DiscoveryMessage")
.field("source", &self.source)
.field("someip_header", &self.someip_header)
.field("sd_header", &self.sd_header)
.finish()
}
}
pub enum ClientUpdate<P: PayloadWireFormat> {
DiscoveryUpdated(DiscoveryMessage<P>),
SenderRebooted(SocketAddr),
Unicast {
message: Message<P>,
e2e_status: Option<E2ECheckStatus>,
source: SocketAddr,
},
Error(Error),
}
impl<P: PayloadWireFormat> core::fmt::Debug for ClientUpdate<P> {
fn fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
match self {
Self::DiscoveryUpdated(msg) => f.debug_tuple("DiscoveryUpdated").field(msg).finish(),
Self::SenderRebooted(addr) => f.debug_tuple("SenderRebooted").field(addr).finish(),
Self::Unicast {
message,
e2e_status,
source,
} => f
.debug_struct("Unicast")
.field("message", message)
.field("e2e_status", e2e_status)
.field("source", source)
.finish(),
Self::Error(err) => f.debug_tuple("Error").field(err).finish(),
}
}
}
pub struct ClientUpdates<MessageDefinitions: PayloadWireFormat + 'static, C: ChannelFactory> {
update_receiver: C::UnboundedReceiver<ClientUpdate<MessageDefinitions>>,
}
impl<MessageDefinitions: PayloadWireFormat + 'static, C: ChannelFactory> core::fmt::Debug
for ClientUpdates<MessageDefinitions, C>
{
fn fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
f.debug_struct("ClientUpdates").finish_non_exhaustive()
}
}
impl<MessageDefinitions: PayloadWireFormat + 'static, C: ChannelFactory>
ClientUpdates<MessageDefinitions, C>
{
pub async fn recv(&mut self) -> Option<ClientUpdate<MessageDefinitions>> {
UnboundedRecv::recv(&mut self.update_receiver).await
}
}
pub struct ClientDeps<F, Tm, R, I, Sp, BP>
where
F: TransportFactory,
Tm: Timer,
R: E2ERegistryHandle,
I: InterfaceHandle,
BP: crate::transport::BufferProvider,
{
pub factory: F,
pub timer: Tm,
pub e2e_registry: R,
pub interface: I,
pub spawner: Sp,
pub buffer_provider: BP,
}
#[cfg(feature = "client-tokio")]
impl
ClientDeps<
crate::tokio_transport::TokioTransport,
TokioTimer,
Arc<Mutex<E2ERegistry>>,
Arc<RwLock<Ipv4Addr>>,
TokioSpawner,
crate::tokio_transport::TokioBufferProvider,
>
{
#[must_use]
pub fn tokio(interface: Ipv4Addr) -> Self {
Self {
factory: crate::tokio_transport::TokioTransport,
timer: TokioTimer,
e2e_registry: Arc::new(Mutex::new(E2ERegistry::new())),
interface: Arc::new(RwLock::new(interface)),
spawner: TokioSpawner,
buffer_provider: crate::tokio_transport::TokioBufferProvider::new(),
}
}
}
impl<F, Tm, R, I, Sp, BP> ClientDeps<F, Tm, R, I, Sp, BP>
where
F: TransportFactory,
Tm: Timer,
R: E2ERegistryHandle,
I: InterfaceHandle,
BP: crate::transport::BufferProvider,
{
pub fn with_factory<F2: TransportFactory>(
self,
factory: F2,
) -> ClientDeps<F2, Tm, R, I, Sp, BP> {
ClientDeps {
factory,
timer: self.timer,
e2e_registry: self.e2e_registry,
interface: self.interface,
spawner: self.spawner,
buffer_provider: self.buffer_provider,
}
}
pub fn with_timer<Tm2: Timer>(self, timer: Tm2) -> ClientDeps<F, Tm2, R, I, Sp, BP> {
ClientDeps {
factory: self.factory,
timer,
e2e_registry: self.e2e_registry,
interface: self.interface,
spawner: self.spawner,
buffer_provider: self.buffer_provider,
}
}
pub fn with_e2e_registry<R2: E2ERegistryHandle>(
self,
e2e_registry: R2,
) -> ClientDeps<F, Tm, R2, I, Sp, BP> {
ClientDeps {
factory: self.factory,
timer: self.timer,
e2e_registry,
interface: self.interface,
spawner: self.spawner,
buffer_provider: self.buffer_provider,
}
}
pub fn with_interface<I2: InterfaceHandle>(
self,
interface: I2,
) -> ClientDeps<F, Tm, R, I2, Sp, BP> {
ClientDeps {
factory: self.factory,
timer: self.timer,
e2e_registry: self.e2e_registry,
interface,
spawner: self.spawner,
buffer_provider: self.buffer_provider,
}
}
pub fn with_spawner<Sp2: Spawner>(self, spawner: Sp2) -> ClientDeps<F, Tm, R, I, Sp2, BP> {
ClientDeps {
factory: self.factory,
timer: self.timer,
e2e_registry: self.e2e_registry,
interface: self.interface,
spawner,
buffer_provider: self.buffer_provider,
}
}
pub fn with_local_spawner<Sp2: crate::transport::LocalSpawner>(
self,
spawner: Sp2,
) -> ClientDeps<F, Tm, R, I, Sp2, BP> {
ClientDeps {
factory: self.factory,
timer: self.timer,
e2e_registry: self.e2e_registry,
interface: self.interface,
spawner,
buffer_provider: self.buffer_provider,
}
}
pub fn with_buffer_provider<BP2: crate::transport::BufferProvider>(
self,
buffer_provider: BP2,
) -> ClientDeps<F, Tm, R, I, Sp, BP2> {
ClientDeps {
factory: self.factory,
timer: self.timer,
e2e_registry: self.e2e_registry,
interface: self.interface,
spawner: self.spawner,
buffer_provider,
}
}
}
#[derive(Clone)]
pub struct Client<
MessageDefinitions: PayloadWireFormat + Send + 'static,
R: E2ERegistryHandle,
I: InterfaceHandle,
C: ChannelFactory,
> {
interface: I,
control_sender: C::BoundedSender<inner::ControlMessage<MessageDefinitions, C>, 4>,
e2e_registry: R,
}
impl<MessageDefinitions, R, I, C> core::fmt::Debug for Client<MessageDefinitions, R, I, C>
where
MessageDefinitions: PayloadWireFormat + Send + 'static,
R: E2ERegistryHandle,
I: InterfaceHandle,
C: ChannelFactory,
{
fn fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
f.debug_struct("Client")
.field("interface", &self.interface.get())
.finish_non_exhaustive()
}
}
#[cfg(feature = "client-tokio")]
impl<MessageDefinitions>
Client<MessageDefinitions, Arc<Mutex<E2ERegistry>>, Arc<RwLock<Ipv4Addr>>, TokioChannels>
where
MessageDefinitions: PayloadWireFormat + Clone + core::fmt::Debug + 'static,
{
#[must_use = "the returned run-loop future must be spawned (e.g. tokio::spawn) for the client to make progress"]
pub fn new(
interface: Ipv4Addr,
) -> (
Self,
ClientUpdates<MessageDefinitions, TokioChannels>,
impl core::future::Future<Output = ()> + Send + 'static,
) {
Self::new_with_loopback(interface, false)
}
#[must_use = "the returned run-loop future must be spawned (e.g. tokio::spawn) for the client to make progress"]
pub fn new_with_loopback(
interface: Ipv4Addr,
multicast_loopback: bool,
) -> (
Self,
ClientUpdates<MessageDefinitions, TokioChannels>,
impl core::future::Future<Output = ()> + Send + 'static,
) {
Self::new_with_spawner_and_loopback(interface, multicast_loopback, TokioSpawner)
}
#[must_use = "the returned run-loop future must be spawned (e.g. via the Spawner) for the client to make progress"]
pub fn new_with_spawner_and_loopback<Sp>(
interface: Ipv4Addr,
multicast_loopback: bool,
spawner: Sp,
) -> (
Self,
ClientUpdates<MessageDefinitions, TokioChannels>,
impl core::future::Future<Output = ()> + Send + 'static,
)
where
Sp: Spawner + Send + Sync + 'static,
{
Self::new_with_deps(
ClientDeps {
factory: crate::tokio_transport::TokioTransport,
timer: TokioTimer,
e2e_registry: Arc::new(Mutex::new(E2ERegistry::new())),
interface: Arc::new(RwLock::new(interface)),
spawner,
buffer_provider: crate::tokio_transport::TokioBufferProvider::new(),
},
multicast_loopback,
)
}
}
impl<MessageDefinitions, R, I, C> Client<MessageDefinitions, R, I, C>
where
MessageDefinitions: PayloadWireFormat + Clone + core::fmt::Debug + Send + 'static,
R: E2ERegistryHandle,
I: InterfaceHandle,
C: ChannelFactory,
Result<(), Error>: OneshotPooled<C>,
Result<MessageDefinitions, Error>: OneshotPooled<C>,
Result<protocol::sd::RebootFlag, Error>: OneshotPooled<C>,
ControlMessage<MessageDefinitions, C>: BoundedPooled<C, 4>,
SendMessage<MessageDefinitions, C>: BoundedPooled<C, 16>,
Result<ReceivedMessage<MessageDefinitions>, Error>: BoundedPooled<C, 16>,
ClientUpdate<MessageDefinitions>: UnboundedPooled<C>,
{
#[allow(clippy::type_complexity)]
#[must_use = "the returned run-loop future must be spawned (e.g. via the Spawner) for the client to make progress"]
pub fn new_with_deps<F, Tm, Sp, BP>(
deps: ClientDeps<F, Tm, R, I, Sp, BP>,
multicast_loopback: bool,
) -> (
Self,
ClientUpdates<MessageDefinitions, C>,
impl core::future::Future<Output = ()> + Send + 'static,
)
where
F: TransportFactory + Send + Sync + 'static,
F::Socket: Send + Sync + 'static,
for<'a> F::BindFuture<'a>: Send,
for<'a> <F::Socket as TransportSocket>::SendFuture<'a>: Send,
for<'a> <F::Socket as TransportSocket>::RecvFuture<'a>: Send,
Sp: Spawner + Send + Sync + 'static,
Tm: Timer + Send + Sync + 'static,
for<'a> Tm::SleepFuture<'a>: Send,
BP: crate::transport::BufferProvider,
{
let ClientDeps {
factory,
timer,
e2e_registry,
interface,
spawner,
buffer_provider,
} = deps;
let initial_addr = interface.get();
let dispatch = bind_dispatch::SpawnerDispatch {
factory,
spawner,
buffer_provider,
};
let (control_sender, update_receiver, run_future) = Inner::<
MessageDefinitions,
Tm,
R,
C,
bind_dispatch::SpawnerDispatch<F, Sp, BP>,
>::build(
initial_addr,
e2e_registry.clone(),
multicast_loopback,
dispatch,
timer,
);
let client = Self {
interface,
control_sender,
e2e_registry,
};
let updates = ClientUpdates { update_receiver };
(client, updates, run_future)
}
#[allow(clippy::type_complexity)]
#[must_use = "the returned run-loop future must be spawned (e.g. via the LocalSpawner) for the client to make progress"]
pub fn new_with_deps_local<F, Tm, Sp, BP>(
deps: ClientDeps<F, Tm, R, I, Sp, BP>,
multicast_loopback: bool,
) -> (
Self,
ClientUpdates<MessageDefinitions, C>,
impl core::future::Future<Output = ()> + 'static,
)
where
F: TransportFactory + 'static,
F::Socket: 'static,
Sp: crate::transport::LocalSpawner + 'static,
Tm: Timer + 'static,
BP: crate::transport::BufferProvider,
{
let ClientDeps {
factory,
timer,
e2e_registry,
interface,
spawner,
buffer_provider,
} = deps;
let initial_addr = interface.get();
let dispatch = bind_dispatch::LocalSpawnerDispatch {
factory,
spawner,
buffer_provider,
};
let (control_sender, update_receiver, run_future) = Inner::<
MessageDefinitions,
Tm,
R,
C,
bind_dispatch::LocalSpawnerDispatch<F, Sp, BP>,
>::build(
initial_addr,
e2e_registry.clone(),
multicast_loopback,
dispatch,
timer,
);
let client = Self {
interface,
control_sender,
e2e_registry,
};
let updates = ClientUpdates { update_receiver };
(client, updates, run_future)
}
#[must_use]
pub fn interface(&self) -> Ipv4Addr {
self.interface.get()
}
pub async fn set_interface(&self, interface: Ipv4Addr) -> Result<(), Error> {
let (response, message) = ControlMessage::set_interface(interface);
self.control_sender
.send(message)
.await
.map_err(|()| Error::Shutdown)?;
response.recv().await.map_err(|_| Error::Shutdown)??;
self.interface.set(interface);
Ok(())
}
pub async fn bind_discovery(&self) -> Result<(), Error> {
let (response, message) = ControlMessage::bind_discovery();
self.control_sender
.send(message)
.await
.map_err(|()| Error::Shutdown)?;
response.recv().await.map_err(|_| Error::Shutdown)?
}
pub async fn unbind_discovery(&self) -> Result<(), Error> {
let (response, message) = ControlMessage::unbind_discovery();
self.control_sender
.send(message)
.await
.map_err(|()| Error::Shutdown)?;
response.recv().await.map_err(|_| Error::Shutdown)?
}
pub async fn subscribe(
&self,
key: ServiceEndpointKey,
major_version: u8,
ttl: u32,
event_group_id: u16,
client_port: u16,
) -> Result<(), Error> {
let (response, message) =
ControlMessage::subscribe(key, major_version, ttl, event_group_id, client_port);
self.control_sender
.send(message)
.await
.map_err(|()| Error::Shutdown)?;
response.recv().await.map_err(|_| Error::Shutdown)?
}
pub async fn subscribe_no_wait(
&self,
key: ServiceEndpointKey,
major_version: u8,
ttl: u32,
event_group_id: u16,
client_port: u16,
) {
let (_response, message) =
ControlMessage::subscribe(key, major_version, ttl, event_group_id, client_port);
let _ = self.control_sender.send(message).await;
}
pub async fn reboot_flag(&self) -> Result<protocol::sd::RebootFlag, Error> {
let (response, message) = ControlMessage::query_reboot_flag();
self.control_sender
.send(message)
.await
.map_err(|()| Error::Shutdown)?;
response.recv().await.map_err(|_| Error::Shutdown)?
}
#[cfg(all(test, feature = "client-tokio"))]
pub(crate) async fn force_sd_session_wrapped_for_test(
&self,
wrapped: bool,
) -> Result<(), Error> {
let (response, message) = ControlMessage::force_sd_session_wrapped_for_test(wrapped);
self.control_sender
.send(message)
.await
.map_err(|()| Error::Shutdown)?;
response.recv().await.map_err(|_| Error::Shutdown)?
}
pub async fn send_sd_message(
&self,
target: SocketAddrV4,
sd_header: <MessageDefinitions as PayloadWireFormat>::SdHeader,
) -> Result<(), Error> {
let (response, message) = ControlMessage::send_sd(target, sd_header);
self.control_sender
.send(message)
.await
.map_err(|()| Error::Shutdown)?;
response.recv().await.map_err(|_| Error::Shutdown)?
}
pub async fn add_endpoint(
&self,
key: ServiceEndpointKey,
instance_id: u16,
local_port: u16,
) -> Result<(), Error> {
let (response, message) = ControlMessage::add_endpoint(key, instance_id, local_port);
self.control_sender
.send(message)
.await
.map_err(|()| Error::Shutdown)?;
response.recv().await.map_err(|_| Error::Shutdown)?
}
pub async fn remove_endpoint(&self, key: ServiceEndpointKey) -> Result<(), Error> {
let (response, message) = ControlMessage::remove_endpoint(key);
self.control_sender
.send(message)
.await
.map_err(|()| Error::Shutdown)?;
response.recv().await.map_err(|_| Error::Shutdown)?
}
pub async fn send_to_service(
&self,
key: ServiceEndpointKey,
message: crate::protocol::Message<MessageDefinitions>,
) -> Result<PendingResponse<MessageDefinitions, C>, Error> {
let (send_rx, response_rx, ctrl_msg) = ControlMessage::send_to_service(key, message);
self.control_sender
.send(ctrl_msg)
.await
.map_err(|()| Error::Shutdown)?;
send_rx.recv().await.map_err(|_| Error::Shutdown)??;
Ok(PendingResponse {
receiver: response_rx,
})
}
pub async fn request(
&self,
key: ServiceEndpointKey,
message: crate::protocol::Message<MessageDefinitions>,
) -> Result<MessageDefinitions, Error> {
let (send_rx, response_rx, ctrl_msg) = ControlMessage::send_to_service(key, message);
self.control_sender
.send(ctrl_msg)
.await
.map_err(|()| Error::Shutdown)?;
send_rx.recv().await.map_err(|_| Error::Shutdown)??;
response_rx.recv().await.map_err(|_| Error::Shutdown)?
}
pub fn register_e2e(
&self,
key: E2EKey,
profile: E2EProfile,
) -> Result<(), crate::e2e::E2ERegistryFull> {
self.e2e_registry.register(key, profile)
}
pub fn unregister_e2e(&self, key: &E2EKey) {
self.e2e_registry.unregister(key);
}
pub fn shut_down(self) {
drop(self.control_sender);
info!("Shutting Down SOME/IP client");
}
}
#[cfg(feature = "client-tokio")]
impl<MessageDefinitions, R, I> Client<MessageDefinitions, R, I, TokioChannels>
where
MessageDefinitions: PayloadWireFormat + Clone + core::fmt::Debug + 'static,
R: E2ERegistryHandle,
I: InterfaceHandle,
{
pub fn sd_announcements_loop(
&self,
sd_header: <MessageDefinitions as PayloadWireFormat>::SdHeader,
interval: std::time::Duration,
) -> impl core::future::Future<Output = ()> + Send + 'static
where
<MessageDefinitions as PayloadWireFormat>::SdHeader: Send + 'static,
{
use crate::protocol::sd;
use crate::transport::OneshotRecv;
let weak_sender = self.control_sender.downgrade();
let target = SocketAddrV4::new(sd::MULTICAST_IP, sd::MULTICAST_PORT);
let interval = interval.max(std::time::Duration::from_millis(100));
async move {
let timer = TokioTimer;
let mut count = 0u64;
loop {
timer.sleep(interval).await;
let (flag_rx, flag_msg) =
ControlMessage::<MessageDefinitions, TokioChannels>::query_reboot_flag();
let Some(sender) = weak_sender.upgrade() else {
crate::log::info!("Client shut down, stopping SD announcements");
break;
};
let enqueue_ok = sender.send(flag_msg).await.is_ok();
drop(sender);
if !enqueue_ok {
crate::log::warn!("SD announcement channel closed, stopping");
break;
}
let reboot = match flag_rx.recv().await {
Ok(Ok(flag)) => flag,
Ok(Err(e)) => {
crate::log::warn!(
"SD announcement reboot-flag query returned error ({:?}), skipping tick",
e
);
continue;
}
Err(_) => {
crate::log::warn!("SD announcement reboot-flag query dropped, stopping");
break;
}
};
let mut header = sd_header.clone();
MessageDefinitions::set_reboot_flag(&mut header, reboot);
let (response, message) =
ControlMessage::<MessageDefinitions, TokioChannels>::send_sd(target, header);
let Some(sender) = weak_sender.upgrade() else {
crate::log::info!("Client shut down, stopping SD announcements");
break;
};
let send_ok = sender.send(message).await.is_ok();
drop(sender);
if !send_ok {
crate::log::warn!("SD announcement channel closed, stopping");
break;
}
match response.recv().await {
Ok(Ok(())) => {
count += 1;
if count == 1 {
crate::log::info!("Sent first client SD announcement");
} else {
crate::log::trace!("Sent {count} client SD announcements");
}
}
Ok(Err(e)) => {
crate::log::error!("Failed to send SD announcement: {e:?}");
}
Err(_) => {
crate::log::warn!("SD announcement response dropped, stopping");
break;
}
}
}
}
}
}
#[cfg(all(test, feature = "client-tokio"))]
mod tests {
use super::*;
use crate::protocol::sd::test_support::{TestPayload, empty_sd_header};
use crate::traits::WireFormat;
use std::format;
type TestClient =
Client<TestPayload, Arc<Mutex<E2ERegistry>>, Arc<RwLock<Ipv4Addr>>, TokioChannels>;
#[tokio::test]
async fn test_client_new_and_interface() {
let (client, _updates, run_fut) = TestClient::new(Ipv4Addr::LOCALHOST);
let _run_handle = tokio::spawn(run_fut);
assert_eq!(client.interface(), Ipv4Addr::LOCALHOST);
client.shut_down();
}
#[tokio::test]
async fn test_client_debug() {
let (client, _updates, run_fut) = TestClient::new(Ipv4Addr::LOCALHOST);
let _run_handle = tokio::spawn(run_fut);
let debug_str = format!("{client:?}");
assert!(debug_str.contains("Client"));
assert!(debug_str.contains("127.0.0.1"));
client.shut_down();
}
#[tokio::test]
async fn test_client_update_debug() {
use std::net::SocketAddr;
let sd_header = empty_sd_header();
let someip_header = crate::protocol::Header::new_sd(1, sd_header.required_size());
let discovery_msg = DiscoveryMessage {
source: SocketAddr::new(Ipv4Addr::LOCALHOST.into(), 30490),
someip_header,
sd_header,
};
let update: ClientUpdate<TestPayload> = ClientUpdate::DiscoveryUpdated(discovery_msg);
let debug_str = format!("{update:?}");
assert!(debug_str.contains("DiscoveryUpdated"));
let update: ClientUpdate<TestPayload> =
ClientUpdate::SenderRebooted(SocketAddr::new(Ipv4Addr::LOCALHOST.into(), 30490));
let debug_str = format!("{update:?}");
assert!(debug_str.contains("SenderRebooted"));
let msg = crate::protocol::Message::new_sd(1, &empty_sd_header());
let update: ClientUpdate<TestPayload> = ClientUpdate::Unicast {
message: msg,
e2e_status: None,
source: SocketAddr::new(Ipv4Addr::LOCALHOST.into(), 30640),
};
let debug_str = format!("{update:?}");
assert!(debug_str.contains("Unicast"));
let update: ClientUpdate<TestPayload> = ClientUpdate::Error(Error::ServiceNotFound);
let debug_str = format!("{update:?}");
assert!(debug_str.contains("Error"));
}
#[test]
fn unicast_update_carries_source() {
let src = SocketAddr::new(Ipv4Addr::new(192, 168, 11, 101).into(), 30640);
let msg = crate::protocol::Message::new_sd(1, &empty_sd_header());
let update: ClientUpdate<TestPayload> = ClientUpdate::Unicast {
message: msg,
e2e_status: None,
source: src,
};
match update {
ClientUpdate::Unicast { source, .. } => assert_eq!(source, src),
_ => panic!("expected Unicast"),
}
}
#[tokio::test]
async fn test_subscribe_unknown_service_returns_error() {
let (client, _updates, run_fut) = TestClient::new(Ipv4Addr::LOCALHOST);
let _run_handle = tokio::spawn(run_fut);
let result = client
.subscribe(
ServiceEndpointKey::udp(
0xFFFF,
SocketAddr::V4(SocketAddrV4::new(Ipv4Addr::LOCALHOST, 1)),
),
1,
3,
0x01,
0,
)
.await;
assert!(
matches!(result, Err(Error::ServiceNotFound)),
"expected ServiceNotFound, got {result:?}"
);
client.shut_down();
}
#[tokio::test]
async fn test_subscribe_no_wait_unknown_service_does_not_panic() {
let (client, _updates, run_fut) = TestClient::new(Ipv4Addr::LOCALHOST);
let _run_handle = tokio::spawn(run_fut);
client
.subscribe_no_wait(
ServiceEndpointKey::udp(
0xFFFF,
SocketAddr::V4(SocketAddrV4::new(Ipv4Addr::LOCALHOST, 1)),
),
1,
3,
0x01,
0,
)
.await;
client.shut_down();
}
#[tokio::test]
async fn test_subscribe_no_wait_fire_and_forget_stress() {
let (client, _updates, run_fut) = TestClient::new(Ipv4Addr::LOCALHOST);
let _run_handle = tokio::spawn(run_fut);
for _ in 0..200 {
client
.subscribe_no_wait(
ServiceEndpointKey::udp(
0xFFFF,
SocketAddr::V4(SocketAddrV4::new(Ipv4Addr::LOCALHOST, 1)),
),
1,
3,
0x01,
0,
)
.await;
}
let msg = crate::protocol::Message::new_sd(1, &empty_sd_header());
let result = tokio::time::timeout(
std::time::Duration::from_secs(2),
client.request(
ServiceEndpointKey::udp(
0xFFFF,
SocketAddr::V4(SocketAddrV4::new(Ipv4Addr::LOCALHOST, 1)),
),
msg,
),
)
.await
.expect("inner loop unresponsive after 200 subscribe_no_wait calls");
assert!(
matches!(result, Err(Error::ServiceNotFound)),
"expected ServiceNotFound, got {result:?}"
);
client.shut_down();
}
#[tokio::test]
async fn test_bind_discovery_and_unbind() {
let (client, _updates, run_fut) = TestClient::new(Ipv4Addr::LOCALHOST);
let _run_handle = tokio::spawn(run_fut);
client.bind_discovery().await.unwrap();
client.unbind_discovery().await.unwrap();
client.shut_down();
}
#[tokio::test]
async fn test_set_interface() {
let (client, _updates, run_fut) = TestClient::new(Ipv4Addr::LOCALHOST);
let _run_handle = tokio::spawn(run_fut);
let new_addr = Ipv4Addr::LOCALHOST;
client.set_interface(new_addr).await.unwrap();
assert_eq!(client.interface(), new_addr);
client.shut_down();
}
#[tokio::test]
async fn test_add_endpoint_succeeds() {
let (client, _updates, run_fut) = TestClient::new(Ipv4Addr::LOCALHOST);
let _run_handle = tokio::spawn(run_fut);
let addr = SocketAddrV4::new(Ipv4Addr::new(192, 168, 1, 1), 30000);
client
.add_endpoint(
ServiceEndpointKey::udp(0x1234, SocketAddr::V4(addr)),
0x0001,
0,
)
.await
.unwrap();
client.shut_down();
}
#[tokio::test]
async fn test_send_to_service_unknown_returns_error() {
let (client, _updates, run_fut) = TestClient::new(Ipv4Addr::LOCALHOST);
let _run_handle = tokio::spawn(run_fut);
let msg = crate::protocol::Message::new_sd(1, &empty_sd_header());
let result = client
.send_to_service(
ServiceEndpointKey::udp(
0xFFFF,
SocketAddr::V4(SocketAddrV4::new(Ipv4Addr::LOCALHOST, 1)),
),
msg,
)
.await;
assert!(
matches!(result, Err(Error::ServiceNotFound)),
"expected ServiceNotFound, got {result:?}"
);
client.shut_down();
}
#[tokio::test]
async fn test_remove_endpoint_succeeds() {
let (client, _updates, run_fut) = TestClient::new(Ipv4Addr::LOCALHOST);
let _run_handle = tokio::spawn(run_fut);
let ep_ip = Ipv4Addr::new(192, 168, 1, 1);
let addr = SocketAddrV4::new(ep_ip, 30000);
client
.add_endpoint(
ServiceEndpointKey::udp(0x1234, SocketAddr::V4(addr)),
0x0001,
0,
)
.await
.unwrap();
client
.remove_endpoint(ServiceEndpointKey::udp(0x1234, SocketAddr::V4(addr)))
.await
.unwrap();
client.shut_down();
}
#[test]
fn test_pending_response_debug() {
let (_tx, rx) = TokioChannels::oneshot::<Result<TestPayload, Error>>();
let pending: PendingResponse<TestPayload, TokioChannels> = PendingResponse { receiver: rx };
let s = format!("{pending:?}");
assert!(s.contains("PendingResponse"));
}
#[tokio::test]
async fn test_pending_response_resolves_ok() {
let (tx, rx) = TokioChannels::oneshot::<Result<TestPayload, Error>>();
let pending: PendingResponse<TestPayload, TokioChannels> = PendingResponse { receiver: rx };
let payload = TestPayload {
header: empty_sd_header(),
};
tx.send(Ok(payload.clone())).unwrap();
let result = pending.response().await;
assert_eq!(result.unwrap(), payload);
}
#[tokio::test]
async fn test_pending_response_resolves_err() {
let (tx, rx) = TokioChannels::oneshot::<Result<TestPayload, Error>>();
let pending: PendingResponse<TestPayload, TokioChannels> = PendingResponse { receiver: rx };
tx.send(Err(Error::ServiceNotFound)).unwrap();
let result = pending.response().await;
assert!(
matches!(result, Err(Error::ServiceNotFound)),
"expected ServiceNotFound, got {result:?}"
);
}
#[tokio::test]
async fn test_send_sd_message() {
let (client, _updates, run_fut) = TestClient::new(Ipv4Addr::LOCALHOST);
let _run_handle = tokio::spawn(run_fut);
client.bind_discovery().await.unwrap();
let target = SocketAddrV4::new(Ipv4Addr::LOCALHOST, 30490);
let sd_header = empty_sd_header();
client.send_sd_message(target, sd_header).await.unwrap();
client.shut_down();
}
#[tokio::test]
async fn test_send_to_service_success_returns_pending_response() {
let (client, _updates, run_fut) = TestClient::new(Ipv4Addr::LOCALHOST);
let _run_handle = tokio::spawn(run_fut);
let addr = SocketAddrV4::new(Ipv4Addr::LOCALHOST, 30000);
client
.add_endpoint(
ServiceEndpointKey::udp(0x1234, SocketAddr::V4(addr)),
0x0001,
0,
)
.await
.unwrap();
let msg = crate::protocol::Message::new_sd(1, &empty_sd_header());
let pending = client
.send_to_service(ServiceEndpointKey::udp(0x1234, SocketAddr::V4(addr)), msg)
.await;
assert!(pending.is_ok());
client.shut_down();
}
#[tokio::test]
async fn test_recv_returns_none_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;
assert!(result.is_ok());
assert!(result.unwrap().is_none());
}
#[tokio::test]
async fn test_register_and_unregister_e2e() {
let (client, _updates, run_fut) = TestClient::new(Ipv4Addr::LOCALHOST);
let _run_handle = tokio::spawn(run_fut);
let key = E2EKey {
service_id: 0x1234,
method_or_event_id: 0x0001,
};
let profile = E2EProfile::Profile4(crate::e2e::Profile4Config::new(42, 10));
client
.register_e2e(key, profile)
.expect("E2E registry has capacity for one entry");
client.unregister_e2e(&key);
client.shut_down();
}
#[tokio::test]
async fn test_client_is_clone() {
let (client, _updates, run_fut) = TestClient::new(Ipv4Addr::LOCALHOST);
let _run_handle = tokio::spawn(run_fut);
let client2 = client.clone();
assert_eq!(client.interface(), client2.interface());
client.shut_down();
}
#[tokio::test]
async fn test_client_updates_debug() {
let (_client, updates, run_fut) = TestClient::new(Ipv4Addr::LOCALHOST);
let _run_handle = tokio::spawn(run_fut);
let debug_str = format!("{updates:?}");
assert!(debug_str.contains("ClientUpdates"));
}
#[tokio::test]
async fn test_request_unknown_service_returns_error() {
let (client, _updates, run_fut) = TestClient::new(Ipv4Addr::LOCALHOST);
let _run_handle = tokio::spawn(run_fut);
let msg = crate::protocol::Message::new_sd(1, &empty_sd_header());
let result = client
.request(
ServiceEndpointKey::udp(
0xFFFF,
SocketAddr::V4(SocketAddrV4::new(Ipv4Addr::LOCALHOST, 1)),
),
msg,
)
.await;
assert!(
matches!(result, Err(Error::ServiceNotFound)),
"expected ServiceNotFound, got {result:?}"
);
client.shut_down();
}
#[tokio::test]
async fn test_sd_announcements_loop_does_not_panic() {
let (client, _updates, run_fut) = TestClient::new(Ipv4Addr::LOCALHOST);
let _run_handle = tokio::spawn(run_fut);
client.bind_discovery().await.unwrap();
let sd_header = empty_sd_header();
let handle = tokio::spawn(
client.sd_announcements_loop(sd_header, std::time::Duration::from_millis(100)),
);
tokio::time::sleep(std::time::Duration::from_millis(250)).await;
handle.abort();
let result = handle.await;
let err = result.unwrap_err();
assert!(
err.is_cancelled(),
"task should have been cancelled, not panicked"
);
client.shut_down();
}
#[tokio::test]
async fn test_sd_announcements_loop_without_discovery_bound() {
let (client, _updates, run_fut) = TestClient::new(Ipv4Addr::LOCALHOST);
let _run_handle = tokio::spawn(run_fut);
let sd_header = empty_sd_header();
let handle = tokio::spawn(
client.sd_announcements_loop(sd_header, std::time::Duration::from_millis(100)),
);
tokio::time::sleep(std::time::Duration::from_millis(250)).await;
handle.abort();
let result = handle.await;
let err = result.unwrap_err();
assert!(
err.is_cancelled(),
"task should have been cancelled, not panicked"
);
client.shut_down();
}
#[tokio::test]
async fn test_sd_announcements_loop_abort_stops_task() {
let (client, _updates, run_fut) = TestClient::new(Ipv4Addr::LOCALHOST);
let _run_handle = tokio::spawn(run_fut);
client.bind_discovery().await.unwrap();
let sd_header = empty_sd_header();
let handle = tokio::spawn(
client.sd_announcements_loop(sd_header, std::time::Duration::from_millis(100)),
);
handle.abort();
let result = handle.await;
let err = result.unwrap_err();
assert!(
err.is_cancelled(),
"task should have been cancelled, not panicked"
);
client.shut_down();
}
#[tokio::test]
async fn test_sd_announcements_loop_overrides_caller_reboot_flag() {
let (client, mut updates, run_fut) =
TestClient::new_with_loopback(Ipv4Addr::LOCALHOST, true);
let _run_handle = tokio::spawn(run_fut);
client.bind_discovery().await.unwrap();
let mut sd_header = empty_sd_header();
sd_header.flags =
crate::protocol::sd::Flags::new_sd(crate::protocol::sd::RebootFlag::Continuous);
let handle = tokio::spawn(
client.sd_announcements_loop(sd_header, std::time::Duration::from_millis(100)),
);
let received = tokio::time::timeout(std::time::Duration::from_secs(2), async {
loop {
match updates.recv().await {
Some(ClientUpdate::DiscoveryUpdated(msg)) => return Some(msg),
Some(_) => {}
None => return None,
}
}
})
.await
.expect("timed out waiting for SD announcement")
.expect("update stream closed");
assert_eq!(
received.sd_header.flags.reboot(),
crate::protocol::sd::RebootFlag::RecentlyRebooted,
"announcer should have overridden the caller-supplied Continuous \
flag with the client's tracked RecentlyRebooted state"
);
handle.abort();
let _ = handle.await;
client.shut_down();
}
#[tokio::test]
async fn test_reboot_flag_uses_persisted_wrap_state_when_unbound() {
let (client, _updates, run_fut) = TestClient::new(Ipv4Addr::LOCALHOST);
let _run_handle = tokio::spawn(run_fut);
assert_eq!(
client.reboot_flag().await.expect("reboot_flag"),
crate::protocol::sd::RebootFlag::RecentlyRebooted
);
client
.force_sd_session_wrapped_for_test(true)
.await
.expect("force_sd_session_wrapped_for_test");
assert_eq!(
client.reboot_flag().await.expect("reboot_flag"),
crate::protocol::sd::RebootFlag::Continuous,
"reboot_flag must report Continuous from persisted state while \
discovery is unbound"
);
client.bind_discovery().await.unwrap();
assert_eq!(
client.reboot_flag().await.expect("reboot_flag"),
crate::protocol::sd::RebootFlag::Continuous,
"seeded socket must report Continuous after wrapped rebind"
);
client.shut_down();
}
#[tokio::test]
async fn test_reboot_flag_defaults_to_recently_rebooted() {
let (client, _updates, run_fut) = TestClient::new(Ipv4Addr::LOCALHOST);
let _run_handle = tokio::spawn(run_fut);
assert_eq!(
client.reboot_flag().await.expect("reboot_flag"),
crate::protocol::sd::RebootFlag::RecentlyRebooted
);
client.bind_discovery().await.unwrap();
assert_eq!(
client.reboot_flag().await.expect("reboot_flag"),
crate::protocol::sd::RebootFlag::RecentlyRebooted
);
client.shut_down();
}
#[tokio::test]
async fn reboot_flag_returns_shutdown_error_when_run_loop_dropped() {
let (client, _updates, run_fut) = TestClient::new(Ipv4Addr::LOCALHOST);
drop(run_fut);
let err = client
.reboot_flag()
.await
.expect_err("reboot_flag must return an error after run loop is dropped");
assert!(
matches!(err, Error::Shutdown),
"expected Shutdown, got {err:?}"
);
}
#[tokio::test]
async fn test_sd_announcements_loop_stops_on_shutdown() {
let (client, _updates, run_fut) = TestClient::new(Ipv4Addr::LOCALHOST);
let _run_handle = tokio::spawn(run_fut);
client.bind_discovery().await.unwrap();
let sd_header = empty_sd_header();
let handle = tokio::spawn(
client.sd_announcements_loop(sd_header, std::time::Duration::from_millis(100)),
);
client.shut_down();
let join_result = tokio::time::timeout(std::time::Duration::from_secs(2), handle)
.await
.expect("task should have exited within timeout");
assert!(
join_result.is_ok() || join_result.as_ref().unwrap_err().is_cancelled(),
"task should have exited cleanly, not panicked"
);
}
#[tokio::test]
async fn dropping_run_future_without_spawn_returns_shutdown_error() {
let (client, _updates, run_fut) = TestClient::new(Ipv4Addr::LOCALHOST);
drop(run_fut);
let err = client
.bind_discovery()
.await
.expect_err("must surface a typed error, not Ok or panic");
assert!(
matches!(err, Error::Shutdown),
"expected Error::Shutdown after run-loop drop, got {err:?}",
);
}
#[tokio::test]
async fn cancelling_run_future_closes_control_channel_returns_shutdown_error() {
let (client, _updates, run_fut) = TestClient::new(Ipv4Addr::LOCALHOST);
let handle = tokio::spawn(run_fut);
tokio::time::sleep(std::time::Duration::from_millis(50)).await;
handle.abort();
tokio::time::sleep(std::time::Duration::from_millis(50)).await;
let err = client
.bind_discovery()
.await
.expect_err("must surface a typed error, not Ok or panic");
assert!(
matches!(err, Error::Shutdown),
"expected Error::Shutdown after run-loop cancel, got {err:?}",
);
}
#[ignore = "requires MULTICAST on the loopback interface; dev \
machines where `lo` lacks the MULTICAST flag will not \
deliver loopback multicast and this test will fail. \
Runs in any environment where loopback multicast is \
available (e.g. CI)."]
#[tokio::test]
async fn sd_announcements_loop_cadence_stays_close_to_requested() {
use crate::protocol::sd;
use socket2::{Domain, Protocol, Socket, Type};
let iface = Ipv4Addr::LOCALHOST;
let recv = {
let s = Socket::new(Domain::IPV4, Type::DGRAM, Some(Protocol::UDP)).unwrap();
s.set_reuse_address(true).unwrap();
#[cfg(unix)]
s.set_reuse_port(true).unwrap();
s.bind(&std::net::SocketAddr::from((iface, sd::MULTICAST_PORT)).into())
.unwrap();
s.set_nonblocking(true).unwrap();
let std_s: std::net::UdpSocket = s.into();
let rs = tokio::net::UdpSocket::from_std(std_s).unwrap();
rs.join_multicast_v4(sd::MULTICAST_IP, iface).unwrap();
rs
};
let (client, _updates, run_fut) = TestClient::new_with_loopback(iface, true);
let _run_handle = tokio::spawn(run_fut);
client.bind_discovery().await.unwrap();
let interval = std::time::Duration::from_millis(100);
let loop_handle = tokio::spawn(client.sd_announcements_loop(empty_sd_header(), interval));
let start = std::time::Instant::now();
let mut count = 0u32;
let mut buf = [0u8; 1500];
while start.elapsed() < std::time::Duration::from_millis(550) {
if tokio::time::timeout(
std::time::Duration::from_millis(200),
recv.recv_from(&mut buf),
)
.await
.map(|r| r.is_ok())
.unwrap_or(false)
{
count += 1;
}
}
loop_handle.abort();
client.shut_down();
assert!(
count >= 3,
"expected >= 3 announcements in 550ms at 100ms interval, got {count} — \
cadence may have regressed"
);
}
#[ignore = "requires MULTICAST on the loopback interface; same \
constraint as `sd_announcements_loop_cadence_stays_close_to_requested`. \
Runs in any environment where loopback multicast is \
available (e.g. CI)."]
#[tokio::test]
async fn sd_announcements_loop_first_emit_within_one_interval() {
use crate::protocol::sd;
use socket2::{Domain, Protocol, Socket, Type};
let iface = Ipv4Addr::LOCALHOST;
let recv = {
let s = Socket::new(Domain::IPV4, Type::DGRAM, Some(Protocol::UDP)).unwrap();
s.set_reuse_address(true).unwrap();
#[cfg(unix)]
s.set_reuse_port(true).unwrap();
s.bind(&std::net::SocketAddr::from((iface, sd::MULTICAST_PORT)).into())
.unwrap();
s.set_nonblocking(true).unwrap();
let std_s: std::net::UdpSocket = s.into();
let rs = tokio::net::UdpSocket::from_std(std_s).unwrap();
rs.join_multicast_v4(sd::MULTICAST_IP, iface).unwrap();
rs
};
let (client, _updates, run_fut) = TestClient::new_with_loopback(iface, true);
let _run_handle = tokio::spawn(run_fut);
client.bind_discovery().await.unwrap();
let interval = std::time::Duration::from_millis(100);
let start = std::time::Instant::now();
let loop_handle = tokio::spawn(client.sd_announcements_loop(empty_sd_header(), interval));
let mut buf = [0u8; 1500];
let first = tokio::time::timeout(
std::time::Duration::from_millis(500),
recv.recv_from(&mut buf),
)
.await
.expect("first SD announcement did not arrive within 500ms")
.expect("recv_from errored");
let first_emit_elapsed = start.elapsed();
let _ = first;
loop_handle.abort();
client.shut_down();
assert!(
first_emit_elapsed < std::time::Duration::from_millis(250),
"first announcement took {first_emit_elapsed:?}, expected < 250ms at 100ms interval — \
likely double-sleep regression"
);
}
#[test]
fn client_new_run_future_is_send_static() {
let (_client, _updates, run_fut) = TestClient::new(Ipv4Addr::LOCALHOST);
let handle = std::thread::spawn(move || drop(run_fut));
handle.join().unwrap();
}
#[tokio::test]
async fn client_new_with_spawner_routes_socket_spawns_through_it() {
use core::sync::atomic::{AtomicUsize, Ordering};
use std::sync::Arc;
#[derive(Clone)]
struct CountingSpawner {
count: Arc<AtomicUsize>,
}
impl Spawner for CountingSpawner {
fn spawn(&self, future: impl core::future::Future<Output = ()> + Send + 'static) {
self.count.fetch_add(1, Ordering::SeqCst);
let _run_handle = tokio::spawn(future);
}
}
let count = Arc::new(AtomicUsize::new(0));
let spawner = CountingSpawner {
count: Arc::clone(&count),
};
let (client, _updates, run_fut) =
TestClient::new_with_spawner_and_loopback(Ipv4Addr::LOCALHOST, false, spawner);
let _run_handle = tokio::spawn(run_fut);
client
.bind_discovery()
.await
.expect("bind_discovery must succeed");
client
.bind_discovery()
.await
.expect("second bind_discovery is idempotent");
assert_eq!(
count.load(Ordering::SeqCst),
2,
"expected two spawns (multicast + unicast SD socket loops), \
got {}",
count.load(Ordering::SeqCst)
);
client.shut_down();
}
const TOKIO_CLIENT_RUN_FUTURE_BUDGET: usize = 132736; const TOKIO_CLIENT_SOCKET_LOOP_BUDGET: usize = 8768;
#[tokio::test]
async fn future_size_witness_tokio_client() {
use core::sync::atomic::{AtomicUsize, Ordering};
use std::sync::Arc;
#[derive(Clone)]
struct SizeRecordingSpawner {
max_spawned: Arc<AtomicUsize>,
}
impl Spawner for SizeRecordingSpawner {
fn spawn(&self, future: impl core::future::Future<Output = ()> + Send + 'static) {
self.max_spawned
.fetch_max(core::mem::size_of_val(&future), Ordering::SeqCst);
let _run_handle = tokio::spawn(future);
}
}
let max_spawned = Arc::new(AtomicUsize::new(0));
let spawner = SizeRecordingSpawner {
max_spawned: Arc::clone(&max_spawned),
};
let (client, _updates, run_fut) =
TestClient::new_with_spawner_and_loopback(Ipv4Addr::LOCALHOST, false, spawner);
let run_size = core::mem::size_of_val(&run_fut);
let _run_handle = tokio::spawn(run_fut);
client.bind_discovery().await.expect("bind_discovery");
let loop_size = max_spawned.load(Ordering::SeqCst);
std::println!("FUTURE_SIZE tokio_client_run_future {run_size}");
std::println!("FUTURE_SIZE tokio_client_socket_loop {loop_size}");
assert!(loop_size > 0, "spawner never received the socket loop");
assert!(
run_size <= TOKIO_CLIENT_RUN_FUTURE_BUDGET,
"Inner::run_future grew: {run_size} B > budget {TOKIO_CLIENT_RUN_FUTURE_BUDGET} B"
);
assert!(
loop_size <= TOKIO_CLIENT_SOCKET_LOOP_BUDGET,
"socket loop future grew: {loop_size} B > budget {TOKIO_CLIENT_SOCKET_LOOP_BUDGET} B"
);
client.shut_down();
}
}