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
404 loop {
405 while let Some(message) = decoder.next_message().map_err(ClientError::from)? {
406 match message {
407 AttachMessage::Data(bytes) => {
408 flush_pending_render_state(
409 &mut output,
410 &mut pending_render,
411 &mut pending_render_started_at,
412 )?;
413 if matches!(
414 handle_attach_data_payload(
415 &mut output,
416 &locked,
417 &closed,
418 &mut stop_detector,
419 &bytes,
420 )?,
421 AttachDataPayloadOutcome::Stop
422 ) {
423 return Ok(());
424 }
425 }
426 AttachMessage::Render(bytes) => {
427 if locked.load(Ordering::SeqCst) {
428 continue;
429 }
430 if pending_render.is_none() {
431 pending_render_started_at = Some(Instant::now());
432 pending_render_drained_after_deadline = false;
433 }
434 pending_render = Some(bytes);
435 if !painted_render_frame
436 && flush_pending_render_state(
437 &mut output,
438 &mut pending_render,
439 &mut pending_render_started_at,
440 )?
441 {
442 pending_render_drained_after_deadline = false;
443 painted_render_frame = true;
444 }
445 }
446 AttachMessage::KeyDispatched(_) => {}
447 AttachMessage::Resize(_) | AttachMessage::ResizeGeometry(_) => {
448 flush_pending_render_state(
449 &mut output,
450 &mut pending_render,
451 &mut pending_render_started_at,
452 )?;
453 return Err(ClientError::Protocol(RmuxError::Decode(
454 "received unexpected resize message from attach stream".to_owned(),
455 )));
456 }
457 AttachMessage::Lock(command) => {
458 flush_pending_render_state(
459 &mut output,
460 &mut pending_render,
461 &mut pending_render_started_at,
462 )?;
463 locked.store(true, Ordering::SeqCst);
464 send_attach_action(&event_tx, ClientAttachAction::Lock(command))?;
465 }
466 AttachMessage::LockShellCommand(command) => {
467 flush_pending_render_state(
468 &mut output,
469 &mut pending_render,
470 &mut pending_render_started_at,
471 )?;
472 locked.store(true, Ordering::SeqCst);
473 send_attach_action(
474 &event_tx,
475 ClientAttachAction::Lock(command.command().to_owned()),
476 )?;
477 }
478 AttachMessage::Suspend => {
479 flush_pending_render_state(
480 &mut output,
481 &mut pending_render,
482 &mut pending_render_started_at,
483 )?;
484 locked.store(true, Ordering::SeqCst);
485 send_attach_action(&event_tx, ClientAttachAction::Suspend)?;
486 }
487 AttachMessage::DetachKill => {
488 flush_pending_render_state(
489 &mut output,
490 &mut pending_render,
491 &mut pending_render_started_at,
492 )?;
493 closed.store(true, Ordering::SeqCst);
494 send_attach_action(&event_tx, ClientAttachAction::DetachKill)?;
495 return Ok(());
496 }
497 AttachMessage::DetachExec(command) => {
498 flush_pending_render_state(
499 &mut output,
500 &mut pending_render,
501 &mut pending_render_started_at,
502 )?;
503 closed.store(true, Ordering::SeqCst);
504 send_attach_action(&event_tx, ClientAttachAction::DetachExec(command))?;
505 return Ok(());
506 }
507 AttachMessage::DetachExecShellCommand(command) => {
508 flush_pending_render_state(
509 &mut output,
510 &mut pending_render,
511 &mut pending_render_started_at,
512 )?;
513 closed.store(true, Ordering::SeqCst);
514 send_attach_action(
515 &event_tx,
516 ClientAttachAction::DetachExec(command.command().to_owned()),
517 )?;
518 return Ok(());
519 }
520 AttachMessage::Unlock => {
521 flush_pending_render_state(
522 &mut output,
523 &mut pending_render,
524 &mut pending_render_started_at,
525 )?;
526 return Err(ClientError::Protocol(RmuxError::Decode(
527 "received unexpected unlock message from attach stream".to_owned(),
528 )));
529 }
530 AttachMessage::Keystroke(_) => {
531 flush_pending_render_state(
532 &mut output,
533 &mut pending_render,
534 &mut pending_render_started_at,
535 )?;
536 return Err(ClientError::Protocol(RmuxError::Decode(
537 "received unexpected keystroke message from attach stream".to_owned(),
538 )));
539 }
540 }
541 }
542
543 if pending_render.is_some() {
544 let pending_expired = pending_render_expired(pending_render_started_at);
545 if (!pending_expired || !pending_render_drained_after_deadline)
546 && drain_available_attach_stream(&mut stream, &mut decoder, &mut read_buffer)?
547 {
548 if pending_expired {
549 pending_render_drained_after_deadline = true;
550 }
551 continue;
552 }
553 }
554 if pending_render.is_some() && !pending_render_expired(pending_render_started_at) {
555 sleep_until_pending_render_deadline(pending_render_started_at);
556 if drain_available_attach_stream(&mut stream, &mut decoder, &mut read_buffer)? {
557 continue;
558 }
559 }
560 if flush_pending_render_state(
561 &mut output,
562 &mut pending_render,
563 &mut pending_render_started_at,
564 )? {
565 pending_render_drained_after_deadline = false;
566 painted_render_frame = true;
567 }
568
569 let bytes_read = match stream.read(&mut read_buffer) {
570 Ok(0) => {
571 closed.store(true, Ordering::SeqCst);
572 if screen_tracker.was_stopped() {
573 return Ok(());
574 }
575 return Err(ClientError::Io(io::Error::new(
576 io::ErrorKind::UnexpectedEof,
577 "attach stream closed before attach-stop sequence",
578 )));
579 }
580 Ok(bytes_read) => bytes_read,
581 Err(error) if error.kind() == io::ErrorKind::Interrupted => continue,
582 Err(error)
583 if screen_tracker.was_stopped()
584 && matches!(
585 error.kind(),
586 io::ErrorKind::ConnectionReset | io::ErrorKind::BrokenPipe
587 ) =>
588 {
589 return Ok(());
590 }
591 Err(error) => return Err(ClientError::Io(error)),
592 };
593
594 let mut consumed = 0;
595 if decoder.is_empty() {
596 while consumed < bytes_read {
597 let Some(frame) = decode_attach_data_frame(&read_buffer[consumed..])
598 .map_err(ClientError::from)?
599 else {
600 break;
601 };
602 if matches!(
603 handle_attach_data_payload(
604 &mut output,
605 &locked,
606 &closed,
607 &mut stop_detector,
608 frame.payload(),
609 )?,
610 AttachDataPayloadOutcome::Stop
611 ) {
612 return Ok(());
613 }
614 consumed += frame.frame_len();
615 }
616 }
617 if consumed < bytes_read {
618 decoder.push_bytes(&read_buffer[consumed..bytes_read]);
619 }
620 }
621}
622
623enum AttachDataPayloadOutcome {
624 Continue,
625 Stop,
626}
627
628fn handle_attach_data_payload<Output>(
629 output: &mut Output,
630 locked: &Arc<AtomicBool>,
631 closed: &Arc<AtomicBool>,
632 stop_detector: &mut AttachStopDetector,
633 bytes: &[u8],
634) -> std::result::Result<AttachDataPayloadOutcome, ClientError>
635where
636 Output: Write,
637{
638 let observation = stop_detector.observe(bytes);
639 let attach_done = observation.attach_done();
640 if locked.load(Ordering::SeqCst) {
641 return Ok(AttachDataPayloadOutcome::Continue);
642 }
643 output.write_all(bytes).map_err(ClientError::Io)?;
644 output.flush().map_err(ClientError::Io)?;
645 if attach_done {
646 closed.store(true, Ordering::SeqCst);
647 return Ok(AttachDataPayloadOutcome::Stop);
648 }
649 Ok(AttachDataPayloadOutcome::Continue)
650}
651
652fn pending_render_expired(started_at: Option<Instant>) -> bool {
653 started_at.is_some_and(|started_at| started_at.elapsed() >= RENDER_MAX_PENDING)
654}
655
656fn sleep_until_pending_render_deadline(started_at: Option<Instant>) {
657 let Some(started_at) = started_at else {
658 return;
659 };
660 let Some(remaining) = RENDER_MAX_PENDING.checked_sub(started_at.elapsed()) else {
661 return;
662 };
663 if !remaining.is_zero() {
664 thread::sleep(remaining);
665 }
666}
667
668fn flush_pending_render_state<Output>(
669 output: &mut Output,
670 pending_render: &mut Option<Vec<u8>>,
671 pending_render_started_at: &mut Option<Instant>,
672) -> std::result::Result<bool, ClientError>
673where
674 Output: Write,
675{
676 let flushed = pending_render.is_some();
677 flush_pending_render(output, pending_render)?;
678 *pending_render_started_at = None;
679 Ok(flushed)
680}
681
682fn wait_for_output_thread(
683 output_thread: thread::JoinHandle<std::result::Result<(), ClientError>>,
684 raw_terminal: Option<&RawTerminal>,
685 lock_stream: &mut UnixStream,
686 locked: &Arc<AtomicBool>,
687 event_rx: mpsc::Receiver<ClientAttachEvent>,
688) -> std::result::Result<std::result::Result<(), ClientError>, ClientError> {
689 while let Ok(ClientAttachEvent::Action(action)) = event_rx.recv() {
690 handle_attach_action(raw_terminal, lock_stream, locked, action)?;
691 }
692
693 while let Ok(event) = event_rx.try_recv() {
694 match event {
695 ClientAttachEvent::Action(action) => {
696 handle_attach_action(raw_terminal, lock_stream, locked, action)?;
697 }
698 ClientAttachEvent::OutputDone => {}
699 }
700 }
701
702 join_attach_thread(output_thread)
703}
704
705fn send_attach_action(
706 event_tx: &mpsc::Sender<ClientAttachEvent>,
707 action: ClientAttachAction,
708) -> std::result::Result<(), ClientError> {
709 event_tx
710 .send(ClientAttachEvent::Action(action))
711 .map_err(|_| ClientError::Io(io::Error::other("attach event receiver closed")))
712}
713
714fn handle_attach_action(
715 raw_terminal: Option<&RawTerminal>,
716 lock_stream: &mut UnixStream,
717 locked: &Arc<AtomicBool>,
718 action: ClientAttachAction,
719) -> std::result::Result<(), ClientError> {
720 match action {
721 ClientAttachAction::Lock(command) => {
722 let Some(raw_terminal) = raw_terminal else {
723 locked.store(false, Ordering::SeqCst);
724 return Err(ClientError::Protocol(RmuxError::Decode(
725 "received unexpected lock request without a managed terminal".to_owned(),
726 )));
727 };
728 raw_terminal
729 .run_lock_command(&command)
730 .map_err(ClientError::from)?;
731 write_attach_message(lock_stream, AttachMessage::Unlock)?;
732 locked.store(false, Ordering::SeqCst);
733 Ok(())
734 }
735 ClientAttachAction::Suspend => {
736 let Some(raw_terminal) = raw_terminal else {
737 locked.store(false, Ordering::SeqCst);
738 return Err(ClientError::Protocol(RmuxError::Decode(
739 "received unexpected suspend request without a managed terminal".to_owned(),
740 )));
741 };
742 raw_terminal.suspend_self().map_err(ClientError::from)?;
743 write_attach_message(lock_stream, AttachMessage::Unlock)?;
744 locked.store(false, Ordering::SeqCst);
745 Ok(())
746 }
747 ClientAttachAction::DetachKill => {
748 if let Some(raw_terminal) = raw_terminal {
749 raw_terminal.restore().map_err(ClientError::from)?;
750 }
751 kill_process(current_process_pid().map_err(ClientError::Io)?, Signal::HUP)
752 .map_err(|error| ClientError::Io(error.into()))?;
753 Ok(())
754 }
755 ClientAttachAction::DetachExec(command) => {
756 let Some(raw_terminal) = raw_terminal else {
757 return Err(ClientError::Protocol(RmuxError::Decode(
758 "received unexpected detach exec request without a managed terminal".to_owned(),
759 )));
760 };
761 raw_terminal
762 .run_detach_exec_command(&command)
763 .map_err(ClientError::from)
764 }
765 }
766}
767
768fn drain_resize_events(
769 stream: &mut UnixStream,
770 resize_events: &mpsc::Receiver<TerminalGeometry>,
771 resize_geometry_enabled: bool,
772) -> std::result::Result<(), ClientError> {
773 while let Ok(geometry) = resize_events.try_recv() {
774 let message = if resize_geometry_enabled && geometry.pixels.is_some() {
775 AttachMessage::ResizeGeometry(geometry)
776 } else {
777 AttachMessage::Resize(geometry.size)
778 };
779 write_attach_message(stream, message)?;
780 }
781
782 Ok(())
783}
784
785fn write_attach_message(
786 stream: &mut UnixStream,
787 message: AttachMessage,
788) -> std::result::Result<(), ClientError> {
789 let frame = encode_attach_message(&message).map_err(ClientError::from)?;
790 stream.write_all(&frame).map_err(ClientError::Io)
791}
792
793fn write_attach_data(
794 stream: &mut UnixStream,
795 bytes: &[u8],
796) -> std::result::Result<(), ClientError> {
797 if bytes.len() <= STACK_ATTACH_DATA_PAYLOAD {
798 let mut frame = [0_u8; STACK_ATTACH_DATA_PAYLOAD + ATTACH_DATA_HEADER_LEN];
799 let len = encode_attach_data_into_slice(bytes, &mut frame).map_err(ClientError::from)?;
800 return stream.write_all(&frame[..len]).map_err(ClientError::Io);
801 }
802
803 let frame = encode_attach_data(bytes).map_err(ClientError::from)?;
804 stream.write_all(&frame).map_err(ClientError::Io)
805}
806
807fn join_attach_thread(
808 thread: thread::JoinHandle<std::result::Result<(), ClientError>>,
809) -> std::result::Result<std::result::Result<(), ClientError>, ClientError> {
810 thread
811 .join()
812 .map_err(|_| ClientError::Io(io::Error::other("attach thread panicked")))
813}
814
815fn shutdown_attach_writes(stream: &UnixStream) -> std::result::Result<(), ClientError> {
816 match stream.shutdown(Shutdown::Write) {
817 Ok(()) => Ok(()),
818 Err(error) if error.kind() == io::ErrorKind::NotConnected => Ok(()),
819 Err(error) => Err(ClientError::Io(error)),
820 }
821}
822
823#[derive(Debug)]
824enum ClientAttachAction {
825 Lock(String),
826 Suspend,
827 DetachKill,
828 DetachExec(String),
829}
830
831#[derive(Debug)]
832enum ClientAttachEvent {
833 Action(ClientAttachAction),
834 OutputDone,
835}
836
837#[cfg(test)]
838mod tests;