1use crate::codec::{
9 changed_runs, fit_to_width, replace_cached_range, sanitize_text, sanitize_to_width,
10};
11use crate::config::{
12 DisplaySettings, MAX_MARQUEE_CHARS, MAX_MARQUEE_CPS, MAX_QUEUED_RAW_BYTES, VfdConfig,
13};
14use crate::error::{ConfigError, Result, VfdError};
15use crate::vfd::Vfd;
16use serialport::SerialPort;
17use std::io::Write;
18use std::sync::mpsc::{self, Receiver, SyncSender};
19use std::thread::{self, JoinHandle};
20use std::time::{Duration, Instant};
21
22type Ack = mpsc::SyncSender<Result<()>>;
23
24#[derive(Debug)]
25enum Cmd {
26 Clear {
27 ack: Ack,
28 },
29 PrintLine {
30 line: u8,
31 text: String,
32 ack: Ack,
33 },
34 PrintLineDiff {
35 line: u8,
36 text: String,
37 ack: Ack,
38 },
39 PrintAt {
40 x: u8,
41 y: u8,
42 text: String,
43 ack: Ack,
44 },
45 WriteRaw {
46 bytes: Vec<u8>,
47 ack: Ack,
48 },
49 SetMarqueeText {
50 text: String,
51 ack: Ack,
52 },
53 StartMarquee {
54 line: u8,
55 cps: u32,
56 end_pause: Duration,
57 ack: Ack,
58 },
59 StopMarquee {
60 ack: Ack,
61 },
62 SetBrightness {
63 level: u8,
64 ack: Ack,
65 },
66 Shutdown {
67 ack: Ack,
68 },
69}
70
71#[derive(Clone)]
78pub struct VfdHandle {
79 tx: SyncSender<Cmd>,
80 columns: usize,
81}
82
83impl VfdHandle {
84 pub fn clear(&self) -> Result<()> {
92 self.call(|ack| Cmd::Clear { ack })
93 }
94
95 pub fn set_brightness(&self, level: u8) -> Result<()> {
104 self.call(|ack| Cmd::SetBrightness { level, ack })
105 }
106
107 pub fn print_line(&self, line: u8, text: impl Into<String>) -> Result<()> {
115 let text = text.into();
116 let text = prepare_text_for_queue(&text, self.columns);
117 self.call(|ack| Cmd::PrintLine { line, text, ack })
118 }
119
120 pub fn print_line_diff(&self, line: u8, text: impl Into<String>) -> Result<()> {
133 let text = text.into();
134 let text = prepare_text_for_queue(&text, self.columns);
135 self.call(|ack| Cmd::PrintLineDiff { line, text, ack })
136 }
137
138 pub fn print_at(&self, x: u8, y: u8, text: impl Into<String>) -> Result<()> {
147 let text = text.into();
148 let remaining = if x == 0 {
149 0
150 } else {
151 self.columns
152 .checked_sub(usize::from(x))
153 .map_or(0, |remaining| remaining + 1)
154 };
155 let text = prepare_text_for_queue(&text, remaining);
156 self.call(|ack| Cmd::PrintAt { x, y, text, ack })
157 }
158
159 pub fn write_raw(&self, bytes: impl Into<Vec<u8>>) -> Result<()> {
170 let bytes = bytes.into();
171 if bytes.len() > MAX_QUEUED_RAW_BYTES {
172 return Err(VfdError::RawPayloadTooLarge {
173 length: bytes.len(),
174 max: MAX_QUEUED_RAW_BYTES,
175 });
176 }
177 self.call(|ack| Cmd::WriteRaw { bytes, ack })
178 }
179
180 pub fn set_marquee_text(&self, text: impl Into<String>) -> Result<()> {
190 let text = text.into();
191 let text = prepare_marquee_text(&text)?;
192 self.call(|ack| Cmd::SetMarqueeText { text, ack })
193 }
194
195 pub fn start_marquee(&self, line: u8, cps: u32, end_pause: Duration) -> Result<()> {
205 validate_marquee_speed(cps)?;
206 self.call(|ack| Cmd::StartMarquee {
207 line,
208 cps,
209 end_pause,
210 ack,
211 })
212 }
213
214 pub fn stop_marquee(&self) -> Result<()> {
223 self.call(|ack| Cmd::StopMarquee { ack })
224 }
225
226 pub fn shutdown(&self) -> Result<()> {
236 self.call(|ack| Cmd::Shutdown { ack })
237 }
238
239 fn call(&self, build: impl FnOnce(Ack) -> Cmd) -> Result<()> {
240 let (ack_tx, ack_rx) = mpsc::sync_channel(1);
241 self.tx.send(build(ack_tx))?;
242 ack_rx.recv()?
243 }
244}
245
246pub struct VfdWorker<T: Write + Send + 'static = Box<dyn SerialPort>> {
252 handle: VfdHandle,
253 join: Option<JoinHandle<Result<Vfd<T>>>>,
254}
255
256impl VfdWorker<Box<dyn SerialPort>> {
257 pub fn start(cfg: VfdConfig) -> Result<Self> {
268 cfg.validate()?;
269 let capacity = cfg.queue_capacity;
270 let vfd = Vfd::open(cfg)?;
271 Self::from_vfd(vfd, capacity)
272 }
273}
274
275impl<T: Write + Send + 'static> VfdWorker<T> {
276 pub fn from_vfd(vfd: Vfd<T>, queue_capacity: usize) -> Result<Self> {
285 if queue_capacity == 0 {
286 return Err(ConfigError::ZeroQueueCapacity.into());
287 }
288 let (tx, rx) = mpsc::sync_channel::<Cmd>(queue_capacity);
289 let columns = vfd.columns();
290 let handle = VfdHandle { tx, columns };
291 let join = thread::spawn(move || writer_loop(vfd, rx));
292
293 Ok(Self {
294 handle,
295 join: Some(join),
296 })
297 }
298
299 pub fn from_transport(
309 transport: T,
310 display: DisplaySettings,
311 queue_capacity: usize,
312 ) -> Result<Self> {
313 if queue_capacity == 0 {
314 return Err(ConfigError::ZeroQueueCapacity.into());
315 }
316 let vfd = Vfd::from_transport(transport, display)?;
317 Self::from_vfd(vfd, queue_capacity)
318 }
319
320 pub fn handle(&self) -> VfdHandle {
325 self.handle.clone()
326 }
327
328 pub fn shutdown(mut self) -> Result<Vfd<T>> {
338 let shutdown_result = self.handle.shutdown();
339 let worker_result = self
340 .join
341 .take()
342 .expect("join handle exists")
343 .join()
344 .map_err(|_| VfdError::WorkerPanicked)?;
345
346 let vfd = worker_result?;
347 shutdown_result?;
348 Ok(vfd)
349 }
350}
351
352impl<T: Write + Send + 'static> Drop for VfdWorker<T> {
353 fn drop(&mut self) {
354 let _ = self.handle.shutdown();
355 if let Some(j) = self.join.take() {
356 let _ = j.join();
357 }
358 }
359}
360
361#[derive(Debug, Clone)]
362struct MarqueeState {
363 active: bool,
364 line: u8,
365 cps: u32,
366 end_pause: Duration,
367 text: String,
368 stream: Vec<char>,
369 offset: usize,
370 paused_until: Option<Instant>,
371 next_step: Option<Instant>,
372}
373
374impl MarqueeState {
375 fn new() -> Self {
376 Self {
377 active: false,
378 line: 1,
379 cps: 5,
380 end_pause: Duration::from_millis(1500),
381 text: String::new(),
382 stream: Vec::new(),
383 offset: 0,
384 paused_until: None,
385 next_step: None,
386 }
387 }
388
389 fn rebuild_stream(&mut self, width: usize) {
390 let text = sanitize_text(&self.text);
391 self.stream.clear();
392 self.stream.reserve(width * 2 + text.chars().count());
393 self.stream.extend(std::iter::repeat_n(' ', width));
394 self.stream.extend(text.chars());
395 self.stream.extend(std::iter::repeat_n(' ', width));
396 self.offset = 0;
397 self.paused_until = None;
398 self.next_step = Some(Instant::now() + self.step_interval());
399 }
400
401 fn step_interval(&self) -> Duration {
402 let cps = u64::from(self.cps.clamp(1, MAX_MARQUEE_CPS));
403 Duration::from_nanos((1_000_000_000 / cps).max(1))
404 }
405
406 fn next_deadline(&self) -> Option<Instant> {
407 if !self.active {
408 return None;
409 }
410 self.paused_until.or(self.next_step)
411 }
412}
413
414fn send_ack(ack: Ack, result: Result<()>) {
415 let _ = ack.send(result);
416}
417
418fn writer_loop<T: Write + Send + 'static>(mut vfd: Vfd<T>, rx: Receiver<Cmd>) -> Result<Vfd<T>> {
419 vfd.clear()?;
420
421 let rows = vfd.rows();
422 let mut marquee = MarqueeState::new();
423 let mut last_lines = vec![String::new(); rows];
424
425 loop {
426 let timeout = marquee
427 .next_deadline()
428 .map(|deadline| deadline.saturating_duration_since(Instant::now()));
429
430 let command = match timeout {
431 Some(delay) => match rx.recv_timeout(delay) {
432 Ok(cmd) => Some(cmd),
433 Err(mpsc::RecvTimeoutError::Timeout) => None,
434 Err(mpsc::RecvTimeoutError::Disconnected) => break,
435 },
436 None => match rx.recv() {
437 Ok(cmd) => Some(cmd),
438 Err(_) => break,
439 },
440 };
441
442 if let Some(cmd) = command {
443 if handle_command(cmd, &mut vfd, &mut marquee, &mut last_lines)? {
444 break;
445 }
446 } else {
447 render_marquee(&mut vfd, &mut marquee, &mut last_lines)?;
448 }
449 }
450
451 Ok(vfd)
452}
453
454fn handle_command<T: Write>(
455 cmd: Cmd,
456 vfd: &mut Vfd<T>,
457 marquee: &mut MarqueeState,
458 last_lines: &mut [String],
459) -> Result<bool> {
460 let width = vfd.columns();
461 let rows = vfd.rows();
462 match cmd {
463 Cmd::Clear { ack } => {
464 let result = vfd.clear();
465 if result.is_ok() {
466 last_lines.fill(String::new());
467 }
468 send_ack(ack, result);
469 }
470 Cmd::SetBrightness { level, ack } => {
471 send_ack(ack, vfd.set_brightness(level));
472 }
473 Cmd::PrintLine { line, text, ack } => {
474 let result = vfd.print_line(line, &text);
475 if result.is_ok() {
476 last_lines[(line - 1) as usize] = fit_to_width(&sanitize_text(&text), width);
477 }
478 send_ack(ack, result);
479 }
480 Cmd::PrintLineDiff { line, text, ack } => {
481 let result = print_line_diff(vfd, marquee, last_lines, line, &text);
482 send_ack(ack, result);
483 }
484 Cmd::PrintAt { x, y, text, ack } => {
485 let result = print_at_cached(vfd, marquee, last_lines, x, y, &text);
486 send_ack(ack, result);
487 }
488 Cmd::WriteRaw { bytes, ack } => {
489 send_ack(ack, vfd.write_raw(&bytes));
490 }
491 Cmd::SetMarqueeText { text, ack } => {
492 marquee.text = text;
493 if marquee.active {
494 marquee.rebuild_stream(width);
495 }
496 send_ack(ack, Ok(()));
497 }
498 Cmd::StartMarquee {
499 line,
500 cps,
501 end_pause,
502 ack,
503 } => {
504 let result = if cps > MAX_MARQUEE_CPS {
505 Err(VfdError::InvalidMarqueeSpeed {
506 cps,
507 max: MAX_MARQUEE_CPS,
508 })
509 } else if line == 0 || usize::from(line) > rows {
510 Err(VfdError::InvalidLine { line, rows })
511 } else {
512 last_lines[(line - 1) as usize].clear();
513 marquee.active = true;
514 marquee.line = line;
515 marquee.cps = cps.max(1);
516 marquee.end_pause = end_pause;
517 marquee.rebuild_stream(width);
518 Ok(())
519 };
520 send_ack(ack, result);
521 }
522 Cmd::StopMarquee { ack } => {
523 if marquee.active && usize::from(marquee.line) <= last_lines.len() {
524 last_lines[(marquee.line - 1) as usize].clear();
525 }
526 marquee.active = false;
527 marquee.paused_until = None;
528 marquee.next_step = None;
529 send_ack(ack, Ok(()));
530 }
531 Cmd::Shutdown { ack } => {
532 send_ack(ack, Ok(()));
533 return Ok(true);
534 }
535 }
536 Ok(false)
537}
538
539fn print_line_diff<T: Write>(
540 vfd: &mut Vfd<T>,
541 marquee: &MarqueeState,
542 last_lines: &mut [String],
543 line: u8,
544 text: &str,
545) -> Result<()> {
546 if line == 0 || usize::from(line) > vfd.rows() {
547 return Err(VfdError::InvalidLine {
548 line,
549 rows: vfd.rows(),
550 });
551 }
552 if marquee.active && marquee.line == line {
553 return Ok(());
554 }
555
556 let next = fit_to_width(&sanitize_text(text), vfd.columns());
557 let idx = (line - 1) as usize;
558 if last_lines[idx] == next {
559 return Ok(());
560 }
561 if last_lines[idx].is_empty() {
562 vfd.print_line(line, &next)?;
563 last_lines[idx] = next;
564 return Ok(());
565 }
566 for (x, text) in changed_runs(&last_lines[idx], &next) {
567 vfd.print_at_prepared(x, line, &text)?;
568 }
569 last_lines[idx] = next;
570 Ok(())
571}
572
573fn print_at_cached<T: Write>(
574 vfd: &mut Vfd<T>,
575 marquee: &MarqueeState,
576 last_lines: &mut [String],
577 x: u8,
578 y: u8,
579 text: &str,
580) -> Result<()> {
581 if x == 0 || y == 0 || usize::from(x) > vfd.columns() || usize::from(y) > vfd.rows() {
582 return Err(VfdError::InvalidCoordinate {
583 x,
584 y,
585 columns: vfd.columns(),
586 rows: vfd.rows(),
587 });
588 }
589 if marquee.active && marquee.line == y {
590 return Ok(());
591 }
592
593 let remaining = vfd.columns() - usize::from(x) + 1;
594 let text: String = sanitize_text(text).chars().take(remaining).collect();
595 vfd.print_at_prepared(x, y, &text)?;
596 replace_cached_range(&mut last_lines[(y - 1) as usize], x, &text, vfd.columns());
597 Ok(())
598}
599
600fn render_marquee<T: Write>(
601 vfd: &mut Vfd<T>,
602 marquee: &mut MarqueeState,
603 last_lines: &mut [String],
604) -> Result<()> {
605 if !marquee.active {
606 return Ok(());
607 }
608 let now = Instant::now();
609 if let Some(until) = marquee.paused_until {
610 if now < until {
611 return Ok(());
612 }
613 marquee.paused_until = None;
614 marquee.next_step = Some(now + marquee.step_interval());
615 return Ok(());
616 }
617 if marquee.next_step.is_some_and(|deadline| now < deadline) {
618 return Ok(());
619 }
620
621 let width = vfd.columns();
622 if marquee.stream.len() < width {
623 marquee.rebuild_stream(width);
624 }
625 let max_off = marquee.stream.len().saturating_sub(width);
626 let start = marquee.offset.min(max_off);
627 let end = (start + width).min(marquee.stream.len());
628 let frame: String = marquee.stream[start..end].iter().collect();
629
630 vfd.print_at_prepared(1, marquee.line, &frame)?;
631 if usize::from(marquee.line) <= last_lines.len() {
632 last_lines[(marquee.line - 1) as usize] = frame;
633 }
634
635 if marquee.offset >= max_off {
636 marquee.offset = 0;
637 marquee.paused_until = Some(now + marquee.end_pause);
638 marquee.next_step = None;
639 } else {
640 marquee.offset += 1;
641 marquee.next_step = Some(now + marquee.step_interval());
642 }
643 Ok(())
644}
645
646fn prepare_text_for_queue(text: &str, max_chars: usize) -> String {
647 sanitize_to_width(text, max_chars)
648}
649
650fn prepare_marquee_text(text: &str) -> Result<String> {
651 if text.chars().nth(MAX_MARQUEE_CHARS).is_some() {
652 return Err(VfdError::TextTooLong {
653 max: MAX_MARQUEE_CHARS,
654 });
655 }
656 Ok(sanitize_to_width(text, MAX_MARQUEE_CHARS))
657}
658
659fn validate_marquee_speed(cps: u32) -> Result<()> {
660 if cps > MAX_MARQUEE_CPS {
661 return Err(VfdError::InvalidMarqueeSpeed {
662 cps,
663 max: MAX_MARQUEE_CPS,
664 });
665 }
666 Ok(())
667}
668
669#[cfg(test)]
670mod tests {
671 use super::*;
672 use crate::config::{DisplaySettings, TextEncoding};
673 use std::io;
674 use std::sync::{
675 Arc,
676 atomic::{AtomicBool, Ordering},
677 };
678
679 struct FailsAfterWrites {
680 writes_left: usize,
681 failed: Arc<AtomicBool>,
682 }
683
684 impl Write for FailsAfterWrites {
685 fn write(&mut self, buf: &[u8]) -> io::Result<usize> {
686 if self.writes_left == 0 {
687 self.failed.store(true, Ordering::SeqCst);
688 return Err(io::Error::other("forced write failure"));
689 }
690 self.writes_left -= 1;
691 Ok(buf.len())
692 }
693
694 fn flush(&mut self) -> io::Result<()> {
695 Ok(())
696 }
697 }
698
699 #[test]
700 fn marquee_interval_never_collapses_to_zero() {
701 let mut marquee = MarqueeState::new();
702 marquee.cps = u32::MAX;
703 assert_eq!(
704 marquee.step_interval(),
705 Duration::from_nanos(1_000_000_000 / u64::from(MAX_MARQUEE_CPS))
706 );
707 }
708
709 #[test]
710 fn worker_returns_io_ack_and_dynamic_rows() {
711 let display = DisplaySettings::new(8, 3, TextEncoding::Ascii);
712 let worker = VfdWorker::from_transport(Vec::<u8>::new(), display, 2).unwrap();
713 let handle = worker.handle();
714
715 handle.print_line(3, "abc").unwrap();
716 assert!(matches!(
717 handle.print_line(4, "bad"),
718 Err(VfdError::InvalidLine { .. })
719 ));
720
721 let vfd = worker.shutdown().unwrap();
722 let bytes = vfd.into_inner();
723 assert_eq!(&bytes[..2], &[0x1B, 0x40]);
724 assert!(bytes.windows(4).any(|window| window == [0x1F, 0x24, 1, 3]));
725 }
726
727 #[test]
728 fn worker_rejects_zero_queue_capacity() {
729 let display = DisplaySettings::new(8, 2, TextEncoding::Ascii);
730 assert!(matches!(
731 VfdWorker::from_transport(Vec::<u8>::new(), display, 0),
732 Err(VfdError::Config(ConfigError::ZeroQueueCapacity))
733 ));
734 }
735
736 #[test]
737 fn worker_print_at_validates_coordinates_before_marquee_skip() {
738 let display = DisplaySettings::new(5, 2, TextEncoding::Ascii);
739 let worker = VfdWorker::from_transport(Vec::<u8>::new(), display, 2).unwrap();
740 let handle = worker.handle();
741
742 handle
743 .start_marquee(2, 10, Duration::from_millis(100))
744 .unwrap();
745
746 assert!(matches!(
747 handle.print_at(0, 2, "bad"),
748 Err(VfdError::InvalidCoordinate { x: 0, y: 2, .. })
749 ));
750 assert!(matches!(
751 handle.print_at(6, 2, "bad"),
752 Err(VfdError::InvalidCoordinate { x: 6, y: 2, .. })
753 ));
754 handle.print_at(1, 2, "skipped").unwrap();
755
756 worker.shutdown().unwrap();
757 }
758
759 #[test]
760 fn worker_shutdown_prefers_startup_io_error_over_closed_queue() {
761 let display = DisplaySettings::new(5, 2, TextEncoding::Ascii);
762 let failed = Arc::new(AtomicBool::new(false));
763 let transport = FailsAfterWrites {
764 writes_left: 1,
765 failed: Arc::clone(&failed),
766 };
767 let worker = VfdWorker::from_transport(transport, display, 2).unwrap();
768
769 while !failed.load(Ordering::SeqCst) {
770 thread::yield_now();
771 }
772
773 assert!(matches!(worker.shutdown(), Err(VfdError::Io(_))));
774 }
775
776 #[test]
777 fn worker_rejects_oversized_queued_payloads_and_marquee_rate() {
778 let display = DisplaySettings::new(20, 2, TextEncoding::Ascii);
779 let worker = VfdWorker::from_transport(Vec::<u8>::new(), display, 2).unwrap();
780 let handle = worker.handle();
781
782 assert!(matches!(
783 handle.set_marquee_text("x".repeat(MAX_MARQUEE_CHARS + 1)),
784 Err(VfdError::TextTooLong { .. })
785 ));
786 assert!(matches!(
787 handle.write_raw(vec![0; MAX_QUEUED_RAW_BYTES + 1]),
788 Err(VfdError::RawPayloadTooLarge { .. })
789 ));
790 assert!(matches!(
791 handle.start_marquee(1, MAX_MARQUEE_CPS + 1, Duration::ZERO),
792 Err(VfdError::InvalidMarqueeSpeed { .. })
793 ));
794
795 worker.shutdown().unwrap();
796 }
797
798 #[test]
799 fn queued_line_text_is_prepared_to_display_width() {
800 assert_eq!(prepare_text_for_queue("ab\u{1b}@long", 4), "ab @");
801 }
802}