Skip to main content

cliban_server/
shell.rs

1//! Board-over-SSH: bridge one russh channel to the blocking cliban-tui
2//! event loop on a dedicated blocking task.
3//!
4//! Wiring: the connection handler forwards channel data bytes and
5//! window-change resizes into an mpsc ([`RemoteInput`]); the board task
6//! renders through a [`RemoteBackend`] whose writer sends each flush to the
7//! client via the session [`Handle`]. Teardown is symmetric: user quit ends
8//! the loop (exit 0); client disconnect / channel close drops the sender,
9//! the session errors out, and the task dies — no leaked tasks either way.
10
11use std::io::{self, Write};
12use std::sync::mpsc::Receiver;
13use std::sync::Arc;
14
15use cliban_tenancy::Tenant;
16use cliban_tui::app::App;
17use cliban_tui::data::Data;
18use cliban_tui::remote::{ChannelSession, RemoteBackend, RemoteInput};
19use cliban_tui::{picker, runtime};
20use ratatui::Terminal;
21use russh::server::Handle;
22use russh::ChannelId;
23use tokio::sync::broadcast;
24
25use crate::server::AppState;
26
27/// Sent when a shell is requested on a channel that never did pty-req.
28pub const NO_TTY: &str = "cliband: the board requires a TTY; use ssh -t\r\n";
29
30/// `io::Write` onto an SSH channel: buffer locally, send one channel-data
31/// message per flush. `Handle::data` completes once the message enters the
32/// session's mpsc (window-stalled data is buffered by russh server-side),
33/// so the only backpressure is that 100-message queue. Known limitation: a
34/// live client that stops reading TCP can park the session loop and, once
35/// the queue fills, park this blocking thread until the TCP connection
36/// itself dies — an accepted russh-level flow-control gap.
37struct ChannelWriter {
38    rt: tokio::runtime::Handle,
39    handle: Handle,
40    channel: ChannelId,
41    buf: Vec<u8>,
42}
43
44impl Write for ChannelWriter {
45    fn write(&mut self, data: &[u8]) -> io::Result<usize> {
46        self.buf.extend_from_slice(data);
47        Ok(data.len())
48    }
49
50    fn flush(&mut self) -> io::Result<()> {
51        if self.buf.is_empty() {
52            return Ok(());
53        }
54        let data = std::mem::take(&mut self.buf);
55        self.rt
56            .block_on(self.handle.data(self.channel, data))
57            .map_err(|_| io::Error::new(io::ErrorKind::BrokenPipe, "ssh channel closed"))
58    }
59}
60
61/// Everything the blocking board task needs, captured on the async side
62/// (`Handle::current()` is only guaranteed there).
63pub struct BoardTask {
64    pub rt: tokio::runtime::Handle,
65    pub state: Arc<AppState>,
66    pub handle: Handle,
67    pub channel: ChannelId,
68    /// Initial pty size from pty-req.
69    pub size: (u16, u16),
70    /// The user's memberships; picker shown when more than one.
71    pub tenants: Vec<Tenant>,
72    pub input: Receiver<RemoteInput>,
73}
74
75/// Drive picker + board to completion, then hang up. Runs on a blocking
76/// thread (`tokio::task::spawn_blocking`). All sends after the loop are
77/// best-effort: the client may already be gone.
78pub fn run_board(task: BoardTask) {
79    let rt = task.rt.clone();
80    let handle = task.handle.clone();
81    let channel = task.channel;
82    let exit = match board_session(task) {
83        Ok(()) => 0,
84        Err(e) => {
85            // A vanished client surfaces as an io error on the input channel
86            // or the writer — a normal ending, not worth a server log line.
87            let disconnect = e.downcast_ref::<io::Error>().is_some_and(|io| {
88                matches!(
89                    io.kind(),
90                    io::ErrorKind::UnexpectedEof | io::ErrorKind::BrokenPipe
91                )
92            });
93            if !disconnect {
94                eprintln!("cliband: board session ended: {e}");
95            }
96            1
97        }
98    };
99    rt.block_on(async {
100        let _ = handle.exit_status_request(channel, exit).await;
101        let _ = handle.eof(channel).await;
102        let _ = handle.close(channel).await;
103    });
104}
105
106fn board_session(task: BoardTask) -> Result<(), Box<dyn std::error::Error>> {
107    let writer = ChannelWriter {
108        rt: task.rt,
109        handle: task.handle,
110        channel: task.channel,
111        buf: Vec::new(),
112    };
113    let (backend, size) = RemoteBackend::new(writer, task.size.0, task.size.1)?;
114    let mut terminal = Terminal::new(backend)?;
115    let mut session = ChannelSession::new(task.input, size);
116
117    let res = (|| -> Result<(), Box<dyn std::error::Error>> {
118        let picked = if task.tenants.len() == 1 {
119            Some(0)
120        } else {
121            let slugs: Vec<String> = task.tenants.iter().map(|t| t.slug.clone()).collect();
122            picker::pick(&mut terminal, &mut session, " pick a board ", &slugs)?
123        };
124        let Some(i) = picked else { return Ok(()) };
125        // Shared cached Store: every session for this tenant funnels through
126        // one writer thread, and each Data write is a single store.call —
127        // no multi-call write sequence here, so the tenant write_lock is
128        // not needed.
129        let tenant_handle = task.state.manager.handle(&task.tenants[i].id)?;
130
131        // Live updates. Subscribe to the tenant's change feed and poll it
132        // from the session (100ms cadence via the event loop): draining
133        // try_recv coalesces any burst of writes into one Refresh, and no
134        // extra task means teardown semantics are untouched. Publish after
135        // every local commit so sibling sessions refresh too.
136        let mut feed = tenant_handle.changes.subscribe();
137        session.set_dirty_check(move || {
138            let mut dirty = false;
139            loop {
140                match feed.try_recv() {
141                    Ok(()) => dirty = true,
142                    // Lagged means "you missed some": still just a refresh.
143                    Err(broadcast::error::TryRecvError::Lagged(_)) => dirty = true,
144                    Err(_) => break, // Empty (or Closed): nothing pending
145                }
146            }
147            dirty
148        });
149        let publish = tenant_handle.changes.clone();
150        let mut data = Data::from_store(tenant_handle.store.clone())?;
151        data.set_on_mutate(move || {
152            // Err only means no subscriber is listening right now.
153            let _ = publish.send(());
154        });
155        let mut app = App::new();
156        runtime::reload(&data, &mut app)?;
157        runtime::event_loop(&mut terminal, &mut session, &data, &mut app)
158    })();
159
160    // Restore the client's screen whatever happened; if the client is gone
161    // this send just fails quietly.
162    let _ = terminal.backend_mut().leave_screen();
163    res
164}