Skip to main content

muxr_client/
runtime.rs

1use std::io::Read;
2use std::path::Path;
3use std::thread;
4use std::time::Duration;
5
6use muxr_config::KeybindingMode;
7use muxr_config::KeybindingsConfig;
8use muxr_config::LocalKeybindingAction;
9use muxr_config::MuxrConfig;
10use muxr_core::ClientKey;
11use muxr_core::ClientKeyCode;
12use muxr_core::ClientKeyModifiers;
13use muxr_core::ClientMouseEvent;
14use muxr_core::ClientRequest;
15use muxr_core::ServerEvent;
16use muxr_core::SessionName;
17use muxr_core::TerminalSize;
18use muxr_transport::ClientRequestWriter;
19use rootcause::prelude::ResultExt;
20use rootcause::report;
21
22use crate::copy_selection::SelectionEdgeScrollRequest;
23use crate::input::DecodedInput;
24use crate::input::InputDecoder;
25use crate::input::InputIdleTimeout;
26use crate::renderer::ClientPresentationSnapshot;
27use crate::renderer::ClientRenderOutcome;
28use crate::renderer::ClientRenderer;
29use crate::session::attach::AttachedSession;
30use crate::stdout_worker::StdoutSender;
31use crate::stdout_worker::StdoutWorker;
32use crate::terminal::TerminalGuard;
33
34const RESIZE_POLL_INTERVAL: Duration = Duration::from_millis(250);
35const AMBIGUOUS_INPUT_TIMEOUT: Duration = Duration::from_millis(50);
36const SELECTION_EDGE_SCROLL_INTERVAL: Duration = Duration::from_millis(50);
37const STDIN_BUFFER_SIZE: usize = 8192;
38const CONTROL_REQUEST_CHANNEL_LIMIT: usize = 128;
39const INPUT_REQUEST_CHANNEL_LIMIT: usize = 1024;
40
41#[derive(Clone, Debug, Eq, PartialEq)]
42enum StdinRead {
43    Bytes(Vec<u8>),
44    Eof,
45}
46
47#[derive(Clone, Debug, Eq, PartialEq)]
48enum ClientInputAction {
49    CopySelection,
50    CopySelectionInline,
51    ClearSelection,
52    Mouse(ClientMouseEvent),
53}
54
55#[derive(Debug)]
56enum ClientInputCmd {
57    Action(ClientInputAction),
58    Barrier(std::sync::mpsc::SyncSender<()>),
59}
60
61#[derive(Clone, Copy, Debug, Eq, PartialEq)]
62enum InputCmdReceiverState {
63    Open,
64    Closed,
65}
66
67#[derive(Clone, Copy, Debug, Eq, PartialEq)]
68enum LocalActionCompletion {
69    Wait,
70    #[cfg(test)]
71    Skip,
72}
73
74#[derive(Clone, Copy, Debug, Eq, PartialEq)]
75pub enum ClientInputSend {
76    Accepted,
77    Closed,
78}
79
80#[derive(Clone, Copy, Debug, Eq, PartialEq)]
81enum InteractiveFlow {
82    Continue,
83    Stop,
84}
85
86struct RenderCoordinator {
87    committed: Option<ClientPresentationSnapshot>,
88    in_flight: Option<ClientPresentationSnapshot>,
89}
90
91trait RenderSink {
92    fn send_render(&self, transaction: Vec<u8>) -> rootcause::Result<()>;
93}
94
95impl RenderSink for StdoutSender {
96    fn send_render(&self, transaction: Vec<u8>) -> rootcause::Result<()> {
97        Self::send_render(self, transaction)
98    }
99}
100
101impl RenderCoordinator {
102    const fn new() -> Self {
103        Self {
104            committed: None,
105            in_flight: None,
106        }
107    }
108
109    fn complete(&mut self, renderer: &mut ClientRenderer, output: &impl RenderSink) -> rootcause::Result<()> {
110        let completed = self
111            .in_flight
112            .take()
113            .ok_or_else(|| report!("muxr client stdout completed an unknown render"))?;
114        renderer.acknowledge_presentation(&completed);
115        self.committed = Some(completed);
116        self.submit(renderer, output)
117    }
118
119    fn submit(&mut self, renderer: &mut ClientRenderer, output: &impl RenderSink) -> rootcause::Result<()> {
120        if self.in_flight.is_some() {
121            return Ok(());
122        }
123        let snapshot = renderer.presentation_snapshot();
124        let Some(transaction) = renderer.presentation_transaction(self.committed.as_ref())? else {
125            return Ok(());
126        };
127        output.send_render(transaction)?;
128        self.in_flight = Some(snapshot);
129        Ok(())
130    }
131}
132
133/// Start or attach to a muxr session and run an interactive client.
134///
135/// # Errors
136/// - The session paths cannot be resolved.
137/// - The server cannot be started or attached.
138/// - The current terminal size cannot be read.
139/// - Terminal input/output or protocol IO fails.
140pub fn start(session: &SessionName, server_executable: &Path, external_layout: Option<&Path>) -> rootcause::Result<()> {
141    tokio::runtime::Runtime::new()
142        .context("failed to build muxr tokio runtime")?
143        .block_on(async {
144            let muxr_config = MuxrConfig::default();
145            let terminal_size = crate::terminal::current_terminal_size()?;
146            let pane_size = crate::terminal::pane_size_for_terminal(muxr_config.tab_bar.width, &terminal_size)?;
147            let attached_session =
148                crate::session::attach::open_session(session, pane_size.clone(), server_executable, external_layout)
149                    .await?;
150            self::run_interactive(&muxr_config, attached_session, pane_size).await
151        })
152}
153
154async fn run_interactive(
155    muxr_config: &MuxrConfig,
156    mut attached_session: AttachedSession,
157    initial_size: TerminalSize,
158) -> rootcause::Result<()> {
159    let _terminal_guard = TerminalGuard::enable_if_terminal()?;
160
161    let (control_sender, control_receiver) = tokio::sync::mpsc::channel(CONTROL_REQUEST_CHANNEL_LIMIT);
162
163    let (input_cmd_sender, mut input_cmd_receiver) = tokio::sync::mpsc::channel(INPUT_REQUEST_CHANNEL_LIMIT);
164
165    let (input_request_sender, input_receiver) = tokio::sync::mpsc::channel(INPUT_REQUEST_CHANNEL_LIMIT);
166
167    let stdin_handle = self::spawn_stdin_forwarder(
168        muxr_config.keybindings.clone(),
169        input_cmd_sender,
170        input_request_sender.clone(),
171    );
172    let resize_handle = self::spawn_resize_forwarder(control_sender.clone(), muxr_config.tab_bar.width, initial_size);
173
174    let writer = attached_session.writer;
175    let writer_handle =
176        tokio::spawn(async move { self::forward_client_requests(writer, control_receiver, input_receiver).await });
177    let (stdout_sender, _stdout_worker, mut stdout_failure_receiver, mut stdout_completion_receiver) =
178        StdoutWorker::spawn();
179
180    let mut renderer = ClientRenderer::new(muxr_config, attached_session.layout, attached_session.pane_regions);
181    renderer.sync_mouse_capture_logical();
182    let mut render_coordinator = RenderCoordinator::new();
183
184    let edge_scroll_tick_start = tokio::time::Instant::now()
185        .checked_add(SELECTION_EDGE_SCROLL_INTERVAL)
186        .ok_or_else(|| report!("muxr selection edge scroll interval overflowed"))?;
187    let mut edge_scroll_tick = tokio::time::interval_at(edge_scroll_tick_start, SELECTION_EDGE_SCROLL_INTERVAL);
188    edge_scroll_tick.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
189
190    let mut input_cmd_receiver_state = InputCmdReceiverState::Open;
191
192    loop {
193        tokio::select! {
194            event = attached_session.reader.recv_event() => {
195                let Some(event) = event? else {
196                    break;
197                };
198                if self::handle_server_event(
199                    event,
200                    &control_sender,
201                    &mut renderer,
202                    &mut render_coordinator,
203                    &stdout_sender,
204                ).await?
205                    == InteractiveFlow::Stop
206                {
207                    break;
208                }
209            },
210            cmd = input_cmd_receiver.recv(), if input_cmd_receiver_state == InputCmdReceiverState::Open => {
211                let Some(cmd) = cmd else {
212                    input_cmd_receiver_state = InputCmdReceiverState::Closed;
213                    continue;
214                };
215                let action = match cmd {
216                    ClientInputCmd::Action(action) => action,
217                    ClientInputCmd::Barrier(completed) => {
218                        let _sent = completed.send(());
219                        continue;
220                    }
221                };
222                if self::handle_client_input_action(action, muxr_config, &input_request_sender, &mut renderer).await? == ClientInputSend::Closed {
223                    break;
224                }
225                render_coordinator.submit(&mut renderer, &stdout_sender)?;
226            },
227            _ = edge_scroll_tick.tick(), if renderer.selection_edge_drag() == crate::renderer::SelectionEdgeDrag::Active => {
228                if self::send_selection_edge_scroll_request(&input_request_sender, &mut renderer) == ClientInputSend::Closed {
229                    break;
230                }
231            },
232            stdout_failure = &mut stdout_failure_receiver => {
233                let error = stdout_failure.unwrap_or_else(|_| "stdout worker stopped unexpectedly".to_owned());
234                return Err(report!("muxr client stdout worker failed").attach(error));
235            },
236            completed = stdout_completion_receiver.recv() => {
237                if completed.is_none() {
238                    return Err(report!("muxr client stdout completion worker stopped unexpectedly"));
239                }
240                render_coordinator.complete(&mut renderer, &stdout_sender)?;
241            },
242            else => {
243                if input_cmd_receiver_state == InputCmdReceiverState::Closed {
244                    break;
245                }
246            }
247        }
248    }
249
250    writer_handle.abort();
251    drop(writer_handle.await);
252    drop(stdin_handle);
253    drop(resize_handle);
254    Ok(())
255}
256
257async fn handle_server_event(
258    event: ServerEvent,
259    control_sender: &tokio::sync::mpsc::Sender<ClientRequest>,
260    renderer: &mut ClientRenderer,
261    render_coordinator: &mut RenderCoordinator,
262    stdout_sender: &StdoutSender,
263) -> rootcause::Result<InteractiveFlow> {
264    match event {
265        ServerEvent::Deleted | ServerEvent::Detached => Ok(InteractiveFlow::Stop),
266        ServerEvent::Error(error) => Err(report!("muxr server returned error")
267            .attach(format!("code={}", error.code()))
268            .attach(format!("msg={}", error.msg()))),
269        ServerEvent::Ping => Ok(if control_sender.send(ClientRequest::Pong).await.is_ok() {
270            InteractiveFlow::Continue
271        } else {
272            InteractiveFlow::Stop
273        }),
274        ServerEvent::Layout(next_layout) => {
275            renderer.apply_layout(next_layout);
276            render_coordinator.submit(renderer, stdout_sender)?;
277            Ok(InteractiveFlow::Continue)
278        }
279        ServerEvent::SidebarLayout(next_layout) => {
280            renderer.apply_sidebar_layout_logical(next_layout);
281            render_coordinator.submit(renderer, stdout_sender)?;
282            Ok(InteractiveFlow::Continue)
283        }
284        ServerEvent::PaneRegions(next_regions) => {
285            renderer.apply_pane_regions_logical(next_regions);
286            render_coordinator.submit(renderer, stdout_sender)?;
287            Ok(InteractiveFlow::Continue)
288        }
289        ServerEvent::Render(update) => {
290            self::handle_render_event(update, control_sender, renderer, render_coordinator, stdout_sender).await
291        }
292        ServerEvent::ScrollPaneLineResult {
293            position,
294            direction,
295            movement,
296        } => {
297            renderer.apply_scroll_pane_line_result(position, direction, movement);
298            Ok(InteractiveFlow::Continue)
299        }
300        ServerEvent::Attached(_) | ServerEvent::Pong => Ok(InteractiveFlow::Continue),
301    }
302}
303
304async fn handle_render_event(
305    update: muxr_core::RenderUpdate,
306    control_sender: &tokio::sync::mpsc::Sender<ClientRequest>,
307    renderer: &mut ClientRenderer,
308    render_coordinator: &mut RenderCoordinator,
309    stdout_sender: &StdoutSender,
310) -> rootcause::Result<InteractiveFlow> {
311    match renderer.apply_render_logical(update)? {
312        ClientRenderOutcome::Drawn => {
313            render_coordinator.submit(renderer, stdout_sender)?;
314            Ok(InteractiveFlow::Continue)
315        }
316        ClientRenderOutcome::NeedsResync => Ok(if control_sender.send(ClientRequest::RenderResync).await.is_ok() {
317            InteractiveFlow::Continue
318        } else {
319            InteractiveFlow::Stop
320        }),
321    }
322}
323
324async fn handle_client_input_action(
325    action: ClientInputAction,
326    muxr_config: &MuxrConfig,
327    input_sender: &tokio::sync::mpsc::Sender<ClientRequest>,
328    renderer: &mut ClientRenderer,
329) -> rootcause::Result<ClientInputSend> {
330    match action {
331        ClientInputAction::CopySelection => {
332            renderer.copy_selection()?;
333            Ok(ClientInputSend::Accepted)
334        }
335        ClientInputAction::CopySelectionInline => {
336            renderer.copy_selection_inline()?;
337            Ok(ClientInputSend::Accepted)
338        }
339        ClientInputAction::ClearSelection => {
340            renderer.clear_selection();
341            Ok(ClientInputSend::Accepted)
342        }
343        ClientInputAction::Mouse(event) => {
344            crate::pane::mouse::handle_mouse_input_action(muxr_config, event, input_sender, renderer).await
345        }
346    }
347}
348
349#[derive(Clone, Copy, Debug, Eq, PartialEq)]
350pub enum DroppableSendOutcome {
351    Closed,
352    Dropped,
353    Sent,
354}
355
356pub fn send_droppable_request(
357    input_sender: &tokio::sync::mpsc::Sender<ClientRequest>,
358    request: ClientRequest,
359) -> DroppableSendOutcome {
360    match input_sender.try_send(request) {
361        Ok(()) => DroppableSendOutcome::Sent,
362        Err(tokio::sync::mpsc::error::TrySendError::Full(request)) => {
363            drop(request);
364            DroppableSendOutcome::Dropped
365        }
366        Err(tokio::sync::mpsc::error::TrySendError::Closed(request)) => {
367            drop(request);
368            DroppableSendOutcome::Closed
369        }
370    }
371}
372
373pub fn send_edge_scroll_request(
374    input_sender: &tokio::sync::mpsc::Sender<ClientRequest>,
375    renderer: &mut ClientRenderer,
376    request: SelectionEdgeScrollRequest,
377) -> ClientInputSend {
378    let (pending, request) = request.into_parts();
379    match self::send_droppable_request(input_sender, request) {
380        DroppableSendOutcome::Sent => {
381            // One queued edge-scroll request must be paired with one moved viewport and its render before another
382            // request is queued; otherwise coalesced renders can skip selected content rows.
383            renderer.mark_selection_edge_scroll_sent(pending);
384            ClientInputSend::Accepted
385        }
386        DroppableSendOutcome::Dropped => ClientInputSend::Accepted,
387        DroppableSendOutcome::Closed => ClientInputSend::Closed,
388    }
389}
390
391fn send_selection_edge_scroll_request(
392    input_sender: &tokio::sync::mpsc::Sender<ClientRequest>,
393    renderer: &mut ClientRenderer,
394) -> ClientInputSend {
395    let Some(request) = renderer.selection_edge_scroll_request() else {
396        return ClientInputSend::Accepted;
397    };
398    self::send_edge_scroll_request(input_sender, renderer, request)
399}
400
401async fn forward_client_requests(
402    mut writer: ClientRequestWriter,
403    mut control_receiver: tokio::sync::mpsc::Receiver<ClientRequest>,
404    mut input_receiver: tokio::sync::mpsc::Receiver<ClientRequest>,
405) -> rootcause::Result<()> {
406    let mut control_closed = false;
407    let mut input_closed = false;
408
409    loop {
410        if control_closed && input_closed {
411            break;
412        }
413
414        tokio::select! {
415            biased;
416            request = control_receiver.recv(), if !control_closed => match request {
417                Some(request) => {
418                    if writer.send_request(&request).await.is_err() {
419                        break;
420                    }
421                }
422                None => control_closed = true,
423            },
424            request = input_receiver.recv(), if !input_closed => match request {
425                Some(request) => {
426                    if writer.send_request(&request).await.is_err() {
427                        break;
428                    }
429                }
430                None => input_closed = true,
431            },
432        }
433    }
434
435    Ok(())
436}
437
438fn spawn_stdin_forwarder(
439    keybindings: KeybindingsConfig,
440    cmd_sender: tokio::sync::mpsc::Sender<ClientInputCmd>,
441    request_sender: tokio::sync::mpsc::Sender<ClientRequest>,
442) -> thread::JoinHandle<()> {
443    thread::spawn(move || {
444        let (read_sender, read_receiver) = std::sync::mpsc::channel();
445        drop(self::spawn_stdin_reader(read_sender));
446        let mut decoder = InputDecoder::with_keybindings(keybindings.clone());
447
448        loop {
449            // Ambiguous escape prefixes need an idle timeout. Bracketed paste waits for its terminator so slow
450            // multi-chunk paste cannot leak raw paste markers into the PTY.
451            let read = if decoder.idle_timeout() == InputIdleTimeout::Needed {
452                match read_receiver.recv_timeout(AMBIGUOUS_INPUT_TIMEOUT) {
453                    Ok(read) => read,
454                    Err(std::sync::mpsc::RecvTimeoutError::Timeout) => {
455                        if self::send_decoded_input_with_ordering(
456                            &keybindings,
457                            &cmd_sender,
458                            &request_sender,
459                            decoder.finalize(),
460                            LocalActionCompletion::Wait,
461                        ) == ClientInputSend::Closed
462                        {
463                            break;
464                        }
465                        continue;
466                    }
467                    Err(std::sync::mpsc::RecvTimeoutError::Disconnected) => StdinRead::Eof,
468                }
469            } else {
470                read_receiver.recv().unwrap_or(StdinRead::Eof)
471            };
472
473            match read {
474                StdinRead::Bytes(bytes) => {
475                    if self::send_decoded_input_with_ordering(
476                        &keybindings,
477                        &cmd_sender,
478                        &request_sender,
479                        decoder.decode(&bytes),
480                        LocalActionCompletion::Wait,
481                    ) == ClientInputSend::Closed
482                    {
483                        break;
484                    }
485                }
486                StdinRead::Eof => {
487                    if self::send_decoded_input_with_ordering(
488                        &keybindings,
489                        &cmd_sender,
490                        &request_sender,
491                        decoder.finalize(),
492                        LocalActionCompletion::Wait,
493                    ) == ClientInputSend::Closed
494                    {
495                        break;
496                    }
497                    // EOF detach follows any queued stdin bytes so piped cmds like `exit\n` reach the shell first.
498                    drop(request_sender.blocking_send(ClientRequest::Detach));
499                    break;
500                }
501            }
502        }
503    })
504}
505
506fn spawn_stdin_reader(sender: std::sync::mpsc::Sender<StdinRead>) -> thread::JoinHandle<()> {
507    thread::spawn(move || {
508        let mut stdin = std::io::stdin();
509        let mut buffer = [0; STDIN_BUFFER_SIZE];
510
511        loop {
512            match stdin.read(&mut buffer) {
513                Ok(0) | Err(_) => {
514                    drop(sender.send(StdinRead::Eof));
515                    break;
516                }
517                Ok(bytes_read) => {
518                    let Some(bytes) = buffer.get(..bytes_read) else {
519                        drop(sender.send(StdinRead::Eof));
520                        break;
521                    };
522                    if sender.send(StdinRead::Bytes(bytes.to_vec())).is_err() {
523                        break;
524                    }
525                }
526            }
527        }
528    })
529}
530
531#[cfg(test)]
532fn send_decoded_input(
533    cmd_sender: &tokio::sync::mpsc::Sender<ClientInputCmd>,
534    request_sender: &tokio::sync::mpsc::Sender<ClientRequest>,
535    decoded: Vec<DecodedInput>,
536) -> ClientInputSend {
537    let keybindings = MuxrConfig::default().keybindings;
538    self::send_decoded_input_with_ordering(
539        &keybindings,
540        cmd_sender,
541        request_sender,
542        decoded,
543        LocalActionCompletion::Skip,
544    )
545}
546
547fn send_decoded_input_with_ordering(
548    keybindings: &KeybindingsConfig,
549    cmd_sender: &tokio::sync::mpsc::Sender<ClientInputCmd>,
550    request_sender: &tokio::sync::mpsc::Sender<ClientRequest>,
551    decoded: Vec<DecodedInput>,
552    local_action_completion: LocalActionCompletion,
553) -> ClientInputSend {
554    let mut selection_reset_sent = false;
555    let mut pending_input = Vec::new();
556
557    for decoded in decoded {
558        if self::send_decoded_event(
559            keybindings,
560            cmd_sender,
561            request_sender,
562            decoded,
563            &mut pending_input,
564            local_action_completion,
565            &mut selection_reset_sent,
566        ) == ClientInputSend::Closed
567        {
568            return ClientInputSend::Closed;
569        }
570    }
571
572    self::send_pending_input(
573        cmd_sender,
574        request_sender,
575        &mut pending_input,
576        local_action_completion,
577        &mut selection_reset_sent,
578    )
579}
580
581fn send_decoded_event(
582    keybindings: &KeybindingsConfig,
583    cmd_sender: &tokio::sync::mpsc::Sender<ClientInputCmd>,
584    request_sender: &tokio::sync::mpsc::Sender<ClientRequest>,
585    decoded: DecodedInput,
586    pending_input: &mut Vec<u8>,
587    local_action_completion: LocalActionCompletion,
588    selection_reset_sent: &mut bool,
589) -> ClientInputSend {
590    match decoded {
591        DecodedInput::Input(bytes) => pending_input.extend(bytes),
592        DecodedInput::Key(key) => {
593            return self::send_key_input(
594                keybindings,
595                cmd_sender,
596                request_sender,
597                key,
598                pending_input,
599                local_action_completion,
600                selection_reset_sent,
601            );
602        }
603        DecodedInput::Mouse(event) => {
604            return self::send_mouse_input(
605                cmd_sender,
606                event,
607                pending_input,
608                request_sender,
609                local_action_completion,
610                selection_reset_sent,
611            );
612        }
613        DecodedInput::Paste(bytes) => {
614            return self::send_paste_input(
615                cmd_sender,
616                request_sender,
617                bytes,
618                pending_input,
619                local_action_completion,
620                selection_reset_sent,
621            );
622        }
623    }
624    ClientInputSend::Accepted
625}
626
627fn send_key_input(
628    keybindings: &KeybindingsConfig,
629    cmd_sender: &tokio::sync::mpsc::Sender<ClientInputCmd>,
630    request_sender: &tokio::sync::mpsc::Sender<ClientRequest>,
631    key: ClientKey,
632    pending_input: &mut Vec<u8>,
633    local_action_completion: LocalActionCompletion,
634    selection_reset_sent: &mut bool,
635) -> ClientInputSend {
636    if let Some(action) = keybindings.resolve_local(&key) {
637        if self::send_pending_input(
638            cmd_sender,
639            request_sender,
640            pending_input,
641            local_action_completion,
642            selection_reset_sent,
643        ) == ClientInputSend::Closed
644        {
645            return ClientInputSend::Closed;
646        }
647        let action = match action {
648            LocalKeybindingAction::CopySelection => ClientInputAction::CopySelection,
649            LocalKeybindingAction::CopySelectionInline => ClientInputAction::CopySelectionInline,
650        };
651        return self::send_input_action(cmd_sender, action, local_action_completion);
652    }
653    if let Some(byte) = self::plain_input_byte(keybindings, &key) {
654        pending_input.push(byte);
655        return ClientInputSend::Accepted;
656    }
657    if self::send_pending_input(
658        cmd_sender,
659        request_sender,
660        pending_input,
661        local_action_completion,
662        selection_reset_sent,
663    ) == ClientInputSend::Closed
664    {
665        return ClientInputSend::Closed;
666    }
667    if self::clear_selection_before_input(cmd_sender, local_action_completion, selection_reset_sent)
668        == ClientInputSend::Closed
669    {
670        return ClientInputSend::Closed;
671    }
672    if request_sender.blocking_send(ClientRequest::Key(key)).is_err() {
673        return ClientInputSend::Closed;
674    }
675    ClientInputSend::Accepted
676}
677
678fn send_mouse_input(
679    cmd_sender: &tokio::sync::mpsc::Sender<ClientInputCmd>,
680    event: ClientMouseEvent,
681    pending_input: &mut Vec<u8>,
682    request_sender: &tokio::sync::mpsc::Sender<ClientRequest>,
683    local_action_completion: LocalActionCompletion,
684    selection_reset_sent: &mut bool,
685) -> ClientInputSend {
686    if self::send_pending_input(
687        cmd_sender,
688        request_sender,
689        pending_input,
690        local_action_completion,
691        selection_reset_sent,
692    ) == ClientInputSend::Closed
693    {
694        return ClientInputSend::Closed;
695    }
696    *selection_reset_sent = false;
697    let action = ClientInputAction::Mouse(event);
698    if crate::pane::mouse::MouseEventDrop::from(event) == crate::pane::mouse::MouseEventDrop::Droppable {
699        self::send_droppable_input_action(cmd_sender, action, local_action_completion)
700    } else {
701        self::send_input_action(cmd_sender, action, local_action_completion)
702    }
703}
704
705fn send_paste_input(
706    cmd_sender: &tokio::sync::mpsc::Sender<ClientInputCmd>,
707    request_sender: &tokio::sync::mpsc::Sender<ClientRequest>,
708    bytes: Vec<u8>,
709    pending_input: &mut Vec<u8>,
710    local_action_completion: LocalActionCompletion,
711    selection_reset_sent: &mut bool,
712) -> ClientInputSend {
713    if self::send_pending_input(
714        cmd_sender,
715        request_sender,
716        pending_input,
717        local_action_completion,
718        selection_reset_sent,
719    ) == ClientInputSend::Closed
720    {
721        return ClientInputSend::Closed;
722    }
723    if self::clear_selection_before_input(cmd_sender, local_action_completion, selection_reset_sent)
724        == ClientInputSend::Closed
725    {
726        return ClientInputSend::Closed;
727    }
728    if request_sender.blocking_send(ClientRequest::Paste(bytes)).is_err() {
729        return ClientInputSend::Closed;
730    }
731    ClientInputSend::Accepted
732}
733
734fn send_pending_input(
735    cmd_sender: &tokio::sync::mpsc::Sender<ClientInputCmd>,
736    request_sender: &tokio::sync::mpsc::Sender<ClientRequest>,
737    pending_input: &mut Vec<u8>,
738    local_action_completion: LocalActionCompletion,
739    selection_reset_sent: &mut bool,
740) -> ClientInputSend {
741    if pending_input.is_empty() {
742        return ClientInputSend::Accepted;
743    }
744
745    if self::clear_selection_before_input(cmd_sender, local_action_completion, selection_reset_sent)
746        == ClientInputSend::Closed
747    {
748        return ClientInputSend::Closed;
749    }
750    if request_sender
751        .blocking_send(ClientRequest::Input(std::mem::take(pending_input)))
752        .is_err()
753    {
754        return ClientInputSend::Closed;
755    }
756    ClientInputSend::Accepted
757}
758
759fn plain_input_byte(keybindings: &KeybindingsConfig, key: &ClientKey) -> Option<u8> {
760    if keybindings.resolve(KeybindingMode::Normal, key).is_some()
761        || keybindings.resolve(KeybindingMode::Resize, key).is_some()
762    {
763        return None;
764    }
765
766    if key.modifiers != ClientKeyModifiers::NONE {
767        return None;
768    }
769    let ClientKeyCode::Char(character) = key.code else {
770        return None;
771    };
772    if !character.is_ascii() || character.is_ascii_control() {
773        return None;
774    }
775    let byte = u8::try_from(u32::from(character)).ok()?;
776    (key.raw_bytes.as_slice() == std::slice::from_ref(&byte)).then_some(byte)
777}
778
779fn clear_selection_before_input(
780    cmd_sender: &tokio::sync::mpsc::Sender<ClientInputCmd>,
781    local_action_completion: LocalActionCompletion,
782    selection_reset_sent: &mut bool,
783) -> ClientInputSend {
784    if *selection_reset_sent {
785        return ClientInputSend::Accepted;
786    }
787
788    let result = self::send_input_action(cmd_sender, ClientInputAction::ClearSelection, local_action_completion);
789    if result == ClientInputSend::Accepted {
790        *selection_reset_sent = true;
791    }
792    result
793}
794
795fn send_droppable_input_action(
796    cmd_sender: &tokio::sync::mpsc::Sender<ClientInputCmd>,
797    action: ClientInputAction,
798    local_action_completion: LocalActionCompletion,
799) -> ClientInputSend {
800    match cmd_sender.try_send(ClientInputCmd::Action(action)) {
801        Ok(()) => match local_action_completion {
802            LocalActionCompletion::Wait => self::send_input_action_barrier(cmd_sender),
803            #[cfg(test)]
804            LocalActionCompletion::Skip => ClientInputSend::Accepted,
805        },
806        Err(tokio::sync::mpsc::error::TrySendError::Full(action)) => {
807            drop(action);
808            ClientInputSend::Accepted
809        }
810        Err(tokio::sync::mpsc::error::TrySendError::Closed(action)) => {
811            drop(action);
812            ClientInputSend::Closed
813        }
814    }
815}
816
817fn send_input_action(
818    cmd_sender: &tokio::sync::mpsc::Sender<ClientInputCmd>,
819    action: ClientInputAction,
820    local_action_completion: LocalActionCompletion,
821) -> ClientInputSend {
822    if cmd_sender.blocking_send(ClientInputCmd::Action(action)).is_err() {
823        return ClientInputSend::Closed;
824    }
825    match local_action_completion {
826        LocalActionCompletion::Wait => self::send_input_action_barrier(cmd_sender),
827        #[cfg(test)]
828        LocalActionCompletion::Skip => ClientInputSend::Accepted,
829    }
830}
831
832fn send_input_action_barrier(cmd_sender: &tokio::sync::mpsc::Sender<ClientInputCmd>) -> ClientInputSend {
833    let (completed_sender, completed_receiver) = std::sync::mpsc::sync_channel(0);
834    if cmd_sender
835        .blocking_send(ClientInputCmd::Barrier(completed_sender))
836        .is_err()
837    {
838        return ClientInputSend::Closed;
839    }
840    if completed_receiver.recv().is_err() {
841        return ClientInputSend::Closed;
842    }
843    ClientInputSend::Accepted
844}
845
846fn spawn_resize_forwarder(
847    sender: tokio::sync::mpsc::Sender<ClientRequest>,
848    tab_bar_width: u16,
849    initial_size: TerminalSize,
850) -> thread::JoinHandle<()> {
851    thread::spawn(move || {
852        let mut last_size = initial_size;
853
854        loop {
855            if sender.is_closed() {
856                break;
857            }
858
859            thread::sleep(RESIZE_POLL_INTERVAL);
860            let Ok(next_terminal_size) = crate::terminal::current_terminal_size() else {
861                break;
862            };
863            // Resize requests use the pane viewport, because left-side host-terminal columns are reserved for tab UI.
864            let Ok(next_size) = crate::terminal::pane_size_for_terminal(tab_bar_width, &next_terminal_size) else {
865                break;
866            };
867            if next_size == last_size {
868                continue;
869            }
870
871            if sender.blocking_send(ClientRequest::Resize(next_size.clone())).is_err() {
872                break;
873            }
874            last_size = next_size;
875        }
876    })
877}
878
879#[cfg(test)]
880mod tests {
881    use std::cell::RefCell;
882    use std::fs;
883    use std::path::Path;
884
885    use muxr_core::ClientKey;
886    use muxr_core::ClientKeyCode;
887    use muxr_core::ClientKeyModifiers;
888    use muxr_core::ClientMouseEventPhase;
889    use muxr_core::ClientMousePosition;
890    use muxr_core::LayoutSnapshot;
891    use muxr_core::PaneId;
892    use muxr_core::PaneRegionsSnapshot;
893    use muxr_core::PaneScrollDirection;
894    use muxr_core::PaneSnapshot;
895    use muxr_core::SessionPaths;
896    use muxr_core::TabId;
897    use muxr_core::TabSnapshot;
898    use muxr_transport::ClientConnection;
899    use muxr_transport::ServerListener;
900    use test_that::prelude::*;
901
902    use super::*;
903    use crate::copy_selection::SelectionInput;
904    use crate::copy_selection::test_helpers as copy_selection_test_helpers;
905    use crate::terminal::SynchronizedOutput;
906
907    #[derive(Default)]
908    struct RecordingRenderSink {
909        transactions: RefCell<Vec<Vec<u8>>>,
910    }
911
912    impl RecordingRenderSink {
913        fn len(&self) -> usize {
914            self.transactions.borrow().len()
915        }
916
917        fn transaction(&self, index: usize) -> rootcause::Result<Vec<u8>> {
918            self.transactions
919                .borrow()
920                .get(index)
921                .cloned()
922                .ok_or_else(|| report!("muxr coordinator test transaction is missing").attach(format!("index={index}")))
923        }
924    }
925
926    impl RenderSink for RecordingRenderSink {
927        fn send_render(&self, transaction: Vec<u8>) -> rootcause::Result<()> {
928            self.transactions.borrow_mut().push(transaction);
929            Ok(())
930        }
931    }
932
933    #[test]
934    fn test_render_coordinator_when_logical_state_changes_in_flight_emits_one_delta_after_completion()
935    -> rootcause::Result<()> {
936        let mut renderer = ClientRenderer::with_synchronized_output(
937            layout_snapshot()?,
938            pane_regions_snapshot()?,
939            SynchronizedOutput::Csi,
940        );
941        let mut coordinator = RenderCoordinator::new();
942        let output = RecordingRenderSink::default();
943
944        renderer.apply_render_logical(muxr_core::RenderUpdate::Baseline(render_baseline()?))?;
945        coordinator.submit(&mut renderer, &output)?;
946        assert_that!(output.len(), eq(1));
947
948        renderer.apply_selection_input_logical(SelectionInput::Start(ClientMousePosition { row: 0, col: 0 }))?;
949        renderer.apply_selection_input_logical(SelectionInput::Update(ClientMousePosition { row: 0, col: 1 }))?;
950        coordinator.submit(&mut renderer, &output)?;
951        assert_that!(output.len(), eq(1));
952
953        coordinator.complete(&mut renderer, &output)?;
954        assert_that!(output.len(), eq(2));
955        let initial = String::from_utf8(output.transaction(0)?)?;
956        let delta = String::from_utf8(output.transaction(1)?)?;
957        assert_that!(initial, contains_substring("/tmp"));
958        assert_that!(delta, not(contains_substring("/tmp")));
959        assert_that!(delta, not(contains_substring("\x1b[2J")));
960
961        coordinator.complete(&mut renderer, &output)?;
962        assert_that!(output.len(), eq(2));
963        Ok(())
964    }
965
966    #[test]
967    fn test_forward_client_requests_when_input_queue_is_ready_sends_control_first() -> rootcause::Result<()> {
968        self::runtime()?.block_on(async {
969            let tempdir = tempfile::tempdir()?;
970            let (_, paths) = self::session_paths(tempdir.path(), "work")?;
971            fs::create_dir_all(&paths.root)?;
972            let listener = ServerListener::bind(&paths.socket)?;
973            let server_handle = tokio::spawn(async move {
974                let mut connection = listener.accept().await?;
975                let Some(request) = connection.recv_request().await? else {
976                    return Err(report!("expected forwarded client request"));
977                };
978                Ok::<ClientRequest, rootcause::Report>(request)
979            });
980
981            let connection = ClientConnection::connect(&paths.socket).await?;
982            let (_reader, writer) = connection.split();
983            let (control_sender, control_receiver) = tokio::sync::mpsc::channel(1);
984            let (input_sender, input_receiver) = tokio::sync::mpsc::channel(1);
985            assert_that!(input_sender.try_send(ClientRequest::Input(vec![b'a'])), ok(eq(())));
986            assert_that!(input_sender.try_send(ClientRequest::Input(vec![b'b'])), err(anything()));
987            assert_that!(control_sender.try_send(ClientRequest::Pong), ok(eq(())));
988
989            let writer_handle = tokio::spawn(self::forward_client_requests(writer, control_receiver, input_receiver));
990            let first_request = server_handle
991                .await
992                .map_err(|error| report!("muxr forward test socket task panicked").attach(format!("{error}")))??;
993
994            assert_that!(first_request, eq(ClientRequest::Pong));
995            drop(control_sender);
996            drop(input_sender);
997            writer_handle
998                .await
999                .map_err(|error| report!("muxr forward test writer task panicked").attach(format!("{error}")))??;
1000            Ok(())
1001        })
1002    }
1003
1004    #[test]
1005    fn test_forward_client_requests_when_stdin_requests_are_mixed_sends_input_queue_in_order() -> rootcause::Result<()>
1006    {
1007        self::runtime()?.block_on(async {
1008            let tempdir = tempfile::tempdir()?;
1009            let (_, paths) = self::session_paths(tempdir.path(), "work")?;
1010            fs::create_dir_all(&paths.root)?;
1011            let listener = ServerListener::bind(&paths.socket)?;
1012            let server_handle = tokio::spawn(async move {
1013                let mut connection = listener.accept().await?;
1014                let mut requests = Vec::new();
1015                for _ in 0..3 {
1016                    let Some(request) = connection.recv_request().await? else {
1017                        return Err(report!("expected forwarded stdin request"));
1018                    };
1019                    requests.push(request);
1020                }
1021                Ok::<Vec<ClientRequest>, rootcause::Report>(requests)
1022            });
1023
1024            let connection = ClientConnection::connect(&paths.socket).await?;
1025            let (_reader, writer) = connection.split();
1026            let (control_sender, control_receiver) = tokio::sync::mpsc::channel(1);
1027            let (input_sender, input_receiver) = tokio::sync::mpsc::channel(3);
1028            let key = ClientKey {
1029                code: ClientKeyCode::Char('E'),
1030                modifiers: ClientKeyModifiers::SHIFT_ALT,
1031                raw_bytes: b"\x1bE".to_vec(),
1032            };
1033            assert_that!(input_sender.try_send(ClientRequest::Input(b"a".to_vec())), ok(eq(())));
1034            assert_that!(input_sender.try_send(ClientRequest::Key(key.clone())), ok(eq(())));
1035            assert_that!(input_sender.try_send(ClientRequest::Input(b"b".to_vec())), ok(eq(())));
1036            drop(control_sender);
1037            drop(input_sender);
1038
1039            let writer_handle = tokio::spawn(self::forward_client_requests(writer, control_receiver, input_receiver));
1040            let requests = server_handle.await.map_err(|error| {
1041                report!("muxr forward order test socket task panicked").attach(format!("{error}"))
1042            })??;
1043
1044            assert_that!(
1045                requests,
1046                eq(vec![
1047                    ClientRequest::Input(b"a".to_vec()),
1048                    ClientRequest::Key(key),
1049                    ClientRequest::Input(b"b".to_vec()),
1050                ])
1051            );
1052            writer_handle.await.map_err(|error| {
1053                report!("muxr forward order test writer task panicked").attach(format!("{error}"))
1054            })??;
1055            Ok(())
1056        })
1057    }
1058
1059    #[test]
1060    fn test_forward_client_requests_when_stdin_detach_follows_input_sends_input_before_detach() -> rootcause::Result<()>
1061    {
1062        self::runtime()?.block_on(async {
1063            let tempdir = tempfile::tempdir()?;
1064            let (_, paths) = self::session_paths(tempdir.path(), "work")?;
1065            fs::create_dir_all(&paths.root)?;
1066            let listener = ServerListener::bind(&paths.socket)?;
1067            let server_handle = tokio::spawn(async move {
1068                let mut connection = listener.accept().await?;
1069                let mut requests = Vec::new();
1070                for _ in 0..2 {
1071                    let Some(request) = connection.recv_request().await? else {
1072                        return Err(report!("expected forwarded stdin detach request"));
1073                    };
1074                    requests.push(request);
1075                }
1076                Ok::<Vec<ClientRequest>, rootcause::Report>(requests)
1077            });
1078
1079            let connection = ClientConnection::connect(&paths.socket).await?;
1080            let (_reader, writer) = connection.split();
1081            let (control_sender, control_receiver) = tokio::sync::mpsc::channel(1);
1082            let (input_sender, input_receiver) = tokio::sync::mpsc::channel(2);
1083            assert_that!(
1084                input_sender.try_send(ClientRequest::Input(b"exit\n".to_vec())),
1085                ok(eq(()))
1086            );
1087            assert_that!(input_sender.try_send(ClientRequest::Detach), ok(eq(())));
1088            drop(control_sender);
1089            drop(input_sender);
1090
1091            let writer_handle = tokio::spawn(self::forward_client_requests(writer, control_receiver, input_receiver));
1092            let requests = server_handle
1093                .await
1094                .map_err(|error| report!("muxr forward EOF test socket task panicked").attach(format!("{error}")))??;
1095
1096            assert_that!(
1097                requests,
1098                eq(vec![ClientRequest::Input(b"exit\n".to_vec()), ClientRequest::Detach])
1099            );
1100            writer_handle
1101                .await
1102                .map_err(|error| report!("muxr forward EOF test writer task panicked").attach(format!("{error}")))??;
1103            Ok(())
1104        })
1105    }
1106
1107    #[test]
1108    fn test_send_decoded_input_when_key_arrives_clears_selection_once_and_preserves_request_order() {
1109        let (cmd_sender, mut cmd_receiver) = tokio::sync::mpsc::channel(1);
1110        let (request_sender, mut request_receiver) = tokio::sync::mpsc::channel(3);
1111        let key = ClientKey {
1112            code: ClientKeyCode::Char('E'),
1113            modifiers: ClientKeyModifiers::SHIFT_ALT,
1114            raw_bytes: b"\x1bE".to_vec(),
1115        };
1116
1117        assert_that!(
1118            send_decoded_input(
1119                &cmd_sender,
1120                &request_sender,
1121                vec![
1122                    DecodedInput::Input(b"a".to_vec()),
1123                    DecodedInput::Key(key.clone()),
1124                    DecodedInput::Input(b"b".to_vec()),
1125                ],
1126            ),
1127            eq(ClientInputSend::Accepted)
1128        );
1129
1130        assert_that!(
1131            request_receiver.blocking_recv(),
1132            eq(Some(ClientRequest::Input(b"a".to_vec())))
1133        );
1134        assert_that!(request_receiver.blocking_recv(), eq(Some(ClientRequest::Key(key))));
1135        assert_that!(
1136            request_receiver.blocking_recv(),
1137            eq(Some(ClientRequest::Input(b"b".to_vec())))
1138        );
1139        assert_that!(
1140            matches!(
1141                cmd_receiver.blocking_recv(),
1142                Some(ClientInputCmd::Action(ClientInputAction::ClearSelection))
1143            ),
1144            eq(true)
1145        );
1146        assert_that!(cmd_receiver.try_recv().is_err(), eq(true));
1147    }
1148
1149    #[test]
1150    fn test_send_decoded_input_when_contiguous_plain_keys_arrive_batches_input() {
1151        let (cmd_sender, _cmd_receiver) = tokio::sync::mpsc::channel(1);
1152        let (request_sender, mut request_receiver) = tokio::sync::mpsc::channel(1);
1153        let key = |character| ClientKey {
1154            code: ClientKeyCode::Char(character),
1155            modifiers: ClientKeyModifiers::NONE,
1156            raw_bytes: vec![character as u8],
1157        };
1158
1159        assert_that!(
1160            send_decoded_input(
1161                &cmd_sender,
1162                &request_sender,
1163                vec![
1164                    DecodedInput::Key(key('a')),
1165                    DecodedInput::Key(key('b')),
1166                    DecodedInput::Key(key('c')),
1167                ],
1168            ),
1169            eq(ClientInputSend::Accepted)
1170        );
1171
1172        assert_that!(
1173            request_receiver.blocking_recv(),
1174            eq(Some(ClientRequest::Input(b"abc".to_vec())))
1175        );
1176    }
1177
1178    #[test]
1179    fn test_send_decoded_input_when_server_bound_plain_key_arrives_preserves_key_request() {
1180        let (cmd_sender, _cmd_receiver) = tokio::sync::mpsc::channel(1);
1181        let (request_sender, mut request_receiver) = tokio::sync::mpsc::channel(1);
1182        let key = ClientKey {
1183            code: ClientKeyCode::Char('h'),
1184            modifiers: ClientKeyModifiers::NONE,
1185            raw_bytes: b"h".to_vec(),
1186        };
1187
1188        assert_that!(
1189            send_decoded_input(&cmd_sender, &request_sender, vec![DecodedInput::Key(key.clone())]),
1190            eq(ClientInputSend::Accepted)
1191        );
1192
1193        assert_that!(request_receiver.blocking_recv(), eq(Some(ClientRequest::Key(key))));
1194    }
1195
1196    #[test]
1197    fn test_send_decoded_input_when_copy_precedes_key_keeps_copy_before_selection_reset() {
1198        let (cmd_sender, mut cmd_receiver) = tokio::sync::mpsc::channel(2);
1199        let (request_sender, mut request_receiver) = tokio::sync::mpsc::channel(1);
1200        let key = ClientKey {
1201            code: ClientKeyCode::Char('E'),
1202            modifiers: ClientKeyModifiers::SHIFT_ALT,
1203            raw_bytes: b"\x1bE".to_vec(),
1204        };
1205
1206        assert_that!(
1207            send_decoded_input(
1208                &cmd_sender,
1209                &request_sender,
1210                vec![
1211                    DecodedInput::Key(ClientKey {
1212                        code: ClientKeyCode::Char('C'),
1213                        modifiers: ClientKeyModifiers::SHIFT_ALT,
1214                        raw_bytes: b"\x1bC".to_vec(),
1215                    }),
1216                    DecodedInput::Key(key.clone()),
1217                ],
1218            ),
1219            eq(ClientInputSend::Accepted)
1220        );
1221
1222        assert_that!(
1223            matches!(
1224                cmd_receiver.blocking_recv(),
1225                Some(ClientInputCmd::Action(ClientInputAction::CopySelection))
1226            ),
1227            eq(true)
1228        );
1229        assert_that!(
1230            matches!(
1231                cmd_receiver.blocking_recv(),
1232                Some(ClientInputCmd::Action(ClientInputAction::ClearSelection))
1233            ),
1234            eq(true)
1235        );
1236        assert_that!(request_receiver.blocking_recv(), eq(Some(ClientRequest::Key(key))));
1237    }
1238
1239    #[test]
1240    fn test_send_decoded_input_when_scrollback_editor_shortcut_arrives_sends_key_request() {
1241        let (cmd_sender, _cmd_receiver) = tokio::sync::mpsc::channel(1);
1242        let (request_sender, mut request_receiver) = tokio::sync::mpsc::channel(1);
1243        let key = ClientKey {
1244            code: ClientKeyCode::Char('S'),
1245            modifiers: ClientKeyModifiers::SHIFT_ALT,
1246            raw_bytes: b"\x1bS".to_vec(),
1247        };
1248
1249        assert_that!(
1250            send_decoded_input(&cmd_sender, &request_sender, vec![DecodedInput::Key(key.clone())]),
1251            eq(ClientInputSend::Accepted)
1252        );
1253
1254        assert_that!(request_receiver.blocking_recv(), eq(Some(ClientRequest::Key(key))));
1255    }
1256
1257    #[test]
1258    fn test_send_decoded_input_when_paste_arrives_uses_input_queue() {
1259        let (cmd_sender, _cmd_receiver) = tokio::sync::mpsc::channel(1);
1260        let (request_sender, mut request_receiver) = tokio::sync::mpsc::channel(1);
1261
1262        assert_that!(
1263            send_decoded_input(
1264                &cmd_sender,
1265                &request_sender,
1266                vec![DecodedInput::Paste(b"one\ntwo\n".to_vec())],
1267            ),
1268            eq(ClientInputSend::Accepted)
1269        );
1270
1271        assert_that!(
1272            request_receiver.blocking_recv(),
1273            eq(Some(ClientRequest::Paste(b"one\ntwo\n".to_vec())))
1274        );
1275    }
1276
1277    #[test]
1278    fn test_send_decoded_input_when_mouse_arrives_emits_local_mouse_action() {
1279        let (cmd_sender, mut cmd_receiver) = tokio::sync::mpsc::channel(1);
1280        let (request_sender, _request_receiver) = tokio::sync::mpsc::channel(1);
1281        let event = ClientMouseEvent {
1282            button: 0,
1283            phase: ClientMouseEventPhase::Press,
1284            position: muxr_core::ClientMousePosition { row: 4, col: 9 },
1285        };
1286
1287        assert_that!(
1288            send_decoded_input(&cmd_sender, &request_sender, vec![DecodedInput::Mouse(event)]),
1289            eq(ClientInputSend::Accepted)
1290        );
1291
1292        assert_that!(
1293            matches!(
1294                cmd_receiver.blocking_recv(),
1295                Some(ClientInputCmd::Action(ClientInputAction::Mouse(actual))) if actual == event
1296            ),
1297            eq(true)
1298        );
1299    }
1300
1301    #[test]
1302    fn test_send_decoded_input_when_mouse_motion_action_queue_is_full_drops_without_blocking() -> rootcause::Result<()>
1303    {
1304        let (cmd_sender, mut cmd_receiver) = tokio::sync::mpsc::channel(1);
1305        let (request_sender, _request_receiver) = tokio::sync::mpsc::channel(1);
1306        assert_that!(
1307            cmd_sender.try_send(ClientInputCmd::Action(ClientInputAction::CopySelection)),
1308            ok(eq(()))
1309        );
1310        let event = ClientMouseEvent {
1311            button: 32,
1312            phase: ClientMouseEventPhase::Press,
1313            position: muxr_core::ClientMousePosition { row: 4, col: 9 },
1314        };
1315        let (result_sender, result_receiver) = std::sync::mpsc::channel();
1316        let handle = thread::spawn(move || {
1317            let _ = result_sender.send(send_decoded_input(
1318                &cmd_sender,
1319                &request_sender,
1320                vec![DecodedInput::Mouse(event)],
1321            ));
1322        });
1323        let result = match result_receiver.recv_timeout(Duration::from_secs(1)) {
1324            Ok(result) => result,
1325            Err(error) => {
1326                drop(cmd_receiver);
1327                handle
1328                    .join()
1329                    .map_err(|error| report!("muxr mouse input test thread panicked").attach(format!("{error:?}")))?;
1330                return Err(report!("muxr mouse motion blocked on full input-action queue").attach(format!("{error}")));
1331            }
1332        };
1333
1334        assert_that!(result, eq(ClientInputSend::Accepted));
1335        assert_that!(
1336            matches!(
1337                cmd_receiver.try_recv(),
1338                Ok(ClientInputCmd::Action(ClientInputAction::CopySelection))
1339            ),
1340            eq(true)
1341        );
1342        assert_that!(cmd_receiver.try_recv(), err(anything()));
1343        handle
1344            .join()
1345            .map_err(|error| report!("muxr mouse input test thread panicked").attach(format!("{error:?}")))?;
1346        Ok(())
1347    }
1348
1349    #[test]
1350    fn test_send_decoded_input_when_mouse_wheel_action_queue_is_full_waits_for_queue_space() -> rootcause::Result<()> {
1351        let (cmd_sender, mut cmd_receiver) = tokio::sync::mpsc::channel(1);
1352        let (request_sender, _request_receiver) = tokio::sync::mpsc::channel(1);
1353        assert_that!(
1354            cmd_sender.try_send(ClientInputCmd::Action(ClientInputAction::CopySelection)),
1355            ok(eq(()))
1356        );
1357        let event = ClientMouseEvent {
1358            button: 64,
1359            phase: ClientMouseEventPhase::Press,
1360            position: muxr_core::ClientMousePosition { row: 4, col: 9 },
1361        };
1362        let (result_sender, result_receiver) = std::sync::mpsc::channel();
1363        let handle = thread::spawn(move || {
1364            let _ = result_sender.send(send_decoded_input(
1365                &cmd_sender,
1366                &request_sender,
1367                vec![DecodedInput::Mouse(event)],
1368            ));
1369        });
1370
1371        assert_that!(result_receiver.recv_timeout(Duration::from_millis(50)), err(anything()));
1372        assert_that!(
1373            matches!(
1374                cmd_receiver.blocking_recv(),
1375                Some(ClientInputCmd::Action(ClientInputAction::CopySelection))
1376            ),
1377            eq(true)
1378        );
1379        assert_that!(
1380            result_receiver.recv_timeout(Duration::from_secs(1)),
1381            eq(Ok(ClientInputSend::Accepted))
1382        );
1383        assert_that!(
1384            matches!(
1385                cmd_receiver.blocking_recv(),
1386                Some(ClientInputCmd::Action(ClientInputAction::Mouse(actual))) if actual == event
1387            ),
1388            eq(true)
1389        );
1390        handle
1391            .join()
1392            .map_err(|error| report!("muxr mouse input test thread panicked").attach(format!("{error:?}")))?;
1393        Ok(())
1394    }
1395
1396    #[test]
1397    fn test_send_decoded_input_when_copy_selection_key_arrives_emits_local_action() {
1398        let (cmd_sender, mut cmd_receiver) = tokio::sync::mpsc::channel(1);
1399        let (request_sender, _request_receiver) = tokio::sync::mpsc::channel(1);
1400        let key = ClientKey {
1401            code: ClientKeyCode::Char('C'),
1402            modifiers: ClientKeyModifiers::SHIFT_ALT,
1403            raw_bytes: b"\x1bC".to_vec(),
1404        };
1405
1406        assert_that!(
1407            send_decoded_input(&cmd_sender, &request_sender, vec![DecodedInput::Key(key)]),
1408            eq(ClientInputSend::Accepted)
1409        );
1410
1411        assert_that!(
1412            matches!(
1413                cmd_receiver.blocking_recv(),
1414                Some(ClientInputCmd::Action(ClientInputAction::CopySelection))
1415            ),
1416            eq(true)
1417        );
1418    }
1419
1420    #[test]
1421    fn test_send_decoded_input_when_inline_copy_selection_key_arrives_emits_local_action() {
1422        let (cmd_sender, mut cmd_receiver) = tokio::sync::mpsc::channel(1);
1423        let (request_sender, _request_receiver) = tokio::sync::mpsc::channel(1);
1424        let key = ClientKey {
1425            code: ClientKeyCode::Char('X'),
1426            modifiers: ClientKeyModifiers::SHIFT_ALT,
1427            raw_bytes: b"\x1bX".to_vec(),
1428        };
1429
1430        assert_that!(
1431            send_decoded_input(&cmd_sender, &request_sender, vec![DecodedInput::Key(key)]),
1432            eq(ClientInputSend::Accepted)
1433        );
1434
1435        assert_that!(
1436            matches!(
1437                cmd_receiver.blocking_recv(),
1438                Some(ClientInputCmd::Action(ClientInputAction::CopySelectionInline))
1439            ),
1440            eq(true)
1441        );
1442    }
1443
1444    #[test]
1445    fn test_send_selection_edge_scroll_request_when_scroll_is_pending_waits_for_render_ack() -> rootcause::Result<()> {
1446        let (input_sender, mut input_receiver) = tokio::sync::mpsc::channel(2);
1447        let mut renderer = ClientRenderer::with_synchronized_output(
1448            layout_snapshot()?,
1449            pane_regions_snapshot()?,
1450            SynchronizedOutput::Csi,
1451        );
1452        renderer.apply_render_logical(muxr_core::RenderUpdate::Baseline(render_baseline()?))?;
1453        renderer.apply_selection_input_logical(SelectionInput::Start(ClientMousePosition { row: 0, col: 0 }))?;
1454        let initial = renderer
1455            .set_selection_edge_drag(ClientMousePosition { row: 2, col: 1 }, None)
1456            .ok_or_else(|| report!("expected initial muxr edge scroll request"))?;
1457        let expected = ClientRequest::ScrollPaneLineAt {
1458            direction: PaneScrollDirection::Down,
1459            position: ClientMousePosition { row: 0, col: 1 },
1460        };
1461        assert_that!(
1462            copy_selection_test_helpers::edge_scroll_request(&initial),
1463            eq(&expected)
1464        );
1465        assert_that!(
1466            send_edge_scroll_request(&input_sender, &mut renderer, initial),
1467            eq(ClientInputSend::Accepted)
1468        );
1469        assert_that!(input_receiver.blocking_recv(), eq(Some(expected.clone())));
1470        assert_that!(
1471            send_selection_edge_scroll_request(&input_sender, &mut renderer),
1472            eq(ClientInputSend::Accepted)
1473        );
1474        assert_that!(
1475            input_receiver.try_recv(),
1476            err(matches_pattern!(tokio::sync::mpsc::error::TryRecvError::Empty))
1477        );
1478
1479        renderer.apply_pane_regions_logical(pane_regions_snapshot_with_visible_top_row(1)?);
1480        renderer.apply_render_logical(muxr_core::RenderUpdate::Baseline(render_baseline()?))?;
1481        let flushed = renderer.presentation_snapshot();
1482        renderer.acknowledge_presentation(&flushed);
1483        assert_that!(
1484            send_selection_edge_scroll_request(&input_sender, &mut renderer),
1485            eq(ClientInputSend::Accepted)
1486        );
1487
1488        assert_that!(input_receiver.try_recv(), eq(Ok(expected)));
1489        Ok(())
1490    }
1491
1492    #[test]
1493    fn test_send_edge_scroll_request_when_queue_is_full_does_not_mark_scroll_pending() -> rootcause::Result<()> {
1494        let (input_sender, mut input_receiver) = tokio::sync::mpsc::channel(1);
1495        assert_that!(input_sender.try_send(ClientRequest::Pong), ok(eq(())));
1496        let mut renderer = ClientRenderer::with_synchronized_output(
1497            layout_snapshot()?,
1498            pane_regions_snapshot()?,
1499            SynchronizedOutput::Csi,
1500        );
1501        renderer.apply_render_logical(muxr_core::RenderUpdate::Baseline(render_baseline()?))?;
1502        renderer.apply_selection_input_logical(SelectionInput::Start(ClientMousePosition { row: 0, col: 0 }))?;
1503        let request = renderer
1504            .set_selection_edge_drag(ClientMousePosition { row: 2, col: 1 }, None)
1505            .ok_or_else(|| report!("expected muxr edge scroll request"))?;
1506
1507        assert_that!(
1508            send_edge_scroll_request(&input_sender, &mut renderer, request),
1509            eq(ClientInputSend::Accepted)
1510        );
1511        assert_that!(input_receiver.try_recv(), eq(Ok(ClientRequest::Pong)));
1512        assert_that!(
1513            send_selection_edge_scroll_request(&input_sender, &mut renderer),
1514            eq(ClientInputSend::Accepted)
1515        );
1516
1517        assert_that!(
1518            input_receiver.blocking_recv(),
1519            eq(Some(ClientRequest::ScrollPaneLineAt {
1520                direction: PaneScrollDirection::Down,
1521                position: ClientMousePosition { row: 0, col: 1 },
1522            }))
1523        );
1524        Ok(())
1525    }
1526
1527    fn session_paths(base: &Path, raw: &str) -> rootcause::Result<(SessionName, SessionPaths)> {
1528        let session = raw.parse()?;
1529        let root = base.join("sessions").join(raw);
1530
1531        Ok((
1532            session,
1533            SessionPaths {
1534                socket: root.join("server.sock"),
1535                pid: root.join("server.pid"),
1536                layout: root.join("layout.json"),
1537                panes: root.join("panes"),
1538                root,
1539            },
1540        ))
1541    }
1542
1543    fn layout_snapshot() -> rootcause::Result<LayoutSnapshot> {
1544        let active_tab = TabId::new(1)?;
1545        let active_pane = PaneId::new(1)?;
1546        let pane = PaneSnapshot {
1547            tracked_process_state: muxr_core::TrackedProcessState::None,
1548            cwd: "/tmp".to_owned(),
1549            cmd_label: None,
1550            focus_seq: 1,
1551            id: active_pane,
1552            title: "shell".to_owned(),
1553        };
1554        let tab = TabSnapshot::new(active_tab, "default", active_pane, vec![pane])?;
1555        LayoutSnapshot::new(active_tab, vec![tab])
1556    }
1557
1558    fn pane_regions_snapshot() -> rootcause::Result<PaneRegionsSnapshot> {
1559        self::pane_regions_snapshot_with_visible_top_row(0)
1560    }
1561
1562    fn pane_regions_snapshot_with_visible_top_row(visible_top_row: u64) -> rootcause::Result<PaneRegionsSnapshot> {
1563        PaneRegionsSnapshot::new(vec![muxr_core::PaneRegionSnapshot::new(
1564            muxr_core::PaneId::new(1)?,
1565            0,
1566            0,
1567            2,
1568            1,
1569            muxr_core::PaneMouseMode::None,
1570            visible_top_row,
1571        )?])
1572    }
1573
1574    fn render_baseline() -> rootcause::Result<muxr_core::RenderBaseline> {
1575        muxr_core::RenderBaseline::new(
1576            1,
1577            TerminalSize::new(2, 1)?,
1578            muxr_core::RenderCursor {
1579                row: 0,
1580                col: 1,
1581                shape: muxr_core::RenderCursorShape::Default,
1582                visibility: muxr_core::RenderCursorVisibility::Visible,
1583            },
1584            vec![muxr_core::RenderRowSpan::new(
1585                0,
1586                0,
1587                vec![render_cell("a"), render_cell("b")],
1588            )?],
1589        )
1590    }
1591
1592    fn render_cell(text: &str) -> muxr_core::RenderCell {
1593        muxr_core::RenderCell::narrow(text, muxr_core::RenderStyle::default())
1594    }
1595
1596    fn runtime() -> rootcause::Result<tokio::runtime::Runtime> {
1597        Ok(tokio::runtime::Runtime::new().context("failed to build muxr client test runtime")?)
1598    }
1599}