use crate::protocols::idn::dac::{Addressed, ServerInfo};
use crate::protocols::idn::error::{CommunicationError, ProtocolError, ResponseError, Result};
use crate::protocols::idn::protocol::{
AcknowledgeResponse, ChannelConfigHeader, ChannelMessageHeader, GroupRequest, GroupResponse,
PacketHeader, ParameterGetRequest, ParameterResponse, ParameterSetRequest, Point, ReadBytes,
ReadFromBytes, SampleChunkHeader, SizeBytes, WriteBytes, EXTENDED_SAMPLE_SIZE,
IDNCMD_GROUP_REQUEST, IDNCMD_GROUP_RESPONSE, IDNCMD_PING_REQUEST, IDNCMD_PING_RESPONSE,
IDNCMD_RT_ACKNOWLEDGE, IDNCMD_RT_CNLMSG, IDNCMD_RT_CNLMSG_ACKREQ, IDNCMD_RT_CNLMSG_CLOSE,
IDNCMD_RT_CNLMSG_CLOSE_ACKREQ, IDNCMD_SERVICE_PARAMS_REQUEST, IDNCMD_SERVICE_PARAMS_RESPONSE,
IDNCMD_UNIT_PARAMS_REQUEST, IDNCMD_UNIT_PARAMS_RESPONSE, IDNFLG_CHNCFG_CLOSE,
IDNFLG_CHNCFG_ROUTING, IDNFLG_CONTENTID_CHANNELMSG, IDNFLG_CONTENTID_CONFIG_LSTFRG,
IDNVAL_CNKTYPE_LPGRF_FRAME, IDNVAL_CNKTYPE_LPGRF_FRAME_FIRST,
IDNVAL_CNKTYPE_LPGRF_FRAME_SEQUEL, IDNVAL_CNKTYPE_LPGRF_WAVE, IDNVAL_CNKTYPE_VOID,
IDNVAL_SMOD_LPGRF_CONTINUOUS, IDNVAL_SMOD_LPGRF_DISCRETE, MAX_UDP_PAYLOAD, XYRGBI_SAMPLE_SIZE,
XYRGB_HIGHRES_SAMPLE_SIZE,
};
use log::{debug, trace, warn};
use std::io;
use std::net::UdpSocket;
use std::time::{Duration, Instant};
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Default)]
pub enum FrameMode {
#[default]
Wave,
Frame,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
pub enum PointFormat {
Xyrgbi,
XyrgbHighRes,
Extended,
}
impl PointFormat {
pub fn size_bytes(&self) -> usize {
match self {
PointFormat::Xyrgbi => XYRGBI_SAMPLE_SIZE,
PointFormat::XyrgbHighRes => XYRGB_HIGHRES_SAMPLE_SIZE,
PointFormat::Extended => EXTENDED_SAMPLE_SIZE,
}
}
fn word_count(&self) -> u8 {
(self.descriptors().len() / 2) as u8
}
fn descriptors(&self) -> &'static [u16] {
use crate::protocols::idn::protocol::channel_descriptors;
match self {
PointFormat::Xyrgbi => channel_descriptors::XYRGBI,
PointFormat::XyrgbHighRes => channel_descriptors::XYRGB_HIGHRES,
PointFormat::Extended => channel_descriptors::EXTENDED,
}
}
}
pub const LINK_TIMEOUT: Duration = Duration::from_secs(1);
pub const KEEPALIVE_INTERVAL: Duration = Duration::from_millis(500);
const MIN_SAMPLES_PER_FRAME: usize = 20;
const CONFIG_REFRESH_INTERVAL: Duration = Duration::from_millis(200);
const STARTUP_LEAD_US: u64 = 30_000;
const REANCHOR_MIN_GAP_US: u64 = 50_000;
pub struct Stream {
dac: Addressed,
socket: UdpSocket,
client_group: u8,
sequence: u16,
timestamp_anchor_us: u64,
points_since_anchor: u64,
anchor_pps: u32,
last_chunk_duration_us: u64,
scan_speed: u32,
point_format: PointFormat,
service_data_match: u8,
last_config_time: Option<Instant>,
previous_format: Option<PointFormat>,
frame_count: u64,
packet_buffer: Vec<u8>,
recv_buffer: [u8; 64],
last_send_time: Option<Instant>,
last_data_send_time: Option<Instant>,
connect_time: Instant,
frame_mode: FrameMode,
closed: bool,
}
pub fn connect(server: &ServerInfo, service_id: u8) -> io::Result<Stream> {
connect_with_group(server, service_id, 0)
}
pub fn connect_with_group(
server: &ServerInfo,
service_id: u8,
client_group: u8,
) -> io::Result<Stream> {
let address = server
.primary_address()
.ok_or_else(|| io::Error::new(io::ErrorKind::NotFound, "server has no addresses"))?;
debug!(
"IDN stream: connecting to {} (service_id={}, group={}, address={})",
server.hostname, service_id, client_group, address
);
let service = server
.find_service(service_id)
.ok_or_else(|| io::Error::new(io::ErrorKind::NotFound, "service not found"))?;
debug!(
"IDN stream: found service '{}' type={:?} flags=0x{:02x}",
service.name, service.service_type, service.flags
);
let socket = UdpSocket::bind("0.0.0.0:0")?;
let local_addr = socket.local_addr().ok();
socket.connect(address)?;
debug!(
"IDN stream: UDP socket bound to {:?}, connected to {}",
local_addr, address
);
let dac = Addressed::new(server.clone(), service.clone(), *address);
Ok(Stream {
dac,
socket,
client_group: client_group & 0x0F,
sequence: 0,
timestamp_anchor_us: 0,
points_since_anchor: 0,
anchor_pps: 30000,
last_chunk_duration_us: 0,
scan_speed: 30000, point_format: PointFormat::Xyrgbi,
service_data_match: 0,
last_config_time: None,
previous_format: None,
frame_count: 0,
packet_buffer: Vec::with_capacity(MAX_UDP_PAYLOAD * 2),
recv_buffer: [0u8; 64],
last_send_time: None,
last_data_send_time: None,
connect_time: Instant::now(),
frame_mode: FrameMode::Wave,
closed: false,
})
}
impl Stream {
fn needs_config(
frame_count: u64,
previous_format: Option<PointFormat>,
point_format: PointFormat,
last_config_time: Option<Instant>,
now: Instant,
) -> bool {
frame_count == 0
|| previous_format != Some(point_format)
|| match last_config_time {
None => true,
Some(t) => now.duration_since(t) >= CONFIG_REFRESH_INTERVAL,
}
}
fn current_timestamp_us(&self) -> u64 {
self.timestamp_anchor_us
+ self.points_since_anchor * 1_000_000 / self.anchor_pps.max(1) as u64
}
pub fn dac(&self) -> &Addressed {
&self.dac
}
pub fn scan_speed(&self) -> u32 {
self.scan_speed
}
pub fn set_scan_speed(&mut self, pps: u32) {
self.scan_speed = pps;
}
pub fn point_format(&self) -> PointFormat {
self.point_format
}
pub fn set_point_format(&mut self, format: PointFormat) {
self.point_format = format;
}
pub fn frame_mode(&self) -> FrameMode {
self.frame_mode
}
pub fn set_frame_mode(&mut self, mode: FrameMode) {
self.frame_mode = mode;
}
pub fn needs_keepalive(&self) -> bool {
match self.last_send_time {
Some(last) => last.elapsed() >= KEEPALIVE_INTERVAL,
None => self.frame_count > 0, }
}
pub fn time_since_last_send(&self) -> Option<Duration> {
self.last_send_time.map(|t| t.elapsed())
}
pub fn send_keepalive(&mut self) -> Result<()> {
let content_id =
IDNFLG_CONTENTID_CHANNELMSG | self.channel_id() | IDNVAL_CNKTYPE_VOID as u16;
let timestamp = self.current_timestamp_us();
trace!(
"IDN stream: sending keepalive (seq={}, timestamp={}, time_since_last={:?})",
self.sequence,
timestamp,
self.last_send_time.map(|t| t.elapsed()),
);
self.packet_buffer.clear();
let packet_header = PacketHeader {
command: IDNCMD_RT_CNLMSG,
flags: self.client_group,
sequence: self.next_sequence(),
};
self.packet_buffer.write_bytes(packet_header)?;
let channel_msg = ChannelMessageHeader {
total_size: ChannelMessageHeader::SIZE_BYTES as u16,
content_id,
timestamp: (timestamp & 0xFFFF_FFFF) as u32,
};
self.packet_buffer.write_bytes(channel_msg)?;
let sent_bytes = self.socket.send(&self.packet_buffer)?;
self.last_send_time = Some(Instant::now());
trace!(
"IDN stream: keepalive sent ({} bytes, content_id=0x{:04x})",
sent_bytes,
content_id
);
Ok(())
}
pub fn write_frame<P: Point>(&mut self, points: &[P]) -> Result<()> {
if points.is_empty() {
trace!("IDN stream: write_frame called with empty points, skipping");
return Ok(());
}
let mut padded = Vec::new();
let points = self.prepare_points(points, &mut padded)?;
match self.frame_mode {
FrameMode::Wave => {
let (seq, points_to_send) =
self.build_and_send_first_packet(points, IDNCMD_RT_CNLMSG)?;
trace!(
"IDN stream: sent frame #{} packet - seq={}, {} points",
self.frame_count - 1,
seq,
points_to_send,
);
if points_to_send < points.len() {
self.write_frame_continuation(&points[points_to_send..])?;
}
}
FrameMode::Frame => {
self.write_frame_fragmented(points)?;
}
}
Ok(())
}
fn write_frame_fragmented<P: Point>(&mut self, points: &[P]) -> Result<()> {
let now = Instant::now();
let bytes_per_sample = P::SIZE_BYTES;
let service_id = self.dac.service_id();
let channel_id = self.channel_id();
let needs_config = Self::needs_config(
self.frame_count,
self.previous_format,
self.point_format,
self.last_config_time,
now,
);
let config_size = if needs_config {
ChannelConfigHeader::SIZE_BYTES + self.point_format.descriptors().len() * 2
} else {
0
};
let ts_start = self.current_timestamp_us();
let total_points = points.len() as u64;
let points_after = self.points_since_anchor + total_points;
let ts_end =
self.timestamp_anchor_us + points_after * 1_000_000 / self.anchor_pps.max(1) as u64;
let frame_duration_us = (ts_end - ts_start).min(0x00FF_FFFF) as u32;
let first_header = PacketHeader::SIZE_BYTES
+ ChannelMessageHeader::SIZE_BYTES
+ config_size
+ SampleChunkHeader::SIZE_BYTES;
let max_first = ((MAX_UDP_PAYLOAD - first_header) / bytes_per_sample).max(1);
let num_packets = points.len().div_ceil(max_first).max(1);
let base = points.len() / num_packets;
let rem = points.len() % num_packets;
let mut offset = 0usize;
for i in 0..num_packets {
let count = base + if i < rem { 1 } else { 0 };
let is_first = i == 0;
let is_last = i == num_packets - 1;
let is_single = num_packets == 1;
self.send_frame_fragment(
&points[offset..offset + count],
is_first,
is_last,
is_single,
needs_config && is_first,
service_id,
channel_id,
ts_start,
frame_duration_us,
now,
)?;
offset += count;
}
self.points_since_anchor = points_after;
self.last_chunk_duration_us = frame_duration_us as u64;
self.last_send_time = Some(now);
self.last_data_send_time = Some(now);
self.frame_count += 1;
Ok(())
}
#[allow(clippy::too_many_arguments)]
fn send_frame_fragment<P: Point>(
&mut self,
points: &[P],
is_first: bool,
is_last: bool,
is_single: bool,
write_cfg: bool,
service_id: u8,
channel_id: u16,
timestamp_us: u64,
frame_duration_us: u32,
now: Instant,
) -> Result<()> {
let bytes_per_sample = P::SIZE_BYTES;
let cnk_type = if is_single {
IDNVAL_CNKTYPE_LPGRF_FRAME
} else if is_first {
IDNVAL_CNKTYPE_LPGRF_FRAME_FIRST
} else {
IDNVAL_CNKTYPE_LPGRF_FRAME_SEQUEL
};
let mut content_id = IDNFLG_CONTENTID_CHANNELMSG | channel_id | cnk_type as u16;
if write_cfg {
content_id |= IDNFLG_CONTENTID_CONFIG_LSTFRG;
}
if is_last && !is_single {
content_id |= IDNFLG_CONTENTID_CONFIG_LSTFRG;
}
let has_chunk_header = is_first;
let config_size = if write_cfg {
ChannelConfigHeader::SIZE_BYTES + self.point_format.descriptors().len() * 2
} else {
0
};
let chunk_header_size = if has_chunk_header {
SampleChunkHeader::SIZE_BYTES
} else {
0
};
let msg_size = ChannelMessageHeader::SIZE_BYTES
+ config_size
+ chunk_header_size
+ points.len() * bytes_per_sample;
self.packet_buffer.clear();
let seq = self.next_sequence();
self.packet_buffer.write_bytes(PacketHeader {
command: IDNCMD_RT_CNLMSG,
flags: self.client_group,
sequence: seq,
})?;
self.packet_buffer.write_bytes(ChannelMessageHeader {
total_size: msg_size as u16,
content_id,
timestamp: (timestamp_us & 0xFFFF_FFFF) as u32,
})?;
if write_cfg {
self.write_config(service_id, now)?;
}
if has_chunk_header {
let chunk_header = SampleChunkHeader::new(self.sdm_flags(), frame_duration_us);
self.packet_buffer.write_bytes(chunk_header)?;
}
P::encode_batch_into(points, &mut self.packet_buffer);
self.socket.send(&self.packet_buffer)?;
Ok(())
}
fn write_frame_continuation<P: Point>(&mut self, points: &[P]) -> Result<()> {
if points.is_empty() {
return Ok(());
}
let bytes_per_sample = P::SIZE_BYTES;
let content_id = IDNFLG_CONTENTID_CHANNELMSG | self.channel_id();
let header_size = PacketHeader::SIZE_BYTES
+ ChannelMessageHeader::SIZE_BYTES
+ SampleChunkHeader::SIZE_BYTES;
let max_points_per_packet = (MAX_UDP_PAYLOAD - header_size) / bytes_per_sample;
let num_packets = points.len().div_ceil(max_points_per_packet);
let points_to_send = points.len() / num_packets;
let ts_start = self.current_timestamp_us();
let points_after = self.points_since_anchor + points_to_send as u64;
let ts_end =
self.timestamp_anchor_us + points_after * 1_000_000 / self.anchor_pps.max(1) as u64;
let duration_us = (ts_end - ts_start).min(0x00FF_FFFF) as u32;
trace!(
"IDN stream: continuation - {} points remaining, sending {}, duration={}us, timestamp={}",
points.len(),
points_to_send,
duration_us,
ts_start,
);
self.packet_buffer.clear();
let seq = self.next_sequence();
let packet_header = PacketHeader {
command: IDNCMD_RT_CNLMSG,
flags: self.client_group,
sequence: seq,
};
self.packet_buffer.write_bytes(packet_header)?;
let msg_size = ChannelMessageHeader::SIZE_BYTES
+ SampleChunkHeader::SIZE_BYTES
+ points_to_send * bytes_per_sample;
let cnk_type = self.chunk_type(false, false);
let channel_msg = ChannelMessageHeader {
total_size: msg_size as u16,
content_id: content_id | cnk_type as u16,
timestamp: (ts_start & 0xFFFF_FFFF) as u32,
};
self.packet_buffer.write_bytes(channel_msg)?;
let chunk_header = SampleChunkHeader::new(self.sdm_flags(), duration_us);
self.packet_buffer.write_bytes(chunk_header)?;
P::encode_batch_into(&points[..points_to_send], &mut self.packet_buffer);
let now = Instant::now();
let sent_bytes = self.socket.send(&self.packet_buffer)?;
self.last_send_time = Some(now);
self.last_data_send_time = Some(now);
trace!(
"IDN stream: continuation sent - seq={}, {} bytes, {} points",
seq,
sent_bytes,
points_to_send,
);
self.points_since_anchor = points_after;
self.last_chunk_duration_us = duration_us as u64;
if points_to_send < points.len() {
self.write_frame_continuation(&points[points_to_send..])?;
}
Ok(())
}
pub fn write_frame_with_ack<P: Point>(
&mut self,
points: &[P],
timeout: Duration,
) -> Result<AcknowledgeResponse> {
if points.is_empty() {
return Err(CommunicationError::Protocol(ProtocolError::BufferTooSmall));
}
let mut padded = Vec::new();
let points = self.prepare_points(points, &mut padded)?;
let (ack_seq, points_to_send) =
self.build_and_send_first_packet(points, IDNCMD_RT_CNLMSG_ACKREQ)?;
let ack = self.recv_acknowledge(timeout, ack_seq)?;
if points_to_send < points.len() {
self.write_frame_continuation(&points[points_to_send..])?;
}
Ok(ack)
}
fn prepare_points<'a, P: Point>(
&mut self,
points: &'a [P],
padded: &'a mut Vec<P>,
) -> Result<&'a [P]> {
let points = if points.len() < MIN_SAMPLES_PER_FRAME {
trace!(
"IDN stream: padding {} points to minimum {}",
points.len(),
MIN_SAMPLES_PER_FRAME
);
*padded = Self::pad_points(points);
padded.as_slice()
} else {
points
};
if self.scan_speed == 0 {
warn!("IDN stream: scan_speed is 0, cannot send frame");
return Err(CommunicationError::Protocol(
ProtocolError::InvalidScanSpeed,
));
}
if self.previous_format != Some(self.point_format) {
self.service_data_match = self.service_data_match.wrapping_add(1);
}
let elapsed_us = self.connect_time.elapsed().as_micros() as u64;
if self.frame_count == 0 {
self.timestamp_anchor_us = elapsed_us.saturating_sub(STARTUP_LEAD_US);
self.points_since_anchor = 0;
self.anchor_pps = self.scan_speed.max(1);
} else {
if let Some(last) = self.last_data_send_time {
let gap_us = last.elapsed().as_micros() as u64;
let threshold = (2 * self.last_chunk_duration_us).max(REANCHOR_MIN_GAP_US);
if gap_us > threshold {
self.timestamp_anchor_us = elapsed_us;
self.points_since_anchor = 0;
self.anchor_pps = self.scan_speed.max(1);
}
}
if self.anchor_pps != self.scan_speed.max(1) {
self.timestamp_anchor_us = self.current_timestamp_us();
self.points_since_anchor = 0;
self.anchor_pps = self.scan_speed.max(1);
}
}
Ok(points)
}
fn build_and_send_first_packet<P: Point>(
&mut self,
points: &[P],
command: u8,
) -> Result<(u16, usize)> {
let now = Instant::now();
let needs_config = Self::needs_config(
self.frame_count,
self.previous_format,
self.point_format,
self.last_config_time,
now,
);
let bytes_per_sample = P::SIZE_BYTES;
let service_id = self.dac.service_id();
let mut content_id = IDNFLG_CONTENTID_CHANNELMSG | self.channel_id();
let config_size = if needs_config {
content_id |= IDNFLG_CONTENTID_CONFIG_LSTFRG;
ChannelConfigHeader::SIZE_BYTES + self.point_format.descriptors().len() * 2
} else {
0
};
let header_size = PacketHeader::SIZE_BYTES
+ ChannelMessageHeader::SIZE_BYTES
+ config_size
+ SampleChunkHeader::SIZE_BYTES;
let max_points_per_packet = (MAX_UDP_PAYLOAD - header_size) / bytes_per_sample;
let num_packets = points.len().div_ceil(max_points_per_packet);
let points_to_send = points.len() / num_packets;
let ts_start = self.current_timestamp_us();
let points_after = self.points_since_anchor + points_to_send as u64;
let ts_end =
self.timestamp_anchor_us + points_after * 1_000_000 / self.anchor_pps.max(1) as u64;
let duration_us = (ts_end - ts_start).min(0x00FF_FFFF) as u32;
let is_only = num_packets == 1;
let cnk_type = self.chunk_type(true, is_only);
self.packet_buffer.clear();
let seq = self.next_sequence();
let packet_header = PacketHeader {
command,
flags: self.client_group,
sequence: seq,
};
self.packet_buffer.write_bytes(packet_header)?;
let msg_size = ChannelMessageHeader::SIZE_BYTES
+ config_size
+ SampleChunkHeader::SIZE_BYTES
+ points_to_send * bytes_per_sample;
let channel_msg = ChannelMessageHeader {
total_size: msg_size as u16,
content_id: content_id | cnk_type as u16,
timestamp: (ts_start & 0xFFFF_FFFF) as u32,
};
self.packet_buffer.write_bytes(channel_msg)?;
if needs_config {
self.write_config(service_id, now)?;
}
let chunk_header = SampleChunkHeader::new(self.sdm_flags(), duration_us);
self.packet_buffer.write_bytes(chunk_header)?;
P::encode_batch_into(&points[..points_to_send], &mut self.packet_buffer);
self.socket.send(&self.packet_buffer)?;
self.last_send_time = Some(now);
self.last_data_send_time = Some(now);
self.points_since_anchor = points_after;
self.last_chunk_duration_us = duration_us as u64;
self.frame_count += 1;
Ok((seq, points_to_send))
}
pub fn recv_acknowledge(
&mut self,
timeout: Duration,
expected_seq: u16,
) -> Result<AcknowledgeResponse> {
trace!(
"IDN stream: waiting for ack (expected_seq={}, timeout={:?})",
expected_seq,
timeout,
);
let ack: AcknowledgeResponse =
self.recv_response(timeout, IDNCMD_RT_ACKNOWLEDGE, expected_seq)?;
trace!(
"IDN stream: received ack - result_code={}, seq={}",
ack.result_code,
expected_seq,
);
if let Some(error) = ResponseError::from_ack_code(ack.result_code) {
warn!("IDN stream: ack error: {:?}", error);
return Err(CommunicationError::Response(error));
}
Ok(ack)
}
pub fn ping(&mut self, timeout: Duration) -> Result<Duration> {
let start = Instant::now();
let seq = self.send_request(IDNCMD_PING_REQUEST, |_| Ok(()))?;
self.recv_response_header(timeout, IDNCMD_PING_RESPONSE, seq)?;
Ok(start.elapsed())
}
pub fn get_client_group_mask(&mut self, timeout: Duration) -> Result<GroupResponse> {
let seq = self.send_request(IDNCMD_GROUP_REQUEST, |buf| {
buf.write_bytes(GroupRequest::get())
})?;
self.recv_response(timeout, IDNCMD_GROUP_RESPONSE, seq)
}
pub fn set_client_group_mask(&mut self, mask: u16, timeout: Duration) -> Result<GroupResponse> {
let seq = self.send_request(IDNCMD_GROUP_REQUEST, |buf| {
buf.write_bytes(GroupRequest::set(mask))
})?;
self.recv_response(timeout, IDNCMD_GROUP_RESPONSE, seq)
}
pub fn get_parameter(
&mut self,
service_id: u8,
param_id: u16,
timeout: Duration,
) -> Result<ParameterResponse> {
let (request_cmd, response_cmd) = param_commands(service_id);
let seq = self.send_request(request_cmd, |buf| {
buf.write_bytes(ParameterGetRequest {
service_id,
reserved: 0,
param_id,
})
})?;
let response: ParameterResponse = self.recv_response(timeout, response_cmd, seq)?;
check_parameter_response(&response)?;
Ok(response)
}
pub fn set_parameter(
&mut self,
service_id: u8,
param_id: u16,
value: u32,
timeout: Duration,
) -> Result<ParameterResponse> {
let (request_cmd, response_cmd) = param_commands(service_id);
let seq = self.send_request(request_cmd, |buf| {
buf.write_bytes(ParameterSetRequest {
service_id,
reserved: 0,
param_id,
value,
})
})?;
let response: ParameterResponse = self.recv_response(timeout, response_cmd, seq)?;
check_parameter_response(&response)?;
Ok(response)
}
pub fn close(&mut self) -> Result<()> {
if self.closed {
return Ok(());
}
self.closed = true;
debug!(
"IDN stream: closing connection (frames sent: {}, timestamp: {})",
self.frame_count,
self.current_timestamp_us()
);
self.send_channel_close()?;
self.packet_buffer.clear();
let close_header = PacketHeader {
command: IDNCMD_RT_CNLMSG_CLOSE,
flags: self.client_group,
sequence: self.next_sequence(),
};
self.packet_buffer.write_bytes(close_header)?;
self.socket.send(&self.packet_buffer)?;
debug!("IDN stream: close complete");
Ok(())
}
pub fn close_with_ack(&mut self, timeout: Duration) -> Result<AcknowledgeResponse> {
self.closed = true;
self.send_channel_close()?;
self.packet_buffer.clear();
let close_seq = self.next_sequence();
let close_header = PacketHeader {
command: IDNCMD_RT_CNLMSG_CLOSE_ACKREQ,
flags: self.client_group,
sequence: close_seq,
};
self.packet_buffer.write_bytes(close_header)?;
self.socket.send(&self.packet_buffer)?;
self.recv_acknowledge(timeout, close_seq)
}
fn send_channel_close(&mut self) -> Result<()> {
debug!("IDN stream: sending channel close");
let service_id = self.dac.service_id();
let channel_id = self.channel_id();
self.packet_buffer.clear();
let packet_header = PacketHeader {
command: IDNCMD_RT_CNLMSG,
flags: self.client_group,
sequence: self.next_sequence(),
};
self.packet_buffer.write_bytes(packet_header)?;
let content_id = IDNFLG_CONTENTID_CHANNELMSG
| IDNFLG_CONTENTID_CONFIG_LSTFRG
| channel_id
| IDNVAL_CNKTYPE_VOID as u16;
let msg_size = ChannelMessageHeader::SIZE_BYTES + ChannelConfigHeader::SIZE_BYTES;
let channel_msg = ChannelMessageHeader {
total_size: msg_size as u16,
content_id,
timestamp: (self.current_timestamp_us() & 0xFFFF_FFFF) as u32,
};
self.packet_buffer.write_bytes(channel_msg)?;
let config = ChannelConfigHeader {
word_count: 0,
flags: IDNFLG_CHNCFG_CLOSE,
service_id,
service_mode: 0,
};
self.packet_buffer.write_bytes(config)?;
self.socket.send(&self.packet_buffer)?;
Ok(())
}
fn send_request(
&mut self,
command: u8,
write_body: impl FnOnce(&mut Vec<u8>) -> io::Result<()>,
) -> Result<u16> {
self.packet_buffer.clear();
let seq = self.next_sequence();
let header = PacketHeader {
command,
flags: self.client_group,
sequence: seq,
};
self.packet_buffer.write_bytes(header)?;
write_body(&mut self.packet_buffer)?;
self.socket.send(&self.packet_buffer)?;
self.last_send_time = Some(Instant::now());
Ok(seq)
}
fn recv_response_header(
&mut self,
timeout: Duration,
expected_cmd: u8,
expected_seq: u16,
) -> Result<PacketHeader> {
let len = self.recv_matching(timeout, expected_cmd, expected_seq)?;
let mut cursor = &self.recv_buffer[..len];
let header: PacketHeader = cursor.read_bytes()?;
Ok(header)
}
fn recv_response<T: ReadFromBytes + SizeBytes>(
&mut self,
timeout: Duration,
expected_cmd: u8,
expected_seq: u16,
) -> Result<T> {
let len = self.recv_matching(timeout, expected_cmd, expected_seq)?;
if len < PacketHeader::SIZE_BYTES + T::SIZE_BYTES {
return Err(CommunicationError::Protocol(ProtocolError::BufferTooSmall));
}
let mut cursor = &self.recv_buffer[..len];
let _header: PacketHeader = cursor.read_bytes()?;
let response: T = cursor.read_bytes()?;
Ok(response)
}
fn recv_matching(
&mut self,
timeout: Duration,
expected_cmd: u8,
expected_seq: u16,
) -> Result<usize> {
let deadline = Instant::now() + timeout;
loop {
let remaining = deadline.saturating_duration_since(Instant::now());
if remaining.is_zero() {
return Err(CommunicationError::Response(ResponseError::Timeout));
}
let len = self.recv_into_buffer(remaining)?;
if len < PacketHeader::SIZE_BYTES {
trace!("IDN stream: discarding undersized datagram ({} bytes)", len);
continue;
}
let mut cursor = &self.recv_buffer[..len];
let header: PacketHeader = match cursor.read_bytes() {
Ok(h) => h,
Err(_) => continue,
};
if header.command != expected_cmd || header.sequence != expected_seq {
trace!(
"IDN stream: discarding non-matching datagram (cmd=0x{:02x} seq={}, \
want cmd=0x{:02x} seq={})",
header.command,
header.sequence,
expected_cmd,
expected_seq,
);
continue;
}
return Ok(len);
}
}
fn recv_into_buffer(&mut self, timeout: Duration) -> Result<usize> {
trace!("IDN stream: waiting for response (timeout={:?})", timeout);
self.socket.set_read_timeout(Some(timeout))?;
match self.socket.recv(&mut self.recv_buffer) {
Ok(len) => {
trace!(
"IDN stream: received {} bytes: {:02x?}",
len,
&self.recv_buffer[..len.min(32)]
);
Ok(len)
}
Err(e)
if e.kind() == io::ErrorKind::WouldBlock || e.kind() == io::ErrorKind::TimedOut =>
{
debug!("IDN stream: receive timeout ({:?})", timeout);
Err(CommunicationError::Response(ResponseError::Timeout))
}
Err(e) => {
warn!("IDN stream: receive error: {}", e);
Err(CommunicationError::Io(e))
}
}
}
fn write_config(&mut self, service_id: u8, now: Instant) -> Result<()> {
let flags = IDNFLG_CHNCFG_ROUTING | self.sdm_flags();
let service_mode = match self.frame_mode {
FrameMode::Wave => IDNVAL_SMOD_LPGRF_CONTINUOUS,
FrameMode::Frame => IDNVAL_SMOD_LPGRF_DISCRETE,
};
let config = ChannelConfigHeader {
word_count: self.point_format.word_count(),
flags,
service_id,
service_mode,
};
trace!(
"IDN stream: write_config - service_id={}, word_count={}, flags=0x{:02x}, \
service_mode=0x{:02x}, format={:?}, descriptors={:?}",
service_id,
config.word_count,
flags,
service_mode,
self.point_format,
self.point_format
.descriptors()
.iter()
.map(|d| format!("0x{:04x}", d))
.collect::<Vec<_>>(),
);
self.packet_buffer.write_bytes(config)?;
for &desc in self.point_format.descriptors() {
self.packet_buffer.push((desc >> 8) as u8);
self.packet_buffer.push(desc as u8);
}
self.last_config_time = Some(now);
self.previous_format = Some(self.point_format);
Ok(())
}
fn chunk_type(&self, is_first: bool, is_only: bool) -> u8 {
match self.frame_mode {
FrameMode::Wave => IDNVAL_CNKTYPE_LPGRF_WAVE,
FrameMode::Frame if is_only => IDNVAL_CNKTYPE_LPGRF_FRAME,
FrameMode::Frame if is_first => IDNVAL_CNKTYPE_LPGRF_FRAME_FIRST,
FrameMode::Frame => IDNVAL_CNKTYPE_LPGRF_FRAME_SEQUEL,
}
}
fn channel_id(&self) -> u16 {
let service_id = self.dac.service_id();
((service_id.saturating_sub(1)) as u16 & 0x3F) << 8
}
fn sdm_flags(&self) -> u8 {
((self.service_data_match & 1) | 2) << 4
}
fn next_sequence(&mut self) -> u16 {
let seq = self.sequence;
self.sequence = self.sequence.wrapping_add(1);
seq
}
fn pad_points<P: Point>(points: &[P]) -> Vec<P> {
let last = *points.last().unwrap();
let mut padded = Vec::with_capacity(MIN_SAMPLES_PER_FRAME);
padded.extend_from_slice(points);
padded.resize(MIN_SAMPLES_PER_FRAME, last);
padded
}
}
impl Drop for Stream {
fn drop(&mut self) {
debug!(
"IDN stream: dropping (frames sent: {}, timestamp: {})",
self.frame_count,
self.current_timestamp_us()
);
let _ = self.close();
}
}
fn param_commands(service_id: u8) -> (u8, u8) {
if service_id == 0 {
(IDNCMD_UNIT_PARAMS_REQUEST, IDNCMD_UNIT_PARAMS_RESPONSE)
} else {
(
IDNCMD_SERVICE_PARAMS_REQUEST,
IDNCMD_SERVICE_PARAMS_RESPONSE,
)
}
}
fn check_parameter_response(response: &ParameterResponse) -> Result<()> {
if !response.is_success() {
if let Some(error) = ResponseError::from_ack_code(response.result_code) {
return Err(CommunicationError::Response(error));
}
}
Ok(())
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn point_format_size_bytes() {
assert_eq!(PointFormat::Xyrgbi.size_bytes(), 8);
assert_eq!(PointFormat::XyrgbHighRes.size_bytes(), 10);
assert_eq!(PointFormat::Extended.size_bytes(), 20);
}
#[test]
fn point_format_word_count() {
assert_eq!(PointFormat::Xyrgbi.word_count(), 4);
assert_eq!(PointFormat::XyrgbHighRes.word_count(), 5);
assert_eq!(PointFormat::Extended.word_count(), 10);
}
#[test]
fn point_format_descriptors_length() {
assert_eq!(PointFormat::Xyrgbi.descriptors().len(), 8);
assert_eq!(PointFormat::XyrgbHighRes.descriptors().len(), 10);
assert_eq!(PointFormat::Extended.descriptors().len(), 20);
}
#[test]
fn point_format_descriptors_not_empty() {
for &desc in PointFormat::Xyrgbi.descriptors() {
assert_ne!(desc, 0);
}
for &desc in PointFormat::XyrgbHighRes.descriptors() {
assert_ne!(desc, 0);
}
for &desc in PointFormat::Extended.descriptors() {
assert_ne!(desc, 0);
}
}
#[test]
fn duration_calculation() {
let points = 1000u64;
let scan_speed = 30000u64;
let duration = (points * 1_000_000) / scan_speed;
assert_eq!(duration, 33333);
let points = 30000u64;
let duration = (points * 1_000_000) / scan_speed;
assert_eq!(duration, 1_000_000);
}
#[test]
fn max_points_per_packet_with_config() {
let header_size = PacketHeader::SIZE_BYTES
+ ChannelMessageHeader::SIZE_BYTES
+ ChannelConfigHeader::SIZE_BYTES
+ PointFormat::Xyrgbi.descriptors().len() * 2
+ SampleChunkHeader::SIZE_BYTES;
assert_eq!(header_size, 36);
let max_points = (MAX_UDP_PAYLOAD - header_size) / 8;
assert_eq!(max_points, 177);
}
#[test]
fn max_points_per_packet_without_config() {
let header_size = PacketHeader::SIZE_BYTES
+ ChannelMessageHeader::SIZE_BYTES
+ SampleChunkHeader::SIZE_BYTES;
assert_eq!(header_size, 16);
let max_points_xyrgbi = (MAX_UDP_PAYLOAD - header_size) / 8;
assert_eq!(max_points_xyrgbi, 179);
let max_points_highres = (MAX_UDP_PAYLOAD - header_size) / 10;
assert_eq!(max_points_highres, 143);
let max_points_extended = (MAX_UDP_PAYLOAD - header_size) / 20;
assert_eq!(max_points_extended, 71);
}
#[test]
fn content_id_construction() {
let service_id: u8 = 1;
let channel_id = ((service_id.saturating_sub(1)) as u16 & 0x3F) << 8;
assert_eq!(channel_id, 0x0000);
let service_id: u8 = 2;
let channel_id = ((service_id.saturating_sub(1)) as u16 & 0x3F) << 8;
assert_eq!(channel_id, 0x0100);
let content_id = IDNFLG_CONTENTID_CHANNELMSG | channel_id;
assert_eq!(content_id, 0x8100);
}
#[test]
fn sequence_number_wrapping() {
let mut seq: u16 = u16::MAX - 1;
seq = seq.wrapping_add(1);
assert_eq!(seq, u16::MAX);
seq = seq.wrapping_add(1);
assert_eq!(seq, 0);
seq = seq.wrapping_add(1);
assert_eq!(seq, 1);
}
#[test]
fn service_data_match_wrapping() {
let mut sdm: u8 = 254;
sdm = sdm.wrapping_add(1);
assert_eq!(sdm, 255);
sdm = sdm.wrapping_add(1);
assert_eq!(sdm, 0);
}
#[test]
fn config_needed_on_first_frame() {
assert!(Stream::needs_config(
0,
None,
PointFormat::Xyrgbi,
Some(Instant::now()),
Instant::now(),
));
}
#[test]
fn config_needed_on_format_change() {
assert!(Stream::needs_config(
42,
Some(PointFormat::Xyrgbi),
PointFormat::XyrgbHighRes,
Some(Instant::now()),
Instant::now(),
));
}
#[test]
fn config_not_needed_when_recently_sent() {
let now = Instant::now();
assert!(!Stream::needs_config(
42,
Some(PointFormat::Xyrgbi),
PointFormat::Xyrgbi,
Some(now),
now,
));
}
#[test]
fn config_needed_after_refresh_interval() {
let now = Instant::now();
let stale = now - CONFIG_REFRESH_INTERVAL - Duration::from_millis(1);
assert!(Stream::needs_config(
42,
Some(PointFormat::Xyrgbi),
PointFormat::Xyrgbi,
Some(stale),
now,
));
}
#[test]
fn config_needed_when_never_sent() {
assert!(Stream::needs_config(
42,
Some(PointFormat::Xyrgbi),
PointFormat::Xyrgbi,
None,
Instant::now(),
));
}
#[test]
fn timestamp_truncation_to_u32() {
let timestamp: u64 = 0x1_0000_ABCD; let truncated = (timestamp & 0xFFFF_FFFF) as u32;
assert_eq!(truncated, 0x0000_ABCD);
let timestamp: u64 = 0xFFFF_FFFF;
let truncated = (timestamp & 0xFFFF_FFFF) as u32;
assert_eq!(truncated, 0xFFFF_FFFF);
}
use crate::protocols::idn::protocol::PointXyrgbi;
#[test]
fn test_pad_points_pads_to_minimum() {
let points: Vec<PointXyrgbi> = (0..5)
.map(|i| PointXyrgbi::new(i as i16, i as i16, 255, 0, 0, 255))
.collect();
let padded = Stream::pad_points(&points);
assert_eq!(padded.len(), MIN_SAMPLES_PER_FRAME);
for i in 0..5 {
assert_eq!(padded[i], points[i]);
}
let last = points[4];
for p in &padded[5..] {
assert_eq!(*p, last);
}
}
#[test]
fn test_pad_points_single_point() {
let points = vec![PointXyrgbi::new(100, -200, 128, 64, 32, 255)];
let padded = Stream::pad_points(&points);
assert_eq!(padded.len(), MIN_SAMPLES_PER_FRAME);
for p in &padded {
assert_eq!(*p, points[0]);
}
}
#[test]
fn test_pad_points_nineteen() {
let points: Vec<PointXyrgbi> = (0..19)
.map(|i| PointXyrgbi::new(i as i16, 0, 0, 0, 0, 0))
.collect();
let padded = Stream::pad_points(&points);
assert_eq!(padded.len(), MIN_SAMPLES_PER_FRAME);
for i in 0..19 {
assert_eq!(padded[i], points[i]);
}
assert_eq!(padded[19], points[18]);
}
#[test]
fn test_even_distribution_300_points() {
let header_size = PacketHeader::SIZE_BYTES
+ ChannelMessageHeader::SIZE_BYTES
+ ChannelConfigHeader::SIZE_BYTES
+ PointFormat::Xyrgbi.descriptors().len() * 2
+ SampleChunkHeader::SIZE_BYTES;
let max_points_per_packet = (MAX_UDP_PAYLOAD - header_size) / PointXyrgbi::SIZE_BYTES;
assert_eq!(max_points_per_packet, 177);
let total = 300usize;
let num_packets = total.div_ceil(max_points_per_packet);
let points_to_send = total / num_packets;
assert_eq!(num_packets, 2);
assert_eq!(points_to_send, 150);
}
#[test]
fn test_even_distribution_500_points() {
let header_size = PacketHeader::SIZE_BYTES
+ ChannelMessageHeader::SIZE_BYTES
+ ChannelConfigHeader::SIZE_BYTES
+ PointFormat::Xyrgbi.descriptors().len() * 2
+ SampleChunkHeader::SIZE_BYTES;
let max_points_per_packet = (MAX_UDP_PAYLOAD - header_size) / PointXyrgbi::SIZE_BYTES;
let total = 500usize;
let num_packets = total.div_ceil(max_points_per_packet);
let points_to_send = total / num_packets;
assert_eq!(num_packets, 3);
assert_eq!(points_to_send, 166);
}
#[test]
fn test_even_distribution_small_frame() {
let header_size = PacketHeader::SIZE_BYTES
+ ChannelMessageHeader::SIZE_BYTES
+ ChannelConfigHeader::SIZE_BYTES
+ PointFormat::Xyrgbi.descriptors().len() * 2
+ SampleChunkHeader::SIZE_BYTES;
let max_points_per_packet = (MAX_UDP_PAYLOAD - header_size) / PointXyrgbi::SIZE_BYTES;
let total = 50usize;
let num_packets = total.div_ceil(max_points_per_packet);
let points_to_send = total / num_packets;
assert_eq!(num_packets, 1);
assert_eq!(points_to_send, 50);
}
#[test]
fn test_even_distribution_exact_max() {
let header_size = PacketHeader::SIZE_BYTES
+ ChannelMessageHeader::SIZE_BYTES
+ SampleChunkHeader::SIZE_BYTES;
let max_points_per_packet = (MAX_UDP_PAYLOAD - header_size) / PointXyrgbi::SIZE_BYTES;
assert_eq!(max_points_per_packet, 179);
let total = 179usize;
let num_packets = total.div_ceil(max_points_per_packet);
let points_to_send = total / num_packets;
assert_eq!(num_packets, 1);
assert_eq!(points_to_send, 179);
}
#[test]
fn timestamp_accumulator_avoids_floor_drift() {
let pps = 30_000u64;
let points_per_chunk = 179u64;
let chunks = 10_000u64;
let per_chunk_floored = points_per_chunk * 1_000_000 / pps;
let naive_sum = per_chunk_floored * chunks;
let total_points = points_per_chunk * chunks;
let exact = total_points * 1_000_000 / pps;
assert!(exact >= naive_sum);
assert!(
exact - naive_sum > 1_000,
"expected measurable drift avoided, got {} us",
exact - naive_sum
);
}
use crate::protocols::idn::dac::{ServiceInfo, ServiceType};
use std::net::UdpSocket;
fn connected_stream() -> (Stream, UdpSocket) {
let receiver = UdpSocket::bind("127.0.0.1:0").unwrap();
let addr = receiver.local_addr().unwrap();
let mut server = ServerInfo::new([0u8; 16], "test".to_string(), (1, 0), 0);
server.addresses.push(addr);
server.services.push(ServiceInfo {
service_id: 1,
service_type: ServiceType::LaserProjector,
name: "laser".to_string(),
flags: 0,
relay_number: 0,
});
let stream = connect(&server, 1).unwrap();
(stream, receiver)
}
fn read_headers(buf: &[u8]) -> (PacketHeader, ChannelMessageHeader) {
let mut cursor = buf;
let ph: PacketHeader = cursor.read_bytes().unwrap();
let cmh: ChannelMessageHeader = cursor.read_bytes().unwrap();
(ph, cmh)
}
#[test]
fn frame_mode_fragmentation_header_layout() {
let (mut stream, receiver) = connected_stream();
receiver
.set_read_timeout(Some(Duration::from_millis(500)))
.unwrap();
stream.set_frame_mode(FrameMode::Frame);
stream.set_scan_speed(30_000);
stream.frame_count = 1;
stream.previous_format = Some(PointFormat::Xyrgbi);
stream.last_config_time = Some(Instant::now());
stream.timestamp_anchor_us = 1_000_000;
stream.points_since_anchor = 0;
stream.anchor_pps = 30_000;
let points: Vec<PointXyrgbi> = (0..400)
.map(|i| PointXyrgbi::new(i as i16, 0, 0, 0, 0, 0))
.collect();
stream.write_frame(&points).unwrap();
let mut buf = [0u8; 2048];
let (n, _) = receiver.recv_from(&mut buf).unwrap();
let (ph0, cmh0) = read_headers(&buf[..n]);
assert_eq!(ph0.command, IDNCMD_RT_CNLMSG);
let chan = IDNFLG_CONTENTID_CHANNELMSG; assert_eq!(
cmh0.content_id,
chan | IDNVAL_CNKTYPE_LPGRF_FRAME_FIRST as u16
);
assert_eq!(cmh0.timestamp, 1_000_000);
let mut cursor = &buf[PacketHeader::SIZE_BYTES + ChannelMessageHeader::SIZE_BYTES..n];
let sch: SampleChunkHeader = cursor.read_bytes().unwrap();
assert_eq!(sch.duration_us(), 13_333);
let (n, _) = receiver.recv_from(&mut buf).unwrap();
let (_, cmh1) = read_headers(&buf[..n]);
assert_eq!(
cmh1.content_id,
chan | IDNVAL_CNKTYPE_LPGRF_FRAME_SEQUEL as u16
);
assert_eq!(cmh1.timestamp, 1_000_000, "sequel shares frame timestamp");
assert_eq!(cmh1.content_id & IDNFLG_CONTENTID_CONFIG_LSTFRG, 0);
assert_eq!(
cmh1.total_size as usize,
ChannelMessageHeader::SIZE_BYTES + 133 * PointXyrgbi::SIZE_BYTES
);
let (n, _) = receiver.recv_from(&mut buf).unwrap();
let (_, cmh2) = read_headers(&buf[..n]);
assert_eq!(
cmh2.content_id,
chan | IDNVAL_CNKTYPE_LPGRF_FRAME_SEQUEL as u16 | IDNFLG_CONTENTID_CONFIG_LSTFRG
);
assert_eq!(cmh2.timestamp, 1_000_000);
}
#[test]
fn frame_mode_single_fragment_uses_frame_type() {
let (mut stream, receiver) = connected_stream();
receiver
.set_read_timeout(Some(Duration::from_millis(500)))
.unwrap();
stream.set_frame_mode(FrameMode::Frame);
stream.set_scan_speed(30_000);
stream.frame_count = 1;
stream.previous_format = Some(PointFormat::Xyrgbi);
stream.last_config_time = Some(Instant::now());
stream.timestamp_anchor_us = 0;
stream.points_since_anchor = 0;
stream.anchor_pps = 30_000;
let points: Vec<PointXyrgbi> = (0..60)
.map(|i| PointXyrgbi::new(i as i16, 0, 0, 0, 0, 0))
.collect();
stream.write_frame(&points).unwrap();
let mut buf = [0u8; 2048];
let (n, _) = receiver.recv_from(&mut buf).unwrap();
let (_, cmh) = read_headers(&buf[..n]);
assert_eq!(
cmh.content_id,
IDNFLG_CONTENTID_CHANNELMSG | IDNVAL_CNKTYPE_LPGRF_FRAME as u16
);
assert_eq!(cmh.content_id & IDNFLG_CONTENTID_CONFIG_LSTFRG, 0);
let mut cursor = &buf[PacketHeader::SIZE_BYTES + ChannelMessageHeader::SIZE_BYTES..n];
let sch: SampleChunkHeader = cursor.read_bytes().unwrap();
assert_eq!(sch.duration_us(), 2000);
}
#[test]
fn wave_wire_timestamps_advance_exactly() {
let (mut stream, receiver) = connected_stream();
receiver
.set_read_timeout(Some(Duration::from_millis(500)))
.unwrap();
stream.set_scan_speed(30_000);
let pts: Vec<PointXyrgbi> = (0..100)
.map(|i| PointXyrgbi::new(i as i16, 0, 0, 0, 0, 0))
.collect();
let mut buf = [0u8; 2048];
stream.write_frame(&pts).unwrap();
let (n, _) = receiver.recv_from(&mut buf).unwrap();
let (_, cmh0) = read_headers(&buf[..n]);
stream.write_frame(&pts).unwrap();
let (n, _) = receiver.recv_from(&mut buf).unwrap();
let (_, cmh1) = read_headers(&buf[..n]);
assert_eq!(cmh1.timestamp.wrapping_sub(cmh0.timestamp), 3333);
}
#[test]
fn config_resent_after_refresh_interval_on_wire() {
let (mut stream, receiver) = connected_stream();
receiver
.set_read_timeout(Some(Duration::from_millis(500)))
.unwrap();
stream.set_scan_speed(30_000);
let pts: Vec<PointXyrgbi> = (0..40)
.map(|i| PointXyrgbi::new(i as i16, 0, 0, 0, 0, 0))
.collect();
let mut buf = [0u8; 2048];
stream.write_frame(&pts).unwrap();
let (n, _) = receiver.recv_from(&mut buf).unwrap();
let (_, cmh0) = read_headers(&buf[..n]);
assert_ne!(cmh0.content_id & IDNFLG_CONTENTID_CONFIG_LSTFRG, 0);
stream.write_frame(&pts).unwrap();
let (n, _) = receiver.recv_from(&mut buf).unwrap();
let (_, cmh1) = read_headers(&buf[..n]);
assert_eq!(cmh1.content_id & IDNFLG_CONTENTID_CONFIG_LSTFRG, 0);
stream.last_config_time =
Some(Instant::now() - CONFIG_REFRESH_INTERVAL - Duration::from_millis(5));
stream.write_frame(&pts).unwrap();
let (n, _) = receiver.recv_from(&mut buf).unwrap();
let (_, cmh2) = read_headers(&buf[..n]);
assert_ne!(
cmh2.content_id & IDNFLG_CONTENTID_CONFIG_LSTFRG,
0,
"config must be resent after the refresh interval"
);
}
#[test]
fn highres_format_writes_16bit_config_and_samples() {
use crate::protocols::idn::protocol::PointXyrgbHighRes;
let (mut stream, receiver) = connected_stream();
receiver
.set_read_timeout(Some(Duration::from_millis(500)))
.unwrap();
stream.set_point_format(PointFormat::XyrgbHighRes);
stream.set_scan_speed(30_000);
let pts: Vec<PointXyrgbHighRes> = (0..40)
.map(|i| PointXyrgbHighRes::new(i, 0, 0x1234, 0x5678, 0x9abc))
.collect();
stream.write_frame(&pts).unwrap();
let mut buf = [0u8; 2048];
let (n, _) = receiver.recv_from(&mut buf).unwrap();
let cfg_off = PacketHeader::SIZE_BYTES + ChannelMessageHeader::SIZE_BYTES;
let word_count = buf[cfg_off] as usize;
assert_eq!(word_count, 5, "hi-res XYRGB descriptor has 5 words");
let samples_off = cfg_off + 4 + word_count * 4 + SampleChunkHeader::SIZE_BYTES;
let sample_bytes = n - samples_off;
assert_eq!(
sample_bytes % PointXyrgbHighRes::SIZE_BYTES,
0,
"sample region must be a whole number of 10-byte hi-res samples"
);
assert_eq!(sample_bytes / PointXyrgbHighRes::SIZE_BYTES, 40);
}
#[test]
fn recv_matching_discards_stale_datagrams() {
let (mut stream, server) = connected_stream();
let client_addr = stream.socket.local_addr().unwrap();
let mut stale = Vec::new();
stale
.write_bytes(PacketHeader {
command: 0xFF,
flags: 0,
sequence: 1,
})
.unwrap();
server.send_to(&stale, client_addr).unwrap();
let mut good = Vec::new();
good.write_bytes(PacketHeader {
command: IDNCMD_PING_RESPONSE,
flags: 0,
sequence: 7,
})
.unwrap();
server.send_to(&good, client_addr).unwrap();
let header = stream
.recv_response_header(Duration::from_millis(500), IDNCMD_PING_RESPONSE, 7)
.expect("should skip the stale datagram and match the ping response");
assert_eq!(header.command, IDNCMD_PING_RESPONSE);
assert_eq!(header.sequence, 7);
}
}