diff --git a/Packages/macOS/CmuxSettings/Sources/CmuxSettings/Keys/BetaFeaturesCatalogSection.swift b/Packages/macOS/CmuxSettings/Sources/CmuxSettings/Keys/BetaFeaturesCatalogSection.swift index 0aa23f9734ff..d3caa04ceeab 100644 --- a/Packages/macOS/CmuxSettings/Sources/CmuxSettings/Keys/BetaFeaturesCatalogSection.swift +++ b/Packages/macOS/CmuxSettings/Sources/CmuxSettings/Keys/BetaFeaturesCatalogSection.swift @@ -102,6 +102,19 @@ public struct BetaFeaturesCatalogSection: SettingCatalogSection { userDefaultsKey: "terminal.beta.tuiBackend.enabled" ) + /// Manual-IO variant of the cmux-tui terminal backend: instead of + /// running `cmux-tui attach` as the surface command (exec bridge), the + /// surface runs in Ghostty's manual-mirror IO mode and an app-side pump + /// relays bytes to/from the daemon through `attach --pipe-io`. The + /// surface parses the daemon terminal's own byte stream, so input + /// encodings always match the daemon's mode state. Requires + /// `tuiTerminalBackend`; defaults off (exec bridge). + public let tuiTerminalBackendManualIO = DefaultsKey( + id: "terminal.beta.tuiBackend.manualIO", + defaultValue: false, + userDefaultsKey: "terminal.beta.tuiBackend.manualIO" + ) + /// Path to the cmux-tui binary used by the `tuiTerminalBackend` spike. /// Dev-only setting with a spike-only default pointing at a locally /// installed npm binary; there is no bundled artifact yet (that is build diff --git a/Packages/macOS/CmuxSettingsUI/Sources/CmuxSettingsUI/Sections/BetaFeaturesSection.swift b/Packages/macOS/CmuxSettingsUI/Sources/CmuxSettingsUI/Sections/BetaFeaturesSection.swift index f6ea913053d1..cc724a263e8d 100644 --- a/Packages/macOS/CmuxSettingsUI/Sources/CmuxSettingsUI/Sections/BetaFeaturesSection.swift +++ b/Packages/macOS/CmuxSettingsUI/Sources/CmuxSettingsUI/Sections/BetaFeaturesSection.swift @@ -14,6 +14,7 @@ public struct BetaFeaturesSection: View { @State private var customSidebars: DefaultsValueModel @State private var remoteTmux: DefaultsValueModel @State private var tuiTerminalBackend: DefaultsValueModel + @State private var tuiTerminalBackendManualIO: DefaultsValueModel @State private var workspaceTodoControls: DefaultsValueModel @State private var workspaceTodosChecklistStyle: DefaultsValueModel @@ -25,6 +26,7 @@ public struct BetaFeaturesSection: View { _customSidebars = State(initialValue: DefaultsValueModel(store: defaultsStore, key: catalog.betaFeatures.customSidebars)) _remoteTmux = State(initialValue: DefaultsValueModel(store: defaultsStore, key: catalog.betaFeatures.remoteTmux)) _tuiTerminalBackend = State(initialValue: DefaultsValueModel(store: defaultsStore, key: catalog.betaFeatures.tuiTerminalBackend)) + _tuiTerminalBackendManualIO = State(initialValue: DefaultsValueModel(store: defaultsStore, key: catalog.betaFeatures.tuiTerminalBackendManualIO)) _workspaceTodoControls = State(initialValue: DefaultsValueModel(store: defaultsStore, key: catalog.betaFeatures.workspaceTodoControls)) _workspaceTodosChecklistStyle = State(initialValue: DefaultsValueModel(store: defaultsStore, key: catalog.betaFeatures.workspaceTodosChecklistStyle)) } @@ -51,6 +53,8 @@ public struct BetaFeaturesSection: View { SettingsCardDivider() tuiTerminalBackendRow SettingsCardDivider() + tuiTerminalBackendManualIORow + SettingsCardDivider() workspaceTodoControlsRow SettingsCardDivider() workspaceTodosChecklistStyleRow @@ -68,6 +72,7 @@ public struct BetaFeaturesSection: View { customSidebars, remoteTmux, tuiTerminalBackend, + tuiTerminalBackendManualIO, workspaceTodoControls, workspaceTodosChecklistStyle, ] @@ -234,6 +239,23 @@ public struct BetaFeaturesSection: View { } } + @ViewBuilder + private var tuiTerminalBackendManualIORow: some View { + SettingsCardRow( + configurationReview: .settingsOnly, + searchAnchorID: "setting:betaFeatures:tuiTerminalBackendManualIO", + String(localized: "settings.betaFeatures.tuiTerminalBackendManualIO", defaultValue: "cmux-tui Manual-IO Data Path"), + subtitle: tuiTerminalBackendManualIO.current + ? String(localized: "settings.betaFeatures.tuiTerminalBackendManualIO.subtitleOn", defaultValue: "Daemon-backed terminals feed bytes straight into the Ghostty surface (no attach subprocess in the pane). Requires the cmux-tui Terminal Backend toggle.") + : String(localized: "settings.betaFeatures.tuiTerminalBackendManualIO.subtitleOff", defaultValue: "Daemon-backed terminals keep running the attach client as the surface command.") + ) { + Toggle("", isOn: Binding(get: { tuiTerminalBackendManualIO.current }, set: { tuiTerminalBackendManualIO.set($0) })) + .labelsHidden() + .controlSize(.small) + .accessibilityIdentifier("SettingsBetaTuiTerminalBackendManualIOToggle") + } + } + } /// Small warning callout with a yellow triangle, used at the top of diff --git a/Packages/macOS/CmuxTerminal/Sources/CmuxTerminal/Surface/TerminalSurface+Input.swift b/Packages/macOS/CmuxTerminal/Sources/CmuxTerminal/Surface/TerminalSurface+Input.swift index 7b36d518c95a..787b29cf2d36 100644 --- a/Packages/macOS/CmuxTerminal/Sources/CmuxTerminal/Surface/TerminalSurface+Input.swift +++ b/Packages/macOS/CmuxTerminal/Sources/CmuxTerminal/Surface/TerminalSurface+Input.swift @@ -660,7 +660,8 @@ extension TerminalSurface { manualIONoReflow = value } - /// Enqueues remote tmux `%output` for the terminal parser. + /// Enqueues remote output (tmux `%output`, tui pipe-io bytes) for the + /// terminal parser. /// /// The native parser runs on the surface generation's FIFO output lane and /// this method returns without waiting for Ghostty's renderer-state mutex. diff --git a/Resources/Localizable.xcstrings b/Resources/Localizable.xcstrings index 22b411375c1e..65219ac212e6 100644 --- a/Resources/Localizable.xcstrings +++ b/Resources/Localizable.xcstrings @@ -174903,6 +174903,27 @@ } } }, + "settings.betaFeatures.tuiTerminalBackendManualIO.subtitleOn": { + "extractionState": "manual", + "localizations": { + "en": { "stringUnit": { "state": "translated", "value": "Daemon-backed terminals feed bytes straight into the Ghostty surface (no attach subprocess in the pane). Requires the cmux-tui Terminal Backend toggle." } }, + "ja": { "stringUnit": { "state": "translated", "value": "デーモン管理のターミナルがGhosttyサーフェスへ直接バイトを送ります(ペイン内にattachサブプロセスなし)。cmux-tuiターミナルバックエンドの有効化が必要です。" } } + } + }, + "settings.betaFeatures.tuiTerminalBackendManualIO.subtitleOff": { + "extractionState": "manual", + "localizations": { + "en": { "stringUnit": { "state": "translated", "value": "Daemon-backed terminals keep running the attach client as the surface command." } }, + "ja": { "stringUnit": { "state": "translated", "value": "デーモン管理のターミナルは引き続きattachクライアントをサーフェスコマンドとして実行します。" } } + } + }, + "settings.betaFeatures.tuiTerminalBackendManualIO": { + "extractionState": "manual", + "localizations": { + "en": { "stringUnit": { "state": "translated", "value": "cmux-tui Manual-IO Data Path" } }, + "ja": { "stringUnit": { "state": "translated", "value": "cmux-tui マニュアルIOデータパス" } } + } + }, "settings.betaFeatures.warning": { "extractionState": "manual", "localizations": { @@ -254478,6 +254499,48 @@ } } }, + "tui.overlay.reconnecting.title": { + "extractionState": "manual", + "localizations": { + "en": { "stringUnit": { "state": "translated", "value": "Reconnecting to the terminal session" } }, + "ja": { "stringUnit": { "state": "translated", "value": "ターミナルセッションに再接続中" } } + } + }, + "tui.overlay.reconnecting.detail": { + "extractionState": "manual", + "localizations": { + "en": { "stringUnit": { "state": "translated", "value": "The session daemon is unreachable (attempt %lld). The terminal shows its last state; input is paused." } }, + "ja": { "stringUnit": { "state": "translated", "value": "セッションデーモンに接続できません(試行 %lld 回目)。ターミナルは最後の状態を表示しており、入力は一時停止中です。" } } + } + }, + "tui.overlay.failed.title": { + "extractionState": "manual", + "localizations": { + "en": { "stringUnit": { "state": "translated", "value": "Terminal backend unavailable" } }, + "ja": { "stringUnit": { "state": "translated", "value": "ターミナルバックエンドを利用できません" } } + } + }, + "tui.overlay.failed.detail": { + "extractionState": "manual", + "localizations": { + "en": { "stringUnit": { "state": "translated", "value": "The relay for this daemon-backed terminal keeps failing. Check the cmux-tui binary path in Settings, then retry." } }, + "ja": { "stringUnit": { "state": "translated", "value": "デーモン管理ターミナルのリレーが繰り返し失敗しています。設定でcmux-tuiバイナリのパスを確認してから再試行してください。" } } + } + }, + "tui.overlay.ended.title": { + "extractionState": "manual", + "localizations": { + "en": { "stringUnit": { "state": "translated", "value": "Terminal session ended" } }, + "ja": { "stringUnit": { "state": "translated", "value": "ターミナルセッションが終了しました" } } + } + }, + "tui.overlay.ended.detail": { + "extractionState": "manual", + "localizations": { + "en": { "stringUnit": { "state": "translated", "value": "The daemon-backed terminal has ended. Close the tab, or keep it to read the final screen." } }, + "ja": { "stringUnit": { "state": "translated", "value": "デーモン管理のターミナルが終了しました。タブを閉じるか、最後の画面を確認するために残せます。" } } + } + }, "uiTest.bonsplit.action.title": { "localizations": { "en": { diff --git a/Sources/TuiManualIOPump.swift b/Sources/TuiManualIOPump.swift new file mode 100644 index 000000000000..2edd933a00ec --- /dev/null +++ b/Sources/TuiManualIOPump.swift @@ -0,0 +1,634 @@ +import CmuxTerminal +import Foundation +#if DEBUG +import CMUXDebugLog +#endif + +/// Pure decision logic for the manual-IO tui pump: relay exit +/// classification, reconnect pacing, stdin line encoding, and the per-pane +/// overlay presentation. Everything here is deterministic and unit-testable; +/// process and pipe I/O live in `TuiManualIOPump`. +enum TuiManualIOPumpPolicy { + /// How one relay process ended, from its exit status and the final JSON + /// line it printed to stderr (`{"exit":{"reason":...}}`). + enum RelayExit: Equatable { + /// The daemon terminal ended; respawning is wrong. + case terminalEnded + /// The daemon connection was lost; respawning reattaches and + /// resyncs from a fresh replay. + case daemonLost + /// The pump closed the relay's stdin itself (teardown/respawn). + case parentClosed + /// Anything else (spawn failure, protocol error, unknown status). + case failure + } + + static func relayExit(status: Int32, stderrText: String?) -> RelayExit { + let reason = stderrText? + .split(separator: "\n") + .reversed() + .compactMap { line -> String? in + guard let data = line.trimmingCharacters(in: .whitespaces).data(using: .utf8), + let object = try? JSONSerialization.jsonObject(with: data) as? [String: Any], + let exit = object["exit"] as? [String: Any] + else { return nil } + return exit["reason"] as? String + } + .first + switch (status, reason) { + case (0, "terminal-ended"): return .terminalEnded + case (0, "parent-closed"): return .parentClosed + // Exit 2 is daemon-lost ONLY with the relay's own reason line: a + // binary without --pipe-io support also exits 2 (usage error), and + // that must not read as an endlessly-retryable daemon outage. + case (2, "daemon-lost"): return .daemonLost + case (0, _): return .terminalEnded + default: return .failure + } + } + + enum NextAction: Equatable { + case end + case retry + case ignore + } + + static func nextAction(after exit: RelayExit) -> NextAction { + switch exit { + case .terminalEnded: return .end + case .daemonLost, .failure: return .retry + // The pump closed stdin itself; its own state machine already + // decided what comes next. + case .parentClosed: return .ignore + } + } + + /// Unexplained relay failures (no exit-reason line: bad binary, usage + /// error, crash) stop retrying after this many consecutive attempts + /// without a live interlude; explained daemon-lost exits retry forever. + static let maxConsecutiveUnexplainedFailures = 5 + + /// Reconnect backoff: fast first retries absorb a daemon restart or + /// handoff, the 30s cap keeps a long outage cheap while the pane shows + /// its last state. Attempts are 1-based. + static func retryDelay(attempt: Int) -> Duration { + let schedule: [Duration] = [ + .milliseconds(500), .seconds(1), .seconds(2), .seconds(4), .seconds(8), .seconds(16), + ] + guard attempt >= 1 else { return schedule[0] } + guard attempt <= schedule.count else { return .seconds(30) } + return schedule[attempt - 1] + } + + /// One stdin line forwarding raw input bytes to the relay. + static func inputLine(bytes: Data) -> Data { + var line = Data(#"{"input":""#.utf8) + line.append(Data(bytes.base64EncodedString().utf8)) + line.append(Data(#""}"#.utf8)) + line.append(0x0A) + return line + } + + /// One stdin line driving the daemon-side viewer size. + static func resizeLine(cols: Int, rows: Int) -> Data { + Data(#"{"resize":{"cols":\#(max(1, cols)),"rows":\#(max(1, rows))}}"#.utf8 + [0x0A]) + } + + /// Relay argv (no shell, no quoting: the pump spawns the binary + /// directly, so the exec-wrapper and env(1) classes of the exec bridge + /// cannot exist here). + static func relayArguments( + sessionName: String, + terminalID: String, + cols: Int, + rows: Int + ) -> [String] { + [ + "attach", + "--session", sessionName, + "--terminal", terminalID, + "--pipe-io", + "--cols", String(max(1, cols)), + "--rows", String(max(1, rows)), + ] + } + + /// Emitted into the surface before a respawned relay's replay: the + /// replay REPLACES terminal state, and only the pump knows whether this + /// surface already rendered a previous attach. + static let resyncReset = Data([0x1B, 0x63, 0x1B, 0x5B, 0x33, 0x4A]) // ESC c, CSI 3 J + + /// Overlay presentation for one pump state, in the same visual family + /// as the cloud terminal reconnect overlay. + static func overlayPresentation( + state: TuiManualIOPump.State + ) -> CloudTerminalReconnectOverlayPolicy.Presentation? { + switch state { + case .connecting, .live: + return nil + case .reconnecting(let attempt): + return CloudTerminalReconnectOverlayPolicy.Presentation( + title: String( + localized: "tui.overlay.reconnecting.title", + defaultValue: "Reconnecting to the terminal session" + ), + detail: String( + localized: "tui.overlay.reconnecting.detail", + defaultValue: "The session daemon is unreachable (attempt \(attempt)). The terminal shows its last state; input is paused." + ), + showsProgress: true, + showsReconnectButton: true + ) + case .ended: + return CloudTerminalReconnectOverlayPolicy.Presentation( + title: String( + localized: "tui.overlay.ended.title", + defaultValue: "Terminal session ended" + ), + detail: String( + localized: "tui.overlay.ended.detail", + defaultValue: "The daemon-backed terminal has ended. Close the tab, or keep it to read the final screen." + ), + showsProgress: false, + showsReconnectButton: false + ) + case .failed: + return CloudTerminalReconnectOverlayPolicy.Presentation( + title: String( + localized: "tui.overlay.failed.title", + defaultValue: "Terminal backend unavailable" + ), + detail: String( + localized: "tui.overlay.failed.detail", + defaultValue: "The relay for this daemon-backed terminal keeps failing. Check the cmux-tui binary path in Settings, then retry." + ), + showsProgress: false, + showsReconnectButton: true + ) + } + } +} + +/// Cross-thread input channel between the Ghostty IO thread (which delivers +/// the surface's encoded input bytes) and the pump's relay stdin writer. +/// Lock-guarded because `manualInputHandler` must not touch the main actor. +final class TuiManualIOInputChannel: @unchecked Sendable { + private let lock = NSLock() + private var handle: FileHandle? + private let queue = DispatchQueue(label: "cmux.tuiManualIO.stdin", qos: .userInitiated) + + /// Swaps the live relay stdin. `nil` pauses input (dropped, never + /// queued: replaying stale input into a shell after a reconnect is + /// worse than losing keystrokes typed into a dead pane). + func setHandle(_ newHandle: FileHandle?) { + lock.lock() + handle = newHandle + lock.unlock() + } + + func send(_ line: Data) { + lock.lock() + let target = handle + lock.unlock() + guard let target else { return } + queue.async { + // A failed write means the relay just died; the pump's + // termination handler owns that transition. + try? target.write(contentsOf: line) + } + } + + /// Closes and detaches the current handle (relay stdin EOF = clean + /// detach on the relay side). + func closeHandle() { + lock.lock() + let target = handle + handle = nil + lock.unlock() + guard let target else { return } + queue.async { + try? target.close() + } + } +} + +/// Owns the `cmux-tui attach --terminal --pipe-io` relay for one +/// manual-mirror terminal panel: pumps relay stdout into the Ghostty +/// surface, forwards the surface's encoded input to relay stdin, drives +/// daemon-side sizing from applied surface resizes, and runs the reconnect +/// state machine when the relay dies. +/// +/// Data-path invariant: the surface parses the same byte stream the daemon +/// terminal emitted, so the surface's input encodings always match the +/// daemon terminal's mode state. The mode-mirroring/re-encoding bug class +/// of the exec attach bridge cannot exist on this path. +@MainActor +final class TuiManualIOPump { + enum State: Equatable { + /// First attach in flight; the pane is empty, no overlay. + case connecting + /// Relay bytes are flowing. + case live + /// The relay died recoverably; a respawn is scheduled (1-based + /// attempt count shown in the overlay). + case reconnecting(attempt: Int) + /// The daemon terminal ended; the pane keeps its final frame. + case ended + /// The relay failed repeatedly with no explanation (bad binary, + /// crash loop); automatic retries stopped, manual retry offered. + case failed + } + + private(set) var state: State = .connecting { + didSet { + guard oldValue != state else { return } + log("state \(oldValue) -> \(state)") + onStateChange?() + } + } + + /// Posted on every state change (wired to the workspace's overlay + /// resynchronization). + var onStateChange: (() -> Void)? + + let terminalID: String + private let binaryPath: String + private let sessionName: String + private let environment: [String: String] + private let inputChannel = TuiManualIOInputChannel() + /// Injected so tests can run the backoff deterministically. Production + /// is a cancellation-checking `Task.sleep`. + private let sleep: @Sendable (Duration) async throws -> Void + + private weak var surface: TerminalSurface? + private var process: Process? + private var stdoutReader: RemoteTmuxProcessOutputReader? + private var stdoutTask: Task? + private var stderrBox = TuiManualIOStderrBox() + private var retryTask: Task? + /// Process termination and stderr EOF race (termination is not EOF); + /// classification needs the relay's final stderr line, so it runs only + /// once BOTH have arrived for the current generation. + private var pendingExitStatus: Int32? + private var stderrDrainedGeneration = 0 + /// Fences callbacks from an old relay after a respawn or stop. + private var generation = 0 + private var stopped = false + private var everRenderedAttach = false + private var consecutiveUnexplainedFailures = 0 + private var lastKnownGrid: (cols: Int, rows: Int)? + private var lastSentGrid: (cols: Int, rows: Int)? + + init( + binaryPath: String, + sessionName: String, + terminalID: String, + environment: [String: String], + sleep: @escaping @Sendable (Duration) async throws -> Void = { try await Task.sleep(for: $0) } + ) { + self.binaryPath = binaryPath + self.sessionName = sessionName + self.terminalID = terminalID + self.environment = environment + self.sleep = sleep + } + + /// The surface's `manualInputHandler`: runs on Ghostty's IO thread, so + /// it only touches the lock-guarded input channel. + nonisolated func makeManualInputHandler() -> @Sendable (TerminalManualInput) -> Void { + let channel = inputChannel + return { input in + switch input { + case .bytes(let bytes): + channel.send(TuiManualIOPumpPolicy.inputLine(bytes: bytes)) + case .namedKey: + // No key-name resolver is installed on this surface, so + // Ghostty encodes every key itself; a named key here is a + // wiring bug, and dropping it is safer than guessing bytes. + break + } + } + } + + /// Binds the pump to its surface and starts the relay at the first + /// known grid size. Applied resizes are the steady-state signal, but a + /// session-restored panel can mount at exactly the runtime's initial + /// grid (no apply ever fires), so the pump also samples the grid + /// directly when the runtime is (or becomes) ready and attaches + /// eagerly; a later mount just resizes the daemon terminal. + func start(surface: TerminalSurface) { + self.surface = surface + surface.onManualSizeApplied = { [weak self] sample in + self?.handleSizingSample(cols: sample.columns, rows: sample.rows) + } + surface.onRuntimeReady = { [weak self, weak surface] in + guard let self, let surface else { return } + surface.flushPendingManualSizeReportIfAttached() + self.sampleCurrentGrid(of: surface) + } + surface.flushPendingManualSizeReportIfAttached() + sampleCurrentGrid(of: surface) + } + + private func sampleCurrentGrid(of surface: TerminalSurface) { + guard let sample = surface.rawSizingSample(), + sample.columns > 1, sample.rows > 1 else { return } + handleSizingSample(cols: sample.columns, rows: sample.rows) + } + + /// The overlay's Reconnect button: skip the remaining backoff, or leave + /// the failed state for another round of attempts. + func retryNow() { + switch state { + case .reconnecting, .failed: + guard !stopped else { return } + consecutiveUnexplainedFailures = 0 + retryTask?.cancel() + retryTask = nil + state = .reconnecting(attempt: 1) + spawnRelay() + case .connecting, .live, .ended: + break + } + } + + /// Tears the pump down (panel closed, app quitting, transfer discard). + /// The relay sees stdin EOF and detaches cleanly; the daemon terminal + /// itself is closed (or kept) by the bridge's existing close semantics. + func stop() { + stopped = true + retryTask?.cancel() + retryTask = nil + stdoutTask?.cancel() + stdoutTask = nil + stdoutReader?.close() + stdoutReader = nil + inputChannel.closeHandle() + generation += 1 + process?.terminationHandler = nil + process?.terminate() + process = nil + } + + var overlayPresentation: CloudTerminalReconnectOverlayPolicy.Presentation? { + TuiManualIOPumpPolicy.overlayPresentation(state: state) + } + + // MARK: - Sizing + + private func handleSizingSample(cols: Int, rows: Int) { + guard !stopped else { return } + lastKnownGrid = (cols, rows) + if process == nil, retryTask == nil, case .connecting = state { + spawnRelay() + return + } + if let sent = lastSentGrid, sent == (cols, rows) { return } + lastSentGrid = (cols, rows) + inputChannel.send(TuiManualIOPumpPolicy.resizeLine(cols: cols, rows: rows)) + } + + // MARK: - Relay lifecycle + + private func spawnRelay() { + guard !stopped, let surface else { return } + // Terminal states only leave through retryNow (which transitions + // back to reconnecting first) or teardown; a straggler retry task + // or sizing sample must not resurrect the relay. + if state == .ended || state == .failed { return } + guard FileManager.default.isExecutableFile(atPath: binaryPath) else { + noteUnexplainedFailureThenRetryOrFail() + return + } + generation += 1 + let spawnGeneration = generation + let grid = lastKnownGrid ?? (80, 24) + lastSentGrid = grid + + // A respawned relay replays the daemon terminal from scratch; reset + // the surface first so the replay replaces the stale frame instead + // of stacking on it. + if everRenderedAttach { + surface.processRemoteOutput(TuiManualIOPumpPolicy.resyncReset) + } + + let process = Process() + process.executableURL = URL(fileURLWithPath: binaryPath) + process.arguments = TuiManualIOPumpPolicy.relayArguments( + sessionName: sessionName, + terminalID: terminalID, + cols: grid.cols, + rows: grid.rows + ) + process.environment = environment + + let stdinPipe = Pipe() + let stdoutPipe = Pipe() + let stderrPipe = Pipe() + process.standardInput = stdinPipe + process.standardOutput = stdoutPipe + process.standardError = stderrPipe + + let stderrBox = TuiManualIOStderrBox() + self.stderrBox = stderrBox + pendingExitStatus = nil + // Blocking drain to EOF on a throwaway queue: it continuously + // empties the pipe (a full pipe would wedge the relay) and EOF is + // the only signal that the final exit-reason line has landed. + let stderrHandle = stderrPipe.fileHandleForReading + DispatchQueue(label: "cmux.tuiManualIO.stderr", qos: .utility).async { [weak self] in + stderrBox.append(stderrHandle.readDataToEndOfFile()) + Task { @MainActor [weak self] in + self?.handleStderrDrained(generation: spawnGeneration) + } + } + + let reader = RemoteTmuxProcessOutputReader( + label: "cmux.tuiManualIO.stdout", + maxPendingChunks: 512, + maxPendingBytes: 8 * 1024 * 1024, + onOverflow: { [weak self] in + // The surface stopped consuming; treat like a dead relay so + // a respawn resyncs from a bounded replay. + self?.handleRelayExit(generation: spawnGeneration, forcedExit: .daemonLost) + } + ) + stdoutReader?.close() + stdoutReader = reader + stdoutTask?.cancel() + stdoutTask = Task { @MainActor [weak self] in + for await chunk in reader.stream { + guard let self, self.generation == spawnGeneration, !self.stopped else { break } + reader.release(chunk) + self.surface?.processRemoteOutput(chunk) + self.everRenderedAttach = true + self.consecutiveUnexplainedFailures = 0 + if self.state != .live { + self.state = .live + } + } + } + + process.terminationHandler = { [weak self] finished in + let status = finished.terminationStatus + Task { @MainActor [weak self] in + self?.handleRelayTermination(generation: spawnGeneration, status: status) + } + } + + do { + try process.run() + } catch { + log("spawn failed: \(error)") + reader.close() + noteUnexplainedFailureThenRetryOrFail() + return + } + self.process = process + inputChannel.setHandle(stdinPipe.fileHandleForWriting) + reader.attach(to: stdoutPipe.fileHandleForReading) + log("relay spawned generation=\(spawnGeneration) grid=\(grid.cols)x\(grid.rows)") + } + + private func handleStderrDrained(generation drainedGeneration: Int) { + guard drainedGeneration == generation, !stopped else { return } + stderrDrainedGeneration = drainedGeneration + if let status = pendingExitStatus { + pendingExitStatus = nil + handleRelayExit(generation: drainedGeneration, status: status) + } + } + + private func handleRelayTermination(generation exitedGeneration: Int, status: Int32) { + guard exitedGeneration == generation, !stopped else { return } + guard stderrDrainedGeneration == exitedGeneration else { + pendingExitStatus = status + return + } + handleRelayExit(generation: exitedGeneration, status: status) + } + + private func handleRelayExit( + generation exitedGeneration: Int, + status: Int32? = nil, + forcedExit: TuiManualIOPumpPolicy.RelayExit? = nil + ) { + guard exitedGeneration == generation, !stopped else { return } + inputChannel.setHandle(nil) + process = nil + let exit = forcedExit + ?? TuiManualIOPumpPolicy.relayExit(status: status ?? -1, stderrText: stderrBox.text()) +#if DEBUG + let stderrTail = (stderrBox.text() ?? "").suffix(300).replacingOccurrences(of: "\n", with: " | ") + log("relay exit \(exit) terminal=\(terminalID.prefix(12)) status=\(status.map(String.init) ?? "nil") stderr=\(stderrTail)") +#else + log("relay exit \(exit)") +#endif + switch TuiManualIOPumpPolicy.nextAction(after: exit) { + case .end: + state = .ended + case .retry: + if exit == .failure { + noteUnexplainedFailureThenRetryOrFail() + } else { + scheduleRetry() + } + case .ignore: + break + } + } + + /// Unexplained failures (spawn failure, usage error, crash) stop + /// retrying after a bounded streak so a permanently broken binary + /// converges to a visible failed state instead of a silent retry loop. + private func noteUnexplainedFailureThenRetryOrFail() { + consecutiveUnexplainedFailures += 1 + if consecutiveUnexplainedFailures + >= TuiManualIOPumpPolicy.maxConsecutiveUnexplainedFailures { + retryTask?.cancel() + retryTask = nil + state = .failed + return + } + scheduleRetry() + } + + private func scheduleRetry() { + guard !stopped else { return } + let attempt: Int + if case .reconnecting(let previous) = state { + attempt = previous + 1 + } else { + attempt = 1 + } + state = .reconnecting(attempt: attempt) + let delay = TuiManualIOPumpPolicy.retryDelay(attempt: attempt) + let sleep = sleep + retryTask?.cancel() + retryTask = Task { @MainActor [weak self] in + do { + try await sleep(delay) + } catch { + return + } + guard let self, !self.stopped, !Task.isCancelled else { return } + self.retryTask = nil + self.spawnRelay() + } + } + + private nonisolated func log(_ message: @autoclosure () -> String) { +#if DEBUG + cmuxDebugLog("tuiManualIO.\(message())") +#endif + } +} + +/// Process-wide pump registry keyed by surface id. The surface id is stable +/// across detach transfers between workspaces, so a global registry needs +/// no transfer bookkeeping: overlays and close paths look pumps up by the +/// surface they render. +@MainActor +final class TuiManualIOPumpRegistry { + static let shared = TuiManualIOPumpRegistry() + + private var pumps: [UUID: TuiManualIOPump] = [:] + + func register(_ pump: TuiManualIOPump, surfaceID: UUID) { + pumps[surfaceID]?.stop() + pumps[surfaceID] = pump + } + + func pump(forSurfaceID surfaceID: UUID) -> TuiManualIOPump? { + pumps[surfaceID] + } + + /// Stops and forgets the pump when its panel is discarded for real (not + /// a detach transfer). The relay sees stdin EOF and detaches cleanly. + func stopAndRemove(surfaceID: UUID) { + pumps.removeValue(forKey: surfaceID)?.stop() + } +} + +/// Lock-guarded stderr accumulator: the relay prints its machine-readable +/// exit reason as the final stderr line. +final class TuiManualIOStderrBox: @unchecked Sendable { + private let lock = NSLock() + private var data = Data() + + func append(_ new: Data) { + lock.lock() + // Only the tail matters (final JSON line); bound the buffer. + data.append(new) + if data.count > 64 * 1024 { + data.removeFirst(data.count - 64 * 1024) + } + lock.unlock() + } + + func text() -> String? { + lock.lock() + defer { lock.unlock() } + return data.isEmpty ? nil : String(decoding: data, as: UTF8.self) + } +} diff --git a/Sources/TuiTerminalAttachBridge.swift b/Sources/TuiTerminalAttachBridge.swift index d90a1a6d1905..03b7d8a7c196 100644 --- a/Sources/TuiTerminalAttachBridge.swift +++ b/Sources/TuiTerminalAttachBridge.swift @@ -29,6 +29,28 @@ final class TuiTerminalAttachBridge { ?? key.defaultValue } + /// Manual-IO data-path variant: the daemon terminal feeds a + /// manual-mirror Ghostty surface through a `--pipe-io` relay instead of + /// running the attach client as the surface command. Requires the main + /// backend flag. + nonisolated static var isManualIOEnabled: Bool { + guard isEnabled else { return false } + let key = SettingCatalog().betaFeatures.tuiTerminalBackendManualIO + return Bool.decodeFromUserDefaults(UserDefaults.standard.object(forKey: key.userDefaultsKey)) + ?? key.defaultValue + } + + /// A pump owning the `--pipe-io` relay for one daemon terminal, sharing + /// this bridge's binary, per-app-tag session, and config isolation. + func makeManualIOPump(terminalID: String) -> TuiManualIOPump { + TuiManualIOPump( + binaryPath: Self.binaryPath, + sessionName: sessionName, + terminalID: terminalID, + environment: Self.bridgeEnvironment + ) + } + private nonisolated static var binaryPath: String { let key = SettingCatalog().betaFeatures.tuiTerminalBackendBinaryPath let stored = String.decodeFromUserDefaults(UserDefaults.standard.object(forKey: key.userDefaultsKey)) diff --git a/Sources/Workspace+PanelLifecycle.swift b/Sources/Workspace+PanelLifecycle.swift index 399b54e4b832..250d671509f1 100644 --- a/Sources/Workspace+PanelLifecycle.swift +++ b/Sources/Workspace+PanelLifecycle.swift @@ -496,6 +496,11 @@ extension Workspace { ) { TuiTerminalAttachBridge.shared.closeTerminalForClosedSurface(terminalID: tuiTerminalID) } + // A manual-IO pump follows its surface: a detach transfer keeps the + // surface (and pump) alive; every other discard stops the relay. + if closePanel, !preservesTerminalForTransfer { + TuiManualIOPumpRegistry.shared.stopAndRemove(surfaceID: panelId) + } let removedPanel = panels.removeValue(forKey: panelId) if discardAgentHibernationTracking { diff --git a/Sources/Workspace+RemoteSessionLifecycle.swift b/Sources/Workspace+RemoteSessionLifecycle.swift index 2797f955d724..9f792ea62cb0 100644 --- a/Sources/Workspace+RemoteSessionLifecycle.swift +++ b/Sources/Workspace+RemoteSessionLifecycle.swift @@ -193,6 +193,12 @@ extension Workspace { @discardableResult func reconnectCloudTerminalSurface(surfaceId: UUID) -> Bool { + // Manual-IO tui panes share the reconnect overlay; their Reconnect + // button skips the pump's remaining backoff. + if let pump = TuiManualIOPumpRegistry.shared.pump(forSurfaceID: surfaceId) { + pump.retryNow() + return true + } guard isManagedCloudVMWorkspace, isRemoteTerminalSurface(surfaceId) || remoteDisconnectPlaceholderPanelIds.contains(surfaceId) else { return false diff --git a/Sources/Workspace.swift b/Sources/Workspace.swift index 1c49f156656d..eab40979ad29 100644 --- a/Sources/Workspace.swift +++ b/Sources/Workspace.swift @@ -1687,6 +1687,7 @@ extension Workspace { // spawning fresh or resuming an agent on top of it. let tuiAttachTerminalID: String? let tuiAttachCommand: String? + let tuiManualIOReattachTerminalID: String? if restoredRemotePTYSessionID == nil, restoredTmuxStartCommand == nil, case .reattach(let reattachTerminalID) = TuiTerminalAttachBridge.shared.restoreDecision( snapshotTerminalID: snapshot.terminal?.tuiTerminalID, @@ -1694,11 +1695,19 @@ extension Workspace { hasRemotePTYSessionID: restoredRemotePTYSessionID != nil ) { tuiAttachTerminalID = reattachTerminalID - tuiAttachCommand = TuiTerminalAttachBridge.shared.attachCommand(terminalID: reattachTerminalID) + if TuiTerminalAttachBridge.isManualIOEnabled { + tuiManualIOReattachTerminalID = reattachTerminalID + tuiAttachCommand = nil + } else { + tuiManualIOReattachTerminalID = nil + tuiAttachCommand = TuiTerminalAttachBridge.shared.attachCommand(terminalID: reattachTerminalID) + } } else { tuiAttachTerminalID = nil tuiAttachCommand = nil + tuiManualIOReattachTerminalID = nil } + let tuiReattachActive = tuiAttachCommand != nil || tuiManualIOReattachTerminalID != nil // A crash-restart can leave this exact agent session alive from the // previous launch (or a duplicate panel can reference the same // session in this same restore pass); firing another `codex @@ -1710,7 +1719,7 @@ extension Workspace { // before the freshly spawned process becomes visible to the index. let agentSessionAlreadyActive: Bool = { guard shouldAutoResumeAgent, restoredHibernation == nil, restoredBindingLaunch == nil, - tuiAttachCommand == nil, + !tuiReattachActive, let restorableAgent else { return false } @@ -1748,7 +1757,7 @@ extension Workspace { }() let restoredAgentResumeLaunch: SurfaceResumeStartupLaunch? = if shouldAutoResumeAgent && restoredHibernation == nil && restoredBindingLaunch == nil - && tuiAttachCommand == nil && !agentSessionAlreadyActive { + && !tuiReattachActive && !agentSessionAlreadyActive { if restoresRemoteWorkspaceTerminalSnapshot { restorableAgent?.resumeStartupInput( useLocalRestoreVerb: false, @@ -1770,7 +1779,7 @@ extension Workspace { // the live process it reattaches to. let deferredAgentResumeCandidateInput: String? = if restoreIndexUnavailable, restoredHibernation == nil, - tuiAttachCommand == nil, + !tuiReattachActive, restorableAgent != nil || resumeBinding?.isAgentHookBinding == true { if let restorableAgent { if restoresRemoteWorkspaceTerminalSnapshot { @@ -1809,7 +1818,7 @@ extension Workspace { let deferredAgentResumeAdmission = deferredAgentResumeStartupInput != nil // A daemon reattach redraws its own scrollback; replaying the // persisted copy on top would duplicate it. - let shouldReplayScrollback = tuiAttachCommand == nil && sessionRestorePolicy.shouldReplaySessionScrollback( + let shouldReplayScrollback = !tuiReattachActive && sessionRestorePolicy.shouldReplaySessionScrollback( hasRestorableAgent: restorableAgent != nil, tmuxStartCommand: restoredTmuxStartCommand, hasResumeStartupWork: restoredBindingLaunch != nil || @@ -1825,7 +1834,7 @@ extension Workspace { tuiAttachCommand ?? restoredRemotePTYAttachCommand ?? restoredTmuxStartupScript - let restoredStartupInput = restoredRemotePTYAttachCommand == nil && tuiAttachCommand == nil + let restoredStartupInput = restoredRemotePTYAttachCommand == nil && !tuiReattachActive ? (restoredBindingLaunch?.initialInput ?? restoredAgentResumeLaunch?.initialInput ?? deferredAgentResumeStartupInput) @@ -1927,7 +1936,8 @@ extension Workspace { representedChangeTokens: Set( snapshot.terminal?.fontSizeChangeTokens ?? [] ) - ) + ), + tuiManualIOReattachTerminalID: tuiManualIOReattachTerminalID ) else { if let replayFileURL { try? FileManager.default.removeItem(at: replayFileURL) } // The claim taken above (if any) was for a launch that never @@ -7195,7 +7205,10 @@ final class Workspace: Identifiable, ObservableObject, FilePreviewTabMetadataHos } func cloudTerminalReconnectOverlayPresentation(forSurfaceId surfaceId: UUID) -> CloudTerminalReconnectOverlayPolicy.Presentation? { - CloudTerminalReconnectOverlayPolicy.presentation( + if let pump = TuiManualIOPumpRegistry.shared.pump(forSurfaceID: surfaceId) { + return pump.overlayPresentation + } + return CloudTerminalReconnectOverlayPolicy.presentation( isManagedCloudWorkspace: isManagedCloudVMWorkspace, isRemoteTerminalSurface: isRemoteTerminalSurface(surfaceId) || remoteDisconnectPlaceholderPanelIds.contains(surfaceId), connectionState: remoteConnectionState, @@ -8667,7 +8680,8 @@ final class Workspace: Identifiable, ObservableObject, FilePreviewTabMetadataHos terminalFontSizeCreationPolicy: TerminalFontSizeCreationPolicy = .inherit, inheritWorkingDirectoryFallback: Bool = false, workingDirectoryFallbackSourcePanelId: UUID? = nil, - allowTextBoxFocusDefault: Bool = true + allowTextBoxFocusDefault: Bool = true, + tuiManualIOReattachTerminalID: String? = nil ) -> TerminalPanel? { return newTerminalSurfaceOutcome( inPane: paneId, @@ -8687,7 +8701,8 @@ final class Workspace: Identifiable, ObservableObject, FilePreviewTabMetadataHos terminalFontSizeCreationPolicy: terminalFontSizeCreationPolicy, inheritWorkingDirectoryFallback: inheritWorkingDirectoryFallback, workingDirectoryFallbackSourcePanelId: workingDirectoryFallbackSourcePanelId, - allowTextBoxFocusDefault: allowTextBoxFocusDefault + allowTextBoxFocusDefault: allowTextBoxFocusDefault, + tuiManualIOReattachTerminalID: tuiManualIOReattachTerminalID ).panel } @@ -8712,7 +8727,8 @@ final class Workspace: Identifiable, ObservableObject, FilePreviewTabMetadataHos terminalFontSizeCreationPolicy: TerminalFontSizeCreationPolicy = .inherit, inheritWorkingDirectoryFallback: Bool = false, workingDirectoryFallbackSourcePanelId: UUID? = nil, - allowTextBoxFocusDefault: Bool = true + allowTextBoxFocusDefault: Bool = true, + tuiManualIOReattachTerminalID: String? = nil ) -> TerminalPanelCreationOutcome { guard !isRetiredFromOwningTabManager else { return .failed } // In a remote tmux mirror, a new tab means "create a tmux window"; never @@ -8763,7 +8779,8 @@ final class Workspace: Identifiable, ObservableObject, FilePreviewTabMetadataHos terminalFontSizeCreationPolicy: terminalFontSizeCreationPolicy, inheritWorkingDirectoryFallback: inheritWorkingDirectoryFallback, workingDirectoryFallbackSourcePanelId: workingDirectoryFallbackSourcePanelId, - allowTextBoxFocusDefault: allowTextBoxFocusDefault + allowTextBoxFocusDefault: allowTextBoxFocusDefault, + tuiManualIOReattachTerminalID: tuiManualIOReattachTerminalID ) else { return .failed } return .created(panel) } @@ -8786,7 +8803,8 @@ final class Workspace: Identifiable, ObservableObject, FilePreviewTabMetadataHos terminalFontSizeCreationPolicy: TerminalFontSizeCreationPolicy, inheritWorkingDirectoryFallback: Bool, workingDirectoryFallbackSourcePanelId: UUID?, - allowTextBoxFocusDefault: Bool + allowTextBoxFocusDefault: Bool, + tuiManualIOReattachTerminalID: String? = nil ) -> TerminalPanel? { let shouldFocusNewTab = focus ?? (bonsplitController.focusedPaneId == paneId) let previousFocusedPanelId = focusedPanelId @@ -8803,7 +8821,7 @@ final class Workspace: Identifiable, ObservableObject, FilePreviewTabMetadataHos // terminal with a daemon terminal so it survives quitting the app. // Any provisioning failure falls through to today's spawn path. var tuiProvisionedTerminal: TuiTerminalAttachBridge.ProvisionedTerminal? - if TuiTerminalAttachPolicy.shouldProvisionNewTerminal( + if tuiManualIOReattachTerminalID == nil, TuiTerminalAttachPolicy.shouldProvisionNewTerminal( flagEnabled: TuiTerminalAttachBridge.isEnabled, hasExplicitStartupCommand: startupCommand != nil, hasTmuxStartCommand: tmuxStartCommand != nil, @@ -8812,7 +8830,12 @@ final class Workspace: Identifiable, ObservableObject, FilePreviewTabMetadataHos ) { tuiProvisionedTerminal = TuiTerminalAttachBridge.shared.provisionTerminalForNewSurface() } - let effectiveStartupCommand = tuiProvisionedTerminal?.attachCommand ?? startupCommand + // Manual-IO variant: the surface runs in Ghostty's manual-mirror IO + // mode and a pump relays daemon bytes; there is no surface command. + let tuiManualIOTerminalID: String? = tuiManualIOReattachTerminalID + ?? (TuiTerminalAttachBridge.isManualIOEnabled ? tuiProvisionedTerminal?.terminalID : nil) + let effectiveStartupCommand = (tuiManualIOTerminalID == nil ? tuiProvisionedTerminal?.attachCommand : nil) + ?? startupCommand let remoteStartupCommandForEnvironment = explicitInitialCommand == nil ? remoteTerminalStartupCommand : nil let newPanelID = restoredSurfaceId ?? UUID() let requestedRemotePTYSessionID = normalizedRemotePTYSessionID(remotePTYSessionID) @@ -8849,23 +8872,44 @@ final class Workspace: Identifiable, ObservableObject, FilePreviewTabMetadataHos // surface id (the panel/surface id IS the ghostty surface id, a // Swift-side UUID), so a session's terminal binding survives relaunch // and restore. The caller only passes an id it has verified is free. - let newPanel = TerminalPanel( - id: newPanelID, - workspaceId: id, - context: GHOSTTY_SURFACE_CONTEXT_SPLIT, - configTemplate: inheritedConfig, - workingDirectory: requestedWorkingDirectory, - portOrdinal: portOrdinal, - initialCommand: effectiveStartupCommand, - tmuxStartCommand: tmuxStartCommand, - initialInput: initialInput, - additionalEnvironment: effectiveStartupEnvironment, - runtimeSpawnPolicy: terminalStartupRestoreCoordinator.runtimeSpawnPolicy( - requestedPolicy: runtimeSpawnPolicy, - willRunStartupCommand: false, - willRunStartupInput: startupRestoreAgent != nil && initialInput != nil + let newPanel: TerminalPanel + if let tuiManualIOTerminalID { + let pump = TuiTerminalAttachBridge.shared.makeManualIOPump(terminalID: tuiManualIOTerminalID) + let surface = TerminalSurface( + id: newPanelID, + tabId: id, + context: GHOSTTY_SURFACE_CONTEXT_SPLIT, + configTemplate: inheritedConfig, + workingDirectory: requestedWorkingDirectory, + portOrdinal: portOrdinal, + ioMode: .manualMirror, + manualInputHandler: pump.makeManualInputHandler() ) - ) + pump.start(surface: surface) + pump.onStateChange = { [weak surface] in + surface?.owningWorkspace()?.postRemoteConnectionPresentationDidChange() + } + TuiManualIOPumpRegistry.shared.register(pump, surfaceID: surface.id) + newPanel = TerminalPanel(workspaceId: id, surface: surface) + } else { + newPanel = TerminalPanel( + id: newPanelID, + workspaceId: id, + context: GHOSTTY_SURFACE_CONTEXT_SPLIT, + configTemplate: inheritedConfig, + workingDirectory: requestedWorkingDirectory, + portOrdinal: portOrdinal, + initialCommand: effectiveStartupCommand, + tmuxStartCommand: tmuxStartCommand, + initialInput: initialInput, + additionalEnvironment: effectiveStartupEnvironment, + runtimeSpawnPolicy: terminalStartupRestoreCoordinator.runtimeSpawnPolicy( + requestedPolicy: runtimeSpawnPolicy, + willRunStartupCommand: false, + willRunStartupInput: startupRestoreAgent != nil && initialInput != nil + ) + ) + } configureNewTerminalPanel( newPanel, allowTextBoxFocusDefault: shouldFocusNewTab && allowTextBoxFocusDefault @@ -8874,6 +8918,8 @@ final class Workspace: Identifiable, ObservableObject, FilePreviewTabMetadataHos panelTitles[newPanel.id] = newPanel.displayTitle if let tuiProvisionedTerminal { tuiTerminalIDsByPanelId[newPanel.id] = tuiProvisionedTerminal.terminalID + } else if let tuiManualIOReattachTerminalID { + tuiTerminalIDsByPanelId[newPanel.id] = tuiManualIOReattachTerminalID } let tracksRemoteTerminalSurface = remoteTerminalStartupCommand != nil || effectiveRemotePTYSessionID != nil if let effectiveRemotePTYSessionID { diff --git a/cmux-tui/crates/cmux-tui/src/main.rs b/cmux-tui/crates/cmux-tui/src/main.rs index 0b96e1b107b2..205e9ee83ad8 100644 --- a/cmux-tui/crates/cmux-tui/src/main.rs +++ b/cmux-tui/crates/cmux-tui/src/main.rs @@ -27,6 +27,7 @@ mod machine_provider_client; #[cfg(unix)] mod machine_provider_runtime; mod machine_runtime; +mod pipe_io; mod plugin_manager; mod process_diagnostics; #[cfg(target_os = "linux")] @@ -429,6 +430,12 @@ START OPTIONS --session Session name (default: main). Determines the socket path. --socket Explicit control socket path. --terminal With attach, show only this terminal (use `cmux terminal list`). + --pipe-io With attach --terminal, relay raw bytes over stdio for an + embedder-fed surface instead of rendering (stdout = replay + then live output, stdin = JSON input/resize lines, + exit 0 = terminal ended, exit 2 = daemon connection lost). + --cols Initial viewer columns for --pipe-io (default 80). + --rows Initial viewer rows for --pipe-io (default 24). --state Durable session-state root (default: platform state dir). --ephemeral Keep workspace state in memory for this run only. --machine-provider @@ -492,6 +499,9 @@ struct Args { session: String, socket: Option, terminal: Option, + pipe_io: bool, + pipe_io_cols: Option, + pipe_io_rows: Option, state: Option, ephemeral: bool, machine_provider: Option, @@ -566,6 +576,9 @@ fn parse_args_result(args: impl IntoIterator) -> Result) -> Result { + out.pipe_io = true; + } + "--cols" => { + let value = args.next().ok_or_else(|| "--cols needs a value".to_string())?; + out.pipe_io_cols = + Some(value.parse().map_err(|_| "--cols needs a number".to_string())?); + } + "--rows" => { + let value = args.next().ok_or_else(|| "--rows needs a value".to_string())?; + out.pipe_io_rows = + Some(value.parse().map_err(|_| "--rows needs a number".to_string())?); + } "--machine-provider" => { if out.machine_provider.is_some() { return Err("--machine-provider may be supplied only once".to_string()); @@ -813,6 +839,12 @@ fn parse_args_result(args: impl IntoIterator) -> Result anyhow::Result<()> { cmux_tui_core::terminal_host_runtime::serve_terminal_host_stdio(args, &mut reader, &mut writer) } +/// Ends a `--pipe-io` relay: the machine-readable exit reason is the final +/// stderr line (the embedder localizes what it shows) and the exit code +/// carries the respawn decision. +fn exit_pipe_io(reason: pipe_io::PipeIoExitReason) -> ! { + { + let mut stderr = std::io::stderr().lock(); + let _ = writeln!(stderr, "{}", serde_json::json!({"exit": {"reason": reason.as_str()}})); + let _ = stderr.flush(); + } + std::process::exit(reason.exit_code()); +} + fn run_attach(args: Args, config: config::StartupConfigSnapshot) -> anyhow::Result<()> { let socket_path = match args.socket { Some(path) => path, @@ -1584,27 +1628,71 @@ fn run_attach(args: Args, config: config::StartupConfigSnapshot) -> anyhow::Resu }) .transpose()?; let remote = if terminal.is_some() { - RemoteSession::connect_for_terminal_attach(&socket_path)? + match RemoteSession::connect_for_terminal_attach(&socket_path) { + Ok(remote) => remote, + // A pipe-io embedder retries daemon-lost forever (a daemon + // restart window looks exactly like this) but gives up on + // unexplained failures; a connect failure is the former. + Err(_) if args.pipe_io => exit_pipe_io(pipe_io::PipeIoExitReason::DaemonLost), + Err(error) => return Err(error), + } } else { RemoteSession::connect(&socket_path)? }; let surface_only = if let Some(terminal) = terminal.as_ref() { - let tree = remote.refresh_tree()?; - let surface = tree - .resolve_terminal(terminal) + let tree = match remote.refresh_tree() { + Ok(tree) => tree, + Err(_) if args.pipe_io => exit_pipe_io(pipe_io::PipeIoExitReason::DaemonLost), + Err(error) => return Err(error), + }; + let resolved = tree.resolve_terminal(terminal); + // A pipe-io embedder respawns on daemon loss; a terminal that no + // longer exists must read as terminal-ended (do not respawn), not + // as a startup failure it would retry forever. + if args.pipe_io && resolved.is_none() { + exit_pipe_io(pipe_io::PipeIoExitReason::TerminalEnded); + } + let surface = resolved .ok_or_else(|| anyhow::anyhow!(messages.unknown_terminal(terminal.as_str())))?; if !remote.supports_surface_subscription_filter() { anyhow::bail!(messages.filtered_subscription_unavailable); } - remote.scope_events_to_surface(surface)?; - let tree = remote.refresh_tree()?; + if let Err(error) = remote.scope_events_to_surface(surface) { + if args.pipe_io { + exit_pipe_io(pipe_io::PipeIoExitReason::DaemonLost); + } + return Err(error); + } + let tree = match remote.refresh_tree() { + Ok(tree) => tree, + Err(_) if args.pipe_io => exit_pipe_io(pipe_io::PipeIoExitReason::DaemonLost), + Err(error) => return Err(error), + }; if tree.resolve_terminal(terminal) != Some(surface) { + if args.pipe_io { + exit_pipe_io(pipe_io::PipeIoExitReason::TerminalEnded); + } anyhow::bail!(messages.unknown_terminal(terminal.as_str())); } Some(surface) } else { None }; + if args.pipe_io { + let surface = surface_only.expect("--pipe-io is validated to carry --terminal"); + let terminal = terminal.as_ref().expect("--pipe-io is validated to carry --terminal"); + let session = Session::Remote(remote.clone()); + let reason = pipe_io::run( + &session, + &remote, + &socket_path, + terminal, + surface, + args.pipe_io_cols.unwrap_or(80), + args.pipe_io_rows.unwrap_or(24), + )?; + exit_pipe_io(reason); + } run_connected_session_client( socket_path, args.session, diff --git a/cmux-tui/crates/cmux-tui/src/pipe_io.rs b/cmux-tui/crates/cmux-tui/src/pipe_io.rs new file mode 100644 index 000000000000..e811db7ffa37 --- /dev/null +++ b/cmux-tui/crates/cmux-tui/src/pipe_io.rs @@ -0,0 +1,321 @@ +//! Renderer-less byte relay for embedder-fed (manual-IO) surfaces. +//! +//! `cmux attach --terminal --pipe-io` is the scoped attach client minus +//! the renderer: the embedder (the macOS GUI's manual-mirror Ghostty +//! surface) parses the byte stream itself, so this process only relays. +//! +//! Contract: +//! - stdout: raw VT bytes — the daemon replay first, then live PTY output. +//! A replay that is not the relay's first output is prefixed with a full +//! reset (`ESC c` + erase-scrollback) because a daemon replay REPLACES +//! terminal state while a byte stream can only append. +//! - stdin: JSON lines — `{"input":""}` forwards raw bytes to the +//! daemon terminal's PTY; `{"resize":{"cols":N,"rows":N}}` drives the +//! daemon-side viewer size. Unknown keys are ignored (forward compat); +//! stdin EOF means the embedder is gone and ends the relay cleanly. +//! - stderr: one final JSON line `{"exit":{"reason":...}}`. +//! - exit code: 0 when the terminal ended or the embedder closed stdin (do +//! not respawn), 2 when the daemon connection was lost (respawning +//! reattaches and resyncs from a fresh replay). + +use std::io::{BufRead, Write}; +use std::path::Path; +use std::sync::Arc; +use std::sync::mpsc::{Receiver, SyncSender}; + +use base64::Engine as _; +use cmux_tui_core::SurfaceId; +use cmux_tui_core::resource::TerminalPublicId; + +use crate::session::{PipeIoEvent, RemoteSession, Session, SurfaceAttach, SurfaceHandle}; + +/// The terminal ended, or the embedder walked away: respawning is wrong. +pub const EXIT_DO_NOT_RESPAWN: i32 = 0; +/// The daemon connection was lost: the embedder may respawn to resync. +pub const EXIT_DAEMON_LOST: i32 = 2; + +/// Bounded event queue between the session reader thread and the stdout +/// pump. A full queue means the embedder stopped reading; the session +/// treats that as a lost transport rather than wedging its reader thread. +const EVENT_QUEUE_CAPACITY: usize = 4096; + +/// Emitted before any replay that is not the relay's first output: full +/// reset plus erase-scrollback, so the replacement replay does not stack on +/// top of the embedder's previous terminal state. +const REPLAY_RESET: &[u8] = b"\x1bc\x1b[3J"; + +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum PipeIoExitReason { + TerminalEnded, + DaemonLost, + ParentClosed, +} + +impl PipeIoExitReason { + pub fn as_str(self) -> &'static str { + match self { + Self::TerminalEnded => "terminal-ended", + Self::DaemonLost => "daemon-lost", + Self::ParentClosed => "parent-closed", + } + } + + pub fn exit_code(self) -> i32 { + match self { + Self::TerminalEnded | Self::ParentClosed => EXIT_DO_NOT_RESPAWN, + Self::DaemonLost => EXIT_DAEMON_LOST, + } + } +} + +/// One parsed stdin line. +#[derive(Debug, PartialEq, Eq)] +pub enum PipeIoRequest { + Input(Vec), + Resize { + cols: u16, + rows: u16, + }, + /// A well-formed line carrying no verb this relay knows: ignored, so a + /// newer embedder can talk to an older relay. + Unknown, +} + +pub fn parse_request(line: &str) -> anyhow::Result { + let value: serde_json::Value = serde_json::from_str(line)?; + if let Some(encoded) = value.get("input").and_then(serde_json::Value::as_str) { + let bytes = base64::engine::general_purpose::STANDARD.decode(encoded)?; + return Ok(PipeIoRequest::Input(bytes)); + } + if let Some(resize) = value.get("resize") { + let cols = resize.get("cols").and_then(serde_json::Value::as_u64); + let rows = resize.get("rows").and_then(serde_json::Value::as_u64); + let (Some(cols), Some(rows)) = (cols, rows) else { + anyhow::bail!("resize needs numeric cols and rows"); + }; + let (cols, rows) = (u16::try_from(cols)?, u16::try_from(rows)?); + return Ok(PipeIoRequest::Resize { cols: cols.max(1), rows: rows.max(1) }); + } + Ok(PipeIoRequest::Unknown) +} + +pub fn run( + session: &Session, + remote: &Arc, + socket_path: &Path, + terminal: &TerminalPublicId, + surface: SurfaceId, + cols: u16, + rows: u16, +) -> anyhow::Result { + let (sender, receiver) = std::sync::mpsc::sync_channel(EVENT_QUEUE_CAPACITY); + // Install before attach so the initial replay cannot be missed. + remote.install_pipe_io_tap(surface, sender.clone()); + let handle = match session.try_surface_sized(surface, Some((cols.max(1), rows.max(1))))? { + SurfaceAttach::Attached(handle) => handle, + SurfaceAttach::Retired | SurfaceAttach::Missing => { + return Ok(PipeIoExitReason::TerminalEnded); + } + SurfaceAttach::Deferred => anyhow::bail!("terminal attach was deferred by the server"), + }; + // The daemon resizes a terminal's PTY only for its geometry-authority + // client (the full TUI client claims this for its active surface). The + // relay is the embedder's only viewer of this terminal, so claim the + // authority or every embedder resize is recorded but never applied. + if let Err(error) = session.claim_terminal_geometry(surface) { + eprintln!( + "{}", + serde_json::json!({"diag": {"claim-terminal-geometry": {"error": error.to_string()}}}) + ); + } + spawn_stdin_pump(handle, sender); + let reason = pump_events_to_stdout(&receiver, &mut std::io::stdout().lock())?; + if reason == PipeIoExitReason::DaemonLost { + return Ok(classify_daemon_loss(remote, socket_path, terminal)); + } + Ok(reason) +} + +/// The stream ended without a terminal-exit event, but "stream lost" covers +/// two situations the embedder must tell apart: the daemon really is +/// unreachable (respawn until it returns), or the terminal was closed and +/// the daemon tore the scoped stream down before the exit event reached us +/// (never respawn). One bounded probe of the daemon resolves it: prefer the +/// existing connection, and fall back to a fresh connect when the socket +/// died with the stream. +fn classify_daemon_loss( + remote: &Arc, + socket_path: &Path, + terminal: &TerminalPublicId, +) -> PipeIoExitReason { + let terminal_still_exists = + remote.refresh_tree().ok().map(|tree| tree.resolve_terminal(terminal).is_some()).or_else( + || { + let probe = RemoteSession::connect_for_terminal_attach(socket_path).ok()?; + let tree = probe.refresh_tree().ok()?; + Some(tree.resolve_terminal(terminal).is_some()) + }, + ); + match terminal_still_exists { + Some(false) => PipeIoExitReason::TerminalEnded, + Some(true) | None => PipeIoExitReason::DaemonLost, + } +} + +/// Forwards embedder requests from stdin until EOF, then reports the closed +/// parent through the shared event queue. +fn spawn_stdin_pump(handle: SurfaceHandle, sender: SyncSender) { + std::thread::Builder::new() + .name("pipe-io-stdin".into()) + .spawn(move || { + let stdin = std::io::stdin(); + for line in stdin.lock().lines() { + let Ok(line) = line else { break }; + if line.trim().is_empty() { + continue; + } + match parse_request(&line) { + Ok(PipeIoRequest::Input(bytes)) => { + if handle.write_bytes(&bytes).is_err() { + // The transport owns loss reporting; input can + // only stop early. + break; + } + } + Ok(PipeIoRequest::Resize { cols, rows }) => { + // Diagnostics only; the exit JSON stays the final + // stderr line and embedders skip lines without an + // "exit" key. + match handle.resize(cols, rows) { + Ok(accepted) => eprintln!( + "{}", + serde_json::json!({ + "diag": {"resize": {"cols": cols, "rows": rows, "accepted": accepted}} + }) + ), + Err(error) => eprintln!( + "{}", + serde_json::json!({ + "diag": {"resize": {"cols": cols, "rows": rows, "error": error.to_string()}} + }) + ), + } + } + Ok(PipeIoRequest::Unknown) => {} + // A malformed line means the embedder side is broken; + // stop consuming rather than misinterpreting input. + Err(_) => break, + } + } + // Blocking send: the queue is drained until the main loop + // returns, and a dropped receiver just ends this thread. + let _ = sender.send(PipeIoEvent::StdinClosed); + }) + .expect("spawn pipe-io stdin pump"); +} + +fn pump_events_to_stdout( + receiver: &Receiver, + stdout: &mut impl Write, +) -> anyhow::Result { + let mut emitted_output = false; + loop { + // A dropped sender without a prior event means the session went + // away wholesale: report it as a lost daemon. + let Ok(event) = receiver.recv() else { + return Ok(PipeIoExitReason::DaemonLost); + }; + let write_result = match &event { + PipeIoEvent::Replay { bytes } => { + let reset = if emitted_output { stdout.write_all(REPLAY_RESET) } else { Ok(()) }; + reset.and_then(|()| stdout.write_all(bytes)).and_then(|()| stdout.flush()) + } + PipeIoEvent::Output(bytes) => stdout.write_all(bytes).and_then(|()| stdout.flush()), + PipeIoEvent::SurfaceExited => return Ok(PipeIoExitReason::TerminalEnded), + PipeIoEvent::TransportLost => return Ok(PipeIoExitReason::DaemonLost), + PipeIoEvent::StdinClosed => return Ok(PipeIoExitReason::ParentClosed), + }; + if write_result.is_err() { + // stdout is the embedder; a failed write means it is gone. + return Ok(PipeIoExitReason::ParentClosed); + } + emitted_output = true; + } +} + +#[cfg(test)] +mod tests { + use super::*; + + fn replay(bytes: &[u8]) -> PipeIoEvent { + PipeIoEvent::Replay { bytes: bytes.to_vec() } + } + + #[test] + fn parse_request_decodes_input_resize_and_ignores_unknown_verbs() { + assert_eq!( + parse_request(r#"{"input":"aGk="}"#).unwrap(), + PipeIoRequest::Input(b"hi".to_vec()) + ); + assert_eq!( + parse_request(r#"{"resize":{"cols":100,"rows":30}}"#).unwrap(), + PipeIoRequest::Resize { cols: 100, rows: 30 } + ); + assert_eq!( + parse_request(r#"{"resize":{"cols":0,"rows":0}}"#).unwrap(), + PipeIoRequest::Resize { cols: 1, rows: 1 } + ); + assert_eq!(parse_request(r#"{"future-verb":true}"#).unwrap(), PipeIoRequest::Unknown); + assert!(parse_request("not json").is_err()); + assert!(parse_request(r#"{"input":"@@not-base64@@"}"#).is_err()); + assert!(parse_request(r#"{"resize":{"cols":100}}"#).is_err()); + } + + #[test] + fn exit_reasons_map_to_the_respawn_contract() { + assert_eq!(PipeIoExitReason::TerminalEnded.exit_code(), EXIT_DO_NOT_RESPAWN); + assert_eq!(PipeIoExitReason::ParentClosed.exit_code(), EXIT_DO_NOT_RESPAWN); + assert_eq!(PipeIoExitReason::DaemonLost.exit_code(), EXIT_DAEMON_LOST); + assert_eq!(PipeIoExitReason::TerminalEnded.as_str(), "terminal-ended"); + assert_eq!(PipeIoExitReason::DaemonLost.as_str(), "daemon-lost"); + assert_eq!(PipeIoExitReason::ParentClosed.as_str(), "parent-closed"); + } + + #[test] + fn stdout_pump_prefixes_only_non_initial_replays_with_a_full_reset() { + let (sender, receiver) = std::sync::mpsc::sync_channel(8); + sender.send(replay(b"FIRST")).unwrap(); + sender.send(PipeIoEvent::Output(b"live".to_vec())).unwrap(); + sender.send(replay(b"SECOND")).unwrap(); + sender.send(PipeIoEvent::SurfaceExited).unwrap(); + let mut stdout = Vec::new(); + let reason = pump_events_to_stdout(&receiver, &mut stdout).unwrap(); + assert_eq!(reason, PipeIoExitReason::TerminalEnded); + let mut expected = b"FIRSTlive".to_vec(); + expected.extend_from_slice(REPLAY_RESET); + expected.extend_from_slice(b"SECOND"); + assert_eq!(stdout, expected); + } + + #[test] + fn stdout_pump_maps_lifecycle_events_to_exit_reasons() { + for (event, expected) in [ + (PipeIoEvent::TransportLost, PipeIoExitReason::DaemonLost), + (PipeIoEvent::StdinClosed, PipeIoExitReason::ParentClosed), + ] { + let (sender, receiver) = std::sync::mpsc::sync_channel(8); + sender.send(event).unwrap(); + let mut stdout = Vec::new(); + assert_eq!(pump_events_to_stdout(&receiver, &mut stdout).unwrap(), expected); + assert!(stdout.is_empty()); + } + // Sender dropped without any event: the session vanished wholesale. + let (sender, receiver) = std::sync::mpsc::sync_channel::(8); + drop(sender); + let mut stdout = Vec::new(); + assert_eq!( + pump_events_to_stdout(&receiver, &mut stdout).unwrap(), + PipeIoExitReason::DaemonLost + ); + } +} diff --git a/cmux-tui/crates/cmux-tui/src/session/mod.rs b/cmux-tui/crates/cmux-tui/src/session/mod.rs index cc5de5529048..c4fd82fa3203 100644 --- a/cmux-tui/crates/cmux-tui/src/session/mod.rs +++ b/cmux-tui/crates/cmux-tui/src/session/mod.rs @@ -36,8 +36,8 @@ use serde::Deserialize; use serde_json::{Map, Value, json}; pub use remote::{ - RemoteMessageReader, RemoteMessageWriter, RemoteSession, RemoteSurface, RemoteTransport, - RemoteTransportAbort, + PipeIoEvent, RemoteMessageReader, RemoteMessageWriter, RemoteSession, RemoteSurface, + RemoteTransport, RemoteTransportAbort, }; pub use tree::{TabNotificationView, TreeView, WorkspaceView}; diff --git a/cmux-tui/crates/cmux-tui/src/session/remote.rs b/cmux-tui/crates/cmux-tui/src/session/remote.rs index d89270f847dd..dbb522d9e472 100644 --- a/cmux-tui/crates/cmux-tui/src/session/remote.rs +++ b/cmux-tui/crates/cmux-tui/src/session/remote.rs @@ -1516,6 +1516,25 @@ pub struct RemoteSession { capabilities: Mutex>, provider_workspace_authority: Option, provider_workspaces_guarded: AtomicBool, + pipe_io_tap: Mutex>, +} + +/// One event on the raw byte stream a `--pipe-io` relay serves to its +/// embedder. `Replay` REPLACES all prior terminal state (the relay must +/// emit a full reset before any replay that is not its first output); +/// `Output` appends live PTY bytes. +#[derive(Debug, PartialEq, Eq)] +pub enum PipeIoEvent { + Replay { bytes: Vec }, + Output(Vec), + SurfaceExited, + TransportLost, + StdinClosed, +} + +struct PipeIoTap { + surface: SurfaceId, + sender: std::sync::mpsc::SyncSender, } pub(super) enum RemoteSurfaceAttach { @@ -1872,6 +1891,7 @@ impl RemoteSession { capabilities: Mutex::new(HashSet::new()), provider_workspace_authority, provider_workspaces_guarded: AtomicBool::new(false), + pipe_io_tap: Mutex::new(None), }); let reader_session = Arc::downgrade(&session); @@ -2126,6 +2146,7 @@ impl RemoteSession { id, format_args!("vt-state cols={cols} rows={rows} bytes={}", replay.len()), ); + self.pipe_io_forward(id, || PipeIoEvent::Replay { bytes: replay.clone() }); let Ok(kitty_image_aliases) = parse_kitty_image_aliases(&value) else { self.disconnect_transport(); return; @@ -2190,6 +2211,7 @@ impl RemoteSession { }; let colors = value.get("colors").and_then(parse_terminal_colors); self.log_frame(id, format_args!("output bytes={}", bytes.len())); + self.pipe_io_forward(id, || PipeIoEvent::Output(bytes.clone())); if let Some(surface) = self.surfaces.lock().unwrap().get(&id).cloned() { surface.scan_cursor_provenance(&bytes); let mut term = surface.term.lock().unwrap(); @@ -2239,6 +2261,9 @@ impl RemoteSession { replay.as_ref().map(|bytes| bytes.len()).unwrap_or(0) ), ); + if let Some(replay) = replay.as_ref() { + self.pipe_io_forward(id, || PipeIoEvent::Replay { bytes: replay.clone() }); + } if let Some(surface) = self.surfaces.lock().unwrap().get(&id).cloned() { if surface .apply_stream_resize_with_colors( @@ -2308,6 +2333,11 @@ impl RemoteSession { } Some("detached") => { if let Some(id) = surface_id() { + // A server-side stream detach (e.g. backpressure + // overflow) ends the byte feed without ending the + // terminal. The relay reports a lost connection so its + // embedder respawns and resyncs from a fresh replay. + self.pipe_io_forward(id, || PipeIoEvent::TransportLost); self.surfaces.lock().unwrap().remove(&id); self.emit(MuxEvent::SurfaceOutput(id)); } @@ -2351,6 +2381,7 @@ impl RemoteSession { } Some("surface-exited") => { if let Some(id) = surface_id() { + self.pipe_io_forward(id, || PipeIoEvent::SurfaceExited); // Retire the mirror immediately. The authoritative tree // refresh may lag this event, but input and reattach must // already fail closed for a known-exited surface. @@ -2892,6 +2923,43 @@ impl RemoteSession { drop(state); self.begin_shutdown(); self.interactive_writer.close(); + // A pipe-io relay learns about every terminal-connection loss here + // (reader-thread EOF and every protocol-error disconnect path). + if let Some(tap) = self.pipe_io_tap.lock().unwrap().as_ref() { + let _ = tap.sender.try_send(PipeIoEvent::TransportLost); + } + } + + /// Routes one scoped surface's raw byte stream to a `--pipe-io` relay. + /// Install before `attach-surface` so the initial replay is not missed. + pub fn install_pipe_io_tap( + &self, + surface: SurfaceId, + sender: std::sync::mpsc::SyncSender, + ) { + *self.pipe_io_tap.lock().unwrap() = Some(PipeIoTap { surface, sender }); + } + + /// Forwards one tap event, treating a full or dropped queue as a lost + /// transport: the relay's reader stalled, and silently dropping bytes + /// would corrupt the embedder's terminal state (bounded-backpressure + /// policy; never wedge the session reader thread). + fn pipe_io_forward(&self, surface: SurfaceId, event: impl FnOnce() -> PipeIoEvent) { + use std::sync::mpsc::TrySendError; + let stalled = { + let tap = self.pipe_io_tap.lock().unwrap(); + let Some(tap) = tap.as_ref() else { return }; + if tap.surface != surface { + return; + } + match tap.sender.try_send(event()) { + Ok(()) => false, + Err(TrySendError::Full(_) | TrySendError::Disconnected(_)) => true, + } + }; + if stalled { + self.disconnect_transport(); + } } /// Returns the first reason recorded when the remote reader stopped. @@ -3824,6 +3892,7 @@ fn test_session_with_writer( capabilities: Mutex::new(capabilities), provider_workspace_authority, provider_workspaces_guarded: AtomicBool::new(false), + pipe_io_tap: Mutex::new(None), }) } @@ -5065,6 +5134,7 @@ mod tests { capabilities: Mutex::new(capabilities), provider_workspace_authority, provider_workspaces_guarded: AtomicBool::new(false), + pipe_io_tap: Mutex::new(None), }) } diff --git a/cmux-tui/crates/cmux-tui/tests/cli.rs b/cmux-tui/crates/cmux-tui/tests/cli.rs index 946397d978b8..93962c65dcb2 100644 --- a/cmux-tui/crates/cmux-tui/tests/cli.rs +++ b/cmux-tui/crates/cmux-tui/tests/cli.rs @@ -4088,3 +4088,327 @@ fn unique_temp_dir(name: &str) -> PathBuf { fn bin() -> &'static str { env!("CARGO_BIN_EXE_cmux-tui") } + +// --------------------------------------------------------------------------- +// `attach --pipe-io`: renderer-less byte relay for embedder-fed surfaces. +// +// Contract under test (GUI-frontend migration, manual-IO surface): +// stdout = raw VT bytes (daemon replay first, then live output), +// stdin = JSON lines ({"input":""} | {"resize":{"cols":N,"rows":N}}), +// stderr = one final JSON line {"exit":{"reason":...}}, +// exit 0 = terminal ended (or parent closed stdin) — do not respawn, +// exit 2 = daemon connection lost — the embedder may respawn to resync. +// --------------------------------------------------------------------------- + +#[cfg(unix)] +struct PipeIoRelay { + child: Child, + stdin: Option, + stdout: std::sync::Arc>>, + stderr: std::sync::Arc>>, +} + +#[cfg(unix)] +impl PipeIoRelay { + fn start(server: &HeadlessServer, terminal: &str, cols: u16, rows: u16) -> Self { + let mut child = Command::new(bin()) + .args(["attach", "--socket"]) + .arg(&server.socket) + .args([ + "--terminal", + terminal, + "--pipe-io", + "--cols", + &cols.to_string(), + "--rows", + &rows.to_string(), + ]) + .env_remove("CMUX_TUI_SOCKET") + .stdin(Stdio::piped()) + .stdout(Stdio::piped()) + .stderr(Stdio::piped()) + .spawn() + .unwrap(); + let stdin = child.stdin.take(); + let stdout = std::sync::Arc::new(std::sync::Mutex::new(Vec::new())); + let stderr = std::sync::Arc::new(std::sync::Mutex::new(Vec::new())); + let out = child.stdout.take().unwrap(); + let err = child.stderr.take().unwrap(); + let stdout_sink = stdout.clone(); + std::thread::spawn(move || drain_into(out, stdout_sink)); + let stderr_sink = stderr.clone(); + std::thread::spawn(move || drain_into(err, stderr_sink)); + Self { child, stdin, stdout, stderr } + } + + fn send_input(&mut self, bytes: &[u8]) { + use base64::Engine as _; + let encoded = base64::engine::general_purpose::STANDARD.encode(bytes); + let line = format!("{}\n", serde_json::json!({ "input": encoded })); + self.stdin.as_mut().unwrap().write_all(line.as_bytes()).unwrap(); + } + + fn send_resize(&mut self, cols: u16, rows: u16) { + let line = format!("{}\n", serde_json::json!({ "resize": { "cols": cols, "rows": rows } })); + self.stdin.as_mut().unwrap().write_all(line.as_bytes()).unwrap(); + } + + fn stdout_text(&self) -> String { + String::from_utf8_lossy(&self.stdout.lock().unwrap()).into_owned() + } + + fn stderr_text(&self) -> String { + String::from_utf8_lossy(&self.stderr.lock().unwrap()).into_owned() + } + + fn wait_for_stdout(&mut self, needle: &str, what: &str) { + let deadline = Instant::now() + Duration::from_secs(20); + while Instant::now() < deadline { + if self.stdout_text().contains(needle) { + return; + } + if let Some(status) = self.child.try_wait().unwrap() { + panic!( + "pipe-io relay exited ({status}) before {what}\nstdout:\n{}\nstderr:\n{}", + self.stdout_text(), + String::from_utf8_lossy(&self.stderr.lock().unwrap()) + ); + } + std::thread::sleep(Duration::from_millis(25)); + } + panic!("timed out waiting for {what}\nstdout:\n{}", self.stdout_text()); + } + + fn wait_for_exit(&mut self) -> (i32, serde_json::Value) { + let deadline = Instant::now() + Duration::from_secs(20); + while Instant::now() < deadline { + if let Some(status) = self.child.try_wait().unwrap() { + // Let the drain threads observe EOF. + std::thread::sleep(Duration::from_millis(100)); + let stderr_text = + String::from_utf8_lossy(&self.stderr.lock().unwrap()).into_owned(); + let exit_line = stderr_text + .lines() + .rev() + .find(|line| line.trim_start().starts_with('{')) + .unwrap_or_else(|| { + panic!("pipe-io relay printed no JSON exit line\nstderr:\n{stderr_text}") + }) + .to_string(); + let value: serde_json::Value = serde_json::from_str(&exit_line).unwrap(); + return (status.code().unwrap_or(-1), value); + } + std::thread::sleep(Duration::from_millis(25)); + } + panic!("pipe-io relay did not exit\nstdout:\n{}", self.stdout_text()); + } +} + +#[cfg(unix)] +fn drain_into(mut source: impl Read, sink: std::sync::Arc>>) { + let mut buffer = [0u8; 4096]; + loop { + match source.read(&mut buffer) { + Ok(0) | Err(_) => return, + Ok(count) => sink.lock().unwrap().extend_from_slice(&buffer[..count]), + } + } +} + +/// Reaps the adoptable terminal hosts a deliberate daemon SIGKILL leaves +/// behind. The fixture's Drop treats any surviving host as a leak, but this +/// test kills the daemon on purpose (host adoption after a daemon death is +/// covered elsewhere); the orphaned hosts are part of the arranged scene. +#[cfg(unix)] +fn reap_orphaned_hosts_after_daemon_kill(server: &HeadlessServer) { + let host_root = cmux_tui_core::terminal_host_runtime::terminal_host_root(&server.state, "main"); + for pid in terminal_host_pids(&host_root) { + if let Ok(pid) = libc::pid_t::try_from(pid) { + // SAFETY: terminating a test-owned child process group member. + unsafe { + libc::kill(pid, libc::SIGKILL); + } + } + } + let deadline = Instant::now() + Duration::from_secs(10); + while terminal_host_pids(&host_root).iter().any(|pid| process_exists(*pid)) { + assert!(Instant::now() < deadline, "orphaned terminal hosts survived SIGKILL"); + std::thread::sleep(Duration::from_millis(25)); + } + let _ = fs::remove_dir_all(&host_root); +} + +#[cfg(unix)] +fn pipe_io_terminal(server: &HeadlessServer) -> String { + let created = json_cli(server, &["workspace", "create", "--name", "pipe-io"]); + assert_success(&created); + json_output(&created)["value"]["terminal_id"].as_str().unwrap().to_string() +} + +/// Count literal `^[[;R` runs, the form `cat -v` uses to +/// render a cursor-position report that reached the inner process. +#[cfg(unix)] +fn visible_cursor_position_reports(text: &str) -> usize { + let bytes = text.as_bytes(); + let mut count = 0; + let mut index = 0; + while let Some(offset) = text[index..].find("^[[") { + let mut cursor = index + offset + 3; + let mut digits_then_semicolon = false; + while cursor < bytes.len() && bytes[cursor].is_ascii_digit() { + cursor += 1; + } + if cursor < bytes.len() && bytes[cursor] == b';' { + cursor += 1; + let row_end = cursor; + while cursor < bytes.len() && bytes[cursor].is_ascii_digit() { + cursor += 1; + } + digits_then_semicolon = cursor > row_end; + } + if digits_then_semicolon && cursor < bytes.len() && bytes[cursor] == b'R' { + count += 1; + } + index += offset + 3; + } + count +} + +#[cfg(unix)] +#[test] +fn pipe_io_attach_streams_output_and_replays_on_reconnect() { + let server = HeadlessServer::start("pipe-io-basic"); + let terminal = pipe_io_terminal(&server); + + let mut first = PipeIoRelay::start(&server, &terminal, 80, 24); + first.send_input(b"printf 'PIPEIO-%s\\n' FIRST\n"); + first.wait_for_stdout("PIPEIO-FIRST", "live output of the typed marker"); + + // Dropping stdin tells the relay its embedder is gone: clean exit 0. + first.stdin = None; + let (code, exit) = first.wait_for_exit(); + assert_eq!(code, 0, "stdin EOF must be a clean detach, got {exit}"); + assert_eq!(exit["exit"]["reason"], "parent-closed"); + + // A fresh relay must serve the daemon replay before live bytes: the + // marker typed through the first relay is visible without retyping. + let mut second = PipeIoRelay::start(&server, &terminal, 80, 24); + second.wait_for_stdout("PIPEIO-FIRST", "replayed marker after reconnect"); +} + +#[cfg(unix)] +#[test] +fn pipe_io_exit_distinguishes_terminal_end_from_daemon_loss() { + let mut server = HeadlessServer::start("pipe-io-exits"); + let terminal = pipe_io_terminal(&server); + + let mut relay = PipeIoRelay::start(&server, &terminal, 80, 24); + relay.send_input(b"printf 'PIPEIO-%s\\n' READY\n"); + relay.wait_for_stdout("PIPEIO-READY", "shell readiness marker"); + let closed = json_cli(&server, &["terminal", &terminal, "close"]); + assert_success(&closed); + let (code, exit) = relay.wait_for_exit(); + assert_eq!(code, 0, "terminal close must exit 0, got {exit}"); + assert_eq!(exit["exit"]["reason"], "terminal-ended"); + + let second_terminal = pipe_io_terminal(&server); + let mut survivor = PipeIoRelay::start(&server, &second_terminal, 80, 24); + survivor.send_input(b"printf 'PIPEIO-%s\\n' TWO\n"); + survivor.wait_for_stdout("PIPEIO-TWO", "second shell readiness marker"); + server.child.kill().unwrap(); + let (code, exit) = survivor.wait_for_exit(); + assert_eq!(code, 2, "daemon loss must exit 2, got {exit}"); + assert_eq!(exit["exit"]["reason"], "daemon-lost"); + reap_orphaned_hosts_after_daemon_kill(&server); +} + +#[cfg(unix)] +#[test] +fn pipe_io_startup_connect_failure_reports_daemon_lost() { + // A daemon restart window (dead socket) must read as daemon-lost so + // the embedder keeps retrying; only unexplained failures make it stop. + let dir = unique_temp_dir("pipe-io-dead-socket"); + fs::create_dir_all(&dir).unwrap(); + let dead_socket = dir.join("mux.sock"); + let output = Command::new(bin()) + .args(["attach", "--socket"]) + .arg(&dead_socket) + .args(["--terminal", "term_0123456789abcdef0123456789abcdef", "--pipe-io"]) + .env_remove("CMUX_TUI_SOCKET") + .stdin(Stdio::null()) + .output() + .unwrap(); + assert_eq!( + output.status.code(), + Some(2), + "stderr: {}", + String::from_utf8_lossy(&output.stderr) + ); + let stderr = String::from_utf8_lossy(&output.stderr); + let exit_line = stderr.lines().rev().find(|line| line.trim_start().starts_with('{')).unwrap(); + let value: serde_json::Value = serde_json::from_str(exit_line).unwrap(); + assert_eq!(value["exit"]["reason"], "daemon-lost"); + let _ = fs::remove_dir_all(&dir); +} + +#[cfg(unix)] +#[test] +fn pipe_io_resize_reaches_the_daemon_pty() { + let server = HeadlessServer::start("pipe-io-resize"); + let terminal = pipe_io_terminal(&server); + + let mut relay = PipeIoRelay::start(&server, &terminal, 80, 24); + relay.send_input(b"printf 'PIPEIO-%s\\n' SIZED\n"); + relay.wait_for_stdout("PIPEIO-SIZED", "shell readiness marker"); + relay.send_resize(100, 30); + // The daemon acknowledges the resize before the PTY winsize necessarily + // lands; poll stty until the new geometry is visible. + let deadline = Instant::now() + Duration::from_secs(20); + loop { + relay.send_input(b"stty size\n"); + let poll_deadline = Instant::now() + Duration::from_secs(2); + while Instant::now() < poll_deadline { + if relay.stdout_text().contains("30 100") { + return; + } + std::thread::sleep(Duration::from_millis(50)); + } + assert!( + Instant::now() < deadline, + "timed out waiting for stty to report the resized PTY\nstdout:\n{}\nstderr:\n{}", + relay.stdout_text(), + relay.stderr_text() + ); + } +} + +#[cfg(unix)] +#[test] +fn pipe_io_inner_query_is_answered_exactly_once() { + let server = HeadlessServer::start("pipe-io-reply-authority"); + let terminal = pipe_io_terminal(&server); + + let mut relay = PipeIoRelay::start(&server, &terminal, 80, 24); + // The inner program emits a DSR cursor-position query and then renders + // its own stdin visibly. Exactly one authority (the daemon-side + // terminal) must answer; the relay and any embedder mirror stay silent. + relay.send_input(b"sh -c 'printf \"\\033[6n\"; exec cat -v'\n"); + let deadline = Instant::now() + Duration::from_secs(20); + let mut seen = 0; + while Instant::now() < deadline { + seen = visible_cursor_position_reports(&relay.stdout_text()); + if seen >= 1 { + // Give a duplicated reply time to surface before asserting. + std::thread::sleep(Duration::from_millis(500)); + seen = visible_cursor_position_reports(&relay.stdout_text()); + break; + } + std::thread::sleep(Duration::from_millis(25)); + } + assert_eq!( + seen, + 1, + "inner DSR query must be answered exactly once\nstdout:\n{}", + relay.stdout_text() + ); +} diff --git a/cmux.xcodeproj/project.pbxproj b/cmux.xcodeproj/project.pbxproj index c655cca32e6e..466347678f8d 100644 --- a/cmux.xcodeproj/project.pbxproj +++ b/cmux.xcodeproj/project.pbxproj @@ -2647,6 +2647,8 @@ C0DE71B10000000000000001 /* AppDelegate+AgentChatNotifications.swift in Sources E30760000000000000000002 /* TmuxWorkspacePaneOverlayView.swift in Sources */ = {isa = PBXBuildFile; fileRef = E30760000000000000000001 /* TmuxWorkspacePaneOverlayView.swift */; }; D3284001A1B2C3D4E5F60718 /* TraditionalChineseIMENumpadRegressionTests.swift in Sources */ = {isa = PBXBuildFile; fileRef = D3284002A1B2C3D4E5F60718 /* TraditionalChineseIMENumpadRegressionTests.swift */; }; AB7C4D21E6F90810 /* TuiBridgeReattachMouseUITests.swift in Sources */ = {isa = PBXBuildFile; fileRef = AB7C4D21E6F90811 /* TuiBridgeReattachMouseUITests.swift */; }; + 7A15A1DE0000000000000008 /* TuiManualIOPump.swift in Sources */ = {isa = PBXBuildFile; fileRef = 7A15A1DE0000000000000007 /* TuiManualIOPump.swift */; }; + 7A15A1DE000000000000000A /* TuiManualIOPumpTests.swift in Sources */ = {isa = PBXBuildFile; fileRef = 7A15A1DE0000000000000009 /* TuiManualIOPumpTests.swift */; }; 7A15A1DE0000000000000004 /* TuiTerminalAttachBridge.swift in Sources */ = {isa = PBXBuildFile; fileRef = 7A15A1DE0000000000000003 /* TuiTerminalAttachBridge.swift */; }; 7A15A1DE0000000000000002 /* TuiTerminalAttachPolicy.swift in Sources */ = {isa = PBXBuildFile; fileRef = 7A15A1DE0000000000000001 /* TuiTerminalAttachPolicy.swift */; }; 7A15A1DE0000000000000006 /* TuiTerminalAttachSpikeTests.swift in Sources */ = {isa = PBXBuildFile; fileRef = 7A15A1DE0000000000000005 /* TuiTerminalAttachSpikeTests.swift */; }; @@ -5519,6 +5521,8 @@ B8B056D80000000000000002 /* MobileHostIdentityTests.swift */ = {isa = PBXFileRef E30760000000000000000001 /* TmuxWorkspacePaneOverlayView.swift */ = {isa = PBXFileReference; lastKnownFileType = sourcecode.swift; path = TmuxWorkspacePaneOverlayView.swift; sourceTree = ""; }; D3284002A1B2C3D4E5F60718 /* TraditionalChineseIMENumpadRegressionTests.swift */ = {isa = PBXFileReference; lastKnownFileType = sourcecode.swift; path = TraditionalChineseIMENumpadRegressionTests.swift; sourceTree = ""; }; AB7C4D21E6F90811 /* TuiBridgeReattachMouseUITests.swift */ = {isa = PBXFileReference; lastKnownFileType = sourcecode.swift; path = TuiBridgeReattachMouseUITests.swift; sourceTree = ""; }; + 7A15A1DE0000000000000007 /* TuiManualIOPump.swift */ = {isa = PBXFileReference; lastKnownFileType = sourcecode.swift; path = TuiManualIOPump.swift; sourceTree = ""; }; + 7A15A1DE0000000000000009 /* TuiManualIOPumpTests.swift */ = {isa = PBXFileReference; lastKnownFileType = sourcecode.swift; path = TuiManualIOPumpTests.swift; sourceTree = ""; }; 7A15A1DE0000000000000003 /* TuiTerminalAttachBridge.swift */ = {isa = PBXFileReference; lastKnownFileType = sourcecode.swift; path = TuiTerminalAttachBridge.swift; sourceTree = ""; }; 7A15A1DE0000000000000001 /* TuiTerminalAttachPolicy.swift */ = {isa = PBXFileReference; lastKnownFileType = sourcecode.swift; path = TuiTerminalAttachPolicy.swift; sourceTree = ""; }; 7A15A1DE0000000000000005 /* TuiTerminalAttachSpikeTests.swift */ = {isa = PBXFileReference; lastKnownFileType = sourcecode.swift; path = TuiTerminalAttachSpikeTests.swift; sourceTree = ""; }; @@ -7704,6 +7708,7 @@ B8B056D80000000000000002 /* MobileHostIdentityTests.swift */ = {isa = PBXFileRef 50FE16F529BD6FBA4FBDECFC /* RemoteTmuxController.swift */, 7A15A1DE0000000000000001 /* TuiTerminalAttachPolicy.swift */, 7A15A1DE0000000000000003 /* TuiTerminalAttachBridge.swift */, + 7A15A1DE0000000000000007 /* TuiManualIOPump.swift */, 740600000000000000000001 /* RemoteTmuxController+Attach.swift */, 740600000000000000000010 /* RemoteTmuxAttachWindowTarget.swift */, 736200000000000000000001 /* RemoteTmuxController+Decisions.swift */, @@ -8116,6 +8121,7 @@ B8B056D80000000000000002 /* MobileHostIdentityTests.swift */ = {isa = PBXFileRef C13519000000000000000008 /* ClaudeConfigDirectoryPathTests.swift */, C65930010000000000000001 /* CrashDiagnosticSessionPolicyTests.swift */, 7A15A1DE0000000000000005 /* TuiTerminalAttachSpikeTests.swift */, + 7A15A1DE0000000000000009 /* TuiManualIOPumpTests.swift */, DEBDADADADADADADAD000004 /* DebugDogfoodCredentialResolverTests.swift */, C135190000000000000000B2 /* SidebarWorkspaceGroupConfigOpenerTests.swift */, A11EAE000000000000000001 /* WorkspaceAppearanceConfigResolutionTests.swift */, @@ -10989,6 +10995,7 @@ B8B056D80000000000000002 /* MobileHostIdentityTests.swift */ = {isa = PBXFileRef B35751000000000000000002 /* TmuxResumeParser.swift in Sources */, E30760050000000000000002 /* TmuxWorkspacePaneOverlayRenderState.swift in Sources */, E30760000000000000000002 /* TmuxWorkspacePaneOverlayView.swift in Sources */, + 7A15A1DE0000000000000008 /* TuiManualIOPump.swift in Sources */, 7A15A1DE0000000000000004 /* TuiTerminalAttachBridge.swift in Sources */, 7A15A1DE0000000000000002 /* TuiTerminalAttachPolicy.swift in Sources */, A5001501 /* UITestRecorder.swift in Sources */, @@ -12063,6 +12070,7 @@ B8B056D80000000000000002 /* MobileHostIdentityTests.swift */ = {isa = PBXFileRef A10507000000000000000001 /* TitleUpdateAmplificationRegressionTests.swift in Sources */, E30760060000000000000002 /* TmuxWorkspacePaneOverlayModelTests.swift in Sources */, D3284001A1B2C3D4E5F60718 /* TraditionalChineseIMENumpadRegressionTests.swift in Sources */, + 7A15A1DE000000000000000A /* TuiManualIOPumpTests.swift in Sources */, 7A15A1DE0000000000000006 /* TuiTerminalAttachSpikeTests.swift in Sources */, 468100000000000000000003 /* TypingHotPathRegressionTests.swift in Sources */, F2000000A1B2C3D4E5F60718 /* UpdatePillReleaseVisibilityTests.swift in Sources */, diff --git a/cmuxTests/TuiManualIOPumpTests.swift b/cmuxTests/TuiManualIOPumpTests.swift new file mode 100644 index 000000000000..3257c537fe2b --- /dev/null +++ b/cmuxTests/TuiManualIOPumpTests.swift @@ -0,0 +1,208 @@ +import CmuxTerminal +import Foundation +import Testing + +#if canImport(cmux_DEV) +@testable import cmux_DEV +#elseif canImport(cmux) +@testable import cmux +#endif + +/// Unit coverage for the manual-IO tui pump: relay exit classification, the +/// respawn decision, backoff pacing, the relay stdin wire format, the relay +/// argv, the overlay presentation mapping, and the cross-thread input +/// channel. All pure or pipe-local; no daemon and no app launch involved. +struct TuiManualIOPumpTests { + // MARK: - Relay exit classification + + @Test + func relayExitReadsTheFinalStderrJSONLine() { + // Config warnings and other noise may precede the exit line. + let stderrText = """ + warning: something unrelated + {"exit":{"reason":"terminal-ended"}} + """ + #expect(TuiManualIOPumpPolicy.relayExit(status: 0, stderrText: stderrText) == .terminalEnded) + #expect( + TuiManualIOPumpPolicy.relayExit( + status: 0, + stderrText: "{\"exit\":{\"reason\":\"parent-closed\"}}\n" + ) == .parentClosed + ) + #expect( + TuiManualIOPumpPolicy.relayExit( + status: 2, + stderrText: "{\"exit\":{\"reason\":\"daemon-lost\"}}\n" + ) == .daemonLost + ) + // Exit 2 without the relay's reason line is a usage error (e.g. a + // binary without --pipe-io), not an endlessly-retryable outage. + #expect(TuiManualIOPumpPolicy.relayExit(status: 2, stderrText: nil) == .failure) + #expect( + TuiManualIOPumpPolicy.relayExit(status: 2, stderrText: "cmux: unknown argument") + == .failure + ) + } + + @Test + func relayExitTreatsUnknownStatusAsFailureAndBareZeroAsEnded() { + // Exit 0 with no reason line: the relay contract says 0 means "do + // not respawn", so a missing line must not turn into a retry loop. + #expect(TuiManualIOPumpPolicy.relayExit(status: 0, stderrText: nil) == .terminalEnded) + #expect(TuiManualIOPumpPolicy.relayExit(status: 0, stderrText: "garbage") == .terminalEnded) + #expect(TuiManualIOPumpPolicy.relayExit(status: 1, stderrText: "usage: ...") == .failure) + #expect(TuiManualIOPumpPolicy.relayExit(status: -1, stderrText: nil) == .failure) + } + + @Test + func nextActionRespawnsOnlyForRecoverableExits() { + #expect(TuiManualIOPumpPolicy.nextAction(after: .terminalEnded) == .end) + #expect(TuiManualIOPumpPolicy.nextAction(after: .daemonLost) == .retry) + #expect(TuiManualIOPumpPolicy.nextAction(after: .failure) == .retry) + // The pump closed stdin itself (teardown/respawn): never a state + // transition of its own. + #expect(TuiManualIOPumpPolicy.nextAction(after: .parentClosed) == .ignore) + } + + // MARK: - Backoff + + @Test + func retryDelayGrowsAndCaps() { + #expect(TuiManualIOPumpPolicy.retryDelay(attempt: 1) == .milliseconds(500)) + #expect(TuiManualIOPumpPolicy.retryDelay(attempt: 2) == .seconds(1)) + #expect(TuiManualIOPumpPolicy.retryDelay(attempt: 6) == .seconds(16)) + #expect(TuiManualIOPumpPolicy.retryDelay(attempt: 7) == .seconds(30)) + #expect(TuiManualIOPumpPolicy.retryDelay(attempt: 100) == .seconds(30)) + #expect(TuiManualIOPumpPolicy.retryDelay(attempt: 0) == .milliseconds(500)) + } + + // MARK: - Relay stdin wire format + + @Test + func inputLineIsBase64JSONWithTrailingNewline() throws { + let line = TuiManualIOPumpPolicy.inputLine(bytes: Data("hi".utf8)) + let text = String(decoding: line, as: UTF8.self) + #expect(text.hasSuffix("\n")) + let object = try JSONSerialization.jsonObject( + with: Data(text.dropLast().utf8) + ) as? [String: Any] + #expect(object?["input"] as? String == "aGk=") + } + + @Test + func resizeLineCarriesClampedGrid() throws { + let line = TuiManualIOPumpPolicy.resizeLine(cols: 120, rows: 0) + let text = String(decoding: line, as: UTF8.self) + #expect(text.hasSuffix("\n")) + let object = try JSONSerialization.jsonObject( + with: Data(text.dropLast().utf8) + ) as? [String: Any] + let resize = object?["resize"] as? [String: Any] + #expect(resize?["cols"] as? Int == 120) + #expect(resize?["rows"] as? Int == 1) + } + + // MARK: - Relay argv + + @Test + func relayArgumentsAreDirectArgvWithoutShellQuoting() { + let arguments = TuiManualIOPumpPolicy.relayArguments( + sessionName: "cmux-tag", + terminalID: "term_abc", + cols: 100, + rows: 30 + ) + #expect(arguments == [ + "attach", "--session", "cmux-tag", "--terminal", "term_abc", + "--pipe-io", "--cols", "100", "--rows", "30", + ]) + } + + @Test + func resyncResetIsFullResetPlusScrollbackErase() { + #expect(TuiManualIOPumpPolicy.resyncReset == Data("\u{1B}c\u{1B}[3J".utf8)) + } + + // MARK: - Overlay presentation + + @Test + func overlayIsHiddenWhileConnectingOrLive() { + #expect(TuiManualIOPumpPolicy.overlayPresentation(state: .connecting) == nil) + #expect(TuiManualIOPumpPolicy.overlayPresentation(state: .live) == nil) + } + + @Test + func reconnectingOverlayShowsProgressAndRetryButton() throws { + let presentation = try #require( + TuiManualIOPumpPolicy.overlayPresentation(state: .reconnecting(attempt: 3)) + ) + #expect(presentation.showsProgress) + #expect(presentation.showsReconnectButton) + #expect(presentation.detail.contains("3")) + } + + @Test + func endedOverlayIsInformationalOnly() throws { + let presentation = try #require( + TuiManualIOPumpPolicy.overlayPresentation(state: .ended) + ) + #expect(!presentation.showsProgress) + #expect(!presentation.showsReconnectButton) + } + + @Test + func failedOverlayOffersManualRetryOnly() throws { + let presentation = try #require( + TuiManualIOPumpPolicy.overlayPresentation(state: .failed) + ) + #expect(!presentation.showsProgress) + #expect(presentation.showsReconnectButton) + } + + // MARK: - Input channel + + @Test + func inputChannelWritesToTheLiveHandleAndDropsWhenPaused() throws { + let pipe = Pipe() + let channel = TuiManualIOInputChannel() + channel.setHandle(pipe.fileHandleForWriting) + channel.send(Data("first\n".utf8)) + // Pause (relay died): input is dropped, never queued for a future + // relay, so stale keystrokes cannot replay into a resynced shell. + channel.setHandle(nil) + channel.send(Data("dropped\n".utf8)) + channel.setHandle(pipe.fileHandleForWriting) + channel.send(Data("second\n".utf8)) + channel.closeHandle() + + let data = pipe.fileHandleForReading.readDataToEndOfFile() + #expect(String(decoding: data, as: UTF8.self) == "first\nsecond\n") + } + + // MARK: - Registry + + @MainActor + @Test + func registryReplacesAndRemovesPumps() { + let registry = TuiManualIOPumpRegistry() + let surfaceID = UUID() + let first = TuiManualIOPump( + binaryPath: "/nonexistent", + sessionName: "cmux-test", + terminalID: "term_1", + environment: [:] + ) + let second = TuiManualIOPump( + binaryPath: "/nonexistent", + sessionName: "cmux-test", + terminalID: "term_2", + environment: [:] + ) + registry.register(first, surfaceID: surfaceID) + #expect(registry.pump(forSurfaceID: surfaceID) === first) + registry.register(second, surfaceID: surfaceID) + #expect(registry.pump(forSurfaceID: surfaceID) === second) + registry.stopAndRemove(surfaceID: surfaceID) + #expect(registry.pump(forSurfaceID: surfaceID) == nil) + } +}