Skip to main content

rmux_client/
attach.rs

1//! Raw terminal lifecycle and attach-stream helpers for attach-mode clients.
2
3use 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
53/// Runs the attach loop using the process stdin/stdout streams.
54pub fn attach_terminal(stream: UnixStream) -> std::result::Result<(), ClientError> {
55    attach_terminal_with_initial_bytes(stream, Vec::new())
56}
57
58/// Runs the attach loop using process stdin/stdout and pre-read stream bytes.
59pub 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
66/// Runs the attach loop and sends resize events with pixel geometry.
67///
68/// Call this only after the daemon advertises the
69/// `stream.attach.resize_geometry` capability. Older daemons do not understand
70/// that attach-stream frame and would close the stream on decode.
71pub 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
102/// Runs the attach loop with an explicit terminal file descriptor.
103///
104/// The `terminal` handle is used for raw-mode lifecycle and resize discovery,
105/// while `input` and `output` carry the byte stream.
106pub 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    // This helper runs while the caller's `RawTerminal` guard is still alive,
182    // which keeps termios restoration as the last drop on every return path.
183    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
214/// Drives raw attach-stream byte forwarding over an upgraded Unix socket.
215pub 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;