term_session_client/
remote_pane.rs1use 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 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 assert!(!pane.drain_pushes());
171
172 push_tx.send(b"hello world".to_vec()).unwrap();
174 assert!(pane.drain_pushes());
175
176 assert!(!pane.drain_pushes());
178 }
179}