Skip to main content

term_session_client/
remote_pane.rs

1use std::cell::Cell;
2use std::io;
3use std::sync::{Arc, Mutex};
4
5use crossbeam_channel::{Receiver, TryRecvError};
6use muxio_tokio_rpc_ipc_client::RpcIpcClient;
7use portable_pty::{ExitStatus, PtySize};
8use term_session_muxio_service_definitions::{CloseSession, ResizePty};
9use term_wm_pty_engine::{Pane, PtyResult};
10use tokio::runtime::Handle;
11
12type InputWriter = Box<dyn FnMut(&[u8]) -> io::Result<()> + Send>;
13
14pub struct RemotePane {
15    pub id: u64,
16    client: Option<std::sync::Arc<RpcIpcClient>>,
17    rt: Handle,
18    parser: Arc<Mutex<vt100::Parser>>,
19    exited: Cell<bool>,
20    push_rx: Receiver<Vec<u8>>,
21    input_writer: InputWriter,
22}
23
24impl RemotePane {
25    pub fn new(
26        id: u64,
27        client: Option<std::sync::Arc<RpcIpcClient>>,
28        rt: Handle,
29        cols: u16,
30        rows: u16,
31        push_rx: Receiver<Vec<u8>>,
32        input_writer: InputWriter,
33    ) -> Self {
34        Self {
35            id,
36            client,
37            rt,
38            parser: Arc::new(Mutex::new(vt100::Parser::new(rows, cols, 0))),
39            exited: Cell::new(false),
40            push_rx,
41            input_writer,
42        }
43    }
44
45    /// Drain pushes from the server, updating the parser.
46    /// Returns `true` if at least one chunk was processed (screen may have changed).
47    pub fn drain_pushes(&mut self) -> bool {
48        let mut updated = false;
49        loop {
50            match self.push_rx.try_recv() {
51                Ok(data) => {
52                    let mut parser = self.parser.lock().unwrap();
53                    parser.process(&data);
54                    updated = true;
55                }
56                Err(TryRecvError::Disconnected) => {
57                    self.exited.set(true);
58                    break;
59                }
60                Err(TryRecvError::Empty) => break,
61            }
62        }
63        updated
64    }
65
66    fn rpc_to_pty<E: std::fmt::Display>(e: E) -> Box<dyn std::error::Error + Send + Sync> {
67        Box::new(io::Error::other(format!("{e}")))
68    }
69}
70
71impl Pane for RemotePane {
72    fn exit_status(&self) -> Option<ExitStatus> {
73        None
74    }
75
76    fn resize(&mut self, size: PtySize) -> PtyResult<()> {
77        let (actual_cols, actual_rows) = if let Some(ref client) = self.client {
78            let result = self.rt.block_on(async {
79                use muxio_tokio_rpc_ipc_client::RpcCallPrebuffered;
80                ResizePty::call(&**client, (self.id, size.cols, size.rows)).await
81            });
82            result.map_err(Self::rpc_to_pty)?
83        } else {
84            (size.cols, size.rows)
85        };
86        {
87            let mut parser = self.parser.lock().unwrap();
88            parser.screen_mut().set_size(actual_rows, actual_cols);
89        }
90        Ok(())
91    }
92
93    fn has_exited(&mut self) -> bool {
94        self.exited.get()
95    }
96
97    fn alternate_screen(&mut self) -> bool {
98        let parser = self.parser.lock().unwrap();
99        parser.screen().alternate_screen()
100    }
101
102    fn scrollback(&mut self) -> usize {
103        0
104    }
105
106    fn set_scrollback(&mut self, _rows: usize) {}
107
108    fn scrollback_len(&self) -> usize {
109        0
110    }
111
112    fn write_bytes(&mut self, input: &[u8]) -> io::Result<()> {
113        (self.input_writer)(input)
114    }
115
116    fn shared_parser(&mut self) -> Arc<Mutex<vt100::Parser>> {
117        self.parser.clone()
118    }
119
120    fn max_scrollback(&mut self) -> usize {
121        0
122    }
123
124    fn take_exit_status(&mut self) -> Option<ExitStatus> {
125        None
126    }
127
128    fn bytes_received(&self) -> usize {
129        0
130    }
131
132    fn last_bytes_text(&self) -> String {
133        String::new()
134    }
135
136    fn kill_child(&mut self) -> PtyResult<()> {
137        if let Some(ref client) = self.client {
138            self.rt
139                .block_on(async {
140                    use muxio_tokio_rpc_ipc_client::RpcCallPrebuffered;
141                    CloseSession::call(&**client, self.id).await
142                })
143                .map_err(Self::rpc_to_pty)?;
144        }
145        self.exited.set(true);
146        Ok(())
147    }
148
149    fn take_pending_title(&mut self) -> Option<String> {
150        None
151    }
152}
153
154#[cfg(test)]
155mod tests {
156    use super::*;
157
158    #[test]
159    fn test_drain_pushes_returns_dirty_flag() {
160        let rt = tokio::runtime::Builder::new_current_thread()
161            .build()
162            .unwrap();
163
164        let (push_tx, push_rx) = crossbeam_channel::unbounded();
165        let input_writer: InputWriter = Box::new(|_| Ok(()));
166
167        let mut pane = RemotePane::new(1, None, rt.handle().clone(), 80, 24, push_rx, input_writer);
168
169        // 1. Idle call with no pending messages must return false
170        assert!(!pane.drain_pushes());
171
172        // 2. Ingesting bytes must return true
173        push_tx.send(b"hello world".to_vec()).unwrap();
174        assert!(pane.drain_pushes());
175
176        // 3. Subsequent call on drained buffer must return false
177        assert!(!pane.drain_pushes());
178    }
179}