use crate::codec::{
changed_runs, fit_to_width, replace_cached_range, sanitize_text, sanitize_to_width,
};
use crate::config::{
DisplaySettings, MAX_MARQUEE_CHARS, MAX_MARQUEE_CPS, MAX_QUEUED_RAW_BYTES, VfdConfig,
};
use crate::error::{ConfigError, Result, VfdError};
use crate::vfd::Vfd;
use serialport::SerialPort;
use std::io::Write;
use std::sync::mpsc::{self, Receiver, SyncSender};
use std::thread::{self, JoinHandle};
use std::time::{Duration, Instant};
type Ack = mpsc::SyncSender<Result<()>>;
#[derive(Debug)]
enum Cmd {
Clear {
ack: Ack,
},
PrintLine {
line: u8,
text: String,
ack: Ack,
},
PrintLineDiff {
line: u8,
text: String,
ack: Ack,
},
PrintAt {
x: u8,
y: u8,
text: String,
ack: Ack,
},
WriteRaw {
bytes: Vec<u8>,
ack: Ack,
},
SetMarqueeText {
text: String,
ack: Ack,
},
StartMarquee {
line: u8,
cps: u32,
end_pause: Duration,
ack: Ack,
},
StopMarquee {
ack: Ack,
},
SetBrightness {
level: u8,
ack: Ack,
},
Shutdown {
ack: Ack,
},
}
#[derive(Clone)]
pub struct VfdHandle {
tx: SyncSender<Cmd>,
columns: usize,
}
impl VfdHandle {
pub fn clear(&self) -> Result<()> {
self.call(|ack| Cmd::Clear { ack })
}
pub fn set_brightness(&self, level: u8) -> Result<()> {
self.call(|ack| Cmd::SetBrightness { level, ack })
}
pub fn print_line(&self, line: u8, text: impl Into<String>) -> Result<()> {
let text = text.into();
let text = prepare_text_for_queue(&text, self.columns);
self.call(|ack| Cmd::PrintLine { line, text, ack })
}
pub fn print_line_diff(&self, line: u8, text: impl Into<String>) -> Result<()> {
let text = text.into();
let text = prepare_text_for_queue(&text, self.columns);
self.call(|ack| Cmd::PrintLineDiff { line, text, ack })
}
pub fn print_at(&self, x: u8, y: u8, text: impl Into<String>) -> Result<()> {
let text = text.into();
let remaining = if x == 0 {
0
} else {
self.columns
.checked_sub(usize::from(x))
.map_or(0, |remaining| remaining + 1)
};
let text = prepare_text_for_queue(&text, remaining);
self.call(|ack| Cmd::PrintAt { x, y, text, ack })
}
pub fn write_raw(&self, bytes: impl Into<Vec<u8>>) -> Result<()> {
let bytes = bytes.into();
if bytes.len() > MAX_QUEUED_RAW_BYTES {
return Err(VfdError::RawPayloadTooLarge {
length: bytes.len(),
max: MAX_QUEUED_RAW_BYTES,
});
}
self.call(|ack| Cmd::WriteRaw { bytes, ack })
}
pub fn set_marquee_text(&self, text: impl Into<String>) -> Result<()> {
let text = text.into();
let text = prepare_marquee_text(&text)?;
self.call(|ack| Cmd::SetMarqueeText { text, ack })
}
pub fn start_marquee(&self, line: u8, cps: u32, end_pause: Duration) -> Result<()> {
validate_marquee_speed(cps)?;
self.call(|ack| Cmd::StartMarquee {
line,
cps,
end_pause,
ack,
})
}
pub fn stop_marquee(&self) -> Result<()> {
self.call(|ack| Cmd::StopMarquee { ack })
}
pub fn shutdown(&self) -> Result<()> {
self.call(|ack| Cmd::Shutdown { ack })
}
fn call(&self, build: impl FnOnce(Ack) -> Cmd) -> Result<()> {
let (ack_tx, ack_rx) = mpsc::sync_channel(1);
self.tx.send(build(ack_tx))?;
ack_rx.recv()?
}
}
pub struct VfdWorker<T: Write + Send + 'static = Box<dyn SerialPort>> {
handle: VfdHandle,
join: Option<JoinHandle<Result<Vfd<T>>>>,
}
impl VfdWorker<Box<dyn SerialPort>> {
pub fn start(cfg: VfdConfig) -> Result<Self> {
cfg.validate()?;
let capacity = cfg.queue_capacity;
let vfd = Vfd::open(cfg)?;
Self::from_vfd(vfd, capacity)
}
}
impl<T: Write + Send + 'static> VfdWorker<T> {
pub fn from_vfd(vfd: Vfd<T>, queue_capacity: usize) -> Result<Self> {
if queue_capacity == 0 {
return Err(ConfigError::ZeroQueueCapacity.into());
}
let (tx, rx) = mpsc::sync_channel::<Cmd>(queue_capacity);
let columns = vfd.columns();
let handle = VfdHandle { tx, columns };
let join = thread::spawn(move || writer_loop(vfd, rx));
Ok(Self {
handle,
join: Some(join),
})
}
pub fn from_transport(
transport: T,
display: DisplaySettings,
queue_capacity: usize,
) -> Result<Self> {
if queue_capacity == 0 {
return Err(ConfigError::ZeroQueueCapacity.into());
}
let vfd = Vfd::from_transport(transport, display)?;
Self::from_vfd(vfd, queue_capacity)
}
pub fn handle(&self) -> VfdHandle {
self.handle.clone()
}
pub fn shutdown(mut self) -> Result<Vfd<T>> {
let shutdown_result = self.handle.shutdown();
let worker_result = self
.join
.take()
.expect("join handle exists")
.join()
.map_err(|_| VfdError::WorkerPanicked)?;
let vfd = worker_result?;
shutdown_result?;
Ok(vfd)
}
}
impl<T: Write + Send + 'static> Drop for VfdWorker<T> {
fn drop(&mut self) {
let _ = self.handle.shutdown();
if let Some(j) = self.join.take() {
let _ = j.join();
}
}
}
#[derive(Debug, Clone)]
struct MarqueeState {
active: bool,
line: u8,
cps: u32,
end_pause: Duration,
text: String,
stream: Vec<char>,
offset: usize,
paused_until: Option<Instant>,
next_step: Option<Instant>,
}
impl MarqueeState {
fn new() -> Self {
Self {
active: false,
line: 1,
cps: 5,
end_pause: Duration::from_millis(1500),
text: String::new(),
stream: Vec::new(),
offset: 0,
paused_until: None,
next_step: None,
}
}
fn rebuild_stream(&mut self, width: usize) {
let text = sanitize_text(&self.text);
self.stream.clear();
self.stream.reserve(width * 2 + text.chars().count());
self.stream.extend(std::iter::repeat_n(' ', width));
self.stream.extend(text.chars());
self.stream.extend(std::iter::repeat_n(' ', width));
self.offset = 0;
self.paused_until = None;
self.next_step = Some(Instant::now() + self.step_interval());
}
fn step_interval(&self) -> Duration {
let cps = u64::from(self.cps.clamp(1, MAX_MARQUEE_CPS));
Duration::from_nanos((1_000_000_000 / cps).max(1))
}
fn next_deadline(&self) -> Option<Instant> {
if !self.active {
return None;
}
self.paused_until.or(self.next_step)
}
}
fn send_ack(ack: Ack, result: Result<()>) {
let _ = ack.send(result);
}
fn writer_loop<T: Write + Send + 'static>(mut vfd: Vfd<T>, rx: Receiver<Cmd>) -> Result<Vfd<T>> {
vfd.clear()?;
let rows = vfd.rows();
let mut marquee = MarqueeState::new();
let mut last_lines = vec![String::new(); rows];
loop {
let timeout = marquee
.next_deadline()
.map(|deadline| deadline.saturating_duration_since(Instant::now()));
let command = match timeout {
Some(delay) => match rx.recv_timeout(delay) {
Ok(cmd) => Some(cmd),
Err(mpsc::RecvTimeoutError::Timeout) => None,
Err(mpsc::RecvTimeoutError::Disconnected) => break,
},
None => match rx.recv() {
Ok(cmd) => Some(cmd),
Err(_) => break,
},
};
if let Some(cmd) = command {
if handle_command(cmd, &mut vfd, &mut marquee, &mut last_lines)? {
break;
}
} else {
render_marquee(&mut vfd, &mut marquee, &mut last_lines)?;
}
}
Ok(vfd)
}
fn handle_command<T: Write>(
cmd: Cmd,
vfd: &mut Vfd<T>,
marquee: &mut MarqueeState,
last_lines: &mut [String],
) -> Result<bool> {
let width = vfd.columns();
let rows = vfd.rows();
match cmd {
Cmd::Clear { ack } => {
let result = vfd.clear();
if result.is_ok() {
last_lines.fill(String::new());
}
send_ack(ack, result);
}
Cmd::SetBrightness { level, ack } => {
send_ack(ack, vfd.set_brightness(level));
}
Cmd::PrintLine { line, text, ack } => {
let result = vfd.print_line(line, &text);
if result.is_ok() {
last_lines[(line - 1) as usize] = fit_to_width(&sanitize_text(&text), width);
}
send_ack(ack, result);
}
Cmd::PrintLineDiff { line, text, ack } => {
let result = print_line_diff(vfd, marquee, last_lines, line, &text);
send_ack(ack, result);
}
Cmd::PrintAt { x, y, text, ack } => {
let result = print_at_cached(vfd, marquee, last_lines, x, y, &text);
send_ack(ack, result);
}
Cmd::WriteRaw { bytes, ack } => {
send_ack(ack, vfd.write_raw(&bytes));
}
Cmd::SetMarqueeText { text, ack } => {
marquee.text = text;
if marquee.active {
marquee.rebuild_stream(width);
}
send_ack(ack, Ok(()));
}
Cmd::StartMarquee {
line,
cps,
end_pause,
ack,
} => {
let result = if cps > MAX_MARQUEE_CPS {
Err(VfdError::InvalidMarqueeSpeed {
cps,
max: MAX_MARQUEE_CPS,
})
} else if line == 0 || usize::from(line) > rows {
Err(VfdError::InvalidLine { line, rows })
} else {
last_lines[(line - 1) as usize].clear();
marquee.active = true;
marquee.line = line;
marquee.cps = cps.max(1);
marquee.end_pause = end_pause;
marquee.rebuild_stream(width);
Ok(())
};
send_ack(ack, result);
}
Cmd::StopMarquee { ack } => {
if marquee.active && usize::from(marquee.line) <= last_lines.len() {
last_lines[(marquee.line - 1) as usize].clear();
}
marquee.active = false;
marquee.paused_until = None;
marquee.next_step = None;
send_ack(ack, Ok(()));
}
Cmd::Shutdown { ack } => {
send_ack(ack, Ok(()));
return Ok(true);
}
}
Ok(false)
}
fn print_line_diff<T: Write>(
vfd: &mut Vfd<T>,
marquee: &MarqueeState,
last_lines: &mut [String],
line: u8,
text: &str,
) -> Result<()> {
if line == 0 || usize::from(line) > vfd.rows() {
return Err(VfdError::InvalidLine {
line,
rows: vfd.rows(),
});
}
if marquee.active && marquee.line == line {
return Ok(());
}
let next = fit_to_width(&sanitize_text(text), vfd.columns());
let idx = (line - 1) as usize;
if last_lines[idx] == next {
return Ok(());
}
if last_lines[idx].is_empty() {
vfd.print_line(line, &next)?;
last_lines[idx] = next;
return Ok(());
}
for (x, text) in changed_runs(&last_lines[idx], &next) {
vfd.print_at_prepared(x, line, &text)?;
}
last_lines[idx] = next;
Ok(())
}
fn print_at_cached<T: Write>(
vfd: &mut Vfd<T>,
marquee: &MarqueeState,
last_lines: &mut [String],
x: u8,
y: u8,
text: &str,
) -> Result<()> {
if x == 0 || y == 0 || usize::from(x) > vfd.columns() || usize::from(y) > vfd.rows() {
return Err(VfdError::InvalidCoordinate {
x,
y,
columns: vfd.columns(),
rows: vfd.rows(),
});
}
if marquee.active && marquee.line == y {
return Ok(());
}
let remaining = vfd.columns() - usize::from(x) + 1;
let text: String = sanitize_text(text).chars().take(remaining).collect();
vfd.print_at_prepared(x, y, &text)?;
replace_cached_range(&mut last_lines[(y - 1) as usize], x, &text, vfd.columns());
Ok(())
}
fn render_marquee<T: Write>(
vfd: &mut Vfd<T>,
marquee: &mut MarqueeState,
last_lines: &mut [String],
) -> Result<()> {
if !marquee.active {
return Ok(());
}
let now = Instant::now();
if let Some(until) = marquee.paused_until {
if now < until {
return Ok(());
}
marquee.paused_until = None;
marquee.next_step = Some(now + marquee.step_interval());
return Ok(());
}
if marquee.next_step.is_some_and(|deadline| now < deadline) {
return Ok(());
}
let width = vfd.columns();
if marquee.stream.len() < width {
marquee.rebuild_stream(width);
}
let max_off = marquee.stream.len().saturating_sub(width);
let start = marquee.offset.min(max_off);
let end = (start + width).min(marquee.stream.len());
let frame: String = marquee.stream[start..end].iter().collect();
vfd.print_at_prepared(1, marquee.line, &frame)?;
if usize::from(marquee.line) <= last_lines.len() {
last_lines[(marquee.line - 1) as usize] = frame;
}
if marquee.offset >= max_off {
marquee.offset = 0;
marquee.paused_until = Some(now + marquee.end_pause);
marquee.next_step = None;
} else {
marquee.offset += 1;
marquee.next_step = Some(now + marquee.step_interval());
}
Ok(())
}
fn prepare_text_for_queue(text: &str, max_chars: usize) -> String {
sanitize_to_width(text, max_chars)
}
fn prepare_marquee_text(text: &str) -> Result<String> {
if text.chars().nth(MAX_MARQUEE_CHARS).is_some() {
return Err(VfdError::TextTooLong {
max: MAX_MARQUEE_CHARS,
});
}
Ok(sanitize_to_width(text, MAX_MARQUEE_CHARS))
}
fn validate_marquee_speed(cps: u32) -> Result<()> {
if cps > MAX_MARQUEE_CPS {
return Err(VfdError::InvalidMarqueeSpeed {
cps,
max: MAX_MARQUEE_CPS,
});
}
Ok(())
}
#[cfg(test)]
mod tests {
use super::*;
use crate::config::{DisplaySettings, TextEncoding};
use std::io;
use std::sync::{
Arc,
atomic::{AtomicBool, Ordering},
};
struct FailsAfterWrites {
writes_left: usize,
failed: Arc<AtomicBool>,
}
impl Write for FailsAfterWrites {
fn write(&mut self, buf: &[u8]) -> io::Result<usize> {
if self.writes_left == 0 {
self.failed.store(true, Ordering::SeqCst);
return Err(io::Error::other("forced write failure"));
}
self.writes_left -= 1;
Ok(buf.len())
}
fn flush(&mut self) -> io::Result<()> {
Ok(())
}
}
#[test]
fn marquee_interval_never_collapses_to_zero() {
let mut marquee = MarqueeState::new();
marquee.cps = u32::MAX;
assert_eq!(
marquee.step_interval(),
Duration::from_nanos(1_000_000_000 / u64::from(MAX_MARQUEE_CPS))
);
}
#[test]
fn worker_returns_io_ack_and_dynamic_rows() {
let display = DisplaySettings::new(8, 3, TextEncoding::Ascii);
let worker = VfdWorker::from_transport(Vec::<u8>::new(), display, 2).unwrap();
let handle = worker.handle();
handle.print_line(3, "abc").unwrap();
assert!(matches!(
handle.print_line(4, "bad"),
Err(VfdError::InvalidLine { .. })
));
let vfd = worker.shutdown().unwrap();
let bytes = vfd.into_inner();
assert_eq!(&bytes[..2], &[0x1B, 0x40]);
assert!(bytes.windows(4).any(|window| window == [0x1F, 0x24, 1, 3]));
}
#[test]
fn worker_rejects_zero_queue_capacity() {
let display = DisplaySettings::new(8, 2, TextEncoding::Ascii);
assert!(matches!(
VfdWorker::from_transport(Vec::<u8>::new(), display, 0),
Err(VfdError::Config(ConfigError::ZeroQueueCapacity))
));
}
#[test]
fn worker_print_at_validates_coordinates_before_marquee_skip() {
let display = DisplaySettings::new(5, 2, TextEncoding::Ascii);
let worker = VfdWorker::from_transport(Vec::<u8>::new(), display, 2).unwrap();
let handle = worker.handle();
handle
.start_marquee(2, 10, Duration::from_millis(100))
.unwrap();
assert!(matches!(
handle.print_at(0, 2, "bad"),
Err(VfdError::InvalidCoordinate { x: 0, y: 2, .. })
));
assert!(matches!(
handle.print_at(6, 2, "bad"),
Err(VfdError::InvalidCoordinate { x: 6, y: 2, .. })
));
handle.print_at(1, 2, "skipped").unwrap();
worker.shutdown().unwrap();
}
#[test]
fn worker_shutdown_prefers_startup_io_error_over_closed_queue() {
let display = DisplaySettings::new(5, 2, TextEncoding::Ascii);
let failed = Arc::new(AtomicBool::new(false));
let transport = FailsAfterWrites {
writes_left: 1,
failed: Arc::clone(&failed),
};
let worker = VfdWorker::from_transport(transport, display, 2).unwrap();
while !failed.load(Ordering::SeqCst) {
thread::yield_now();
}
assert!(matches!(worker.shutdown(), Err(VfdError::Io(_))));
}
#[test]
fn worker_rejects_oversized_queued_payloads_and_marquee_rate() {
let display = DisplaySettings::new(20, 2, TextEncoding::Ascii);
let worker = VfdWorker::from_transport(Vec::<u8>::new(), display, 2).unwrap();
let handle = worker.handle();
assert!(matches!(
handle.set_marquee_text("x".repeat(MAX_MARQUEE_CHARS + 1)),
Err(VfdError::TextTooLong { .. })
));
assert!(matches!(
handle.write_raw(vec![0; MAX_QUEUED_RAW_BYTES + 1]),
Err(VfdError::RawPayloadTooLarge { .. })
));
assert!(matches!(
handle.start_marquee(1, MAX_MARQUEE_CPS + 1, Duration::ZERO),
Err(VfdError::InvalidMarqueeSpeed { .. })
));
worker.shutdown().unwrap();
}
#[test]
fn queued_line_text_is_prepared_to_display_width() {
assert_eq!(prepare_text_for_queue("ab\u{1b}@long", 4), "ab @");
}
}