1use std::io::Read;
2use std::path::Path;
3use std::thread;
4use std::time::Duration;
5
6use muxr_config::MuxrConfig;
7use muxr_core::ClientMouseEvent;
8use muxr_core::ClientRequest;
9use muxr_core::ServerEvent;
10use muxr_core::SessionName;
11use muxr_core::TerminalSize;
12use muxr_transport::ClientRequestWriter;
13use rootcause::prelude::ResultExt;
14use rootcause::report;
15
16use crate::copy_selection::SelectionEdgeScrollRequest;
17use crate::input::DecodedInput;
18use crate::input::InputDecoder;
19use crate::input::InputIdleTimeout;
20use crate::renderer::ClientPresentationSnapshot;
21use crate::renderer::ClientRenderOutcome;
22use crate::renderer::ClientRenderer;
23use crate::session::attach::AttachedSession;
24use crate::stdout_worker::StdoutSender;
25use crate::stdout_worker::StdoutWorker;
26use crate::terminal::TerminalGuard;
27
28const RESIZE_POLL_INTERVAL: Duration = Duration::from_millis(250);
29const AMBIGUOUS_INPUT_TIMEOUT: Duration = Duration::from_millis(50);
30const SELECTION_EDGE_SCROLL_INTERVAL: Duration = Duration::from_millis(50);
31const STDIN_BUFFER_SIZE: usize = 8192;
32const CONTROL_REQUEST_CHANNEL_LIMIT: usize = 128;
33const INPUT_REQUEST_CHANNEL_LIMIT: usize = 1024;
34
35#[derive(Clone, Debug, Eq, PartialEq)]
36enum StdinRead {
37 Bytes(Vec<u8>),
38 Eof,
39}
40
41#[derive(Clone, Debug, Eq, PartialEq)]
42enum ClientInputAction {
43 CopySelection,
44 CopySelectionInline,
45 Mouse(ClientMouseEvent),
46}
47
48#[derive(Debug)]
49enum ClientInputCmd {
50 Action(ClientInputAction),
51 Barrier(std::sync::mpsc::SyncSender<()>),
52}
53
54#[derive(Clone, Copy, Debug, Eq, PartialEq)]
55enum InputCmdReceiverState {
56 Open,
57 Closed,
58}
59
60impl InputCmdReceiverState {
61 const fn is_open(self) -> bool {
62 matches!(self, Self::Open)
63 }
64}
65
66#[derive(Clone, Copy, Debug, Eq, PartialEq)]
67enum LocalActionCompletion {
68 Wait,
69 #[cfg(test)]
70 Skip,
71}
72
73impl LocalActionCompletion {
74 const fn waits(self) -> bool {
75 match self {
76 Self::Wait => true,
77 #[cfg(test)]
78 Self::Skip => false,
79 }
80 }
81}
82
83#[derive(Clone, Copy, Debug, Eq, PartialEq)]
84pub enum ClientInputSend {
85 Accepted,
86 Closed,
87}
88
89#[derive(Clone, Copy, Debug, Eq, PartialEq)]
90enum InteractiveFlow {
91 Continue,
92 Stop,
93}
94
95struct RenderCoordinator {
96 committed: Option<ClientPresentationSnapshot>,
97 in_flight: Option<ClientPresentationSnapshot>,
98}
99
100trait RenderSink {
101 fn send_render(&self, transaction: Vec<u8>) -> rootcause::Result<()>;
102}
103
104impl RenderSink for StdoutSender {
105 fn send_render(&self, transaction: Vec<u8>) -> rootcause::Result<()> {
106 Self::send_render(self, transaction)
107 }
108}
109
110impl RenderCoordinator {
111 const fn new() -> Self {
112 Self {
113 committed: None,
114 in_flight: None,
115 }
116 }
117
118 fn complete(&mut self, renderer: &mut ClientRenderer, output: &impl RenderSink) -> rootcause::Result<()> {
119 let completed = self
120 .in_flight
121 .take()
122 .ok_or_else(|| report!("muxr client stdout completed an unknown render"))?;
123 renderer.acknowledge_presentation(&completed);
124 self.committed = Some(completed);
125 self.submit(renderer, output)
126 }
127
128 fn submit(&mut self, renderer: &mut ClientRenderer, output: &impl RenderSink) -> rootcause::Result<()> {
129 if self.in_flight.is_some() {
130 return Ok(());
131 }
132 let snapshot = renderer.presentation_snapshot();
133 let Some(transaction) = renderer.presentation_transaction(self.committed.as_ref())? else {
134 return Ok(());
135 };
136 output.send_render(transaction)?;
137 self.in_flight = Some(snapshot);
138 Ok(())
139 }
140}
141
142pub fn start(session: &SessionName, server_executable: &Path, external_layout: Option<&Path>) -> rootcause::Result<()> {
150 tokio::runtime::Runtime::new()
151 .context("failed to build muxr tokio runtime")?
152 .block_on(async {
153 let muxr_config = MuxrConfig::default();
154 let terminal_size = crate::terminal::current_terminal_size()?;
155 let pane_size = crate::terminal::pane_size_for_terminal(muxr_config.tab_bar.width, &terminal_size)?;
156 let attached_session =
157 crate::session::attach::open_session(session, pane_size.clone(), server_executable, external_layout)
158 .await?;
159 self::run_interactive(&muxr_config, attached_session, pane_size).await
160 })
161}
162
163async fn run_interactive(
164 muxr_config: &MuxrConfig,
165 mut attached_session: AttachedSession,
166 initial_size: TerminalSize,
167) -> rootcause::Result<()> {
168 let _terminal_guard = TerminalGuard::enable_if_terminal()?;
169 let (control_sender, control_receiver) = tokio::sync::mpsc::channel(CONTROL_REQUEST_CHANNEL_LIMIT);
170 let (input_cmd_sender, mut input_cmd_receiver) = tokio::sync::mpsc::channel(INPUT_REQUEST_CHANNEL_LIMIT);
171 let (input_request_sender, input_receiver) = tokio::sync::mpsc::channel(INPUT_REQUEST_CHANNEL_LIMIT);
172 let stdin_handle = self::spawn_stdin_forwarder(input_cmd_sender, input_request_sender.clone());
173 let resize_handle = self::spawn_resize_forwarder(control_sender.clone(), muxr_config.tab_bar.width, initial_size);
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 let mut renderer = ClientRenderer::new(muxr_config, attached_session.layout, attached_session.pane_regions);
180 renderer.sync_mouse_capture_logical();
181 let mut render_coordinator = RenderCoordinator::new();
182 let edge_scroll_tick_start = tokio::time::Instant::now()
183 .checked_add(SELECTION_EDGE_SCROLL_INTERVAL)
184 .ok_or_else(|| report!("muxr selection edge scroll interval overflowed"))?;
185 let mut edge_scroll_tick = tokio::time::interval_at(edge_scroll_tick_start, SELECTION_EDGE_SCROLL_INTERVAL);
186 edge_scroll_tick.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
187 let mut input_cmd_receiver_state = InputCmdReceiverState::Open;
188
189 loop {
190 tokio::select! {
191 event = attached_session.reader.recv_event() => {
192 let Some(event) = event? else {
193 break;
194 };
195 if self::handle_server_event(
196 event,
197 &control_sender,
198 &mut renderer,
199 &mut render_coordinator,
200 &stdout_sender,
201 ).await?
202 == InteractiveFlow::Stop
203 {
204 break;
205 }
206 },
207 cmd = input_cmd_receiver.recv(), if input_cmd_receiver_state.is_open() => {
208 let Some(cmd) = cmd else {
209 input_cmd_receiver_state = InputCmdReceiverState::Closed;
210 continue;
211 };
212 let action = match cmd {
213 ClientInputCmd::Action(action) => action,
214 ClientInputCmd::Barrier(completed) => {
215 let _sent = completed.send(());
216 continue;
217 }
218 };
219 if self::handle_client_input_action(action, muxr_config, &input_request_sender, &mut renderer).await? == ClientInputSend::Closed {
220 break;
221 }
222 render_coordinator.submit(&mut renderer, &stdout_sender)?;
223 },
224 _ = edge_scroll_tick.tick(), if renderer.selection_edge_drag() == crate::renderer::SelectionEdgeDrag::Active => {
225 if self::send_selection_edge_scroll_request(&input_request_sender, &mut renderer) == ClientInputSend::Closed {
226 break;
227 }
228 },
229 stdout_failure = &mut stdout_failure_receiver => {
230 let error = stdout_failure.unwrap_or_else(|_| "stdout worker stopped unexpectedly".to_owned());
231 return Err(report!("muxr client stdout worker failed").attach(error));
232 },
233 completed = stdout_completion_receiver.recv() => {
234 if completed.is_none() {
235 return Err(report!("muxr client stdout completion worker stopped unexpectedly"));
236 }
237 render_coordinator.complete(&mut renderer, &stdout_sender)?;
238 },
239 else => {
240 if input_cmd_receiver_state == InputCmdReceiverState::Closed {
241 break;
242 }
243 }
244 }
245 }
246
247 writer_handle.abort();
248 drop(writer_handle.await);
249 drop(stdin_handle);
250 drop(resize_handle);
251 Ok(())
252}
253
254async fn handle_server_event(
255 event: ServerEvent,
256 control_sender: &tokio::sync::mpsc::Sender<ClientRequest>,
257 renderer: &mut ClientRenderer,
258 render_coordinator: &mut RenderCoordinator,
259 stdout_sender: &StdoutSender,
260) -> rootcause::Result<InteractiveFlow> {
261 match event {
262 ServerEvent::Deleted | ServerEvent::Detached => Ok(InteractiveFlow::Stop),
263 ServerEvent::Error(error) => Err(report!("muxr server returned error")
264 .attach(format!("code={}", error.code()))
265 .attach(format!("msg={}", error.msg()))),
266 ServerEvent::Ping => Ok(if control_sender.send(ClientRequest::Pong).await.is_ok() {
267 InteractiveFlow::Continue
268 } else {
269 InteractiveFlow::Stop
270 }),
271 ServerEvent::Layout(next_layout) => {
272 renderer.apply_layout(next_layout);
273 render_coordinator.submit(renderer, stdout_sender)?;
274 Ok(InteractiveFlow::Continue)
275 }
276 ServerEvent::SidebarLayout(next_layout) => {
277 renderer.apply_sidebar_layout_logical(next_layout);
278 render_coordinator.submit(renderer, stdout_sender)?;
279 Ok(InteractiveFlow::Continue)
280 }
281 ServerEvent::PaneRegions(next_regions) => {
282 renderer.apply_pane_regions_logical(next_regions);
283 render_coordinator.submit(renderer, stdout_sender)?;
284 Ok(InteractiveFlow::Continue)
285 }
286 ServerEvent::Render(update) => {
287 self::handle_render_event(update, control_sender, renderer, render_coordinator, stdout_sender).await
288 }
289 ServerEvent::ScrollPaneLineResult {
290 position,
291 direction,
292 movement,
293 } => {
294 renderer.apply_scroll_pane_line_result(position, direction, movement);
295 Ok(InteractiveFlow::Continue)
296 }
297 ServerEvent::Attached(_) | ServerEvent::Pong => Ok(InteractiveFlow::Continue),
298 }
299}
300
301async fn handle_render_event(
302 update: muxr_core::RenderUpdate,
303 control_sender: &tokio::sync::mpsc::Sender<ClientRequest>,
304 renderer: &mut ClientRenderer,
305 render_coordinator: &mut RenderCoordinator,
306 stdout_sender: &StdoutSender,
307) -> rootcause::Result<InteractiveFlow> {
308 match renderer.apply_render_logical(update)? {
309 ClientRenderOutcome::Drawn => {
310 render_coordinator.submit(renderer, stdout_sender)?;
311 Ok(InteractiveFlow::Continue)
312 }
313 ClientRenderOutcome::NeedsResync => Ok(if control_sender.send(ClientRequest::RenderResync).await.is_ok() {
314 InteractiveFlow::Continue
315 } else {
316 InteractiveFlow::Stop
317 }),
318 }
319}
320
321async fn handle_client_input_action(
322 action: ClientInputAction,
323 muxr_config: &MuxrConfig,
324 input_sender: &tokio::sync::mpsc::Sender<ClientRequest>,
325 renderer: &mut ClientRenderer,
326) -> rootcause::Result<ClientInputSend> {
327 match action {
328 ClientInputAction::CopySelection => {
329 renderer.copy_selection()?;
330 Ok(ClientInputSend::Accepted)
331 }
332 ClientInputAction::CopySelectionInline => {
333 renderer.copy_selection_inline()?;
334 Ok(ClientInputSend::Accepted)
335 }
336 ClientInputAction::Mouse(event) => {
337 crate::pane::mouse::handle_mouse_input_action(muxr_config, event, input_sender, renderer).await
338 }
339 }
340}
341
342#[derive(Clone, Copy, Debug, Eq, PartialEq)]
343pub enum DroppableSendOutcome {
344 Closed,
345 Dropped,
346 Sent,
347}
348
349pub fn send_droppable_request(
350 input_sender: &tokio::sync::mpsc::Sender<ClientRequest>,
351 request: ClientRequest,
352) -> DroppableSendOutcome {
353 match input_sender.try_send(request) {
354 Ok(()) => DroppableSendOutcome::Sent,
355 Err(tokio::sync::mpsc::error::TrySendError::Full(request)) => {
356 drop(request);
357 DroppableSendOutcome::Dropped
358 }
359 Err(tokio::sync::mpsc::error::TrySendError::Closed(request)) => {
360 drop(request);
361 DroppableSendOutcome::Closed
362 }
363 }
364}
365
366pub fn send_edge_scroll_request(
367 input_sender: &tokio::sync::mpsc::Sender<ClientRequest>,
368 renderer: &mut ClientRenderer,
369 request: SelectionEdgeScrollRequest,
370) -> ClientInputSend {
371 let (pending, request) = request.into_parts();
372 match self::send_droppable_request(input_sender, request) {
373 DroppableSendOutcome::Sent => {
374 renderer.mark_selection_edge_scroll_sent(pending);
377 ClientInputSend::Accepted
378 }
379 DroppableSendOutcome::Dropped => ClientInputSend::Accepted,
380 DroppableSendOutcome::Closed => ClientInputSend::Closed,
381 }
382}
383
384fn send_selection_edge_scroll_request(
385 input_sender: &tokio::sync::mpsc::Sender<ClientRequest>,
386 renderer: &mut ClientRenderer,
387) -> ClientInputSend {
388 let Some(request) = renderer.selection_edge_scroll_request() else {
389 return ClientInputSend::Accepted;
390 };
391 self::send_edge_scroll_request(input_sender, renderer, request)
392}
393
394async fn forward_client_requests(
395 mut writer: ClientRequestWriter,
396 mut control_receiver: tokio::sync::mpsc::Receiver<ClientRequest>,
397 mut input_receiver: tokio::sync::mpsc::Receiver<ClientRequest>,
398) -> rootcause::Result<()> {
399 let mut control_closed = false;
400 let mut input_closed = false;
401
402 loop {
403 if control_closed && input_closed {
404 break;
405 }
406
407 tokio::select! {
408 biased;
409 request = control_receiver.recv(), if !control_closed => match request {
410 Some(request) => {
411 if writer.send_request(&request).await.is_err() {
412 break;
413 }
414 }
415 None => control_closed = true,
416 },
417 request = input_receiver.recv(), if !input_closed => match request {
418 Some(request) => {
419 if writer.send_request(&request).await.is_err() {
420 break;
421 }
422 }
423 None => input_closed = true,
424 },
425 }
426 }
427
428 Ok(())
429}
430
431fn spawn_stdin_forwarder(
432 cmd_sender: tokio::sync::mpsc::Sender<ClientInputCmd>,
433 request_sender: tokio::sync::mpsc::Sender<ClientRequest>,
434) -> thread::JoinHandle<()> {
435 thread::spawn(move || {
436 let (read_sender, read_receiver) = std::sync::mpsc::channel();
437 drop(self::spawn_stdin_reader(read_sender));
438 let mut decoder = InputDecoder::default();
439
440 loop {
441 let read = if decoder.idle_timeout() == InputIdleTimeout::Needed {
444 match read_receiver.recv_timeout(AMBIGUOUS_INPUT_TIMEOUT) {
445 Ok(read) => read,
446 Err(std::sync::mpsc::RecvTimeoutError::Timeout) => {
447 if self::send_decoded_input_with_ordering(
448 &cmd_sender,
449 &request_sender,
450 decoder.finalize(),
451 LocalActionCompletion::Wait,
452 ) == ClientInputSend::Closed
453 {
454 break;
455 }
456 continue;
457 }
458 Err(std::sync::mpsc::RecvTimeoutError::Disconnected) => StdinRead::Eof,
459 }
460 } else {
461 read_receiver.recv().unwrap_or(StdinRead::Eof)
462 };
463
464 match read {
465 StdinRead::Bytes(bytes) => {
466 if self::send_decoded_input_with_ordering(
467 &cmd_sender,
468 &request_sender,
469 decoder.decode(&bytes),
470 LocalActionCompletion::Wait,
471 ) == ClientInputSend::Closed
472 {
473 break;
474 }
475 }
476 StdinRead::Eof => {
477 if self::send_decoded_input_with_ordering(
478 &cmd_sender,
479 &request_sender,
480 decoder.finalize(),
481 LocalActionCompletion::Wait,
482 ) == ClientInputSend::Closed
483 {
484 break;
485 }
486 drop(request_sender.blocking_send(ClientRequest::Detach));
488 break;
489 }
490 }
491 }
492 })
493}
494
495fn spawn_stdin_reader(sender: std::sync::mpsc::Sender<StdinRead>) -> thread::JoinHandle<()> {
496 thread::spawn(move || {
497 let mut stdin = std::io::stdin();
498 let mut buffer = [0; STDIN_BUFFER_SIZE];
499
500 loop {
501 match stdin.read(&mut buffer) {
502 Ok(0) | Err(_) => {
503 drop(sender.send(StdinRead::Eof));
504 break;
505 }
506 Ok(bytes_read) => {
507 let Some(bytes) = buffer.get(..bytes_read) else {
508 drop(sender.send(StdinRead::Eof));
509 break;
510 };
511 if sender.send(StdinRead::Bytes(bytes.to_vec())).is_err() {
512 break;
513 }
514 }
515 }
516 }
517 })
518}
519
520#[cfg(test)]
521fn send_decoded_input(
522 cmd_sender: &tokio::sync::mpsc::Sender<ClientInputCmd>,
523 request_sender: &tokio::sync::mpsc::Sender<ClientRequest>,
524 decoded: Vec<DecodedInput>,
525) -> ClientInputSend {
526 self::send_decoded_input_with_ordering(cmd_sender, request_sender, decoded, LocalActionCompletion::Skip)
527}
528
529fn send_decoded_input_with_ordering(
530 cmd_sender: &tokio::sync::mpsc::Sender<ClientInputCmd>,
531 request_sender: &tokio::sync::mpsc::Sender<ClientRequest>,
532 decoded: Vec<DecodedInput>,
533 local_action_completion: LocalActionCompletion,
534) -> ClientInputSend {
535 for decoded in decoded {
536 let action = match decoded {
537 DecodedInput::CopySelection => ClientInputAction::CopySelection,
538 DecodedInput::CopySelectionInline => ClientInputAction::CopySelectionInline,
539 DecodedInput::Input(bytes) => {
542 if request_sender.blocking_send(ClientRequest::Input(bytes)).is_err() {
543 return ClientInputSend::Closed;
544 }
545 continue;
546 }
547 DecodedInput::Key(key) => {
548 if request_sender.blocking_send(ClientRequest::Key(key)).is_err() {
549 return ClientInputSend::Closed;
550 }
551 continue;
552 }
553 DecodedInput::Mouse(event)
554 if crate::pane::mouse::MouseEventDrop::from(event) == crate::pane::mouse::MouseEventDrop::Droppable =>
555 {
556 if self::send_droppable_input_action(
557 cmd_sender,
558 ClientInputAction::Mouse(event),
559 local_action_completion,
560 ) == ClientInputSend::Closed
561 {
562 return ClientInputSend::Closed;
563 }
564 continue;
565 }
566 DecodedInput::Mouse(event) => ClientInputAction::Mouse(event),
567 DecodedInput::Paste(bytes) => {
568 if request_sender.blocking_send(ClientRequest::Paste(bytes)).is_err() {
569 return ClientInputSend::Closed;
570 }
571 continue;
572 }
573 };
574 if self::send_input_action(cmd_sender, action, local_action_completion) == ClientInputSend::Closed {
575 return ClientInputSend::Closed;
576 }
577 }
578
579 ClientInputSend::Accepted
580}
581
582fn send_droppable_input_action(
583 cmd_sender: &tokio::sync::mpsc::Sender<ClientInputCmd>,
584 action: ClientInputAction,
585 local_action_completion: LocalActionCompletion,
586) -> ClientInputSend {
587 match cmd_sender.try_send(ClientInputCmd::Action(action)) {
588 Ok(()) if local_action_completion.waits() => self::send_input_action_barrier(cmd_sender),
589 Ok(()) => ClientInputSend::Accepted,
590 Err(tokio::sync::mpsc::error::TrySendError::Full(action)) => {
591 drop(action);
592 ClientInputSend::Accepted
593 }
594 Err(tokio::sync::mpsc::error::TrySendError::Closed(action)) => {
595 drop(action);
596 ClientInputSend::Closed
597 }
598 }
599}
600
601fn send_input_action(
602 cmd_sender: &tokio::sync::mpsc::Sender<ClientInputCmd>,
603 action: ClientInputAction,
604 local_action_completion: LocalActionCompletion,
605) -> ClientInputSend {
606 if cmd_sender.blocking_send(ClientInputCmd::Action(action)).is_err() {
607 return ClientInputSend::Closed;
608 }
609 if local_action_completion.waits() {
610 self::send_input_action_barrier(cmd_sender)
611 } else {
612 ClientInputSend::Accepted
613 }
614}
615
616fn send_input_action_barrier(cmd_sender: &tokio::sync::mpsc::Sender<ClientInputCmd>) -> ClientInputSend {
617 let (completed_sender, completed_receiver) = std::sync::mpsc::sync_channel(0);
618 if cmd_sender
619 .blocking_send(ClientInputCmd::Barrier(completed_sender))
620 .is_err()
621 {
622 return ClientInputSend::Closed;
623 }
624 if completed_receiver.recv().is_err() {
625 return ClientInputSend::Closed;
626 }
627 ClientInputSend::Accepted
628}
629
630fn spawn_resize_forwarder(
631 sender: tokio::sync::mpsc::Sender<ClientRequest>,
632 tab_bar_width: u16,
633 initial_size: TerminalSize,
634) -> thread::JoinHandle<()> {
635 thread::spawn(move || {
636 let mut last_size = initial_size;
637
638 loop {
639 if sender.is_closed() {
640 break;
641 }
642
643 thread::sleep(RESIZE_POLL_INTERVAL);
644 let Ok(next_terminal_size) = crate::terminal::current_terminal_size() else {
645 break;
646 };
647 let Ok(next_size) = crate::terminal::pane_size_for_terminal(tab_bar_width, &next_terminal_size) else {
649 break;
650 };
651 if next_size == last_size {
652 continue;
653 }
654
655 if sender.blocking_send(ClientRequest::Resize(next_size.clone())).is_err() {
656 break;
657 }
658 last_size = next_size;
659 }
660 })
661}
662
663#[cfg(test)]
664mod tests {
665 use std::cell::RefCell;
666 use std::fs;
667 use std::path::Path;
668
669 use muxr_core::ClientKey;
670 use muxr_core::ClientKeyCode;
671 use muxr_core::ClientKeyModifiers;
672 use muxr_core::ClientMouseEventPhase;
673 use muxr_core::ClientMousePosition;
674 use muxr_core::LayoutSnapshot;
675 use muxr_core::PaneId;
676 use muxr_core::PaneRegionsSnapshot;
677 use muxr_core::PaneScrollDirection;
678 use muxr_core::PaneSnapshot;
679 use muxr_core::SessionPaths;
680 use muxr_core::TabId;
681 use muxr_core::TabSnapshot;
682 use muxr_transport::ClientConnection;
683 use muxr_transport::ServerListener;
684 use test_that::prelude::*;
685
686 use super::*;
687 use crate::copy_selection::SelectionInput;
688 use crate::copy_selection::test_helpers as copy_selection_test_helpers;
689 use crate::terminal::SynchronizedOutput;
690
691 #[derive(Default)]
692 struct RecordingRenderSink {
693 transactions: RefCell<Vec<Vec<u8>>>,
694 }
695
696 impl RecordingRenderSink {
697 fn len(&self) -> usize {
698 self.transactions.borrow().len()
699 }
700
701 fn transaction(&self, index: usize) -> rootcause::Result<Vec<u8>> {
702 self.transactions
703 .borrow()
704 .get(index)
705 .cloned()
706 .ok_or_else(|| report!("muxr coordinator test transaction is missing").attach(format!("index={index}")))
707 }
708 }
709
710 impl RenderSink for RecordingRenderSink {
711 fn send_render(&self, transaction: Vec<u8>) -> rootcause::Result<()> {
712 self.transactions.borrow_mut().push(transaction);
713 Ok(())
714 }
715 }
716
717 #[test]
718 fn test_render_coordinator_when_logical_state_changes_in_flight_emits_one_delta_after_completion()
719 -> rootcause::Result<()> {
720 let mut renderer = ClientRenderer::with_synchronized_output(
721 layout_snapshot()?,
722 pane_regions_snapshot()?,
723 SynchronizedOutput::Csi,
724 );
725 let mut coordinator = RenderCoordinator::new();
726 let output = RecordingRenderSink::default();
727
728 renderer.apply_render_logical(muxr_core::RenderUpdate::Baseline(render_baseline()?))?;
729 coordinator.submit(&mut renderer, &output)?;
730 assert_that!(output.len(), eq(1));
731
732 renderer.apply_selection_input_logical(SelectionInput::Start(ClientMousePosition { row: 0, col: 0 }))?;
733 renderer.apply_selection_input_logical(SelectionInput::Update(ClientMousePosition { row: 0, col: 1 }))?;
734 coordinator.submit(&mut renderer, &output)?;
735 assert_that!(output.len(), eq(1));
736
737 coordinator.complete(&mut renderer, &output)?;
738 assert_that!(output.len(), eq(2));
739 let initial = String::from_utf8(output.transaction(0)?)?;
740 let delta = String::from_utf8(output.transaction(1)?)?;
741 assert_that!(initial, contains_substring("/tmp"));
742 assert_that!(delta, not(contains_substring("/tmp")));
743 assert_that!(delta, not(contains_substring("\x1b[2J")));
744
745 coordinator.complete(&mut renderer, &output)?;
746 assert_that!(output.len(), eq(2));
747 Ok(())
748 }
749
750 #[test]
751 fn test_forward_client_requests_when_input_queue_is_ready_sends_control_first() -> rootcause::Result<()> {
752 self::runtime()?.block_on(async {
753 let tempdir = tempfile::tempdir()?;
754 let (_, paths) = self::session_paths(tempdir.path(), "work")?;
755 fs::create_dir_all(&paths.root)?;
756 let listener = ServerListener::bind(&paths.socket)?;
757 let server_handle = tokio::spawn(async move {
758 let mut connection = listener.accept().await?;
759 let Some(request) = connection.recv_request().await? else {
760 return Err(report!("expected forwarded client request"));
761 };
762 Ok::<ClientRequest, rootcause::Report>(request)
763 });
764
765 let connection = ClientConnection::connect(&paths.socket).await?;
766 let (_reader, writer) = connection.split();
767 let (control_sender, control_receiver) = tokio::sync::mpsc::channel(1);
768 let (input_sender, input_receiver) = tokio::sync::mpsc::channel(1);
769 assert_that!(input_sender.try_send(ClientRequest::Input(vec![b'a'])), ok(eq(())));
770 assert_that!(input_sender.try_send(ClientRequest::Input(vec![b'b'])), err(anything()));
771 assert_that!(control_sender.try_send(ClientRequest::Pong), ok(eq(())));
772
773 let writer_handle = tokio::spawn(self::forward_client_requests(writer, control_receiver, input_receiver));
774 let first_request = server_handle
775 .await
776 .map_err(|error| report!("muxr forward test socket task panicked").attach(format!("{error}")))??;
777
778 assert_that!(first_request, eq(ClientRequest::Pong));
779 drop(control_sender);
780 drop(input_sender);
781 writer_handle
782 .await
783 .map_err(|error| report!("muxr forward test writer task panicked").attach(format!("{error}")))??;
784 Ok(())
785 })
786 }
787
788 #[test]
789 fn test_forward_client_requests_when_stdin_requests_are_mixed_sends_input_queue_in_order() -> rootcause::Result<()>
790 {
791 self::runtime()?.block_on(async {
792 let tempdir = tempfile::tempdir()?;
793 let (_, paths) = self::session_paths(tempdir.path(), "work")?;
794 fs::create_dir_all(&paths.root)?;
795 let listener = ServerListener::bind(&paths.socket)?;
796 let server_handle = tokio::spawn(async move {
797 let mut connection = listener.accept().await?;
798 let mut requests = Vec::new();
799 for _ in 0..3 {
800 let Some(request) = connection.recv_request().await? else {
801 return Err(report!("expected forwarded stdin request"));
802 };
803 requests.push(request);
804 }
805 Ok::<Vec<ClientRequest>, rootcause::Report>(requests)
806 });
807
808 let connection = ClientConnection::connect(&paths.socket).await?;
809 let (_reader, writer) = connection.split();
810 let (control_sender, control_receiver) = tokio::sync::mpsc::channel(1);
811 let (input_sender, input_receiver) = tokio::sync::mpsc::channel(3);
812 let key = ClientKey {
813 code: ClientKeyCode::Char('E'),
814 modifiers: ClientKeyModifiers::SHIFT_ALT,
815 raw_bytes: b"\x1bE".to_vec(),
816 };
817 assert_that!(input_sender.try_send(ClientRequest::Input(b"a".to_vec())), ok(eq(())));
818 assert_that!(input_sender.try_send(ClientRequest::Key(key.clone())), ok(eq(())));
819 assert_that!(input_sender.try_send(ClientRequest::Input(b"b".to_vec())), ok(eq(())));
820 drop(control_sender);
821 drop(input_sender);
822
823 let writer_handle = tokio::spawn(self::forward_client_requests(writer, control_receiver, input_receiver));
824 let requests = server_handle.await.map_err(|error| {
825 report!("muxr forward order test socket task panicked").attach(format!("{error}"))
826 })??;
827
828 assert_that!(
829 requests,
830 eq(vec![
831 ClientRequest::Input(b"a".to_vec()),
832 ClientRequest::Key(key),
833 ClientRequest::Input(b"b".to_vec()),
834 ])
835 );
836 writer_handle.await.map_err(|error| {
837 report!("muxr forward order test writer task panicked").attach(format!("{error}"))
838 })??;
839 Ok(())
840 })
841 }
842
843 #[test]
844 fn test_forward_client_requests_when_stdin_detach_follows_input_sends_input_before_detach() -> rootcause::Result<()>
845 {
846 self::runtime()?.block_on(async {
847 let tempdir = tempfile::tempdir()?;
848 let (_, paths) = self::session_paths(tempdir.path(), "work")?;
849 fs::create_dir_all(&paths.root)?;
850 let listener = ServerListener::bind(&paths.socket)?;
851 let server_handle = tokio::spawn(async move {
852 let mut connection = listener.accept().await?;
853 let mut requests = Vec::new();
854 for _ in 0..2 {
855 let Some(request) = connection.recv_request().await? else {
856 return Err(report!("expected forwarded stdin detach request"));
857 };
858 requests.push(request);
859 }
860 Ok::<Vec<ClientRequest>, rootcause::Report>(requests)
861 });
862
863 let connection = ClientConnection::connect(&paths.socket).await?;
864 let (_reader, writer) = connection.split();
865 let (control_sender, control_receiver) = tokio::sync::mpsc::channel(1);
866 let (input_sender, input_receiver) = tokio::sync::mpsc::channel(2);
867 assert_that!(
868 input_sender.try_send(ClientRequest::Input(b"exit\n".to_vec())),
869 ok(eq(()))
870 );
871 assert_that!(input_sender.try_send(ClientRequest::Detach), ok(eq(())));
872 drop(control_sender);
873 drop(input_sender);
874
875 let writer_handle = tokio::spawn(self::forward_client_requests(writer, control_receiver, input_receiver));
876 let requests = server_handle
877 .await
878 .map_err(|error| report!("muxr forward EOF test socket task panicked").attach(format!("{error}")))??;
879
880 assert_that!(
881 requests,
882 eq(vec![ClientRequest::Input(b"exit\n".to_vec()), ClientRequest::Detach])
883 );
884 writer_handle
885 .await
886 .map_err(|error| report!("muxr forward EOF test writer task panicked").attach(format!("{error}")))??;
887 Ok(())
888 })
889 }
890
891 #[test]
892 fn test_send_decoded_input_when_key_arrives_bypasses_the_renderer_action_queue_in_order() {
893 let (cmd_sender, mut cmd_receiver) = tokio::sync::mpsc::channel(1);
894 let (request_sender, mut request_receiver) = tokio::sync::mpsc::channel(3);
895 let key = ClientKey {
896 code: ClientKeyCode::Char('E'),
897 modifiers: ClientKeyModifiers::SHIFT_ALT,
898 raw_bytes: b"\x1bE".to_vec(),
899 };
900
901 assert_that!(
902 send_decoded_input(
903 &cmd_sender,
904 &request_sender,
905 vec![
906 DecodedInput::Input(b"a".to_vec()),
907 DecodedInput::Key(key.clone()),
908 DecodedInput::Input(b"b".to_vec()),
909 ],
910 ),
911 eq(ClientInputSend::Accepted)
912 );
913
914 assert_that!(
915 request_receiver.blocking_recv(),
916 eq(Some(ClientRequest::Input(b"a".to_vec())))
917 );
918 assert_that!(request_receiver.blocking_recv(), eq(Some(ClientRequest::Key(key))));
919 assert_that!(
920 request_receiver.blocking_recv(),
921 eq(Some(ClientRequest::Input(b"b".to_vec())))
922 );
923 assert_that!(cmd_receiver.try_recv().is_err(), eq(true));
924 }
925
926 #[test]
927 fn test_send_decoded_input_when_scrollback_editor_shortcut_arrives_sends_key_request() {
928 let (cmd_sender, _cmd_receiver) = tokio::sync::mpsc::channel(1);
929 let (request_sender, mut request_receiver) = tokio::sync::mpsc::channel(1);
930 let key = ClientKey {
931 code: ClientKeyCode::Char('S'),
932 modifiers: ClientKeyModifiers::SHIFT_ALT,
933 raw_bytes: b"\x1bS".to_vec(),
934 };
935
936 assert_that!(
937 send_decoded_input(&cmd_sender, &request_sender, vec![DecodedInput::Key(key.clone())]),
938 eq(ClientInputSend::Accepted)
939 );
940
941 assert_that!(request_receiver.blocking_recv(), eq(Some(ClientRequest::Key(key))));
942 }
943
944 #[test]
945 fn test_send_decoded_input_when_paste_arrives_uses_input_queue() {
946 let (cmd_sender, _cmd_receiver) = tokio::sync::mpsc::channel(1);
947 let (request_sender, mut request_receiver) = tokio::sync::mpsc::channel(1);
948
949 assert_that!(
950 send_decoded_input(
951 &cmd_sender,
952 &request_sender,
953 vec![DecodedInput::Paste(b"one\ntwo\n".to_vec())],
954 ),
955 eq(ClientInputSend::Accepted)
956 );
957
958 assert_that!(
959 request_receiver.blocking_recv(),
960 eq(Some(ClientRequest::Paste(b"one\ntwo\n".to_vec())))
961 );
962 }
963
964 #[test]
965 fn test_send_decoded_input_when_mouse_arrives_emits_local_mouse_action() {
966 let (cmd_sender, mut cmd_receiver) = tokio::sync::mpsc::channel(1);
967 let (request_sender, _request_receiver) = tokio::sync::mpsc::channel(1);
968 let event = ClientMouseEvent {
969 button: 0,
970 phase: ClientMouseEventPhase::Press,
971 position: muxr_core::ClientMousePosition { row: 4, col: 9 },
972 };
973
974 assert_that!(
975 send_decoded_input(&cmd_sender, &request_sender, vec![DecodedInput::Mouse(event)]),
976 eq(ClientInputSend::Accepted)
977 );
978
979 assert_that!(
980 matches!(
981 cmd_receiver.blocking_recv(),
982 Some(ClientInputCmd::Action(ClientInputAction::Mouse(actual))) if actual == event
983 ),
984 eq(true)
985 );
986 }
987
988 #[test]
989 fn test_send_decoded_input_when_mouse_motion_action_queue_is_full_drops_without_blocking() -> rootcause::Result<()>
990 {
991 let (cmd_sender, mut cmd_receiver) = tokio::sync::mpsc::channel(1);
992 let (request_sender, _request_receiver) = tokio::sync::mpsc::channel(1);
993 assert_that!(
994 cmd_sender.try_send(ClientInputCmd::Action(ClientInputAction::CopySelection)),
995 ok(eq(()))
996 );
997 let event = ClientMouseEvent {
998 button: 32,
999 phase: ClientMouseEventPhase::Press,
1000 position: muxr_core::ClientMousePosition { row: 4, col: 9 },
1001 };
1002 let (result_sender, result_receiver) = std::sync::mpsc::channel();
1003 let handle = thread::spawn(move || {
1004 let _ = result_sender.send(send_decoded_input(
1005 &cmd_sender,
1006 &request_sender,
1007 vec![DecodedInput::Mouse(event)],
1008 ));
1009 });
1010 let result = match result_receiver.recv_timeout(Duration::from_secs(1)) {
1011 Ok(result) => result,
1012 Err(error) => {
1013 drop(cmd_receiver);
1014 handle
1015 .join()
1016 .map_err(|error| report!("muxr mouse input test thread panicked").attach(format!("{error:?}")))?;
1017 return Err(report!("muxr mouse motion blocked on full input-action queue").attach(format!("{error}")));
1018 }
1019 };
1020
1021 assert_that!(result, eq(ClientInputSend::Accepted));
1022 assert_that!(
1023 matches!(
1024 cmd_receiver.try_recv(),
1025 Ok(ClientInputCmd::Action(ClientInputAction::CopySelection))
1026 ),
1027 eq(true)
1028 );
1029 assert_that!(cmd_receiver.try_recv(), err(anything()));
1030 handle
1031 .join()
1032 .map_err(|error| report!("muxr mouse input test thread panicked").attach(format!("{error:?}")))?;
1033 Ok(())
1034 }
1035
1036 #[test]
1037 fn test_send_decoded_input_when_mouse_wheel_action_queue_is_full_waits_for_queue_space() -> rootcause::Result<()> {
1038 let (cmd_sender, mut cmd_receiver) = tokio::sync::mpsc::channel(1);
1039 let (request_sender, _request_receiver) = tokio::sync::mpsc::channel(1);
1040 assert_that!(
1041 cmd_sender.try_send(ClientInputCmd::Action(ClientInputAction::CopySelection)),
1042 ok(eq(()))
1043 );
1044 let event = ClientMouseEvent {
1045 button: 64,
1046 phase: ClientMouseEventPhase::Press,
1047 position: muxr_core::ClientMousePosition { row: 4, col: 9 },
1048 };
1049 let (result_sender, result_receiver) = std::sync::mpsc::channel();
1050 let handle = thread::spawn(move || {
1051 let _ = result_sender.send(send_decoded_input(
1052 &cmd_sender,
1053 &request_sender,
1054 vec![DecodedInput::Mouse(event)],
1055 ));
1056 });
1057
1058 assert_that!(result_receiver.recv_timeout(Duration::from_millis(50)), err(anything()));
1059 assert_that!(
1060 matches!(
1061 cmd_receiver.blocking_recv(),
1062 Some(ClientInputCmd::Action(ClientInputAction::CopySelection))
1063 ),
1064 eq(true)
1065 );
1066 assert_that!(
1067 result_receiver.recv_timeout(Duration::from_secs(1)),
1068 eq(Ok(ClientInputSend::Accepted))
1069 );
1070 assert_that!(
1071 matches!(
1072 cmd_receiver.blocking_recv(),
1073 Some(ClientInputCmd::Action(ClientInputAction::Mouse(actual))) if actual == event
1074 ),
1075 eq(true)
1076 );
1077 handle
1078 .join()
1079 .map_err(|error| report!("muxr mouse input test thread panicked").attach(format!("{error:?}")))?;
1080 Ok(())
1081 }
1082
1083 #[test]
1084 fn test_send_decoded_input_when_copy_selection_arrives_emits_local_action() {
1085 let (cmd_sender, mut cmd_receiver) = tokio::sync::mpsc::channel(1);
1086 let (request_sender, _request_receiver) = tokio::sync::mpsc::channel(1);
1087
1088 assert_that!(
1089 send_decoded_input(&cmd_sender, &request_sender, vec![DecodedInput::CopySelection]),
1090 eq(ClientInputSend::Accepted)
1091 );
1092
1093 assert_that!(
1094 matches!(
1095 cmd_receiver.blocking_recv(),
1096 Some(ClientInputCmd::Action(ClientInputAction::CopySelection))
1097 ),
1098 eq(true)
1099 );
1100 }
1101
1102 #[test]
1103 fn test_send_decoded_input_when_inline_copy_selection_arrives_emits_local_action() {
1104 let (cmd_sender, mut cmd_receiver) = tokio::sync::mpsc::channel(1);
1105 let (request_sender, _request_receiver) = tokio::sync::mpsc::channel(1);
1106
1107 assert_that!(
1108 send_decoded_input(&cmd_sender, &request_sender, vec![DecodedInput::CopySelectionInline]),
1109 eq(ClientInputSend::Accepted)
1110 );
1111
1112 assert_that!(
1113 matches!(
1114 cmd_receiver.blocking_recv(),
1115 Some(ClientInputCmd::Action(ClientInputAction::CopySelectionInline))
1116 ),
1117 eq(true)
1118 );
1119 }
1120
1121 #[test]
1122 fn test_send_selection_edge_scroll_request_when_scroll_is_pending_waits_for_render_ack() -> rootcause::Result<()> {
1123 let (input_sender, mut input_receiver) = tokio::sync::mpsc::channel(2);
1124 let mut renderer = ClientRenderer::with_synchronized_output(
1125 layout_snapshot()?,
1126 pane_regions_snapshot()?,
1127 SynchronizedOutput::Csi,
1128 );
1129 renderer.apply_render_logical(muxr_core::RenderUpdate::Baseline(render_baseline()?))?;
1130 renderer.apply_selection_input_logical(SelectionInput::Start(ClientMousePosition { row: 0, col: 0 }))?;
1131 let initial = renderer
1132 .set_selection_edge_drag(ClientMousePosition { row: 2, col: 1 }, None)
1133 .ok_or_else(|| report!("expected initial muxr edge scroll request"))?;
1134 let expected = ClientRequest::ScrollPaneLineAt {
1135 direction: PaneScrollDirection::Down,
1136 position: ClientMousePosition { row: 0, col: 1 },
1137 };
1138 assert_that!(
1139 copy_selection_test_helpers::edge_scroll_request(&initial),
1140 eq(&expected)
1141 );
1142 assert_that!(
1143 send_edge_scroll_request(&input_sender, &mut renderer, initial),
1144 eq(ClientInputSend::Accepted)
1145 );
1146 assert_that!(input_receiver.blocking_recv(), eq(Some(expected.clone())));
1147 assert_that!(
1148 send_selection_edge_scroll_request(&input_sender, &mut renderer),
1149 eq(ClientInputSend::Accepted)
1150 );
1151 assert_that!(
1152 input_receiver.try_recv(),
1153 err(matches_pattern!(tokio::sync::mpsc::error::TryRecvError::Empty))
1154 );
1155
1156 renderer.apply_pane_regions_logical(pane_regions_snapshot_with_visible_top_row(1)?);
1157 renderer.apply_render_logical(muxr_core::RenderUpdate::Baseline(render_baseline()?))?;
1158 let flushed = renderer.presentation_snapshot();
1159 renderer.acknowledge_presentation(&flushed);
1160 assert_that!(
1161 send_selection_edge_scroll_request(&input_sender, &mut renderer),
1162 eq(ClientInputSend::Accepted)
1163 );
1164
1165 assert_that!(input_receiver.try_recv(), eq(Ok(expected)));
1166 Ok(())
1167 }
1168
1169 #[test]
1170 fn test_send_edge_scroll_request_when_queue_is_full_does_not_mark_scroll_pending() -> rootcause::Result<()> {
1171 let (input_sender, mut input_receiver) = tokio::sync::mpsc::channel(1);
1172 assert_that!(input_sender.try_send(ClientRequest::Pong), ok(eq(())));
1173 let mut renderer = ClientRenderer::with_synchronized_output(
1174 layout_snapshot()?,
1175 pane_regions_snapshot()?,
1176 SynchronizedOutput::Csi,
1177 );
1178 renderer.apply_render_logical(muxr_core::RenderUpdate::Baseline(render_baseline()?))?;
1179 renderer.apply_selection_input_logical(SelectionInput::Start(ClientMousePosition { row: 0, col: 0 }))?;
1180 let request = renderer
1181 .set_selection_edge_drag(ClientMousePosition { row: 2, col: 1 }, None)
1182 .ok_or_else(|| report!("expected muxr edge scroll request"))?;
1183
1184 assert_that!(
1185 send_edge_scroll_request(&input_sender, &mut renderer, request),
1186 eq(ClientInputSend::Accepted)
1187 );
1188 assert_that!(input_receiver.try_recv(), eq(Ok(ClientRequest::Pong)));
1189 assert_that!(
1190 send_selection_edge_scroll_request(&input_sender, &mut renderer),
1191 eq(ClientInputSend::Accepted)
1192 );
1193
1194 assert_that!(
1195 input_receiver.blocking_recv(),
1196 eq(Some(ClientRequest::ScrollPaneLineAt {
1197 direction: PaneScrollDirection::Down,
1198 position: ClientMousePosition { row: 0, col: 1 },
1199 }))
1200 );
1201 Ok(())
1202 }
1203
1204 fn session_paths(base: &Path, raw: &str) -> rootcause::Result<(SessionName, SessionPaths)> {
1205 let session = raw.parse()?;
1206 let root = base.join("sessions").join(raw);
1207
1208 Ok((
1209 session,
1210 SessionPaths {
1211 socket: root.join("server.sock"),
1212 pid: root.join("server.pid"),
1213 layout: root.join("layout.json"),
1214 panes: root.join("panes"),
1215 root,
1216 },
1217 ))
1218 }
1219
1220 fn layout_snapshot() -> rootcause::Result<LayoutSnapshot> {
1221 let active_tab = TabId::new(1)?;
1222 let active_pane = PaneId::new(1)?;
1223 let pane = PaneSnapshot {
1224 tracked_process_state: muxr_core::TrackedProcessState::None,
1225 cwd: "/tmp".to_owned(),
1226 cmd_label: None,
1227 focus_seq: 1,
1228 id: active_pane,
1229 title: "shell".to_owned(),
1230 };
1231 let tab = TabSnapshot::new(active_tab, "default", active_pane, vec![pane])?;
1232 LayoutSnapshot::new(active_tab, vec![tab])
1233 }
1234
1235 fn pane_regions_snapshot() -> rootcause::Result<PaneRegionsSnapshot> {
1236 self::pane_regions_snapshot_with_visible_top_row(0)
1237 }
1238
1239 fn pane_regions_snapshot_with_visible_top_row(visible_top_row: u64) -> rootcause::Result<PaneRegionsSnapshot> {
1240 PaneRegionsSnapshot::new(vec![muxr_core::PaneRegionSnapshot::new(
1241 muxr_core::PaneId::new(1)?,
1242 0,
1243 0,
1244 2,
1245 1,
1246 muxr_core::PaneMouseMode::None,
1247 visible_top_row,
1248 )?])
1249 }
1250
1251 fn render_baseline() -> rootcause::Result<muxr_core::RenderBaseline> {
1252 muxr_core::RenderBaseline::new(
1253 1,
1254 TerminalSize::new(2, 1)?,
1255 muxr_core::RenderCursor {
1256 row: 0,
1257 col: 1,
1258 shape: muxr_core::RenderCursorShape::Default,
1259 visibility: muxr_core::RenderCursorVisibility::Visible,
1260 },
1261 vec![muxr_core::RenderRowSpan::new(
1262 0,
1263 0,
1264 vec![render_cell("a"), render_cell("b")],
1265 )?],
1266 )
1267 }
1268
1269 fn render_cell(text: &str) -> muxr_core::RenderCell {
1270 muxr_core::RenderCell::narrow(text, muxr_core::RenderStyle::default())
1271 }
1272
1273 fn runtime() -> rootcause::Result<tokio::runtime::Runtime> {
1274 Ok(tokio::runtime::Runtime::new().context("failed to build muxr client test runtime")?)
1275 }
1276}