use core::cell::RefCell;
use core::mem::MaybeUninit;
use core::net::{Ipv4Addr, SocketAddrV4};
use core::sync::atomic::{AtomicBool, AtomicU16, AtomicU32, AtomicUsize, Ordering};
use core::time::Duration;
use embassy_executor::raw::Executor;
use embassy_sync::blocking_mutex::Mutex as BlockingMutex;
use embassy_sync::blocking_mutex::raw::CriticalSectionRawMutex;
use crate::bare_metal_tasks::{SomeipRun, publish_notification, run_someip};
use crate::e2e::{E2ERegistry, Profile5Config};
use crate::protocol::MessageId;
use crate::protocol::sd::RebootFlag;
use crate::sd_codec::{
OfferServiceRequest, SubscribeEventgroupRequest, build_multi_stop_offer_service_datagram,
next_sd_session,
};
use crate::server::{
EventPublisher, SdStateManager, Server, ServerConfig, ServerStorage, StaticSubscriptionHandle,
StaticSubscriptionStorage, SubscriptionManager,
};
use crate::transport::E2ERegistryHandle;
use crate::{E2EKey, E2EProfile, StaticE2EHandle, StaticE2EStorage};
use super::mailbox::RxMailbox;
use super::transport::{CallbackFactory, CallbackSocket, CallbackTimer, NowMsFn, Platform, SendFn};
pub const RX_SLOTS: usize = crate::from_env_or(option_env!("SIMPLE_SOMEIP_RX_SLOTS"), 2);
pub const RX_CAP: usize = 1500;
const MAX_OFFERS: usize = crate::from_env_or(option_env!("SIMPLE_SOMEIP_MAX_OFFERS"), 4);
const MAX_SUBS: usize = crate::from_env_or(option_env!("SIMPLE_SOMEIP_MAX_SUBS"), 1);
const SD_SCRATCH_CAP: usize = 512;
const SUB_SCRATCH_CAP: usize = 128;
const API_SCRATCH_CAP: usize = 512;
const _: () = assert!(
API_SCRATCH_CAP >= SD_SCRATCH_CAP,
"publish_scratch must hold deinit's StopOffer datagram"
);
const DEFAULT_SD_PORT: u16 = 30490;
const DEFAULT_SD_MCAST: u32 = 0xEFFF_00FF;
type Mailbox = RxMailbox<RX_SLOTS, RX_CAP>;
type Sock = CallbackSocket<'static, RX_SLOTS, RX_CAP>;
type Factory = CallbackFactory<'static, RX_SLOTS, RX_CAP>;
type Publisher = EventPublisher<StaticE2EHandle, StaticSubscriptionHandle, &'static Sock, Sock>;
type RtServer = Server<
Factory,
CallbackTimer,
StaticE2EHandle,
StaticSubscriptionHandle,
&'static Sock,
&'static SdStateManager,
&'static Publisher,
>;
pub type DispatchFn = fn(
ctx: usize,
source: core::net::SocketAddrV4,
service_id: u16,
method_id: u16,
payload: &[u8],
e2e_status: u8,
response_out: &mut [u8],
) -> i32;
pub type BindFn = extern "C" fn(port: u16, is_sd: bool, mcast: u32) -> i32;
#[repr(C)]
#[derive(Clone, Copy)]
pub struct OfferEntry {
pub service_id: u16,
pub instance_id: u16,
pub event_group_id: u16,
pub unicast_port: u16,
pub major_version: u8,
pub ttl_seconds: u32,
}
#[repr(C)]
#[derive(Clone, Copy)]
pub struct SubEntry {
pub service_id: u16,
pub instance_id: u16,
pub event_group_id: u16,
pub method_id: u16,
pub local_rx_port: u16,
pub major_version: u8,
pub e2e_enabled: bool,
pub e2e_data_id: u16,
pub e2e_data_length: u16,
pub e2e_max_delta: u16,
}
pub struct RuntimeBuffers {
pub unicast: [u8; RX_CAP],
pub sd: [u8; RX_CAP],
pub rx: [u8; RX_CAP],
pub tx_scratch: [u8; RX_CAP],
pub offer_scratch: [u8; SD_SCRATCH_CAP],
pub sub_scratch: [u8; SUB_SCRATCH_CAP],
pub publish_scratch: [u8; API_SCRATCH_CAP],
}
impl Default for RuntimeBuffers {
fn default() -> Self {
Self::new()
}
}
impl RuntimeBuffers {
#[must_use]
pub const fn new() -> Self {
Self {
unicast: [0; RX_CAP],
sd: [0; RX_CAP],
rx: [0; RX_CAP],
tx_scratch: [0; RX_CAP],
offer_scratch: [0; SD_SCRATCH_CAP],
sub_scratch: [0; SUB_SCRATCH_CAP],
publish_scratch: [0; API_SCRATCH_CAP],
}
}
}
pub struct RuntimeConfig {
pub interface: u32,
pub sd_port: u16,
pub sd_mcast: u32,
pub multicast_loopback: bool,
pub send: SendFn,
pub now_ms: NowMsFn,
pub dispatch: DispatchFn,
pub dispatch_ctx: usize,
pub bind: BindFn,
pub offers: &'static [OfferEntry],
pub subscriptions: &'static [SubEntry],
pub mailbox: &'static Mailbox,
pub buffers: &'static mut RuntimeBuffers,
}
#[allow(clippy::declare_interior_mutable_const)]
const E2E_INIT: StaticE2EStorage =
BlockingMutex::<CriticalSectionRawMutex, RefCell<E2ERegistry>>::new(RefCell::new(
E2ERegistry::new(),
));
static E2E_STORAGE: StaticE2EStorage = E2E_INIT;
static SUBS_STORAGE: StaticSubscriptionStorage = BlockingMutex::<
CriticalSectionRawMutex,
RefCell<SubscriptionManager>,
>::new(RefCell::new(SubscriptionManager::new()));
static SD_STATE: SdStateManager = SdStateManager::new();
static SERVER_STARTED: AtomicBool = AtomicBool::new(false);
static SD_PORT: AtomicU16 = AtomicU16::new(DEFAULT_SD_PORT);
static SD_MCAST: AtomicU32 = AtomicU32::new(DEFAULT_SD_MCAST);
static IFACE: AtomicU32 = AtomicU32::new(0);
static DISPATCH: AtomicUsize = AtomicUsize::new(0); static DISPATCH_CTX: AtomicUsize = AtomicUsize::new(0); static OFFER_SESSION: AtomicU16 = AtomicU16::new(1);
static SUB_SESSION: AtomicU16 = AtomicU16::new(1);
static STOP_SESSION: AtomicU16 = AtomicU16::new(1);
static PUBLISH_SESSION: AtomicU16 = AtomicU16::new(1);
static OFFERS_LEN: AtomicUsize = AtomicUsize::new(0);
static SUBS_LEN: AtomicUsize = AtomicUsize::new(0);
static mut OFFERS: [OfferEntry; MAX_OFFERS] = [OfferEntry {
service_id: 0,
instance_id: 0,
event_group_id: 0,
unicast_port: 0,
major_version: 0,
ttl_seconds: 0,
}; MAX_OFFERS];
static mut SUBS: [SubEntry; MAX_SUBS] = [SubEntry {
service_id: 0,
instance_id: 0,
event_group_id: 0,
method_id: 0,
local_rx_port: 0,
major_version: 0,
e2e_enabled: false,
e2e_data_id: 0,
e2e_data_length: 0,
e2e_max_delta: 0,
}; MAX_SUBS];
static mut UNICAST_SOCK: MaybeUninit<Sock> = MaybeUninit::uninit();
static mut SD_SOCK: MaybeUninit<Sock> = MaybeUninit::uninit();
static mut PUBLISHER: MaybeUninit<Publisher> = MaybeUninit::uninit();
static mut SERVER: MaybeUninit<RtServer> = MaybeUninit::uninit();
static mut BUFS: *mut RuntimeBuffers = core::ptr::null_mut();
static mut MAILBOX: Option<&'static Mailbox> = None;
static SEND_FN: AtomicUsize = AtomicUsize::new(0); static NOW_FN: AtomicUsize = AtomicUsize::new(0);
static mut EXECUTOR: MaybeUninit<Executor> = MaybeUninit::uninit();
static EXECUTOR_INIT: AtomicBool = AtomicBool::new(false);
static INIT_CLAIMED: AtomicBool = AtomicBool::new(false);
static RUN_READY: AtomicBool = AtomicBool::new(false);
#[unsafe(no_mangle)]
pub extern "Rust" fn __pender(_context: *mut ()) {}
fn offers() -> &'static [OfferEntry] {
let len = OFFERS_LEN.load(Ordering::Acquire).min(MAX_OFFERS);
unsafe { core::slice::from_raw_parts(core::ptr::addr_of!(OFFERS).cast::<OfferEntry>(), len) }
}
fn subs() -> &'static [SubEntry] {
let len = SUBS_LEN.load(Ordering::Acquire).min(MAX_SUBS);
unsafe { core::slice::from_raw_parts(core::ptr::addr_of!(SUBS).cast::<SubEntry>(), len) }
}
fn find_offer(
service_id: u16,
instance_id: u16,
event_group_id: u16,
) -> Option<&'static OfferEntry> {
offers().iter().find(|o| {
o.service_id == service_id
&& o.instance_id == instance_id
&& o.event_group_id == event_group_id
})
}
fn node_sd_ttl_secs() -> u64 {
offers()
.first()
.map(|o| u64::from(o.ttl_seconds))
.filter(|&t| t != 0)
.unwrap_or(3)
}
fn platform() -> Platform<'static, RX_SLOTS, RX_CAP> {
let mailbox = unsafe { (*core::ptr::addr_of!(MAILBOX)).expect("mailbox set in init") };
let send: SendFn =
unsafe { core::mem::transmute::<usize, SendFn>(SEND_FN.load(Ordering::Acquire)) };
let now: NowMsFn =
unsafe { core::mem::transmute::<usize, NowMsFn>(NOW_FN.load(Ordering::Acquire)) };
Platform {
send,
now_ms: now,
mailbox,
interface: IFACE.load(Ordering::Acquire),
}
}
fn dispatch(
_ctx: usize,
source: core::net::SocketAddrV4,
service_id: u16,
method_id: u16,
payload: &[u8],
e2e_status: u8,
response_out: &mut [u8],
) -> i32 {
let raw = DISPATCH.load(Ordering::Acquire);
if raw == 0 {
return -1;
}
let f: DispatchFn = unsafe { core::mem::transmute::<usize, DispatchFn>(raw) };
f(
DISPATCH_CTX.load(Ordering::Acquire),
source,
service_id,
method_id,
payload,
e2e_status,
response_out,
)
}
fn offer_requests(out: &mut [OfferServiceRequest; MAX_OFFERS]) -> usize {
let iface = Ipv4Addr::from(IFACE.load(Ordering::Acquire).to_be_bytes());
let o = offers();
for (dst, src) in out.iter_mut().zip(o.iter()) {
*dst = OfferServiceRequest {
service_id: src.service_id,
instance_id: src.instance_id,
major_version: src.major_version,
minor_version: 0,
ttl: src.ttl_seconds,
local_ip: iface,
unicast_port: src.unicast_port,
};
}
o.len()
}
const DEFAULT_OFFER_REQUEST: OfferServiceRequest = OfferServiceRequest {
service_id: 0,
instance_id: 0,
major_version: 0,
minor_version: 0,
ttl: 0,
local_ip: Ipv4Addr::UNSPECIFIED,
unicast_port: 0,
};
fn build_server() -> bool {
let o = offers();
let Some(primary) = o.first() else {
return false;
};
let plat = platform();
let factory = Factory::new(plat);
let unicast_port = primary.unicast_port;
let sd_port = SD_PORT.load(Ordering::Acquire);
unsafe {
(*core::ptr::addr_of_mut!(UNICAST_SOCK)).write(factory.socket(unicast_port));
(*core::ptr::addr_of_mut!(SD_SOCK)).write(factory.socket(sd_port));
let unicast_ref: &'static Sock = (*core::ptr::addr_of!(UNICAST_SOCK)).assume_init_ref();
let e2e = StaticE2EHandle::new(&E2E_STORAGE);
let subs_handle = StaticSubscriptionHandle::new(&SUBS_STORAGE);
(*core::ptr::addr_of_mut!(PUBLISHER)).write(EventPublisher::new(
subs_handle,
unicast_ref,
e2e,
));
}
let iface = Ipv4Addr::from(IFACE.load(Ordering::Acquire).to_be_bytes());
let mut config = ServerConfig::new(primary.service_id, primary.instance_id)
.with_interface(iface)
.with_local_port(primary.unicast_port)
.with_major_version(primary.major_version)
.with_ttl(Duration::from_secs(u64::from(primary.ttl_seconds)))
.with_event_group(primary.event_group_id)
.with_announce(false);
for entry in o {
config = config.with_accepted_offer(
entry.service_id,
entry.instance_id,
entry.major_version,
entry.event_group_id,
);
}
let storage = unsafe {
ServerStorage {
factory,
timer: CallbackTimer::new(platform().now_ms),
e2e_registry: StaticE2EHandle::new(&E2E_STORAGE),
subscriptions: StaticSubscriptionHandle::new(&SUBS_STORAGE),
unicast_socket: (*core::ptr::addr_of!(UNICAST_SOCK)).assume_init_ref(),
sd_socket: (*core::ptr::addr_of!(SD_SOCK)).assume_init_ref(),
sd_state: &SD_STATE,
publisher: (*core::ptr::addr_of!(PUBLISHER)).assume_init_ref(),
started: &SERVER_STARTED,
non_sd_observer: Some((dispatch as crate::server::NonSdRequestCallback, 0)),
}
};
match Server::new_with_handles(storage, config) {
Ok(s) => {
unsafe { (*core::ptr::addr_of_mut!(SERVER)).write(s) };
true
}
Err(_) => false,
}
}
const CLIENT_SUB_OFFSET_SECS: u64 = 1;
#[embassy_executor::task]
async fn someip_task() {
let server: &'static RtServer = unsafe { (*core::ptr::addr_of!(SERVER)).assume_init_ref() };
let unicast: &'static mut [u8; RX_CAP] =
unsafe { &mut *core::ptr::addr_of_mut!((*BUFS).unicast) };
let sd: &'static mut [u8; RX_CAP] = unsafe { &mut *core::ptr::addr_of_mut!((*BUFS).sd) };
let tx_scratch: &'static mut [u8; RX_CAP] =
unsafe { &mut *core::ptr::addr_of_mut!((*BUFS).tx_scratch) };
let offer_scratch: &'static mut [u8; SD_SCRATCH_CAP] =
unsafe { &mut *core::ptr::addr_of_mut!((*BUFS).offer_scratch) };
let sub_scratch: &'static mut [u8; SUB_SCRATCH_CAP] =
unsafe { &mut *core::ptr::addr_of_mut!((*BUFS).sub_scratch) };
let rx: &'static mut [u8; RX_CAP] = unsafe { &mut *core::ptr::addr_of_mut!((*BUFS).rx) };
let recv = server.run_with_buffers(unicast, sd, tx_scratch, &mut []);
let mut reqs = [DEFAULT_OFFER_REQUEST; MAX_OFFERS];
let n_offers = offer_requests(&mut reqs);
let plat = platform();
let sd_socket: &'static Sock = unsafe { (*core::ptr::addr_of!(SD_SOCK)).assume_init_ref() };
let sd_mcast = SocketAddrV4::new(
Ipv4Addr::from(SD_MCAST.load(Ordering::Acquire).to_be_bytes()),
SD_PORT.load(Ordering::Acquire),
);
let timer = CallbackTimer::new(plat.now_ms);
let e2e = StaticE2EHandle::new(&E2E_STORAGE);
let sub = subs().first().copied();
let factory = Factory::new(plat);
let rx_socket = factory.socket(sub.map_or(0, |s| s.local_rx_port));
let subscribe = sub.map(|s| SubscribeEventgroupRequest {
service_id: s.service_id,
instance_id: s.instance_id,
major_version: s.major_version,
event_group_id: s.event_group_id,
ttl: 0x00FF_FFFF,
local_ip: Ipv4Addr::from(IFACE.load(Ordering::Acquire).to_be_bytes()),
local_rx_port: s.local_rx_port,
});
let sub_e2e_enabled = sub.is_some_and(|s| s.e2e_enabled);
let period = Duration::from_secs(node_sd_ttl_secs());
run_someip(
recv,
SomeipRun {
sd_socket,
rx_socket: &rx_socket,
timer: &timer,
e2e: &e2e,
sd_multicast: sd_mcast,
period,
offers: &reqs[..n_offers],
offer_session: &OFFER_SESSION,
offer_scratch,
subscribe,
sub_session: &SUB_SESSION,
sub_scratch,
sub_reboot: RebootFlag::RecentlyRebooted,
sub_offset: Duration::from_secs(CLIENT_SUB_OFFSET_SECS),
sub_e2e_enabled,
rx_buf: rx,
dispatch,
dispatch_ctx: 0,
},
)
.await;
}
#[allow(clippy::missing_panics_doc)]
#[allow(clippy::needless_pass_by_value)]
pub fn init(config: RuntimeConfig) -> i32 {
if INIT_CLAIMED
.compare_exchange(false, true, Ordering::AcqRel, Ordering::Acquire)
.is_err()
{
return -4;
}
IFACE.store(config.interface, Ordering::Release);
SD_PORT.store(
if config.sd_port == 0 {
DEFAULT_SD_PORT
} else {
config.sd_port
},
Ordering::Release,
);
SD_MCAST.store(
if config.sd_mcast == 0 {
DEFAULT_SD_MCAST
} else {
config.sd_mcast
},
Ordering::Release,
);
SEND_FN.store(config.send as usize, Ordering::Release);
NOW_FN.store(config.now_ms as usize, Ordering::Release);
DISPATCH.store(config.dispatch as usize, Ordering::Release);
DISPATCH_CTX.store(config.dispatch_ctx, Ordering::Release);
unsafe {
*core::ptr::addr_of_mut!(MAILBOX) = Some(config.mailbox);
*core::ptr::addr_of_mut!(BUFS) = core::ptr::from_mut(config.buffers);
}
let n_offers = config.offers.len().min(MAX_OFFERS);
for (i, e) in config.offers.iter().take(n_offers).enumerate() {
unsafe { (*core::ptr::addr_of_mut!(OFFERS))[i] = *e };
}
OFFERS_LEN.store(n_offers, Ordering::Release);
let n_subs = config.subscriptions.len().min(MAX_SUBS);
for (i, e) in config.subscriptions.iter().take(n_subs).enumerate() {
unsafe { (*core::ptr::addr_of_mut!(SUBS))[i] = *e };
}
SUBS_LEN.store(n_subs, Ordering::Release);
let e2e = StaticE2EHandle::new(&E2E_STORAGE);
for s in subs() {
if !s.e2e_enabled {
continue;
}
let p5 = E2EProfile::Profile5WithHeader(Profile5Config::new(
s.e2e_data_id,
s.e2e_data_length,
#[allow(clippy::cast_possible_truncation)]
{
s.e2e_max_delta as u8
},
));
let key = E2EKey::from_message_id(MessageId::new_from_service_and_method(
s.service_id,
s.method_id,
));
let _ = e2e.register(key, p5);
}
let bind = config.bind;
bind(
SD_PORT.load(Ordering::Acquire),
true,
SD_MCAST.load(Ordering::Acquire),
);
let mut bound: [u16; MAX_OFFERS + MAX_SUBS] = [0; MAX_OFFERS + MAX_SUBS];
let mut nb = 0usize;
let mut bind_once = |port: u16| {
if port != 0 && !bound[..nb].contains(&port) {
bind(port, false, 0);
bound[nb] = port;
nb += 1;
}
};
for o in offers() {
bind_once(o.unicast_port);
}
for s in subs() {
bind_once(s.local_rx_port);
}
if !build_server() {
INIT_CLAIMED.store(false, Ordering::Release);
return -1;
}
let spawner = unsafe {
let slot = &mut *core::ptr::addr_of_mut!(EXECUTOR);
slot.write(Executor::new(core::ptr::null_mut()));
slot.assume_init_ref().spawner()
};
if spawner.spawn(someip_task()).is_err() {
INIT_CLAIMED.store(false, Ordering::Release);
return -2;
}
EXECUTOR_INIT.store(true, Ordering::Release);
RUN_READY.store(true, Ordering::Release);
0
}
pub fn poll() {
if !EXECUTOR_INIT.load(Ordering::Acquire) {
return;
}
unsafe { (*core::ptr::addr_of!(EXECUTOR)).assume_init_ref().poll() };
}
fn next_publish_session() -> u16 {
loop {
let s = PUBLISH_SESSION.fetch_add(1, Ordering::Relaxed);
if s != 0 {
return s;
}
}
}
pub unsafe fn publish(
service_id: u16,
instance_id: u16,
event_group_id: u16,
method_id: u16,
payload: *const u8,
len: usize,
) -> i32 {
if !RUN_READY.load(Ordering::Acquire) {
return -1;
}
let Some(offer) = find_offer(service_id, instance_id, event_group_id) else {
return -3;
};
let src_port = offer.unicast_port;
if len > API_SCRATCH_CAP - 16 {
return -2;
}
if len > 0 && payload.is_null() {
return -1;
}
let payload_slice = if len == 0 {
&[][..]
} else {
unsafe { core::slice::from_raw_parts(payload, len) }
};
let subs_handle = StaticSubscriptionHandle::new(&SUBS_STORAGE);
let plat = platform();
let scratch = unsafe { &mut *core::ptr::addr_of_mut!((*BUFS).publish_scratch) };
publish_notification(
&subs_handle,
service_id,
instance_id,
event_group_id,
method_id,
next_publish_session(),
payload_slice,
scratch,
|datagram, dst| {
let dst_addr = u32::from_be_bytes(dst.ip().octets());
(plat.send)(
src_port,
datagram.as_ptr(),
datagram.len(),
dst_addr,
dst.port(),
);
},
)
}
pub fn deinit() {
if !RUN_READY.swap(false, Ordering::AcqRel) {
return;
}
let mut reqs = [DEFAULT_OFFER_REQUEST; MAX_OFFERS];
let n = offer_requests(&mut reqs);
if n == 0 {
return;
}
let plat = platform();
let scratch = unsafe { &mut *core::ptr::addr_of_mut!((*BUFS).publish_scratch) };
let session = next_sd_session(&STOP_SESSION);
if let Ok(total) =
build_multi_stop_offer_service_datagram::<MAX_OFFERS>(scratch, &reqs[..n], session)
{
let dst = SD_MCAST.load(Ordering::Acquire);
let sd_port = SD_PORT.load(Ordering::Acquire);
(plat.send)(sd_port, scratch.as_ptr(), total, dst, sd_port);
}
}
pub unsafe fn on_rx(local_port: u16, src_addr: u32, src_port: u16, buf: *const u8, len: usize) {
if buf.is_null() || len == 0 {
return;
}
if let Some(mailbox) = unsafe { *core::ptr::addr_of!(MAILBOX) } {
unsafe {
let _ = mailbox.push(local_port, src_addr, src_port, buf, len);
}
}
}