use std::time::Duration;
use crate::backend::{BackendKind, WriteOutcome};
use crate::device::DacInfo;
use super::super::content_source::{ContentSourceKind, FifoContentSource};
use super::{LoopCtx, OutputModelAdapter, StepOutcome};
const CHUNK_POINTS: usize = 512;
const IDLE_SLEEP: Duration = Duration::from_millis(2);
fn fifo_source<'a>(source: &'a mut ContentSourceKind<'_>) -> &'a mut dyn FifoContentSource {
match source {
ContentSourceKind::Fifo(s) => &mut **s,
ContentSourceKind::Frame(_) => {
unreachable!("BlockingFifoAdapter requires a Fifo content source")
}
}
}
pub(crate) struct BlockingFifoAdapter {
chunk_points: usize,
has_retain: bool,
was_armed: bool,
}
impl BlockingFifoAdapter {
pub fn new(backend: &BackendKind) -> Self {
Self {
chunk_points: CHUNK_POINTS.min(backend.caps().max_points_per_chunk),
has_retain: false,
was_armed: false,
}
}
}
impl OutputModelAdapter for BlockingFifoAdapter {
fn step(&mut self, ctx: &mut LoopCtx<'_>) -> StepOutcome {
if !ctx.is_armed {
self.was_armed = false;
if let Err(stopped) = ctx.sleep_with_control_check(IDLE_SLEEP) {
return stopped;
}
return StepOutcome::Continue;
}
if !self.was_armed {
self.was_armed = true;
if let Err(e) = ctx.backend.reset_device_buffer() {
log::debug!("reset_device_buffer on re-arm failed (non-fatal): {e}");
}
fifo_source(&mut ctx.source).discard_cached();
self.has_retain = false;
}
let pps = ctx.pps;
let chunk = self.chunk_points;
fifo_source(&mut ctx.source).reserve_buf(chunk);
if !self.has_retain {
if fifo_source(&mut ctx.source)
.produce_chunk(chunk, pps, ctx.is_armed)
.is_empty()
{
ctx.sleep_and_mark_activity(Duration::from_millis(1));
return StepOutcome::Continue;
}
self.has_retain = true;
}
let (n, outcome) = match fifo_source(&mut ctx.source).cached_slice() {
Some(slice) => (slice.len(), ctx.backend.try_write(pps, slice)),
None => {
self.has_retain = false;
return StepOutcome::Continue;
}
};
match outcome {
Ok(WriteOutcome::Written) => {
ctx.metrics.mark_write_success();
fifo_source(&mut ctx.source).commit_written(n, ctx.is_armed);
self.has_retain = false;
}
Ok(WriteOutcome::WouldBlock) => {
if ctx.control.is_stop_requested() {
return StepOutcome::Stopped;
}
ctx.sleep_and_mark_activity(Duration::from_micros(100));
}
Err(e) if e.is_stopped() => return StepOutcome::Stopped,
Err(e) if e.is_disconnected() => {
(ctx.error_sink)(e);
return StepOutcome::Disconnected;
}
Err(e) => {
log::warn!("write error, disconnecting backend: {e}");
let _ = ctx.backend.disconnect();
(ctx.error_sink)(e);
return StepOutcome::Disconnected;
}
}
StepOutcome::Continue
}
fn on_reconnect(&mut self, info: &DacInfo, _backend: &mut BackendKind) {
self.chunk_points = CHUNK_POINTS.min(info.caps.max_points_per_chunk);
self.has_retain = false;
self.was_armed = false;
}
fn drain_and_blank(&mut self, ctx: &mut LoopCtx<'_>, timeout: Duration) {
super::drain_via_estimator(ctx, timeout);
let _ = ctx.backend.reset_device_buffer();
super::blank_and_close_shutter(ctx);
}
}
#[cfg(test)]
mod tests {
use std::sync::atomic::{AtomicUsize, Ordering};
use std::sync::mpsc;
use std::sync::{Arc, Mutex};
use crate::backend::{BackendKind, DacBackend, FifoBackend, WriteOutcome};
use crate::buffer_estimate::{BufferEstimator, SoftwareDecayEstimator};
use crate::config::IdlePolicy;
use crate::device::{DacCapabilities, DacType, OutputModel};
use crate::error::Result as DacResult;
use crate::point::LaserPoint;
use crate::presentation::content_source::{ContentSourceKind, FifoContentSource};
use crate::presentation::engine::PresentationEngine;
use crate::presentation::output_model::{
drain_via_estimator, LoopCtx, OutputModelAdapter, StepOutcome, SystemClock,
};
use crate::presentation::session::FrameSessionMetrics;
use crate::presentation::slice_pipeline::SlicePipeline;
use crate::presentation::{Frame, TransitionPlan};
use crate::stream::{ControlMsg, StreamControl};
use super::{BlockingFifoAdapter, CHUNK_POINTS};
struct FakeBlockingFifo {
caps: DacCapabilities,
writes: Arc<Mutex<Vec<Vec<LaserPoint>>>>,
resets: Arc<AtomicUsize>,
estimator: Box<dyn BufferEstimator>,
}
impl DacBackend for FakeBlockingFifo {
fn dac_type(&self) -> DacType {
DacType::Custom("FakeBlockingFifo".into())
}
fn caps(&self) -> &DacCapabilities {
&self.caps
}
fn connect(&mut self) -> DacResult<()> {
Ok(())
}
fn disconnect(&mut self) -> DacResult<()> {
Ok(())
}
fn is_connected(&self) -> bool {
true
}
fn stop(&mut self) -> DacResult<()> {
Ok(())
}
fn set_shutter(&mut self, _open: bool) -> DacResult<()> {
Ok(())
}
}
impl FifoBackend for FakeBlockingFifo {
fn try_write_points(
&mut self,
_pps: u32,
points: &[LaserPoint],
) -> DacResult<WriteOutcome> {
self.writes.lock().unwrap().push(points.to_vec());
Ok(WriteOutcome::Written)
}
fn estimator(&self) -> &dyn BufferEstimator {
&*self.estimator
}
fn reset_device_buffer(&mut self) -> DacResult<()> {
self.resets.fetch_add(1, Ordering::SeqCst);
Ok(())
}
}
fn caps() -> DacCapabilities {
DacCapabilities {
pps_min: 1,
pps_max: 35_000,
max_points_per_chunk: 4096,
output_model: OutputModel::BlockingFifo,
}
}
fn frame_with_points(n: usize) -> Frame {
let pts: Vec<LaserPoint> = (0..n)
.map(|i| LaserPoint::new(i as f32 * 0.0001, 0.0, 1000, 1000, 1000, 1000))
.collect();
Frame::new(pts)
}
struct Harness {
backend: BackendKind,
adapter: BlockingFifoAdapter,
pipeline: SlicePipeline,
control: StreamControl,
rx: mpsc::Receiver<ControlMsg>,
metrics: FrameSessionMetrics,
shutter: bool,
writes: Arc<Mutex<Vec<Vec<LaserPoint>>>>,
resets: Arc<AtomicUsize>,
}
impl Harness {
fn new() -> Self {
let writes = Arc::new(Mutex::new(Vec::new()));
let resets = Arc::new(AtomicUsize::new(0));
let backend = FakeBlockingFifo {
caps: caps(),
writes: Arc::clone(&writes),
resets: Arc::clone(&resets),
estimator: Box::new(SoftwareDecayEstimator::new()),
};
let backend = BackendKind::Fifo(Box::new(backend));
let adapter = BlockingFifoAdapter::new(&backend);
let mut engine =
PresentationEngine::new(Box::new(|_, _| TransitionPlan::Transition(Vec::new())));
engine.set_pending(frame_with_points(474));
let pipeline = SlicePipeline::new(engine, 0, None, IdlePolicy::Blank, 0);
let (tx, rx) = mpsc::channel::<ControlMsg>();
let control = StreamControl::new(tx, std::time::Duration::ZERO, 30_000);
let metrics = FrameSessionMetrics::new(true);
Self {
backend,
adapter,
pipeline,
control,
rx,
metrics,
shutter: false,
writes,
resets,
}
}
fn step(&mut self, is_armed: bool) -> StepOutcome {
let source = ContentSourceKind::Fifo(&mut self.pipeline as &mut dyn FifoContentSource);
let mut ctx = LoopCtx {
backend: &mut self.backend,
source,
control: &self.control,
control_rx: &self.rx,
metrics: &self.metrics,
shutter_open: &mut self.shutter,
error_sink: &mut |_| {},
target_buffer: std::time::Duration::from_millis(20),
pps: 30_000,
is_armed,
clock: &SystemClock,
};
self.adapter.step(&mut ctx)
}
fn drain_and_blank(&mut self, timeout: std::time::Duration) {
let source = ContentSourceKind::Fifo(&mut self.pipeline as &mut dyn FifoContentSource);
let mut ctx = LoopCtx {
backend: &mut self.backend,
source,
control: &self.control,
control_rx: &self.rx,
metrics: &self.metrics,
shutter_open: &mut self.shutter,
error_sink: &mut |_| {},
target_buffer: std::time::Duration::from_millis(20),
pps: 30_000,
is_armed: true,
clock: &SystemClock,
};
self.adapter.drain_and_blank(&mut ctx, timeout);
}
}
#[test]
fn armed_step_writes_fixed_chunk_size() {
let mut h = Harness::new();
assert!(matches!(h.step(true), StepOutcome::Continue));
let writes = h.writes.lock().unwrap();
assert_eq!(writes.len(), 1, "one blocking write per armed step");
assert_eq!(
writes[0].len(),
CHUNK_POINTS,
"chunk should be the fixed CHUNK_POINTS, not an estimator-derived trickle"
);
}
#[test]
fn disarmed_step_issues_no_write() {
let mut h = Harness::new();
assert!(matches!(h.step(false), StepOutcome::Continue));
assert_eq!(
h.writes.lock().unwrap().len(),
0,
"no blocking write while disarmed (would wedge on a disabled device)"
);
}
#[test]
fn rearm_clears_ring_exactly_once_per_edge() {
let mut h = Harness::new();
h.step(true);
assert_eq!(h.resets.load(Ordering::SeqCst), 1);
h.step(true);
assert_eq!(h.resets.load(Ordering::SeqCst), 1);
h.step(false);
h.step(true);
assert_eq!(
h.resets.load(Ordering::SeqCst),
2,
"ring should be cleared once per re-arm, not per armed step"
);
}
#[test]
fn drain_and_blank_clears_ring_before_trailing_blank() {
let mut h = Harness::new();
h.step(true); let before = h.resets.load(Ordering::SeqCst);
h.drain_and_blank(std::time::Duration::ZERO);
assert!(
h.resets.load(Ordering::SeqCst) > before,
"ring must be cleared before the trailing blank write so it can't wedge"
);
}
struct CapturingEstimator {
queried_pps: Arc<Mutex<Vec<u32>>>,
}
impl BufferEstimator for CapturingEstimator {
fn estimated_fullness(&self, _now: std::time::Instant, pps: u32) -> u64 {
self.queried_pps.lock().unwrap().push(pps);
0
}
}
#[test]
fn drain_via_estimator_polls_at_device_clamped_rate() {
const DEVICE_MAX: u32 = 30_000; const CONFIGURED_PPS: u32 = 33_000;
let queried_pps = Arc::new(Mutex::new(Vec::new()));
let mut device_caps = caps();
device_caps.pps_max = DEVICE_MAX;
let backend = FakeBlockingFifo {
caps: device_caps,
writes: Arc::new(Mutex::new(Vec::new())),
resets: Arc::new(AtomicUsize::new(0)),
estimator: Box::new(CapturingEstimator {
queried_pps: Arc::clone(&queried_pps),
}),
};
let mut backend = BackendKind::Fifo(Box::new(backend));
let mut engine =
PresentationEngine::new(Box::new(|_, _| TransitionPlan::Transition(Vec::new())));
engine.set_pending(frame_with_points(64));
let mut pipeline = SlicePipeline::new(engine, 0, None, IdlePolicy::Blank, 0);
let (tx, rx) = mpsc::channel::<ControlMsg>();
let control = StreamControl::new(tx, std::time::Duration::ZERO, CONFIGURED_PPS);
let metrics = FrameSessionMetrics::new(true);
let mut shutter = true;
let source = ContentSourceKind::Fifo(&mut pipeline as &mut dyn FifoContentSource);
let mut ctx = LoopCtx {
backend: &mut backend,
source,
control: &control,
control_rx: &rx,
metrics: &metrics,
shutter_open: &mut shutter,
error_sink: &mut |_| {},
target_buffer: std::time::Duration::from_millis(20),
pps: CONFIGURED_PPS,
is_armed: true,
clock: &SystemClock,
};
drain_via_estimator(&mut ctx, std::time::Duration::from_secs(1));
let polled = queried_pps.lock().unwrap();
assert!(
!polled.is_empty(),
"drain must poll the estimator at least once"
);
assert!(
polled.iter().all(|&p| p == DEVICE_MAX),
"drain must poll at the device-clamped rate {DEVICE_MAX}, not the \
configured {CONFIGURED_PPS}; got {polled:?}"
);
}
}