use std::sync::atomic::{AtomicBool, AtomicU32, AtomicU64, Ordering};
use std::sync::mpsc::{self, Receiver, Sender};
use std::sync::{Arc, Mutex};
use std::time::Duration;
use crate::backend::{BackendKind, Error, Result};
use crate::config::StreamConfig;
use crate::device::{DacCapabilities, DacInfo, DacType};
use crate::discovery::DacDiscovery;
use crate::point::LaserPoint;
use crate::reconnect::{ReconnectPolicy, ReconnectTarget};
pub(crate) mod chunk_producer;
#[cfg(feature = "serde")]
use serde::{Deserialize, Serialize};
#[derive(Clone, Copy, Debug, Default, PartialEq, Eq, PartialOrd, Ord, Hash)]
#[cfg_attr(feature = "serde", derive(Serialize, Deserialize))]
pub struct StreamInstant(pub u64);
impl StreamInstant {
pub fn new(points: u64) -> Self {
Self(points)
}
pub fn points(&self) -> u64 {
self.0
}
pub fn as_seconds(&self, pps: u32) -> f64 {
self.0 as f64 / pps as f64
}
#[inline]
pub fn as_secs_f64(&self, pps: u32) -> f64 {
self.as_seconds(pps)
}
pub fn from_seconds(seconds: f64, pps: u32) -> Self {
Self((seconds * pps as f64) as u64)
}
pub fn add_points(&self, points: u64) -> Self {
Self(self.0.saturating_add(points))
}
pub fn sub_points(&self, points: u64) -> Self {
Self(self.0.saturating_sub(points))
}
}
impl std::ops::Add<u64> for StreamInstant {
type Output = Self;
fn add(self, rhs: u64) -> Self::Output {
self.add_points(rhs)
}
}
impl std::ops::Sub<u64> for StreamInstant {
type Output = Self;
fn sub(self, rhs: u64) -> Self::Output {
self.sub_points(rhs)
}
}
impl std::ops::AddAssign<u64> for StreamInstant {
fn add_assign(&mut self, rhs: u64) {
self.0 = self.0.saturating_add(rhs);
}
}
impl std::ops::SubAssign<u64> for StreamInstant {
fn sub_assign(&mut self, rhs: u64) {
self.0 = self.0.saturating_sub(rhs);
}
}
#[derive(Clone, Debug)]
pub struct ChunkRequest {
pub start: StreamInstant,
pub pps: u32,
pub target_points: usize,
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub enum ChunkResult {
Filled(usize),
Starved,
End,
}
#[derive(Clone, Debug)]
pub struct StreamStatus {
pub connected: bool,
pub scheduled_ahead_points: u64,
pub device_queued_points: Option<u64>,
pub stats: Option<StreamStats>,
}
#[derive(Clone, Debug, Default)]
pub struct StreamStats {
pub underrun_count: u64,
pub late_chunk_count: u64,
pub reconnect_count: u64,
pub chunks_written: u64,
pub points_written: u64,
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub enum RunExit {
Stopped,
ProducerEnded,
Disconnected,
}
#[derive(Debug, Clone, Copy)]
pub(crate) enum ControlMsg {
Arm,
Disarm,
Stop,
}
#[derive(Clone)]
pub struct StreamControl {
inner: Arc<StreamControlInner>,
}
struct StreamControlInner {
armed: AtomicBool,
stop_requested: AtomicBool,
control_tx: Mutex<Sender<ControlMsg>>,
color_delay_micros: AtomicU64,
pps: AtomicU32,
}
impl StreamControl {
pub(crate) fn new(control_tx: Sender<ControlMsg>, color_delay: Duration, pps: u32) -> Self {
Self {
inner: Arc::new(StreamControlInner {
armed: AtomicBool::new(false),
stop_requested: AtomicBool::new(false),
control_tx: Mutex::new(control_tx),
color_delay_micros: AtomicU64::new(color_delay.as_micros() as u64),
pps: AtomicU32::new(pps),
}),
}
}
pub fn arm(&self) -> Result<()> {
self.inner.armed.store(true, Ordering::SeqCst);
if let Ok(tx) = self.inner.control_tx.lock() {
let _ = tx.send(ControlMsg::Arm);
}
Ok(())
}
pub fn disarm(&self) -> Result<()> {
self.inner.armed.store(false, Ordering::SeqCst);
if let Ok(tx) = self.inner.control_tx.lock() {
let _ = tx.send(ControlMsg::Disarm);
}
Ok(())
}
pub fn is_armed(&self) -> bool {
self.inner.armed.load(Ordering::SeqCst)
}
pub fn set_color_delay(&self, delay: Duration) {
self.inner
.color_delay_micros
.store(delay.as_micros() as u64, Ordering::SeqCst);
}
pub fn color_delay(&self) -> Duration {
Duration::from_micros(self.inner.color_delay_micros.load(Ordering::SeqCst))
}
pub fn set_pps(&self, pps: u32) {
self.inner.pps.store(pps, Ordering::SeqCst);
}
pub fn pps(&self) -> u32 {
self.inner.pps.load(Ordering::SeqCst)
}
pub fn stop(&self) -> Result<()> {
self.inner.stop_requested.store(true, Ordering::SeqCst);
if let Ok(tx) = self.inner.control_tx.lock() {
let _ = tx.send(ControlMsg::Stop);
}
Ok(())
}
pub fn is_stop_requested(&self) -> bool {
self.inner.stop_requested.load(Ordering::SeqCst)
}
}
#[derive(Default)]
struct StreamState {
stats: StreamStats,
}
impl StreamState {
fn new() -> Self {
Self::default()
}
}
pub struct Stream {
info: DacInfo,
backend: Option<BackendKind>,
config: StreamConfig,
control: StreamControl,
control_rx: Receiver<ControlMsg>,
state: StreamState,
pub(crate) reconnect_policy: Option<ReconnectPolicy>,
pub(crate) reconnect_target: Option<ReconnectTarget>,
}
impl Stream {
pub(crate) fn with_backend(info: DacInfo, backend: BackendKind, config: StreamConfig) -> Self {
let (control_tx, control_rx) = mpsc::channel();
let color_delay = config.color_delay;
let pps = config.pps;
Self {
info,
backend: Some(backend),
config,
control: StreamControl::new(control_tx, color_delay, pps),
control_rx,
state: StreamState::new(),
reconnect_policy: None,
reconnect_target: None,
}
}
pub fn info(&self) -> &DacInfo {
&self.info
}
pub fn config(&self) -> &StreamConfig {
&self.config
}
pub fn control(&self) -> StreamControl {
self.control.clone()
}
pub fn status(&self) -> Result<StreamStatus> {
let buffered = self.estimate_buffer_points();
Ok(StreamStatus {
connected: self.backend.as_ref().is_some_and(|b| b.is_connected()),
scheduled_ahead_points: buffered,
device_queued_points: Some(buffered),
stats: Some(self.state.stats.clone()),
})
}
fn shutdown_backend(&mut self) {
let _ = self.control.disarm();
let _ = self.control.stop();
if let Some(b) = &mut self.backend {
let _ = b.set_shutter(false);
let _ = b.stop();
}
}
pub fn stop(&mut self) -> Result<()> {
self.shutdown_backend();
if let Some(b) = &mut self.backend {
b.disconnect()?;
}
Ok(())
}
pub fn into_dac(mut self) -> (Dac, StreamStats) {
self.shutdown_backend();
let backend = self.backend.take();
let stats = self.state.stats.clone();
let reconnect_target = self
.reconnect_target
.take()
.or_else(|| self.reconnect_policy.take().map(|p| p.target));
let dac = Dac {
info: self.info.clone(),
backend,
reconnect_target,
};
(dac, stats)
}
pub fn run<F, E>(mut self, producer: F, on_error: E) -> Result<RunExit>
where
F: FnMut(&ChunkRequest, &mut [LaserPoint]) -> ChunkResult + Send + 'static,
E: FnMut(Error) + Send + 'static,
{
use crate::presentation::driver::{self, DriverInputs, SourceOwned};
use crate::presentation::FrameSessionMetrics;
let backend = self
.backend
.take()
.ok_or_else(|| Error::disconnected("backend already consumed"))?;
if backend.is_frame_swap() {
return Err(Error::invalid_config(
"Stream::run is FIFO-only; use start_frame_session for frame-swap DACs",
));
}
let max_points = self.info.caps.max_points_per_chunk;
let chunk_producer = chunk_producer::ChunkProducer::new(
producer,
self.control.clone(),
self.config.idle_policy.clone(),
self.config.startup_blank,
max_points,
);
let validator = Self::build_reconnect_validator();
let metrics = FrameSessionMetrics::new(true);
let (_dummy_tx, dummy_rx) = mpsc::channel();
let control_rx = std::mem::replace(&mut self.control_rx, dummy_rx);
driver::run(DriverInputs {
backend,
source: SourceOwned::Fifo(Box::new(chunk_producer)),
control: self.control.clone(),
control_rx,
metrics,
reconnect_policy: self.reconnect_policy.take(),
validator,
error_sink: Box::new(on_error),
target_buffer: self.config.target_buffer,
drain_timeout: self.config.drain_timeout,
pending_frame: None,
clock: DriverInputs::system_clock(),
})
}
fn build_reconnect_validator() -> crate::presentation::driver::ReconnectValidator {
Box::new(
move |_info: &DacInfo, new_backend: &BackendKind, pps: u32| {
if new_backend.is_frame_swap() {
log::error!("reconnected device is frame-swap, incompatible with streaming");
return Err(RunExit::Disconnected);
}
if Dac::validate_pps(new_backend.caps(), pps).is_err() {
log::error!("reconnected device PPS range incompatible with stream config");
return Err(RunExit::Disconnected);
}
Ok(())
},
)
}
fn estimate_buffer_points(&self) -> u64 {
let pps = self.config.pps;
let now = std::time::Instant::now();
self.backend
.as_ref()
.and_then(|b| b.estimator())
.map_or(0, |e| e.estimated_fullness(now, pps))
}
}
impl Drop for Stream {
fn drop(&mut self) {
let _ = self.stop();
}
}
pub struct Dac {
info: DacInfo,
backend: Option<BackendKind>,
pub(crate) reconnect_target: Option<ReconnectTarget>,
}
impl Dac {
pub fn new(info: DacInfo, backend: BackendKind) -> Self {
Self {
info,
backend: Some(backend),
reconnect_target: None,
}
}
pub fn with_discovery_factory<F>(mut self, factory: F) -> Self
where
F: Fn() -> DacDiscovery + Send + 'static,
{
match self.reconnect_target {
Some(ref mut target) => {
target.discovery_factory = Some(Box::new(factory));
}
None => {
self.reconnect_target = Some(ReconnectTarget {
device_id: self.info.id.clone(),
discovery_factory: Some(Box::new(factory)),
});
}
}
self
}
pub fn info(&self) -> &DacInfo {
&self.info
}
pub fn id(&self) -> &str {
&self.info.id
}
pub fn name(&self) -> &str {
&self.info.name
}
pub fn kind(&self) -> &DacType {
&self.info.kind
}
pub fn caps(&self) -> &DacCapabilities {
&self.info.caps
}
pub fn has_backend(&self) -> bool {
self.backend.is_some()
}
pub(crate) fn into_backend(mut self) -> Option<BackendKind> {
self.backend.take()
}
pub fn is_connected(&self) -> bool {
self.backend.as_ref().is_some_and(|b| b.is_connected())
}
pub fn start_stream(mut self, mut cfg: StreamConfig) -> Result<(Stream, DacInfo)> {
let reconnect_config = cfg.reconnect.take();
let mut backend = self.backend.take().ok_or_else(|| {
Error::invalid_config("device backend has already been used for a stream")
})?;
if backend.is_frame_swap() {
return Err(Error::invalid_config(
"streaming is not supported on frame-swap DACs (e.g. Helios); \
use start_frame_session() instead",
));
}
let cfg = Self::apply_backend_buffer_defaults(&self.info, cfg);
Self::validate_pps(&self.info.caps, cfg.pps)?;
if !backend.is_connected() {
backend.connect()?;
}
let mut stream = Stream::with_backend(self.info.clone(), backend, cfg);
stream.reconnect_target = self.reconnect_target.take();
if let Some(rc) = reconnect_config {
let target = stream.reconnect_target.take().ok_or_else(|| {
Error::invalid_config("reconnect requires a reconnect target — use open_device(), open_device_with(), or Dac::with_discovery_factory()")
})?;
stream.reconnect_policy = Some(ReconnectPolicy::new(rc, target));
}
Ok((stream, self.info))
}
fn apply_backend_buffer_defaults(info: &DacInfo, mut cfg: StreamConfig) -> StreamConfig {
if cfg.target_buffer == StreamConfig::DEFAULT_TARGET_BUFFER {
cfg.target_buffer =
StreamConfig::default_target_buffer_for(&info.kind, &info.caps.output_model);
}
cfg
}
fn validate_pps(caps: &DacCapabilities, pps: u32) -> Result<()> {
if pps < caps.pps_min || pps > caps.pps_max {
return Err(Error::invalid_config(format!(
"PPS {} is outside device range [{}, {}]",
pps, caps.pps_min, caps.pps_max
)));
}
Ok(())
}
pub fn start_frame_session(
mut self,
mut config: crate::presentation::FrameSessionConfig,
) -> Result<(crate::presentation::FrameSession, DacInfo)> {
let reconnect_config = config.reconnect.take();
let backend = self.backend.take().ok_or_else(|| {
Error::invalid_config("device backend has already been used for a session")
})?;
Self::validate_pps(backend.caps(), config.pps)?;
let reconnect_policy = match reconnect_config {
Some(rc) => {
let target = self.reconnect_target.take().ok_or_else(|| {
Error::invalid_config("reconnect requires a reconnect target — use open_device(), open_device_with(), or Dac::with_discovery_factory()")
})?;
Some(ReconnectPolicy::new(rc, target))
}
None => None,
};
let session = crate::presentation::FrameSession::start(backend, config, reconnect_policy)?;
Ok((session, self.info))
}
}
#[cfg(test)]
mod tests;