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
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;