From 789633a56b00f8f58e8e348bacfa1f1701979682 Mon Sep 17 00:00:00 2001 From: Levi Neuwirth Date: Sun, 17 May 2026 21:37:07 -0400 Subject: [PATCH] support fish over daemon --- src/attach.rs | 231 ++++++++++++++++++++++++++++++------------- src/daemon_attach.rs | 72 ++------------ 2 files changed, 168 insertions(+), 135 deletions(-) diff --git a/src/attach.rs b/src/attach.rs index 212aa19..f81f11b 100644 --- a/src/attach.rs +++ b/src/attach.rs @@ -1033,12 +1033,12 @@ fn attach_debug(msg: impl AsRef) { } fn hex_preview(bytes: &[u8]) -> String { + use std::fmt::Write as _; let mut out = String::new(); for (idx, byte) in bytes.iter().take(16).enumerate() { if idx > 0 { out.push(' '); } - use std::fmt::Write as _; let _ = write!(&mut out, "{byte:02x}"); } if bytes.len() > 16 { @@ -1051,14 +1051,18 @@ struct DebugReader { label: &'static str, inner: R, total: u64, + logged_chunks: u8, + max_logged_chunks: u8, } impl DebugReader { - fn new(label: &'static str, inner: R) -> Self { + fn new(label: &'static str, inner: R, max_logged_chunks: u8) -> Self { Self { label, inner, total: 0, + logged_chunks: 0, + max_logged_chunks, } } } @@ -1066,11 +1070,14 @@ impl DebugReader { impl Read for DebugReader { fn read(&mut self, buf: &mut [u8]) -> io::Result { let n = self.inner.read(buf)?; + if n > 0 { + self.total += n as u64; + } if attach_debug_enabled() { if n == 0 { attach_debug(format!("{}: EOF after {} bytes", self.label, self.total)); - } else { - self.total += n as u64; + } else if self.logged_chunks < self.max_logged_chunks { + self.logged_chunks = self.logged_chunks.saturating_add(1); attach_debug(format!( "{}: read chunk {} bytes (total {}), first bytes [{}]", self.label, @@ -1078,6 +1085,12 @@ impl Read for DebugReader { self.total, hex_preview(&buf[..n]) )); + if self.logged_chunks == self.max_logged_chunks { + attach_debug(format!( + "{}: suppressing further chunk logs after {} chunks", + self.label, self.max_logged_chunks + )); + } } } Ok(n) @@ -1090,8 +1103,11 @@ impl Read for DebugReader { /// directly by [`run_attach_ssh`] and by unit tests that want to /// assert argument shape without spawning a real `ssh`. /// -/// Argument order: `[user-flag] host pmacs --daemon-attach -/// [--socket NAME]`. The user flag is `-l USER` if `user` is set; +/// Argument order: `-T [user-flag] host exec [env VAR=val] pmacs +/// --daemon-attach [--socket NAME]`. The leading `exec` (F8b fix) +/// makes the remote login shell exec-replace itself with the bridge +/// so no shell parent lingers holding the SSH channel fds. The user +/// flag is `-l USER` if `user` is set; /// SSH's own `~/.ssh/config` is consulted by the binary, so we /// don't try to second-guess host aliases here. /// @@ -1116,6 +1132,22 @@ pub(crate) fn build_ssh_command(target: &AttachTarget) -> Option { cmd.arg("-l").arg(u); } cmd.arg(host); + // F8b fix: `exec`-replace the remote login shell with the bridge. + // ssh runs the remote command via the user's login shell + // (`$SHELL -c "joined args"`). If that shell lingers as a parent + // process while the long-lived bridge runs, the SSH session's + // stdout/stderr fds have an extra holder, and OpenSSH keeps the + // channel's output buffered until *every* holder exits — so the + // bridge's `Hello` (and all traffic) is withheld until the bridge + // dies, i.e. forever during a session: the F8b hang. Confound-free + // probing (`f8b-shell-probe.sh`, sh==fish==login-shell all + // withheld; an `exec`'d single long-lived process streamed at the + // SSH floor while alive) pinned this exactly. `exec` is a builtin + // honored by every standard login shell ssh dispatches through; + // the shell has nothing to do after the bridge anyway, so + // exec-replacing it is unconditionally correct and removes the + // lingering fd holder. + cmd.arg("exec"); if protocol_stderr_enabled() { cmd.arg("env").arg(format!("{PMACS_ATTACH_PROTOCOL_FD}=2")); } else if attach_debug_enabled() { @@ -1495,10 +1527,112 @@ fn attempt_session( } } +/// Protocol byte source plus the SSH-stderr tee's join handle and +/// captured tail. Returned by [`open_protocol_channel`]; named so the +/// handshake's binding site isn't a four-line inline tuple type. +type ProtocolChannel = ( + Box, + thread::JoinHandle<()>, + Arc>>, +); + +/// Pick the SSH channel the protocol rides on. +/// +/// Default: protocol on stdout, SSH stderr tee'd to our stderr only +/// on the first attempt (`tee_to_stderr`) — on reconnect raw mode is +/// active and live tee'd bytes would corrupt the cell grid, so the +/// tail is still captured for the give-up message but not echoed. +/// F8b workaround (`protocol_over_stderr`): protocol on SSH stderr; +/// remote stderr diagnostics are disabled by the caller, and there +/// is no tee (a no-op join handle / empty tail keep the shape). +fn open_protocol_channel( + child_stdout: std::process::ChildStdout, + child_stderr: std::process::ChildStderr, + protocol_over_stderr: bool, + tee_to_stderr: bool, +) -> ProtocolChannel { + if protocol_over_stderr { + attach_debug("using SSH stderr as protocol stream; remote stderr diagnostics disabled"); + ( + Box::new(DebugReader::new("ssh stderr(protocol)", child_stderr, 4)), + thread::spawn(|| {}), + Arc::new(Mutex::new(VecDeque::new())), + ) + } else { + let (stderr_handle, stderr_tail) = spawn_stderr_tee(child_stderr, tee_to_stderr); + ( + Box::new(DebugReader::new("ssh stdout", child_stdout, 4)), + stderr_handle, + stderr_tail, + ) + } +} + +/// Spawn the F8b watchdog: with attach-debug on, warn if the `Hello` +/// hasn't arrived within 5s (the F8b symptom — remote bridge +/// connected but no daemon bytes reached local stdout). Returns the +/// flag the caller sets once `Hello` is read or fails, silencing it. +fn spawn_hello_watchdog() -> Arc { + let done = Arc::new(AtomicBool::new(false)); + if attach_debug_enabled() { + let done = Arc::clone(&done); + thread::spawn(move || { + thread::sleep(Duration::from_secs(5)); + if !done.load(Ordering::SeqCst) { + eprintln!( + "pmacs attach debug: still waiting for Hello after 5s; \ + remote bridge connected but no daemon bytes reached local stdout" + ); + } + }); + } + done +} + +/// Run the attach pump to completion, then reap the child and drain +/// the stderr tee with the Frontend still alive (none of these steps +/// emit terminal output, so raw mode being on is harmless). The +/// Frontend is NOT dropped — the outer scope owns its lifetime so it +/// persists into the next reconnect attempt. +fn run_pump_and_reap( + mut child: Child, + child_stdin: std::process::ChildStdin, + protocol_reader: Box, + stderr_handle: thread::JoinHandle<()>, + stderr_tail: &Arc>>, + frontend: &mut Frontend, + assigned_frontend_id: FrontendId, +) -> Result<(), AttachError> { + let pid = child.id(); + let io = AttachIo { + reader: protocol_reader, + writer: Box::new(child_stdin), + kick: ssh_kick(pid), + }; + let pump_result = run_attach_pair(io, frontend, assigned_frontend_id); + attach_debug(format!( + "attach pump exited: {:?}", + pump_result.as_ref().err() + )); + let exit_status = child.wait().ok(); + let _ = stderr_handle.join(); + let stderr_text = drain_stderr_tail(stderr_tail); + classify_ssh_exit(pump_result, exit_status, stderr_text) +} + /// Body of [`attempt_session`] expressed in raw `AttachError` so the /// caller can fold the result through [`classify_for_reconnect`]. /// Splitting this out keeps the IO path linear and tests-friendly: /// `attempt_session` is just classify + dispatch. +// The protocol-channel, Hello-watchdog, and pump-and-reap phases are +// already factored out (`open_protocol_channel`, `spawn_hello_watchdog`, +// `run_pump_and_reap`). What remains is a linear Hello/AttachRequest +// handshake whose several early-return error paths each need ownership +// of `child` + the stderr handles to tear down via +// `handshake_error_with_child`; extracting it further would fragment +// that error-ownership flow rather than clarify it. Same precedent as +// the `#[allow(clippy::too_many_lines)]` on `run_attach_pair` above. +#[allow(clippy::too_many_lines)] fn run_one_session( frontend_slot: &mut Option, target: &AttachTarget, @@ -1509,7 +1643,7 @@ fn run_one_session( .stderr(Stdio::piped()); let bin_for_error = ssh_binary(); - attach_debug(format!("spawning ssh command: {:?}", cmd)); + attach_debug(format!("spawning ssh command: {cmd:?}")); let mut child = cmd.spawn().map_err(|source| AttachError::SshSpawnFailed { command: bin_for_error, source, @@ -1530,31 +1664,12 @@ fn run_one_session( .expect("Stdio::piped on stderr guarantees a handle"); let protocol_over_stderr = protocol_stderr_enabled(); - let (mut protocol_reader, stderr_handle, stderr_tail): ( - Box, - thread::JoinHandle<()>, - Arc>>, - ) = if protocol_over_stderr { - attach_debug("using SSH stderr as protocol stream; remote stderr diagnostics disabled"); - ( - Box::new(DebugReader::new("ssh stderr(protocol)", child_stderr)), - thread::spawn(|| {}), - Arc::new(Mutex::new(VecDeque::new())), - ) - } else { - // Tee SSH stderr to our stderr only on the first attempt, - // where raw mode hasn't engaged yet. On reconnect attempts - // (slot is `Some`), raw mode is active and live tee'd bytes - // would corrupt the cell grid; the tail is still captured - // for the give-up message via `AttachError::SshChildExited`. - let tee_to_stderr = frontend_slot.is_none(); - let (stderr_handle, stderr_tail) = spawn_stderr_tee(child_stderr, tee_to_stderr); - ( - Box::new(DebugReader::new("ssh stdout", child_stdout)), - stderr_handle, - stderr_tail, - ) - }; + let (mut protocol_reader, stderr_handle, stderr_tail) = open_protocol_channel( + child_stdout, + child_stderr, + protocol_over_stderr, + frontend_slot.is_none(), + ); // Hello / AttachRequest handshake. On the FIRST attempt // (`frontend_slot` is `None`) raw mode is not engaged yet, so @@ -1563,19 +1678,7 @@ fn run_one_session( // overlay communicates state; handshake errors here are bubbled // up so the M5.8d loop can decide whether to retry. attach_debug("waiting for Hello from remote daemon bridge"); - let hello_wait_done = Arc::new(AtomicBool::new(false)); - if attach_debug_enabled() { - let done = Arc::clone(&hello_wait_done); - thread::spawn(move || { - thread::sleep(Duration::from_secs(5)); - if !done.load(Ordering::SeqCst) { - eprintln!( - "pmacs attach debug: still waiting for Hello after 5s; \ - remote bridge connected but no daemon bytes reached local stdout" - ); - } - }); - } + let hello_wait_done = spawn_hello_watchdog(); let hello: Hello = match read_message(&mut protocol_reader) { Ok(h) => h, Err(e) => { @@ -1668,29 +1771,15 @@ fn run_one_session( .as_mut() .expect("frontend_slot was just initialized or already Some"); - let pid = child.id(); - let io = AttachIo { - reader: protocol_reader, - writer: Box::new(child_stdin), - kick: ssh_kick(pid), - }; - - let pump_result = run_attach_pair(io, frontend, hello.assigned_frontend_id); - attach_debug(format!( - "attach pump exited: {:?}", - pump_result.as_ref().err() - )); - - // Reap the child and drain the stderr tee with the Frontend - // still alive. None of these steps emit terminal output, so raw - // mode being on is harmless. Critically, we do NOT drop the - // Frontend here — the outer scope owns its lifetime so it - // persists into the next reconnect attempt. - let exit_status = child.wait().ok(); - let _ = stderr_handle.join(); - let stderr_text = drain_stderr_tail(&stderr_tail); - - classify_ssh_exit(pump_result, exit_status, stderr_text) + run_pump_and_reap( + child, + child_stdin, + protocol_reader, + stderr_handle, + &stderr_tail, + frontend, + hello.assigned_frontend_id, + ) } /// Helper for the early-exit paths in [`run_attach_ssh`] that want @@ -2181,8 +2270,8 @@ mod tests { let arg_strs: Vec<&str> = args.iter().filter_map(|s| s.to_str()).collect(); assert_eq!( arg_strs, - vec!["-T", "mac-studio", "pmacs", "--daemon-attach"], - "bare host should produce: -T pmacs --daemon-attach", + vec!["-T", "mac-studio", "exec", "pmacs", "--daemon-attach"], + "bare host should produce: -T exec pmacs --daemon-attach", ); } @@ -2202,6 +2291,7 @@ mod tests { "-l", "alice", "workstation", + "exec", "pmacs", "--daemon-attach", ], @@ -2224,6 +2314,7 @@ mod tests { "-l", "bob", "workstation", + "exec", "pmacs", "--daemon-attach", "--socket", diff --git a/src/daemon_attach.rs b/src/daemon_attach.rs index 0354f61..d1cb0d5 100644 --- a/src/daemon_attach.rs +++ b/src/daemon_attach.rs @@ -439,41 +439,6 @@ impl Write for RawProtocolOutput { } } -/// F8 diagnostic byte-capture. Active only when the marker directory -/// `/tmp/pmacs-bridge-capture` exists on the bridge's host (or -/// `PMACS_BRIDGE_CAPTURE_DIR` points at a directory). When active, -/// each direction's bytes are mirrored verbatim to a file there so -/// the exact wire image the daemon sees can be decoded offline. The -/// mirror is best-effort and never alters the forwarded stream; -/// when inactive this is a single `is_dir` syscall and nothing else. -fn capture_sink(name: &str) -> Option { - let dir = std::env::var_os("PMACS_BRIDGE_CAPTURE_DIR") - .map(PathBuf::from) - .filter(|p| p.is_dir()) - .or_else(|| { - let marker = PathBuf::from("/tmp/pmacs-bridge-capture"); - marker.is_dir().then_some(marker) - })?; - std::fs::File::create(dir.join(name)).ok() -} - -/// Reader that mirrors every byte it yields to a side file -/// (best-effort) before returning it to the caller. Diagnostic only; -/// the byte stream the bridge forwards is unchanged. -struct TeeReader { - inner: R, - sink: std::fs::File, -} - -impl Read for TeeReader { - fn read(&mut self, buf: &mut [u8]) -> std::io::Result { - let n = self.inner.read(buf)?; - let _ = self.sink.write_all(&buf[..n]); - let _ = self.sink.flush(); - Ok(n) - } -} - // Owned `PathBuf` matches the public entry-point signature; clippy's // pedantic pass-by-value is wrong for this seam. #[allow(clippy::needless_pass_by_value)] @@ -513,20 +478,8 @@ where let mut local_input = local_input; let mut socket_writer = socket_for_writer; bridge_debug("stdin->daemon copy loop started"); - // Diagnostic: mirror exactly the bytes we forward toward the - // daemon (the client's AttachRequest et seq). - let result = match capture_sink("to_daemon.bin") { - Some(sink) => { - let mut tee = TeeReader { - inner: local_input, - sink, - }; - copy_with_flush_named(&mut tee, &mut socket_writer, Some("stdin->daemon")) - } - None => { - copy_with_flush_named(&mut local_input, &mut socket_writer, Some("stdin->daemon")) - } - }; + let result = + copy_with_flush_named(&mut local_input, &mut socket_writer, Some("stdin->daemon")); match &result { Ok(n) => bridge_debug(format!("stdin->daemon copy loop ended after {n} bytes")), Err(e) => bridge_debug(format!("stdin->daemon copy loop failed: {e}")), @@ -546,23 +499,12 @@ where // swallowed — once we've decided to tear down, the caller cares // only about whether the connect succeeded, not about the exact // shape of the disconnect. - // Diagnostic: mirror exactly the bytes the daemon sent us (its - // Hello et seq), as the client receives them. bridge_debug("daemon->stdout copy loop started"); - let daemon_to_stdout = match capture_sink("from_daemon.bin") { - Some(sink) => { - let mut tee = TeeReader { - inner: socket_reader, - sink, - }; - copy_with_flush_named(&mut tee, &mut local_output, Some("daemon->stdout")) - } - None => copy_with_flush_named( - &mut socket_reader, - &mut local_output, - Some("daemon->stdout"), - ), - }; + let daemon_to_stdout = copy_with_flush_named( + &mut socket_reader, + &mut local_output, + Some("daemon->stdout"), + ); match &daemon_to_stdout { Ok(n) => bridge_debug(format!("daemon->stdout copy loop ended after {n} bytes")), Err(e) => bridge_debug(format!("daemon->stdout copy loop failed: {e}")),