use crate::backend::{DacBackend, FifoBackend, WriteOutcome};
use crate::buffer_estimate::{BufferEstimator, StatusDecayEstimator};
use crate::device::{DacCapabilities, DacType};
use crate::error::{Error, Result};
use crate::point::LaserPoint;
use crate::protocols::ether_dream::dac::stream::{
self, CommunicationError, Nak, ResponseErrorKind,
};
use crate::protocols::ether_dream::dac::{LightEngine, Playback, PlaybackFlags};
use crate::protocols::ether_dream::protocol::{DacBroadcast, DacPoint};
use std::net::IpAddr;
use std::time::{Duration, Instant};
const WARMUP_PING_INTERVAL: Duration = Duration::from_millis(100);
const ESTOP_RETRY_INTERVAL: Duration = Duration::from_secs(1);
pub struct EtherDreamBackend {
broadcast: DacBroadcast,
ip_addr: IpAddr,
stream: Option<stream::Stream>,
caps: DacCapabilities,
last_status_time: Option<Instant>,
last_point_rate: u32,
last_ping_time: Option<Instant>,
last_estop_attempt: Option<Instant>,
point_buffer: Vec<DacPoint>,
estimator: StatusDecayEstimator,
}
impl EtherDreamBackend {
pub fn new(broadcast: DacBroadcast, ip_addr: IpAddr) -> Self {
Self {
broadcast,
ip_addr,
stream: None,
caps: super::default_capabilities(),
last_status_time: None,
last_point_rate: 0,
last_ping_time: None,
last_estop_attempt: None,
point_buffer: Vec::new(),
estimator: StatusDecayEstimator::new(),
}
}
}
impl DacBackend for EtherDreamBackend {
fn dac_type(&self) -> DacType {
DacType::EtherDream
}
fn caps(&self) -> &DacCapabilities {
&self.caps
}
fn connect(&mut self) -> Result<()> {
let stream = stream::connect_timeout(&self.broadcast, self.ip_addr, Duration::from_secs(5))
.map_err(Error::backend)?;
self.stream = Some(stream);
Ok(())
}
fn disconnect(&mut self) -> Result<()> {
if let Some(stream) = &mut self.stream {
let _ = stream.queue_commands().stop().submit();
}
self.stream = None;
Ok(())
}
fn is_connected(&self) -> bool {
self.stream.is_some()
}
fn stop(&mut self) -> Result<()> {
if let Some(stream) = &mut self.stream {
match stream.queue_commands().stop().submit() {
Ok(()) => {}
Err(e) if matches!(nak_of(&e), Some(Nak::Invalid)) => {}
Err(e) => return Err(Error::backend(e)),
}
}
Ok(())
}
fn set_shutter(&mut self, _open: bool) -> Result<()> {
Ok(())
}
}
impl FifoBackend for EtherDreamBackend {
fn try_write_points(&mut self, pps: u32, points: &[LaserPoint]) -> Result<WriteOutcome> {
let stream = self
.stream
.as_mut()
.ok_or_else(|| Error::disconnected("Not connected"))?;
if points.is_empty() {
return Ok(WriteOutcome::WouldBlock);
}
match stream.dac().status.light_engine {
LightEngine::EmergencyStop => {
let now = Instant::now();
let due = self
.last_estop_attempt
.is_none_or(|t| now.duration_since(t) >= ESTOP_RETRY_INTERVAL);
if !due {
return Ok(WriteOutcome::WouldBlock);
}
self.last_estop_attempt = Some(now);
match stream.queue_commands().clear_emergency_stop().submit() {
Ok(()) => {}
Err(e) => match nak_of(&e) {
Some(Nak::StopCondition) => {
log::warn!(
"Ether Dream stuck in emergency stop - check hardware interlock"
);
return Ok(WriteOutcome::WouldBlock);
}
_ => return Err(Error::backend(e)),
},
}
stream
.queue_commands()
.ping()
.submit()
.map_err(Error::backend)?;
if stream.dac().status.light_engine == LightEngine::EmergencyStop {
log::warn!(
"Ether Dream still in emergency stop after clear - check hardware interlock"
);
return Ok(WriteOutcome::WouldBlock);
}
let now = Instant::now();
self.last_status_time = Some(now);
self.estimator
.record_status(now, stream.dac().status.buffer_fullness as u64);
}
LightEngine::Warmup | LightEngine::Cooldown => {
let now = Instant::now();
let due = self
.last_ping_time
.is_none_or(|t| now.duration_since(t) >= WARMUP_PING_INTERVAL);
if due {
self.last_ping_time = Some(now);
stream
.queue_commands()
.ping()
.submit()
.map_err(Error::backend)?;
}
return Ok(WriteOutcome::WouldBlock);
}
LightEngine::Ready => {}
}
let point_rate = if pps > 0 {
pps
} else {
stream.dac().max_point_rate / 16
};
let max_rate = stream.dac().max_point_rate;
let point_rate = if max_rate > 0 {
point_rate.min(max_rate)
} else {
point_rate
};
let playback = stream.dac().status.playback;
let buffer_capacity = stream.dac().buffer_capacity;
let raw_fullness = stream.dac().status.buffer_fullness;
let fullness = decay_fullness(
raw_fullness,
buffer_capacity,
self.last_status_time,
self.last_point_rate,
playback == Playback::Playing,
);
let available = buffer_capacity.saturating_sub(fullness).saturating_sub(1) as usize;
if available < points.len() {
return Ok(WriteOutcome::WouldBlock);
}
self.point_buffer.clear();
self.point_buffer.extend(points.iter().map(DacPoint::from));
let playback_flags = stream.dac().status.playback_flags;
let current_point_rate = stream.dac().status.point_rate;
let needs_prepare =
playback_flags.contains(PlaybackFlags::UNDERFLOWED) || playback == Playback::Idle;
if needs_prepare {
stream
.queue_commands()
.prepare_stream()
.submit()
.map_err(Error::backend)?;
}
let send_result = if playback == Playback::Playing && current_point_rate != point_rate {
stream
.queue_commands()
.update(0, point_rate)
.data(self.point_buffer.iter().copied())
.submit()
} else {
stream
.queue_commands()
.data(self.point_buffer.iter().copied())
.submit()
};
match send_result {
Ok(()) => {}
Err(e) => match nak_of(&e) {
Some(Nak::Full) => {
self.last_status_time = Some(Instant::now());
return Ok(WriteOutcome::WouldBlock);
}
Some(Nak::Invalid) => {
if playback == Playback::Idle {
stream
.queue_commands()
.prepare_stream()
.submit()
.map_err(Error::backend)?;
match stream
.queue_commands()
.data(self.point_buffer.iter().copied())
.submit()
{
Ok(()) => {}
Err(e2) => match nak_of(&e2) {
Some(Nak::Full) | Some(Nak::Invalid) => {
self.last_status_time = Some(Instant::now());
return Ok(WriteOutcome::WouldBlock);
}
_ => return Err(Error::backend(e2)),
},
}
} else {
self.last_status_time = Some(Instant::now());
return Ok(WriteOutcome::WouldBlock);
}
}
_ => return Err(Error::backend(e)),
},
}
let playback_after = stream.dac().status.playback;
let buffer_fullness = stream.dac().status.buffer_fullness;
let needs_begin =
playback_after != Playback::Playing && buffer_fullness >= begin_threshold(point_rate);
if needs_begin {
stream
.queue_commands()
.begin(0, point_rate)
.submit()
.map_err(Error::backend)?;
}
let now = Instant::now();
self.last_status_time = Some(now);
self.last_point_rate = point_rate;
self.estimator
.set_playing(stream.dac().status.playback == Playback::Playing);
self.estimator
.record_status(now, stream.dac().status.buffer_fullness as u64);
Ok(WriteOutcome::Written)
}
fn estimator(&self) -> &dyn BufferEstimator {
&self.estimator
}
}
fn nak_of(err: &CommunicationError) -> Option<Nak> {
match err {
CommunicationError::Response(re) => match &re.kind {
ResponseErrorKind::Nak(nak) => Some(*nak),
_ => None,
},
_ => None,
}
}
fn begin_threshold(point_rate: u32) -> u16 {
let ten_ms = point_rate / 100;
ten_ms.clamp(64, 1700) as u16
}
fn decay_fullness(
raw: u16,
capacity: u16,
last_status_time: Option<Instant>,
point_rate: u32,
playing: bool,
) -> u16 {
if !playing {
return raw.min(capacity);
}
if let Some(last_time) = last_status_time {
let elapsed_secs = last_time.elapsed().as_secs_f64();
let consumed = (elapsed_secs * point_rate as f64) as u16;
raw.saturating_sub(consumed).min(capacity)
} else {
raw
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::protocols::ether_dream::dac::{
Addressed, Dac, DataSource, LightEngineFlags, MacAddress, Status,
};
use crate::protocols::ether_dream::protocol::command::Command as _;
use crate::protocols::ether_dream::protocol::{
self, DacResponse, DacStatus, SizeBytes, WriteToBytes,
};
use std::io::{Read, Write};
use std::net::{TcpListener, TcpStream};
use std::sync::{Arc, Mutex};
use std::thread;
#[test]
fn begin_threshold_is_pps_derived() {
assert_eq!(begin_threshold(1_000), 64); assert_eq!(begin_threshold(6_400), 64); assert_eq!(begin_threshold(30_000), 300); assert_eq!(begin_threshold(100_000), 1000);
assert_eq!(begin_threshold(1_000_000), 1700); }
fn mk_status(le: u8, pb: u8, fullness: u16, rate: u32) -> DacStatus {
DacStatus {
protocol: 0,
light_engine_state: le,
playback_state: pb,
source: 0,
light_engine_flags: 0,
playback_flags: 0,
source_flags: 0,
buffer_fullness: fullness,
point_rate: rate,
point_count: 0,
}
}
fn ack(cmd: u8, status: DacStatus) -> DacResponse {
DacResponse {
response: DacResponse::ACK,
command: cmd,
dac_status: status,
}
}
fn nak(kind: u8, cmd: u8, status: DacStatus) -> DacResponse {
DacResponse {
response: kind,
command: cmd,
dac_status: status,
}
}
fn addressed(le: LightEngine, pb: Playback, fullness: u16, capacity: u16) -> Addressed {
Addressed {
mac_address: MacAddress([0; 6]),
dac: Dac {
hw_revision: 0,
sw_revision: 0,
buffer_capacity: capacity,
max_point_rate: 100_000,
status: Status {
protocol: 0,
light_engine: le,
playback: pb,
data_source: DataSource::NetworkStreaming,
light_engine_flags: LightEngineFlags::empty(),
playback_flags: PlaybackFlags::empty(),
buffer_fullness: fullness,
point_rate: 0,
point_count: 0,
},
},
}
}
fn test_broadcast() -> DacBroadcast {
DacBroadcast {
mac_address: [0; 6],
hw_revision: 0,
sw_revision: 0,
buffer_capacity: 1000,
max_point_rate: 100_000,
dac_status: mk_status(0, 0, 0, 0),
}
}
fn connect_mock(
initial: Addressed,
handler: impl FnMut(u8) -> DacResponse + Send + 'static,
) -> (EtherDreamBackend, Arc<Mutex<Vec<u8>>>) {
let listener = TcpListener::bind("127.0.0.1:0").unwrap();
let addr = listener.local_addr().unwrap();
let commands = Arc::new(Mutex::new(Vec::new()));
let commands_srv = commands.clone();
let mut handler = handler;
thread::spawn(move || {
let Ok((mut sock, _)) = listener.accept() else {
return;
};
loop {
let mut cmd = [0u8; 1];
if sock.read_exact(&mut cmd).is_err() {
break;
}
let cmd = cmd[0];
let extra = match cmd {
protocol::command::Begin::START_BYTE
| protocol::command::Update::START_BYTE => 6,
protocol::command::PointRate::START_BYTE => 4,
protocol::command::Data::START_BYTE => {
let mut n = [0u8; 2];
if sock.read_exact(&mut n).is_err() {
break;
}
u16::from_le_bytes(n) as usize * DacPoint::SIZE_BYTES
}
_ => 0,
};
if extra > 0 {
let mut buf = vec![0u8; extra];
if sock.read_exact(&mut buf).is_err() {
break;
}
}
commands_srv.lock().unwrap().push(cmd);
let resp = handler(cmd);
let mut out = Vec::new();
resp.write_to_bytes(&mut out).unwrap();
if sock.write_all(&out).is_err() {
break;
}
}
});
let client = TcpStream::connect(addr).unwrap();
let stream = stream::Stream::from_tcp_stream_for_test(initial, client).unwrap();
let mut backend = EtherDreamBackend::new(test_broadcast(), "127.0.0.1".parse().unwrap());
backend.stream = Some(stream);
(backend, commands)
}
fn points(n: usize) -> Vec<LaserPoint> {
vec![LaserPoint::new(0.0, 0.0, 0, 0, 0, 0); n]
}
#[test]
fn ack_fullness_is_not_double_counted() {
let (mut backend, _cmds) = connect_mock(
addressed(LightEngine::Ready, Playback::Idle, 0, 1000),
|cmd| match cmd {
protocol::command::PrepareStream::START_BYTE => {
ack(cmd, mk_status(0, DacStatus::PLAYBACK_PREPARED, 0, 30_000))
}
protocol::command::Data::START_BYTE => {
ack(cmd, mk_status(0, DacStatus::PLAYBACK_PREPARED, 500, 30_000))
}
protocol::command::Begin::START_BYTE => {
ack(cmd, mk_status(0, DacStatus::PLAYBACK_PLAYING, 500, 30_000))
}
_ => ack(cmd, mk_status(0, DacStatus::PLAYBACK_PLAYING, 500, 30_000)),
},
);
let outcome = backend.try_write_points(30_000, &points(10)).unwrap();
assert_eq!(outcome, WriteOutcome::Written);
let est = backend
.estimator()
.estimated_fullness(Instant::now(), 30_000);
assert_eq!(est, 500);
}
#[test]
fn partial_fit_blocks_without_writing() {
let (mut backend, cmds) = connect_mock(
addressed(LightEngine::Ready, Playback::Prepared, 95, 100),
|cmd| ack(cmd, mk_status(0, DacStatus::PLAYBACK_PREPARED, 95, 30_000)),
);
let outcome = backend.try_write_points(30_000, &points(10)).unwrap();
assert_eq!(outcome, WriteOutcome::WouldBlock);
assert!(
cmds.lock().unwrap().is_empty(),
"no command should be sent when the chunk does not fit"
);
}
#[test]
fn warmup_pings_rate_limited() {
let (mut backend, cmds) = connect_mock(
addressed(LightEngine::Warmup, Playback::Idle, 0, 1000),
|cmd| {
ack(
cmd,
mk_status(
DacStatus::LIGHT_ENGINE_WARMUP,
DacStatus::PLAYBACK_IDLE,
0,
0,
),
)
},
);
let o1 = backend.try_write_points(30_000, &points(10)).unwrap();
let o2 = backend.try_write_points(30_000, &points(10)).unwrap();
assert_eq!(o1, WriteOutcome::WouldBlock);
assert_eq!(o2, WriteOutcome::WouldBlock);
let log = cmds.lock().unwrap();
assert_eq!(log.as_slice(), &[protocol::command::Ping::START_BYTE]);
}
#[test]
fn nak_full_maps_to_would_block() {
let (mut backend, cmds) = connect_mock(
addressed(LightEngine::Ready, Playback::Prepared, 0, 1000),
|cmd| match cmd {
protocol::command::Data::START_BYTE => nak(
DacResponse::NAK_FULL,
cmd,
mk_status(0, DacStatus::PLAYBACK_PREPARED, 999, 30_000),
),
_ => ack(cmd, mk_status(0, DacStatus::PLAYBACK_PREPARED, 0, 30_000)),
},
);
let outcome = backend.try_write_points(30_000, &points(10)).unwrap();
assert_eq!(outcome, WriteOutcome::WouldBlock);
assert!(cmds
.lock()
.unwrap()
.contains(&protocol::command::Data::START_BYTE));
}
#[test]
fn nak_invalid_while_prepared_maps_to_would_block() {
let (mut backend, _cmds) = connect_mock(
addressed(LightEngine::Ready, Playback::Prepared, 0, 1000),
|cmd| match cmd {
protocol::command::Data::START_BYTE => nak(
DacResponse::NAK_INVALID,
cmd,
mk_status(0, DacStatus::PLAYBACK_PREPARED, 0, 30_000),
),
_ => ack(cmd, mk_status(0, DacStatus::PLAYBACK_PREPARED, 0, 30_000)),
},
);
let outcome = backend.try_write_points(30_000, &points(10)).unwrap();
assert_eq!(outcome, WriteOutcome::WouldBlock);
}
#[test]
fn estop_stop_condition_blocks_and_rate_limits() {
let (mut backend, cmds) = connect_mock(
addressed(LightEngine::EmergencyStop, Playback::Idle, 0, 1000),
|cmd| {
nak(
DacResponse::NAK_STOP_CONDITION,
cmd,
mk_status(
DacStatus::LIGHT_ENGINE_EMERGENCY_STOP,
DacStatus::PLAYBACK_IDLE,
0,
0,
),
)
},
);
let o1 = backend.try_write_points(30_000, &points(10)).unwrap();
let o2 = backend.try_write_points(30_000, &points(10)).unwrap();
assert_eq!(o1, WriteOutcome::WouldBlock);
assert_eq!(o2, WriteOutcome::WouldBlock);
let log = cmds.lock().unwrap();
assert_eq!(
log.as_slice(),
&[protocol::command::ClearEmergencyStop::START_BYTE]
);
}
#[test]
fn io_error_on_send_is_fatal_not_retried() {
let listener = TcpListener::bind("127.0.0.1:0").unwrap();
let addr = listener.local_addr().unwrap();
let commands = Arc::new(Mutex::new(Vec::new()));
let commands_srv = commands.clone();
thread::spawn(move || {
let Ok((mut sock, _)) = listener.accept() else {
return;
};
loop {
let mut cmd = [0u8; 1];
if sock.read_exact(&mut cmd).is_err() {
break;
}
let cmd = cmd[0];
if cmd == protocol::command::Data::START_BYTE {
let mut n = [0u8; 2];
if sock.read_exact(&mut n).is_err() {
break;
}
let extra = u16::from_le_bytes(n) as usize * DacPoint::SIZE_BYTES;
let mut buf = vec![0u8; extra];
let _ = sock.read_exact(&mut buf);
commands_srv.lock().unwrap().push(cmd);
break; }
commands_srv.lock().unwrap().push(cmd);
let resp = ack(cmd, mk_status(0, DacStatus::PLAYBACK_PREPARED, 0, 30_000));
let mut out = Vec::new();
resp.write_to_bytes(&mut out).unwrap();
if sock.write_all(&out).is_err() {
break;
}
}
});
let client = TcpStream::connect(addr).unwrap();
let stream = stream::Stream::from_tcp_stream_for_test(
addressed(LightEngine::Ready, Playback::Idle, 0, 1000),
client,
)
.unwrap();
let mut backend = EtherDreamBackend::new(test_broadcast(), "127.0.0.1".parse().unwrap());
backend.stream = Some(stream);
let result = backend.try_write_points(30_000, &points(10));
assert!(result.is_err(), "IO error on send must be fatal");
let log = commands.lock().unwrap();
assert_eq!(
log.iter()
.filter(|&&c| c == protocol::command::Data::START_BYTE)
.count(),
1
);
}
}