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
133pub 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 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 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 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 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}