use core::future::Future;
use core::net::SocketAddrV4;
use core::pin::pin;
use core::sync::atomic::AtomicU16;
use core::task::{Context, Poll, RawWaker, RawWakerVTable, Waker};
use core::time::Duration;
use crate::E2ECheckStatus;
use crate::protocol::sd::RebootFlag;
use crate::sd_codec::{
OfferServiceRequest, SubscribeEventgroupRequest, build_multi_offer_service_datagram,
build_notification_datagram, build_subscribe_eventgroup_datagram, check_parsed_e2e,
e2e_status_code, next_sd_session, parse_someip_datagram,
};
use crate::server::{ServerConfig, Subscriber, SubscriptionHandle};
use crate::transport::{E2ERegistryHandle, Timer, TransportSocket};
use heapless::Vec as HeaplessVec;
pub type RequestDispatchFn = fn(
ctx: usize,
source: SocketAddrV4,
service_id: u16,
method_id: u16,
payload: &[u8],
e2e_status: u8,
response_out: &mut [u8],
) -> i32;
pub async fn announce_offers_future<'a, S, Tm, const N: usize>(
sd_socket: &'a S,
timer: &'a Tm,
offers: &'a [OfferServiceRequest],
sd_multicast: SocketAddrV4,
session: &'a AtomicU16,
period: Duration,
scratch: &'a mut [u8],
) where
S: TransportSocket,
Tm: Timer,
{
loop {
let s = next_sd_session(session);
if let Ok(len) = build_multi_offer_service_datagram::<N>(scratch, offers, s) {
let _ = sd_socket.send_to(&scratch[..len], sd_multicast).await;
}
timer.sleep(period).await;
}
}
#[allow(clippy::too_many_arguments)]
pub async fn subscribe_announce_future<'a, S, Tm>(
sd_socket: &'a S,
timer: &'a Tm,
request: SubscribeEventgroupRequest,
sd_multicast: SocketAddrV4,
session: &'a AtomicU16,
reboot: RebootFlag,
period: Duration,
scratch: &'a mut [u8],
) where
S: TransportSocket,
Tm: Timer,
{
loop {
let s = next_sd_session(session);
if let Ok(len) = build_subscribe_eventgroup_datagram(scratch, &request, s, reboot) {
let _ = sd_socket.send_to(&scratch[..len], sd_multicast).await;
}
timer.sleep(period).await;
}
}
pub async fn event_rx_dispatch_future<'a, S, R>(
rx_socket: &'a S,
e2e: &'a R,
e2e_enabled: bool,
dispatch: RequestDispatchFn,
ctx: usize,
buf: &'a mut [u8],
) where
S: TransportSocket,
R: E2ERegistryHandle,
{
loop {
let (n, source) = match rx_socket.recv_from(&mut *buf).await {
Ok(d) => (d.bytes_received, d.source),
Err(_) => continue,
};
let Some(parsed) = parse_someip_datagram(&buf[..n]) else {
continue;
};
let (status, body) = if e2e_enabled {
check_parsed_e2e(e2e, core::net::IpAddr::V4(*source.ip()), &parsed)
} else {
(E2ECheckStatus::Unchecked, parsed.payload)
};
let _ = dispatch(
ctx,
source,
parsed.service_id,
parsed.method_id,
body,
e2e_status_code(status),
&mut [],
);
}
}
const RUN_OFFER_CAP: usize = crate::from_env_or(option_env!("SIMPLE_SOMEIP_MAX_OFFERS"), 4);
pub struct SomeipRun<'a, S, Tm, R>
where
S: TransportSocket,
Tm: Timer,
R: E2ERegistryHandle,
{
pub sd_socket: &'a S,
pub rx_socket: &'a S,
pub timer: &'a Tm,
pub e2e: &'a R,
pub sd_multicast: SocketAddrV4,
pub period: Duration,
pub offers: &'a [OfferServiceRequest],
pub offer_session: &'a AtomicU16,
pub offer_scratch: &'a mut [u8],
pub subscribe: Option<SubscribeEventgroupRequest>,
pub sub_session: &'a AtomicU16,
pub sub_scratch: &'a mut [u8],
pub sub_reboot: RebootFlag,
pub sub_offset: Duration,
pub sub_e2e_enabled: bool,
pub rx_buf: &'a mut [u8],
pub dispatch: RequestDispatchFn,
pub dispatch_ctx: usize,
}
pub async fn run_someip<S, Tm, R, RecvFut>(recv: RecvFut, cfg: SomeipRun<'_, S, Tm, R>)
where
S: TransportSocket,
Tm: Timer,
R: E2ERegistryHandle,
RecvFut: Future,
{
let SomeipRun {
sd_socket,
rx_socket,
timer,
e2e,
sd_multicast,
period,
offers,
offer_session,
offer_scratch,
subscribe,
sub_session,
sub_scratch,
sub_reboot,
sub_offset,
sub_e2e_enabled,
rx_buf,
dispatch,
dispatch_ctx,
} = cfg;
let announce = announce_offers_future::<_, _, RUN_OFFER_CAP>(
sd_socket,
timer,
offers,
sd_multicast,
offer_session,
period,
offer_scratch,
);
let subscribe_fut = async move {
if let Some(request) = subscribe {
timer.sleep(sub_offset).await;
subscribe_announce_future(
sd_socket,
timer,
request,
sd_multicast,
sub_session,
sub_reboot,
period,
sub_scratch,
)
.await;
} else {
core::future::pending::<()>().await;
}
};
let rx_fut = async move {
if subscribe.is_some() {
event_rx_dispatch_future(
rx_socket,
e2e,
sub_e2e_enabled,
dispatch,
dispatch_ctx,
rx_buf,
)
.await;
} else {
core::future::pending::<()>().await;
}
};
futures_util::join!(recv, announce, subscribe_fut, rx_fut);
}
const NOOP_RAW_WAKER: RawWaker = {
const VTABLE: RawWakerVTable = RawWakerVTable::new(|_| NOOP_RAW_WAKER, |_| {}, |_| {}, |_| {});
RawWaker::new(core::ptr::null(), &VTABLE)
};
#[allow(clippy::too_many_arguments)]
pub fn publish_notification<Sub, FSend>(
subscriptions: &Sub,
service_id: u16,
instance_id: u16,
event_group_id: u16,
method_id: u16,
session: u16,
payload: &[u8],
scratch: &mut [u8],
mut send: FSend,
) -> i32
where
Sub: SubscriptionHandle,
FSend: FnMut(&[u8], SocketAddrV4),
{
let Ok(total) = build_notification_datagram(scratch, service_id, method_id, session, payload)
else {
return -2;
};
let datagram = &scratch[..total];
let mut targets: HeaplessVec<SocketAddrV4, { ServerConfig::SUBSCRIBERS_PER_GROUP_CAP }> =
HeaplessVec::new();
let visited = {
let mut visit = |sub: &Subscriber| {
let _ = targets.push(sub.address);
};
let fut =
subscriptions.for_each_subscriber(service_id, instance_id, event_group_id, &mut visit);
let mut fut = pin!(fut);
let waker = unsafe { Waker::from_raw(NOOP_RAW_WAKER) };
let mut cx = Context::from_waker(&waker);
match fut.as_mut().poll(&mut cx) {
Poll::Ready(n) => n,
Poll::Pending => 0,
}
};
for target in &targets {
send(datagram, *target);
}
#[allow(clippy::cast_possible_truncation, clippy::cast_possible_wrap)]
{
visited as i32
}
}