1use std::io;
14use std::path::{Path, PathBuf};
15use std::sync::mpsc::{self, Receiver, Sender};
16use std::sync::{Arc, Mutex};
17use std::time::Duration;
18
19use super::protocol::{Answer, Malformed, Question};
20use super::{SCRIPT, SCRIPT_NAME, SOCKET_NAME};
21
22const QUESTION_WAIT: Duration = Duration::from_secs(10);
25
26const ANSWER_WAIT: Duration = Duration::from_secs(100);
30
31const MOST_CONNECTIONS: usize = 16;
34
35#[derive(Debug, Clone)]
37pub struct Call {
38 pub question: Result<Question, Malformed>,
40 reply: Sender<Answer>,
41}
42
43impl Call {
44 #[must_use]
47 pub fn new(question: Result<Question, Malformed>) -> (Self, Receiver<Answer>) {
48 let (reply, answers) = mpsc::channel();
49 (Self { question, reply }, answers)
50 }
51
52 pub fn answer(&self, answer: Answer) {
55 let _ = self.reply.send(answer);
56 }
57}
58
59#[derive(Debug, Clone)]
61pub struct Inbox(Arc<Mutex<Receiver<Call>>>);
62
63impl Inbox {
64 #[must_use]
66 pub fn next(&self) -> Option<Call> {
67 let receiver = self.0.lock().ok()?;
68 receiver.recv().ok()
69 }
70}
71
72#[derive(Debug)]
74pub struct Listener {
75 socket: PathBuf,
76 inbox: Inbox,
77 #[cfg(unix)]
78 open: Arc<std::sync::atomic::AtomicBool>,
79}
80
81impl Listener {
82 pub fn open(folder: &Path) -> io::Result<Self> {
93 #[cfg(unix)]
94 {
95 unix::open(folder)
96 }
97 #[cfg(not(unix))]
98 {
99 let _ = folder;
100 Err(io::Error::new(io::ErrorKind::Unsupported, "this system has no unix sockets"))
101 }
102 }
103
104 #[must_use]
106 pub fn socket(&self) -> &Path {
107 &self.socket
108 }
109
110 #[must_use]
112 pub fn inbox(&self) -> Inbox {
113 self.inbox.clone()
114 }
115}
116
117fn write_script(folder: &Path) -> io::Result<()> {
120 let path = folder.join(SCRIPT_NAME);
121 if std::fs::read(&path).is_ok_and(|written| written == SCRIPT.as_bytes()) {
122 return Ok(());
123 }
124 qframe::storage::atomic_write(&path, SCRIPT.as_bytes())
125}
126
127#[cfg(unix)]
128mod unix {
129 use std::io::{self, BufRead, BufReader, Read, Write};
130 use std::os::unix::fs::PermissionsExt;
131 use std::os::unix::net::{UnixListener, UnixStream};
132 use std::path::{Path, PathBuf};
133 use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
134 use std::sync::mpsc::{self, Sender};
135 use std::sync::{Arc, Mutex};
136
137 use super::super::protocol::{self, MOST_LINE, Malformed};
138 use super::{ANSWER_WAIT, Call, Inbox, Listener, MOST_CONNECTIONS, QUESTION_WAIT, SOCKET_NAME, write_script};
139
140 const MOST_PATH: usize = 107;
143
144 pub(super) fn open(folder: &Path) -> io::Result<Listener> {
145 std::fs::create_dir_all(folder)?;
146 std::fs::set_permissions(folder, std::fs::Permissions::from_mode(0o700))?;
148 write_script(folder)?;
149 let socket = folder.join(SOCKET_NAME);
150 if socket.exists() {
151 if reach(&socket, |path| UnixStream::connect(path)).is_ok() {
152 return Err(io::Error::new(io::ErrorKind::AddrInUse, socket.display().to_string()));
153 }
154 std::fs::remove_file(&socket)?;
155 }
156 let listener = reach(&socket, |path| UnixListener::bind(path))?;
157 let (calls, received) = mpsc::channel();
158 let open = Arc::new(AtomicBool::new(true));
159 let still_open = Arc::clone(&open);
160 std::thread::Builder::new()
161 .name("qcode-bridge".to_owned())
162 .spawn(move || accept(&listener, &calls, &still_open))?;
163 Ok(Listener { socket, inbox: Inbox(Arc::new(Mutex::new(received))), open })
164 }
165
166 fn reach<T>(socket: &Path, with: impl Fn(&Path) -> io::Result<T>) -> io::Result<T> {
169 if socket.as_os_str().len() <= MOST_PATH {
170 return with(socket);
171 }
172 let folder = socket.parent().ok_or_else(|| io::Error::from(io::ErrorKind::InvalidInput))?;
173 let handle = std::fs::File::open(folder)?;
174 let short = PathBuf::from(format!("/proc/self/fd/{}/{SOCKET_NAME}", std::os::fd::AsRawFd::as_raw_fd(&handle)));
175 let result = with(&short);
176 drop(handle);
177 result
178 }
179
180 fn accept(listener: &UnixListener, calls: &Sender<Call>, open: &AtomicBool) {
182 let serving = Arc::new(AtomicUsize::new(0));
183 for stream in listener.incoming() {
184 if !open.load(Ordering::SeqCst) {
185 break;
186 }
187 let Ok(stream) = stream else { continue };
188 if serving.load(Ordering::SeqCst) >= MOST_CONNECTIONS {
189 continue;
190 }
191 serving.fetch_add(1, Ordering::SeqCst);
192 let (calls, done) = (calls.clone(), Arc::clone(&serving));
193 let spawned = std::thread::Builder::new().name("qcode-bridge-call".to_owned()).spawn(move || {
194 serve(stream, &calls);
195 done.fetch_sub(1, Ordering::SeqCst);
196 });
197 if spawned.is_err() {
198 serving.fetch_sub(1, Ordering::SeqCst);
199 }
200 }
201 }
202
203 fn serve(stream: UnixStream, calls: &Sender<Call>) {
205 if stream.set_read_timeout(Some(QUESTION_WAIT)).is_err() {
206 return;
207 }
208 let Ok(reading) = stream.try_clone() else { return };
209 let mut line = String::new();
210 let limit = u64::try_from(MOST_LINE).unwrap_or(u64::MAX) + 1;
211 let read = BufReader::new(reading.take(limit)).read_line(&mut line);
212 let question = match read {
213 Ok(_) if line.ends_with('\n') => protocol::parse(line.trim_end_matches(['\n', '\r'])),
214 _ => Err(Malformed),
215 };
216 let (call, answers) = Call::new(question);
217 if calls.send(call).is_err() {
218 return;
219 }
220 if let Ok(answer) = answers.recv_timeout(ANSWER_WAIT) {
221 let mut writing = stream;
222 let _ = writing.write_all(answer.line().as_bytes());
223 }
224 }
225
226 impl Drop for Listener {
227 fn drop(&mut self) {
230 self.open.store(false, Ordering::SeqCst);
231 let _ = reach(&self.socket, |path| UnixStream::connect(path));
232 let _ = std::fs::remove_file(&self.socket);
233 }
234 }
235}
236
237#[cfg(all(test, unix))]
238mod tests {
239 use std::io::{BufRead, BufReader, Write};
240 use std::os::unix::net::UnixStream;
241 use std::path::PathBuf;
242
243 use super::super::protocol::{Answer, Malformed, Question, Request};
244 use super::*;
245
246 struct Scratch(PathBuf);
247
248 impl Scratch {
249 fn new(name: &str) -> Self {
250 let stamp =
251 std::time::SystemTime::now().duration_since(std::time::UNIX_EPOCH).unwrap_or_default().as_nanos();
252 Self(std::env::temp_dir().join(format!("qcode-bridge-{name}-{stamp}")))
253 }
254 }
255
256 impl Drop for Scratch {
257 fn drop(&mut self) {
258 let _ = std::fs::remove_dir_all(&self.0);
259 }
260 }
261
262 fn ask(socket: &Path, line: &str) -> String {
264 let mut stream = connect(socket);
265 stream.write_all(line.as_bytes()).expect("the question is written");
266 let mut answer = String::new();
267 BufReader::new(stream).read_line(&mut answer).expect("an answer comes");
268 answer
269 }
270
271 fn connect(socket: &Path) -> UnixStream {
272 if socket.as_os_str().len() <= 107 {
273 return UnixStream::connect(socket).expect("the socket answers");
274 }
275 let handle = std::fs::File::open(socket.parent().expect("a folder")).expect("the folder opens");
276 let short = format!("/proc/self/fd/{}/{SOCKET_NAME}", std::os::fd::AsRawFd::as_raw_fd(&handle));
277 UnixStream::connect(short).expect("the socket answers")
278 }
279
280 fn answering(inbox: Inbox, answer: fn(&Result<Question, Malformed>) -> Answer) -> std::thread::JoinHandle<usize> {
282 std::thread::spawn(move || {
283 let mut served = 0;
284 while let Some(call) = inbox.next() {
285 call.answer(answer(&call.question));
286 served += 1;
287 }
288 served
289 })
290 }
291
292 fn echo(question: &Result<Question, Malformed>) -> Answer {
293 match question {
294 Ok(Question { request: Request::Send { text, .. }, .. }) => Answer::done(text.clone()),
295 Ok(Question { request: Request::List, .. }) => {
296 Answer::listed("none".to_owned(), crate::bridge::protocol::You::default(), Vec::new())
297 }
298 Ok(Question { request: Request::Inbox | Request::Peek, .. }) => {
299 Answer::received("empty".to_owned(), Vec::new())
300 }
301 Err(Malformed) => Answer::refused("malformed".to_owned()),
302 }
303 }
304
305 #[test]
306 fn a_question_reaches_the_screen_and_its_answer_comes_back() {
307 let scratch = Scratch::new("round");
308 let listener = Listener::open(&scratch.0).expect("the socket opens");
309 assert!(scratch.0.join(SCRIPT_NAME).exists(), "the server is written beside the socket");
310 let served = answering(listener.inbox(), echo);
311
312 let answer = ask(listener.socket(), "{\"token\":\"t\",\"op\":\"send\",\"tab\":\"2\",\"text\":\"hello\"}\n");
313 assert_eq!(answer, "{\"ok\":true,\"text\":\"hello\"}\n");
314 let answer = ask(listener.socket(), "not json\n");
315 assert_eq!(answer, "{\"ok\":false,\"text\":\"malformed\"}\n");
316
317 let socket = listener.socket().to_path_buf();
318 drop(listener);
319 assert_eq!(served.join().expect("the answering thread ends"), 2, "closing ends the inbox");
320 assert!(!socket.exists(), "closing takes the socket away");
321 }
322
323 #[test]
324 fn a_line_without_its_end_is_malformed_rather_than_waited_for_forever() {
325 let scratch = Scratch::new("cut");
326 let listener = Listener::open(&scratch.0).expect("the socket opens");
327 let served = answering(listener.inbox(), echo);
328 let mut stream = connect(listener.socket());
329 stream.write_all(b"{\"token\":\"t\",\"op\":\"list\"}").expect("written");
330 stream.shutdown(std::net::Shutdown::Write).expect("the writing side closes");
331 let mut answer = String::new();
332 BufReader::new(stream).read_line(&mut answer).expect("an answer comes");
333 assert_eq!(answer, "{\"ok\":false,\"text\":\"malformed\"}\n");
334 drop(listener);
335 let _ = served.join();
336 }
337
338 #[test]
339 fn a_socket_another_qcode_answers_on_is_left_to_it() {
340 let scratch = Scratch::new("twice");
341 let first = Listener::open(&scratch.0).expect("the socket opens");
342 let second = Listener::open(&scratch.0);
343 assert_eq!(second.map(|_| ()).map_err(|error| error.kind()), Err(io::ErrorKind::AddrInUse));
344 let served = answering(first.inbox(), echo);
345 assert!(ask(first.socket(), "{\"token\":\"t\",\"op\":\"list\"}\n").contains("\"ok\":true"));
346 drop(first);
347 let _ = served.join();
348 }
349
350 fn refused(socket: &Path) {
363 let deadline = std::time::Instant::now() + std::time::Duration::from_secs(60);
364 while UnixStream::connect(socket).is_ok() {
365 assert!(std::time::Instant::now() < deadline, "something still answers on {}", socket.display());
366 std::thread::sleep(std::time::Duration::from_millis(10));
367 }
368 }
369
370 #[test]
371 fn a_socket_left_behind_is_replaced() {
372 let scratch = Scratch::new("stale");
373 std::fs::create_dir_all(&scratch.0).expect("a folder");
374 let stale = std::os::unix::net::UnixListener::bind(scratch.0.join(SOCKET_NAME)).expect("a socket");
375 drop(stale);
376 refused(&scratch.0.join(SOCKET_NAME));
377 let listener = Listener::open(&scratch.0).expect("the stale socket is replaced");
378 let served = answering(listener.inbox(), echo);
379 assert!(ask(listener.socket(), "{\"token\":\"t\",\"op\":\"list\"}\n").contains("none"));
380 drop(listener);
381 let _ = served.join();
382 }
383
384 #[test]
385 fn a_store_deep_enough_to_outgrow_a_socket_address_still_gets_its_socket() {
386 let scratch = Scratch::new("deep");
387 let deep = scratch.0.join("a-folder-name-long-enough".repeat(4)).join("Containers").join("MCP");
388 assert!(deep.join(SOCKET_NAME).as_os_str().len() > 108, "{}", deep.display());
389 let listener = Listener::open(&deep).expect("the socket opens");
390 assert_eq!(listener.socket(), deep.join(SOCKET_NAME));
391 assert!(listener.socket().exists());
392 let served = answering(listener.inbox(), echo);
393 assert!(ask(listener.socket(), "{\"token\":\"t\",\"op\":\"list\"}\n").contains("none"));
394 drop(listener);
395 let _ = served.join();
396 }
397
398 #[test]
399 fn the_folder_is_the_persons_alone_and_the_server_is_written_once() {
400 use std::os::unix::fs::PermissionsExt;
401
402 let scratch = Scratch::new("folder");
403 let listener = Listener::open(&scratch.0).expect("the socket opens");
404 let mode = std::fs::metadata(&scratch.0).expect("the folder").permissions().mode() & 0o777;
405 assert_eq!(mode, 0o700);
406 assert_eq!(std::fs::read_to_string(scratch.0.join(SCRIPT_NAME)).expect("the server"), SCRIPT);
407 drop(listener);
408 }
409}