Skip to main content

qcode/bridge/
socket.rs

1//! The workspace's socket: QCode listens on it while the workspace is open, and the server in every
2//! profile container of the workspace asks its questions through it.
3//!
4//! A unix socket rather than a port, because a profile without the network has no network to
5//! reach a port through, while a socket is a file in a folder the container mounts. One thread
6//! waits for connections and one short thread serves each: it reads the question, hands it to
7//! the screen through the [`Inbox`] and writes back whatever the screen answers. The screen
8//! decides everything; this module only carries lines.
9//!
10//! Where the system has no unix sockets, [`Listener::open`] says so and the workspace has no
11//! bridge; nothing else changes.
12
13use 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
22/// How long a connection may take to send its question. The server writes it at once; a
23/// connection that sends nothing is not one of ours.
24const QUESTION_WAIT: Duration = Duration::from_secs(10);
25
26/// How long a connection waits for the screen's answer. The screen answers every question
27/// itself, a question for the person included (it says the message waits for them when they
28/// take long), so this only ends a connection whose answer was lost with the screen.
29const ANSWER_WAIT: Duration = Duration::from_secs(100);
30
31/// The most connections served at once. Every tab asks one question at a time; more than this
32/// is something in a container knocking, and the rest are closed unanswered.
33const MOST_CONNECTIONS: usize = 16;
34
35/// A question that came in, and the way back to the one who asked.
36#[derive(Debug, Clone)]
37pub struct Call {
38    /// The question, or [`Malformed`] when the line was not one.
39    pub question: Result<Question, Malformed>,
40    reply: Sender<Answer>,
41}
42
43impl Call {
44    /// A call carrying `question`, with the receiving end of its answer; for the screen's tests,
45    /// which have no socket.
46    #[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    /// Answers the call. A connection that went away meanwhile is not an error: the agent that
53    /// asked is gone.
54    pub fn answer(&self, answer: Answer) {
55        let _ = self.reply.send(answer);
56    }
57}
58
59/// Where the screen waits for the next call.
60#[derive(Debug, Clone)]
61pub struct Inbox(Arc<Mutex<Receiver<Call>>>);
62
63impl Inbox {
64    /// Waits for the next call, on a background thread. `None` once the listener is closed.
65    #[must_use]
66    pub fn next(&self) -> Option<Call> {
67        let receiver = self.0.lock().ok()?;
68        receiver.recv().ok()
69    }
70}
71
72/// A workspace's socket, listened on for as long as this value lives.
73#[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    /// Listens in `folder`, the workspace's `Containers/MCP/`: makes the folder, writes the
83    /// server beside where the socket goes, and starts waiting for connections.
84    ///
85    /// A socket left behind by a QCode that ended without clearing it is replaced; one another
86    /// QCode still answers on is left to it, and this QCode has no bridge for the workspace.
87    ///
88    /// # Errors
89    ///
90    /// When the folder or the server cannot be written, another QCode answers on the socket
91    /// ([`io::ErrorKind::AddrInUse`]), the socket cannot be made, or the system has none.
92    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    /// Where the socket is.
105    #[must_use]
106    pub fn socket(&self) -> &Path {
107        &self.socket
108    }
109
110    /// Where the screen waits for calls.
111    #[must_use]
112    pub fn inbox(&self) -> Inbox {
113        self.inbox.clone()
114    }
115}
116
117/// Writes the server into `folder` unless it is there already as this QCode carries it, so a
118/// workspace opened again rewrites nothing.
119fn 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    /// The longest socket path the system takes, less the byte its terminator needs. A path the
141    /// store makes longer is reached through the folder's handle instead (see [`reach`]).
142    const MOST_PATH: usize = 107;
143
144    pub(super) fn open(folder: &Path) -> io::Result<Listener> {
145        std::fs::create_dir_all(folder)?;
146        // Only the person's own processes and the containers running as them reach the socket.
147        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    /// Runs `with` on the socket path, or, when the path is too long for a socket address, on the
167    /// same socket named through the folder's open handle in `/proc/self/fd`, which is short.
168    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    /// Waits for connections until the listener is closed, serving each on a thread of its own.
181    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    /// Reads one question, hands it on and writes the answer back.
204    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        /// Stops listening and takes the socket away. The thread waiting for connections is
228        /// woken by one last connection of its own, so it sees the listener is closed and ends.
229        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    /// Asks `socket` one line the way the server in a container does, and reads the answer.
263    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    /// Answers every call the way `answer` says, on a thread, until the listener closes.
281    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    /// Waits until nothing answers on `socket` any more, which is what a socket left behind by a
351    /// QCode that ended is.
352    ///
353    /// Dropping the listener is not yet that, in a test binary. Other tests start processes on
354    /// other threads, and a process being started is a copy of this one until it runs its
355    /// program: for that moment it holds every descriptor this one has, the listener just dropped
356    /// included, and the socket still takes connections. `Listener::open` then finds it answered
357    /// and rightly leaves it alone. Measured: 13 of 2000 opens right after a drop found it
358    /// answered while another thread started `true` over and over, none of 2000 when nothing was
359    /// started, and none of 2000 with this wait in between. A QCode
360    /// that ended has no such copy left, so the product has nothing to wait for; the test waits
361    /// for the state it means to set up, for as long as a loaded machine could take.
362    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}