1use std::fs::File;
4use std::io::{self, Read, Write};
5use std::net::Shutdown;
6use std::os::fd::AsFd;
7use std::os::unix::net::UnixStream;
8use std::sync::atomic::{AtomicBool, Ordering};
9use std::sync::{mpsc, Arc};
10use std::thread;
11use std::time::{Duration, Instant};
12
13use rmux_proto::{
14 decode_attach_data_frame, encode_attach_data, encode_attach_data_into_slice,
15 encode_attach_message, AttachFrameDecoder, AttachMessage, RmuxError, TerminalGeometry,
16 TerminalSize, ATTACH_DATA_HEADER_LEN,
17};
18use rustix::event::{poll, PollFd, PollFlags, Timespec};
19use rustix::process::{kill_process, Signal};
20
21use crate::ClientError;
22
23#[path = "attach/render_drain.rs"]
24mod render_drain;
25#[path = "attach/resize.rs"]
26mod resize;
27#[path = "attach/screen.rs"]
28mod screen;
29#[path = "attach/terminal.rs"]
30mod terminal;
31#[path = "attach/terminal_cleanup.rs"]
32mod terminal_cleanup;
33
34use render_drain::{drain_available_attach_stream, flush_pending_render};
35#[cfg(test)]
36use resize::terminal_size_from_fd;
37use resize::{terminal_geometry_from_fd, ResizeWatcher, SignalMaskGuard};
38use screen::{AttachScreenTracker, AttachStopDetector};
39use terminal::current_process_pid;
40pub use terminal::{AttachError, RawTerminal, Result};
41
42#[cfg(test)]
43use terminal_cleanup::fallback_attach_stop_sequence;
44
45const READ_BUFFER_SIZE: usize = 8192;
46const STACK_ATTACH_DATA_PAYLOAD: usize = 1024;
47const POLL_TIMEOUT: Timespec = Timespec {
48 tv_sec: 0,
49 tv_nsec: 100_000_000,
50};
51const RENDER_MAX_PENDING: Duration = Duration::from_millis(8);
52
53pub fn attach_terminal(stream: UnixStream) -> std::result::Result<(), ClientError> {
55 attach_terminal_with_initial_bytes(stream, Vec::new())
56}
57
58pub fn attach_terminal_with_initial_bytes(
60 stream: UnixStream,
61 initial_bytes: Vec<u8>,
62) -> std::result::Result<(), ClientError> {
63 attach_terminal_with_initial_bytes_and_geometry_flag(stream, initial_bytes, false)
64}
65
66pub fn attach_terminal_with_initial_bytes_and_resize_geometry(
72 stream: UnixStream,
73 initial_bytes: Vec<u8>,
74) -> std::result::Result<(), ClientError> {
75 attach_terminal_with_initial_bytes_and_geometry_flag(stream, initial_bytes, true)
76}
77
78fn attach_terminal_with_initial_bytes_and_geometry_flag(
79 stream: UnixStream,
80 initial_bytes: Vec<u8>,
81 resize_geometry_enabled: bool,
82) -> std::result::Result<(), ClientError> {
83 let terminal = io::stdin();
84 let input = io::stdin();
85 let output = File::from(
86 io::stdout()
87 .as_fd()
88 .try_clone_to_owned()
89 .map_err(AttachError::from)?,
90 );
91
92 attach_with_terminal_with_initial_bytes(
93 stream,
94 initial_bytes,
95 &terminal,
96 input,
97 output,
98 resize_geometry_enabled,
99 )
100}
101
102pub fn attach_with_terminal<Terminal, Input, Output>(
107 stream: UnixStream,
108 terminal: &Terminal,
109 input: Input,
110 output: Output,
111) -> std::result::Result<(), ClientError>
112where
113 Terminal: AsFd,
114 Input: Read + AsFd + Send + 'static,
115 Output: Write + Send + 'static,
116{
117 attach_with_terminal_with_initial_bytes(stream, Vec::new(), terminal, input, output, false)
118}
119
120fn attach_with_terminal_with_initial_bytes<Terminal, Input, Output>(
121 stream: UnixStream,
122 initial_bytes: Vec<u8>,
123 terminal: &Terminal,
124 input: Input,
125 output: Output,
126 resize_geometry_enabled: bool,
127) -> std::result::Result<(), ClientError>
128where
129 Terminal: AsFd,
130 Input: Read + AsFd + Send + 'static,
131 Output: Write + Send + 'static,
132{
133 let raw_terminal = RawTerminal::from_fd(terminal).map_err(ClientError::from)?;
134 let _ = raw_terminal.flush_pending_input();
135 let screen_tracker = AttachScreenTracker::default();
136 let attach_state = AttachTerminalState {
137 stream,
138 initial_bytes,
139 terminal,
140 raw_terminal: &raw_terminal,
141 screen_tracker: &screen_tracker,
142 resize_geometry_enabled,
143 };
144 let result = drive_attach_with_terminal_state(attach_state, input, output);
145 if result.is_err() && !screen_tracker.was_stopped() {
146 let _ = raw_terminal.restore_attach_terminal_state();
147 }
148 let _ = raw_terminal.flush_pending_input();
149 drop(raw_terminal);
150 result
151}
152
153struct AttachTerminalState<'a, Terminal> {
154 stream: UnixStream,
155 initial_bytes: Vec<u8>,
156 terminal: &'a Terminal,
157 raw_terminal: &'a RawTerminal,
158 screen_tracker: &'a AttachScreenTracker,
159 resize_geometry_enabled: bool,
160}
161
162struct AttachStreamState<'a> {
163 stream: UnixStream,
164 initial_bytes: Vec<u8>,
165 raw_terminal: Option<&'a RawTerminal>,
166 screen_tracker: AttachScreenTracker,
167 resize_events: mpsc::Receiver<TerminalGeometry>,
168 resize_geometry_enabled: bool,
169}
170
171fn drive_attach_with_terminal_state<Terminal, Input, Output>(
172 state: AttachTerminalState<'_, Terminal>,
173 input: Input,
174 output: Output,
175) -> std::result::Result<(), ClientError>
176where
177 Terminal: AsFd,
178 Input: Read + AsFd + Send + 'static,
179 Output: Write + Send + 'static,
180{
181 let _signal_mask = SignalMaskGuard::block_winch().map_err(ClientError::from)?;
184 let (resize_tx, resize_rx) = mpsc::channel();
185 let initial_geometry = terminal_geometry_from_fd(state.terminal).map_err(ClientError::from)?;
186 let terminal_fd = state
187 .terminal
188 .as_fd()
189 .try_clone_to_owned()
190 .map_err(AttachError::from)?;
191
192 if let Some(initial_geometry) = initial_geometry {
193 resize_tx.send(initial_geometry).map_err(|_| {
194 ClientError::Io(io::Error::other(
195 "resize channel closed before attach start",
196 ))
197 })?;
198 }
199
200 let resize_watcher = ResizeWatcher::spawn(terminal_fd, resize_tx)?;
201 let stream_state = AttachStreamState {
202 stream: state.stream,
203 initial_bytes: state.initial_bytes,
204 raw_terminal: Some(state.raw_terminal),
205 screen_tracker: state.screen_tracker.clone(),
206 resize_events: resize_rx,
207 resize_geometry_enabled: state.resize_geometry_enabled,
208 };
209 let attach_result = drive_attach_stream_inner(stream_state, input, output);
210 drop(resize_watcher);
211 attach_result
212}
213
214pub fn drive_attach_stream<Input, Output>(
216 stream: UnixStream,
217 input: Input,
218 output: Output,
219 resize_events: mpsc::Receiver<TerminalSize>,
220) -> std::result::Result<(), ClientError>
221where
222 Input: Read + AsFd + Send + 'static,
223 Output: Write + Send + 'static,
224{
225 let resize_events = geometry_resize_events_from_size_events(resize_events);
226 let stream_state = AttachStreamState {
227 stream,
228 initial_bytes: Vec::new(),
229 raw_terminal: None,
230 screen_tracker: AttachScreenTracker::default(),
231 resize_events,
232 resize_geometry_enabled: false,
233 };
234 drive_attach_stream_inner(stream_state, input, output)
235}
236
237fn drive_attach_stream_inner<Input, Output>(
238 state: AttachStreamState<'_>,
239 input: Input,
240 output: Output,
241) -> std::result::Result<(), ClientError>
242where
243 Input: Read + AsFd + Send + 'static,
244 Output: Write + Send + 'static,
245{
246 let control = state.stream.try_clone().map_err(ClientError::Io)?;
247 let mut lock_stream = state.stream.try_clone().map_err(ClientError::Io)?;
248 let input_stream = state.stream.try_clone().map_err(ClientError::Io)?;
249 let (input_wakeup, wake_input_thread) = UnixStream::pair().map_err(ClientError::Io)?;
250 let closed = Arc::new(AtomicBool::new(false));
251 let input_closed = Arc::clone(&closed);
252 let output_closed = Arc::clone(&closed);
253 let locked = Arc::new(AtomicBool::new(false));
254 let input_locked = Arc::clone(&locked);
255 let output_locked = Arc::clone(&locked);
256 let (event_tx, event_rx) = mpsc::channel();
257
258 let input_thread = thread::spawn(move || {
259 input_loop(
260 input_stream,
261 input,
262 state.resize_events,
263 state.resize_geometry_enabled,
264 input_closed,
265 input_locked,
266 wake_input_thread,
267 )
268 });
269 let output_screen_tracker = state.screen_tracker.clone();
270 let output_thread = thread::spawn(move || {
271 let result = output_loop(
272 state.stream,
273 state.initial_bytes,
274 output,
275 output_closed,
276 output_locked,
277 output_screen_tracker,
278 event_tx.clone(),
279 );
280 let _ = event_tx.send(ClientAttachEvent::OutputDone);
281 result
282 });
283
284 let output_result = wait_for_output_thread(
285 output_thread,
286 state.raw_terminal,
287 &mut lock_stream,
288 &locked,
289 event_rx,
290 )?;
291 closed.store(true, Ordering::SeqCst);
292 let _ = control.shutdown(Shutdown::Both);
293 let _ = input_wakeup.shutdown(Shutdown::Both);
294 let input_result = join_attach_thread(input_thread)?;
295
296 output_result?;
297 input_result
298}
299
300fn geometry_resize_events_from_size_events(
301 resize_events: mpsc::Receiver<TerminalSize>,
302) -> mpsc::Receiver<TerminalGeometry> {
303 let (geometry_tx, geometry_rx) = mpsc::channel();
304 let _forwarder = thread::spawn(move || {
305 while let Ok(size) = resize_events.recv() {
306 if geometry_tx.send(TerminalGeometry::from_size(size)).is_err() {
307 break;
308 }
309 }
310 });
311 geometry_rx
312}
313
314fn input_loop<Input>(
315 mut stream: UnixStream,
316 mut input: Input,
317 resize_events: mpsc::Receiver<TerminalGeometry>,
318 resize_geometry_enabled: bool,
319 closed: Arc<AtomicBool>,
320 locked: Arc<AtomicBool>,
321 wakeup: UnixStream,
322) -> std::result::Result<(), ClientError>
323where
324 Input: Read + AsFd,
325{
326 let mut read_buffer = [0_u8; READ_BUFFER_SIZE];
327
328 loop {
329 if closed.load(Ordering::SeqCst) {
330 return Ok(());
331 }
332
333 drain_resize_events(&mut stream, &resize_events, resize_geometry_enabled)?;
334 if locked.load(Ordering::SeqCst) {
335 thread::sleep(Duration::from_millis(20));
336 continue;
337 }
338
339 let mut fds = [
340 PollFd::new(&input, PollFlags::IN | PollFlags::ERR | PollFlags::HUP),
341 PollFd::new(&wakeup, PollFlags::IN | PollFlags::ERR | PollFlags::HUP),
342 ];
343 match poll(&mut fds, Some(&POLL_TIMEOUT)) {
344 Ok(0) => continue,
345 Ok(_) => {}
346 Err(rustix::io::Errno::INTR) => continue,
347 Err(error) => return Err(ClientError::Io(error.into())),
348 }
349
350 if !fds[1].revents().is_empty() {
351 return Ok(());
352 }
353
354 let ready = fds[0].revents();
355 if ready.is_empty() {
356 continue;
357 }
358 if closed.load(Ordering::SeqCst) {
359 return Ok(());
360 }
361 if !ready.contains(PollFlags::IN) {
362 if ready.contains(PollFlags::HUP) || ready.contains(PollFlags::ERR) {
363 shutdown_attach_writes(&stream)?;
364 return Ok(());
365 }
366 continue;
367 }
368
369 let bytes_read = match input.read(&mut read_buffer) {
370 Ok(0) => {
371 shutdown_attach_writes(&stream)?;
372 return Ok(());
373 }
374 Ok(bytes_read) => bytes_read,
375 Err(error) if error.kind() == io::ErrorKind::Interrupted => continue,
376 Err(error) => return Err(ClientError::Io(error)),
377 };
378
379 write_attach_data(&mut stream, &read_buffer[..bytes_read])?;
380 }
381}
382
383fn output_loop<Output>(
384 mut stream: UnixStream,
385 initial_bytes: Vec<u8>,
386 mut output: Output,
387 closed: Arc<AtomicBool>,
388 locked: Arc<AtomicBool>,
389 screen_tracker: AttachScreenTracker,
390 event_tx: mpsc::Sender<ClientAttachEvent>,
391) -> std::result::Result<(), ClientError>
392where
393 Output: Write,
394{
395 let mut decoder = AttachFrameDecoder::new();
396 decoder.push_bytes(&initial_bytes);
397 let mut read_buffer = [0_u8; READ_BUFFER_SIZE];
398 let mut stop_detector = AttachStopDetector::new(screen_tracker.clone());
399 let mut pending_render = None::<Vec<u8>>;
400 let mut pending_render_started_at = None::<Instant>;
401 let mut pending_render_drained_after_deadline = false;
402 let mut painted_render_frame = false;
403 let mut data_scratch = [0_u8; READ_BUFFER_SIZE];
404
405 loop {
406 loop {
407 while let Some(bytes) = decoder
408 .next_data_payload_into(&mut data_scratch)
409 .map_err(ClientError::from)?
410 {
411 flush_pending_render_state(
412 &mut output,
413 &mut pending_render,
414 &mut pending_render_started_at,
415 )?;
416 if matches!(
417 handle_attach_data_payload(
418 &mut output,
419 &locked,
420 &closed,
421 &mut stop_detector,
422 bytes,
423 )?,
424 AttachDataPayloadOutcome::Stop
425 ) {
426 return Ok(());
427 }
428 }
429
430 let Some(message) = decoder.next_message().map_err(ClientError::from)? else {
431 break;
432 };
433 match message {
434 AttachMessage::Data(bytes) => {
435 flush_pending_render_state(
436 &mut output,
437 &mut pending_render,
438 &mut pending_render_started_at,
439 )?;
440 if matches!(
441 handle_attach_data_payload(
442 &mut output,
443 &locked,
444 &closed,
445 &mut stop_detector,
446 &bytes,
447 )?,
448 AttachDataPayloadOutcome::Stop
449 ) {
450 return Ok(());
451 }
452 }
453 AttachMessage::Render(bytes) => {
454 if locked.load(Ordering::SeqCst) {
455 continue;
456 }
457 if pending_render.is_none() {
458 pending_render_started_at = Some(Instant::now());
459 pending_render_drained_after_deadline = false;
460 }
461 pending_render = Some(bytes);
462 if !painted_render_frame
463 && flush_pending_render_state(
464 &mut output,
465 &mut pending_render,
466 &mut pending_render_started_at,
467 )?
468 {
469 pending_render_drained_after_deadline = false;
470 painted_render_frame = true;
471 }
472 }
473 AttachMessage::KeyDispatched(_) => {}
474 AttachMessage::Resize(_) | AttachMessage::ResizeGeometry(_) => {
475 flush_pending_render_state(
476 &mut output,
477 &mut pending_render,
478 &mut pending_render_started_at,
479 )?;
480 return Err(ClientError::Protocol(RmuxError::Decode(
481 "received unexpected resize message from attach stream".to_owned(),
482 )));
483 }
484 AttachMessage::Lock(command) => {
485 flush_pending_render_state(
486 &mut output,
487 &mut pending_render,
488 &mut pending_render_started_at,
489 )?;
490 locked.store(true, Ordering::SeqCst);
491 send_attach_action(&event_tx, ClientAttachAction::Lock(command))?;
492 }
493 AttachMessage::LockShellCommand(command) => {
494 flush_pending_render_state(
495 &mut output,
496 &mut pending_render,
497 &mut pending_render_started_at,
498 )?;
499 locked.store(true, Ordering::SeqCst);
500 send_attach_action(
501 &event_tx,
502 ClientAttachAction::Lock(command.command().to_owned()),
503 )?;
504 }
505 AttachMessage::Suspend => {
506 flush_pending_render_state(
507 &mut output,
508 &mut pending_render,
509 &mut pending_render_started_at,
510 )?;
511 locked.store(true, Ordering::SeqCst);
512 send_attach_action(&event_tx, ClientAttachAction::Suspend)?;
513 }
514 AttachMessage::DetachKill => {
515 flush_pending_render_state(
516 &mut output,
517 &mut pending_render,
518 &mut pending_render_started_at,
519 )?;
520 closed.store(true, Ordering::SeqCst);
521 send_attach_action(&event_tx, ClientAttachAction::DetachKill)?;
522 return Ok(());
523 }
524 AttachMessage::DetachExec(command) => {
525 flush_pending_render_state(
526 &mut output,
527 &mut pending_render,
528 &mut pending_render_started_at,
529 )?;
530 closed.store(true, Ordering::SeqCst);
531 send_attach_action(&event_tx, ClientAttachAction::DetachExec(command))?;
532 return Ok(());
533 }
534 AttachMessage::DetachExecShellCommand(command) => {
535 flush_pending_render_state(
536 &mut output,
537 &mut pending_render,
538 &mut pending_render_started_at,
539 )?;
540 closed.store(true, Ordering::SeqCst);
541 send_attach_action(
542 &event_tx,
543 ClientAttachAction::DetachExec(command.command().to_owned()),
544 )?;
545 return Ok(());
546 }
547 AttachMessage::Unlock => {
548 flush_pending_render_state(
549 &mut output,
550 &mut pending_render,
551 &mut pending_render_started_at,
552 )?;
553 return Err(ClientError::Protocol(RmuxError::Decode(
554 "received unexpected unlock message from attach stream".to_owned(),
555 )));
556 }
557 AttachMessage::Keystroke(_) => {
558 flush_pending_render_state(
559 &mut output,
560 &mut pending_render,
561 &mut pending_render_started_at,
562 )?;
563 return Err(ClientError::Protocol(RmuxError::Decode(
564 "received unexpected keystroke message from attach stream".to_owned(),
565 )));
566 }
567 }
568 }
569
570 if pending_render.is_some() {
571 let pending_expired = pending_render_expired(pending_render_started_at);
572 if (!pending_expired || !pending_render_drained_after_deadline)
573 && drain_available_attach_stream(&mut stream, &mut decoder, &mut read_buffer)?
574 {
575 if pending_expired {
576 pending_render_drained_after_deadline = true;
577 }
578 continue;
579 }
580 }
581 if pending_render.is_some() && !pending_render_expired(pending_render_started_at) {
582 sleep_until_pending_render_deadline(pending_render_started_at);
583 if drain_available_attach_stream(&mut stream, &mut decoder, &mut read_buffer)? {
584 continue;
585 }
586 }
587 if flush_pending_render_state(
588 &mut output,
589 &mut pending_render,
590 &mut pending_render_started_at,
591 )? {
592 pending_render_drained_after_deadline = false;
593 painted_render_frame = true;
594 }
595
596 let bytes_read = match stream.read(&mut read_buffer) {
597 Ok(0) => {
598 closed.store(true, Ordering::SeqCst);
599 if screen_tracker.was_stopped() {
600 return Ok(());
601 }
602 return Err(ClientError::Io(io::Error::new(
603 io::ErrorKind::UnexpectedEof,
604 "attach stream closed before attach-stop sequence",
605 )));
606 }
607 Ok(bytes_read) => bytes_read,
608 Err(error) if error.kind() == io::ErrorKind::Interrupted => continue,
609 Err(error)
610 if screen_tracker.was_stopped()
611 && matches!(
612 error.kind(),
613 io::ErrorKind::ConnectionReset | io::ErrorKind::BrokenPipe
614 ) =>
615 {
616 return Ok(());
617 }
618 Err(error) => return Err(ClientError::Io(error)),
619 };
620
621 let mut consumed = 0;
622 if decoder.is_empty() {
623 while consumed < bytes_read {
624 let Some(frame) = decode_attach_data_frame(&read_buffer[consumed..])
625 .map_err(ClientError::from)?
626 else {
627 break;
628 };
629 if matches!(
630 handle_attach_data_payload(
631 &mut output,
632 &locked,
633 &closed,
634 &mut stop_detector,
635 frame.payload(),
636 )?,
637 AttachDataPayloadOutcome::Stop
638 ) {
639 return Ok(());
640 }
641 consumed += frame.frame_len();
642 }
643 }
644 if consumed < bytes_read {
645 decoder.push_bytes(&read_buffer[consumed..bytes_read]);
646 }
647 }
648}
649
650enum AttachDataPayloadOutcome {
651 Continue,
652 Stop,
653}
654
655fn handle_attach_data_payload<Output>(
656 output: &mut Output,
657 locked: &Arc<AtomicBool>,
658 closed: &Arc<AtomicBool>,
659 stop_detector: &mut AttachStopDetector,
660 bytes: &[u8],
661) -> std::result::Result<AttachDataPayloadOutcome, ClientError>
662where
663 Output: Write,
664{
665 let observation = stop_detector.observe(bytes);
666 let attach_done = observation.attach_done();
667 if locked.load(Ordering::SeqCst) {
668 return Ok(AttachDataPayloadOutcome::Continue);
669 }
670 output.write_all(bytes).map_err(ClientError::Io)?;
671 output.flush().map_err(ClientError::Io)?;
672 if attach_done {
673 closed.store(true, Ordering::SeqCst);
674 return Ok(AttachDataPayloadOutcome::Stop);
675 }
676 Ok(AttachDataPayloadOutcome::Continue)
677}
678
679fn pending_render_expired(started_at: Option<Instant>) -> bool {
680 started_at.is_some_and(|started_at| started_at.elapsed() >= RENDER_MAX_PENDING)
681}
682
683fn sleep_until_pending_render_deadline(started_at: Option<Instant>) {
684 let Some(started_at) = started_at else {
685 return;
686 };
687 let Some(remaining) = RENDER_MAX_PENDING.checked_sub(started_at.elapsed()) else {
688 return;
689 };
690 if !remaining.is_zero() {
691 thread::sleep(remaining);
692 }
693}
694
695fn flush_pending_render_state<Output>(
696 output: &mut Output,
697 pending_render: &mut Option<Vec<u8>>,
698 pending_render_started_at: &mut Option<Instant>,
699) -> std::result::Result<bool, ClientError>
700where
701 Output: Write,
702{
703 let flushed = pending_render.is_some();
704 flush_pending_render(output, pending_render)?;
705 *pending_render_started_at = None;
706 Ok(flushed)
707}
708
709fn wait_for_output_thread(
710 output_thread: thread::JoinHandle<std::result::Result<(), ClientError>>,
711 raw_terminal: Option<&RawTerminal>,
712 lock_stream: &mut UnixStream,
713 locked: &Arc<AtomicBool>,
714 event_rx: mpsc::Receiver<ClientAttachEvent>,
715) -> std::result::Result<std::result::Result<(), ClientError>, ClientError> {
716 while let Ok(ClientAttachEvent::Action(action)) = event_rx.recv() {
717 handle_attach_action(raw_terminal, lock_stream, locked, action)?;
718 }
719
720 while let Ok(event) = event_rx.try_recv() {
721 match event {
722 ClientAttachEvent::Action(action) => {
723 handle_attach_action(raw_terminal, lock_stream, locked, action)?;
724 }
725 ClientAttachEvent::OutputDone => {}
726 }
727 }
728
729 join_attach_thread(output_thread)
730}
731
732fn send_attach_action(
733 event_tx: &mpsc::Sender<ClientAttachEvent>,
734 action: ClientAttachAction,
735) -> std::result::Result<(), ClientError> {
736 event_tx
737 .send(ClientAttachEvent::Action(action))
738 .map_err(|_| ClientError::Io(io::Error::other("attach event receiver closed")))
739}
740
741fn handle_attach_action(
742 raw_terminal: Option<&RawTerminal>,
743 lock_stream: &mut UnixStream,
744 locked: &Arc<AtomicBool>,
745 action: ClientAttachAction,
746) -> std::result::Result<(), ClientError> {
747 match action {
748 ClientAttachAction::Lock(command) => {
749 let Some(raw_terminal) = raw_terminal else {
750 locked.store(false, Ordering::SeqCst);
751 return Err(ClientError::Protocol(RmuxError::Decode(
752 "received unexpected lock request without a managed terminal".to_owned(),
753 )));
754 };
755 raw_terminal
756 .run_lock_command(&command)
757 .map_err(ClientError::from)?;
758 write_attach_message(lock_stream, AttachMessage::Unlock)?;
759 locked.store(false, Ordering::SeqCst);
760 Ok(())
761 }
762 ClientAttachAction::Suspend => {
763 let Some(raw_terminal) = raw_terminal else {
764 locked.store(false, Ordering::SeqCst);
765 return Err(ClientError::Protocol(RmuxError::Decode(
766 "received unexpected suspend request without a managed terminal".to_owned(),
767 )));
768 };
769 raw_terminal.suspend_self().map_err(ClientError::from)?;
770 write_attach_message(lock_stream, AttachMessage::Unlock)?;
771 locked.store(false, Ordering::SeqCst);
772 Ok(())
773 }
774 ClientAttachAction::DetachKill => {
775 if let Some(raw_terminal) = raw_terminal {
776 raw_terminal.restore().map_err(ClientError::from)?;
777 }
778 kill_process(current_process_pid().map_err(ClientError::Io)?, Signal::HUP)
779 .map_err(|error| ClientError::Io(error.into()))?;
780 Ok(())
781 }
782 ClientAttachAction::DetachExec(command) => {
783 let Some(raw_terminal) = raw_terminal else {
784 return Err(ClientError::Protocol(RmuxError::Decode(
785 "received unexpected detach exec request without a managed terminal".to_owned(),
786 )));
787 };
788 raw_terminal
789 .run_detach_exec_command(&command)
790 .map_err(ClientError::from)
791 }
792 }
793}
794
795fn drain_resize_events(
796 stream: &mut UnixStream,
797 resize_events: &mpsc::Receiver<TerminalGeometry>,
798 resize_geometry_enabled: bool,
799) -> std::result::Result<(), ClientError> {
800 while let Ok(geometry) = resize_events.try_recv() {
801 let message = if resize_geometry_enabled && geometry.pixels.is_some() {
802 AttachMessage::ResizeGeometry(geometry)
803 } else {
804 AttachMessage::Resize(geometry.size)
805 };
806 write_attach_message(stream, message)?;
807 }
808
809 Ok(())
810}
811
812fn write_attach_message(
813 stream: &mut UnixStream,
814 message: AttachMessage,
815) -> std::result::Result<(), ClientError> {
816 let frame = encode_attach_message(&message).map_err(ClientError::from)?;
817 stream.write_all(&frame).map_err(ClientError::Io)
818}
819
820fn write_attach_data(
821 stream: &mut UnixStream,
822 bytes: &[u8],
823) -> std::result::Result<(), ClientError> {
824 if bytes.len() <= STACK_ATTACH_DATA_PAYLOAD {
825 let mut frame = [0_u8; STACK_ATTACH_DATA_PAYLOAD + ATTACH_DATA_HEADER_LEN];
826 let len = encode_attach_data_into_slice(bytes, &mut frame).map_err(ClientError::from)?;
827 return stream.write_all(&frame[..len]).map_err(ClientError::Io);
828 }
829
830 let frame = encode_attach_data(bytes).map_err(ClientError::from)?;
831 stream.write_all(&frame).map_err(ClientError::Io)
832}
833
834fn join_attach_thread(
835 thread: thread::JoinHandle<std::result::Result<(), ClientError>>,
836) -> std::result::Result<std::result::Result<(), ClientError>, ClientError> {
837 thread
838 .join()
839 .map_err(|_| ClientError::Io(io::Error::other("attach thread panicked")))
840}
841
842fn shutdown_attach_writes(stream: &UnixStream) -> std::result::Result<(), ClientError> {
843 match stream.shutdown(Shutdown::Write) {
844 Ok(()) => Ok(()),
845 Err(error) if error.kind() == io::ErrorKind::NotConnected => Ok(()),
846 Err(error) => Err(ClientError::Io(error)),
847 }
848}
849
850#[derive(Debug)]
851enum ClientAttachAction {
852 Lock(String),
853 Suspend,
854 DetachKill,
855 DetachExec(String),
856}
857
858#[derive(Debug)]
859enum ClientAttachEvent {
860 Action(ClientAttachAction),
861 OutputDone,
862}
863
864#[cfg(test)]
865mod tests;