use core::time::Duration;
use kinavis_kernel::time::{Instant, Utc};
use crate::error::Nmea2000Error;
use crate::frame::{Frame, Payload, MAX_PAYLOAD_BYTES};
use crate::id::{Pgn, Transport};
pub const MAX_ASSEMBLIES: usize = 8;
pub const DEFAULT_TIMEOUT: Duration = Duration::from_millis(750);
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
struct Key {
pgn: Pgn,
source: u8,
}
#[derive(Debug, Clone, Copy)]
struct Assembly {
key: Key,
sequence: u8,
next: u8,
declared: usize,
started: Instant<Utc>,
payload: Payload,
}
#[derive(Debug, Clone)]
pub struct Assembler<const N: usize = MAX_ASSEMBLIES> {
slots: [Option<Assembly>; N],
timeout: Duration,
}
impl Default for Assembler<MAX_ASSEMBLIES> {
fn default() -> Self {
Self::new()
}
}
impl Assembler<MAX_ASSEMBLIES> {
#[must_use]
pub const fn new() -> Self {
Self::with_slots(DEFAULT_TIMEOUT)
}
#[must_use]
pub const fn with_timeout(timeout: Duration) -> Self {
Self::with_slots(timeout)
}
}
impl<const N: usize> Assembler<N> {
#[must_use]
pub const fn with_slots(timeout: Duration) -> Self {
Self {
slots: [None; N],
timeout,
}
}
pub fn push(
&mut self,
frame: &Frame,
now: Instant<Utc>,
) -> Result<Option<Payload>, Nmea2000Error> {
let pgn = frame.id.pgn();
let transport = pgn.transport().ok_or(Nmea2000Error::UnknownPgn { pgn })?;
self.push_as(frame, transport, now)
}
pub fn push_as(
&mut self,
frame: &Frame,
transport: Transport,
now: Instant<Utc>,
) -> Result<Option<Payload>, Nmea2000Error> {
self.expire(now);
let key = Key {
pgn: frame.id.pgn(),
source: frame.id.source(),
};
match transport {
Transport::SingleFrame => {
Payload::from_bytes(key.pgn, key.source, frame.data()).map(Some)
}
Transport::FastPacket => self.push_fast(frame, key, now),
}
}
fn push_fast(
&mut self,
frame: &Frame,
key: Key,
now: Instant<Utc>,
) -> Result<Option<Payload>, Nmea2000Error> {
let data = frame.data();
let Some(&counter) = data.first() else {
return Err(Nmea2000Error::UnexpectedFrame {
expected: 0,
found: 0,
});
};
let sequence = counter >> 5;
let number = counter & 0x1F;
let slot = self
.slots
.iter_mut()
.find(|slot| slot.is_some_and(|assembly| assembly.key == key));
if number == 0 {
let declared = usize::from(data.get(1).copied().unwrap_or(0));
if declared > MAX_PAYLOAD_BYTES {
return Err(Nmea2000Error::BadLength {
declared: data.get(1).copied().unwrap_or(0),
limit: MAX_PAYLOAD_BYTES,
});
}
let mut assembly = Assembly {
key,
sequence,
next: 1,
declared,
started: now,
payload: Payload::new(key.pgn, key.source),
};
assembly.payload.append(data.get(2..).unwrap_or(&[]))?;
if assembly.payload.len() >= declared {
assembly.payload.truncate(declared);
if let Some(slot) = slot {
*slot = None;
}
return Ok(Some(assembly.payload));
}
let slot = match slot {
Some(slot) => slot,
None => self.free_or_oldest_slot().ok_or(Nmea2000Error::NoSlot)?,
};
*slot = Some(assembly);
return Ok(None);
}
let Some(slot) = slot else {
return Err(Nmea2000Error::UnexpectedFrame {
expected: 0,
found: number,
});
};
let Some(assembly) = slot.as_mut() else {
return Ok(None);
};
if number != assembly.next || sequence != assembly.sequence {
let expected = assembly.next;
*slot = None;
return Err(Nmea2000Error::UnexpectedFrame {
expected,
found: number,
});
}
if let Err(error) = assembly.payload.append(data.get(1..).unwrap_or(&[])) {
*slot = None;
return Err(error);
}
if assembly.payload.len() >= assembly.declared {
let mut payload = assembly.payload;
payload.truncate(assembly.declared);
*slot = None;
return Ok(Some(payload));
}
assembly.next = number.saturating_add(1);
Ok(None)
}
fn free_or_oldest_slot(&mut self) -> Option<&mut Option<Assembly>> {
let index = self.slots.iter().position(Option::is_none).or_else(|| {
self.slots
.iter()
.enumerate()
.min_by_key(|(_, slot)| slot.map(|assembly| assembly.started))
.map(|(index, _)| index)
})?;
self.slots.get_mut(index)
}
pub fn expire(&mut self, now: Instant<Utc>) {
for slot in &mut self.slots {
let stale = slot.is_some_and(|assembly| {
now.checked_duration_since(assembly.started)
.is_none_or(|age| age > self.timeout)
});
if stale {
*slot = None;
}
}
}
#[must_use]
pub fn pending(&self) -> usize {
self.slots.iter().filter(|slot| slot.is_some()).count()
}
#[must_use]
pub const fn timeout(&self) -> Duration {
self.timeout
}
}