android_usb_serial/
reader.rs1use crate::error::{ReadOutcome, Result, TransferError, UsbSerialError};
7use crate::rx_filter::{apply_filters, RxFilter};
8use crate::transport::BulkIn;
9use std::sync::mpsc::{self, Receiver, Sender, TryRecvError};
10use std::sync::{Arc, Mutex};
11use std::thread::{self, JoinHandle};
12use std::time::Duration;
13
14const STOP_JOIN_TIMEOUT: Duration = Duration::from_secs(1);
15
16pub struct SerialReader {
18 rx: Receiver<Vec<u8>>,
19 stop_tx: Option<Sender<()>>,
20 join: Option<JoinHandle<()>>,
21 error: Arc<Mutex<Option<UsbSerialError>>>,
22}
23
24impl SerialReader {
25 pub fn start(
27 bulk_in: Box<dyn BulkIn>,
28 max_packet_size: u16,
29 read_timeout_ms: u32,
30 filters: Vec<Box<dyn RxFilter>>,
31 ) -> Self {
32 let bufsize = (max_packet_size as usize).saturating_mul(4).max(64);
33 let (data_tx, data_rx) = mpsc::channel();
34 let (stop_tx, stop_rx) = mpsc::channel();
35 let error = Arc::new(Mutex::new(None));
36 let err_clone = error.clone();
37 let join = thread::spawn(move || {
38 reader_loop(
39 bulk_in,
40 bufsize,
41 read_timeout_ms,
42 filters,
43 &data_tx,
44 &stop_rx,
45 &err_clone,
46 );
47 });
48 Self {
49 rx: data_rx,
50 stop_tx: Some(stop_tx),
51 join: Some(join),
52 error,
53 }
54 }
55
56 pub fn try_read(&mut self, buf: &mut [u8]) -> Result<usize> {
58 if let Some(err) = self.error.lock().unwrap().take() {
59 return Err(err);
60 }
61 match self.rx.try_recv() {
62 Ok(chunk) => {
63 let n = chunk.len().min(buf.len());
64 buf[..n].copy_from_slice(&chunk[..n]);
65 Ok(n)
66 }
67 Err(TryRecvError::Empty) => Ok(0),
68 Err(TryRecvError::Disconnected) => {
69 if let Some(err) = self.error.lock().unwrap().take() {
70 return Err(err);
71 }
72 Err(UsbSerialError::Disconnected)
73 }
74 }
75 }
76
77 pub fn stop(&mut self) {
79 if let Some(tx) = self.stop_tx.take() {
80 let _ = tx.send(());
81 }
82 if let Some(j) = self.join.take() {
83 let (done_tx, done_rx) = mpsc::channel();
84 thread::spawn(move || {
85 let _ = j.join();
86 let _ = done_tx.send(());
87 });
88 let _ = done_rx.recv_timeout(STOP_JOIN_TIMEOUT);
89 }
90 }
91}
92
93impl Drop for SerialReader {
94 fn drop(&mut self) {
95 self.stop();
96 }
97}
98
99fn is_stall(err: &UsbSerialError) -> bool {
100 matches!(err, UsbSerialError::Io(msg) if msg == "stall")
101}
102
103fn reader_loop(
104 mut bulk_in: Box<dyn BulkIn>,
105 bufsize: usize,
106 timeout_ms: u32,
107 mut filters: Vec<Box<dyn RxFilter>>,
108 data_tx: &Sender<Vec<u8>>,
109 stop_rx: &Receiver<()>,
110 error: &Arc<Mutex<Option<UsbSerialError>>>,
111) {
112 let mut buf = vec![0u8; bufsize];
113 let mut consecutive_stalls = 0u32;
114 loop {
115 if stop_rx.try_recv().is_ok() {
116 bulk_in.cancel_all();
117 break;
118 }
119 match bulk_in.read(&mut buf, timeout_ms) {
120 Ok(ReadOutcome::Data(data)) if !data.is_empty() => {
121 consecutive_stalls = 0;
122 let filtered = if filters.is_empty() {
123 data
124 } else {
125 apply_filters(&mut filters, &data)
126 };
127 if !filtered.is_empty() && data_tx.send(filtered).is_err() {
128 break;
129 }
130 }
131 Ok(ReadOutcome::TimedOut) | Ok(ReadOutcome::Data(_)) => {}
132 Ok(ReadOutcome::Cancelled) => break,
133 Err(e) if is_stall(&e) => {
134 if consecutive_stalls == 0 {
135 consecutive_stalls += 1;
136 if bulk_in.clear_halt().is_err() {
137 *error.lock().unwrap() = Some(e);
138 break;
139 }
140 continue;
141 }
142 *error.lock().unwrap() = Some(UsbSerialError::from(TransferError::Stall));
143 break;
144 }
145 Err(e) => {
146 *error.lock().unwrap() = Some(e);
147 break;
148 }
149 }
150 }
151}
152
153#[cfg(all(test, feature = "fake-transport"))]
154mod tests {
155 use super::*;
156 use crate::fake::FakeTransport;
157 use crate::transport::Transport;
158 use std::time::Instant;
159
160 #[test]
161 fn reader_delivers_injected_rx() {
162 let fake = FakeTransport::cdc_single_iface();
163 fake.push_rx(b"hello");
164 let bulk = fake.open_bulk_in(0x81, 64).unwrap();
165 let mut reader = SerialReader::start(bulk, 64, 100, vec![]);
166 let deadline = Instant::now() + Duration::from_secs(2);
167 let mut out = [0u8; 8];
168 let mut n = 0;
169 while n == 0 && Instant::now() < deadline {
170 n = reader.try_read(&mut out).unwrap();
171 std::thread::sleep(Duration::from_millis(5));
172 }
173 assert_eq!(n, 5);
174 assert_eq!(&out[..5], b"hello");
175 reader.stop();
176 }
177}