diff --git a/CLI/CMUXCLI+SSHStartupScripts.swift b/CLI/CMUXCLI+SSHStartupScripts.swift index 31a87d76d973..0281cbfb398e 100644 --- a/CLI/CMUXCLI+SSHStartupScripts.swift +++ b/CLI/CMUXCLI+SSHStartupScripts.swift @@ -130,7 +130,13 @@ extension CMUXCLI { " if {[string length $cmux_buffer] > 0} { send_user -- $cmux_buffer }", " cmux_relay_session", " }", - " eof { set status [wait]; exit [lindex $status 3] }", + // A session that authenticates and ends inside this window still + // owes the user its output: flush what expect buffered before exiting. + " eof {", + " catch { send_user -- $expect_out(buffer) }", + " set status [wait]", + " exit [lindex $status 3]", + " }", " }", "}", "expect {", @@ -217,7 +223,13 @@ extension CMUXCLI { " if {[string length $cmux_buffer] > 0} { send_user -- $cmux_buffer }", " cmux_relay_session", " }", - " eof { set status [wait]; exit [lindex $status 3] }", + // A session that authenticates and ends inside this window still + // owes the user its output: flush what expect buffered before exiting. + " eof {", + " catch { send_user -- $expect_out(buffer) }", + " set status [wait]", + " exit [lindex $status 3]", + " }", " }", "}", "expect {", diff --git a/cmuxTests/CLIVMDevTests.swift b/cmuxTests/CLIVMDevTests.swift index ab6ebe7a3958..47e3353045c3 100644 --- a/cmuxTests/CLIVMDevTests.swift +++ b/cmuxTests/CLIVMDevTests.swift @@ -152,6 +152,7 @@ extension CLINotifyProcessIntegrationRegressionTests { _ name: String, arguments: [String], home: URL, + expectsSocketConnection: Bool = true, respond: @escaping @Sendable (_ method: String, _ params: [String: Any]) -> [String: Any]? ) throws -> (result: ProcessRunResult, log: VMDevRequestLog) { let cliPath = try bundledCLIPath() @@ -163,14 +164,37 @@ extension CLINotifyProcessIntegrationRegressionTests { Darwin.close(listenerFD) unlink(socketPath) } - let serverHandled = startVMDevMock(listenerFD: listenerFD, state: state, log: log, respond: respond) + let serverHandled: XCTestExpectation? + if expectsSocketConnection { + serverHandled = startVMDevMock(listenerFD: listenerFD, state: state, log: log, respond: respond) + } else { + // Invalid input is rejected before connecting, so the mock's expectation + // would only be fulfilled by the listener closing, which happens after the + // wait. Keep the socket listening instead and check its backlog once the + // process has exited: an accidental connection is still observable. + let flags = fcntl(listenerFD, F_GETFL) + XCTAssertGreaterThanOrEqual(flags, 0) + XCTAssertEqual(fcntl(listenerFD, F_SETFL, flags | O_NONBLOCK), 0) + serverHandled = nil + } let result = runProcess( executablePath: cliPath, arguments: arguments, environment: vmDevEnvironment(socketPath: socketPath, home: home), timeout: 60 ) - wait(for: [serverHandled], timeout: 60) + if let serverHandled { + wait(for: [serverHandled], timeout: 60) + } else { + let unexpectedClient = Darwin.accept(listenerFD, nil, nil) + let acceptError = errno + if unexpectedClient >= 0 { + Darwin.close(unexpectedClient) + XCTFail("Invalid vm dev input must not open a socket connection") + } else { + XCTAssertTrue(acceptError == EAGAIN || acceptError == EWOULDBLOCK) + } + } XCTAssertFalse(result.timedOut, "\(arguments) timed out: \(result.stderr)") return (result, log) } @@ -643,7 +667,8 @@ extension CLINotifyProcessIntegrationRegressionTests { let (invalidLayout, invalidLog) = try runVMDev( "vm-dev-bad-layout", arguments: ["vm", "dev", "brave-otter", fixture.project.path, "--layout", badLayout], - home: fixture.home + home: fixture.home, + expectsSocketConnection: false ) { _, _ in nil } XCTAssertEqual(invalidLayout.status, 2, "stdout=\(invalidLayout.stdout) stderr=\(invalidLayout.stderr)") XCTAssertTrue(invalidLayout.stderr.contains("invalid layout document"), invalidLayout.stderr) @@ -653,7 +678,8 @@ extension CLINotifyProcessIntegrationRegressionTests { let (missingFolder, missingLog) = try runVMDev( "vm-dev-missing-folder", arguments: ["vm", "dev", "brave-otter", fixture.project.appendingPathComponent("nope").path], - home: fixture.home + home: fixture.home, + expectsSocketConnection: false ) { _, _ in nil } XCTAssertNotEqual(missingFolder.status, 0) XCTAssertTrue(missingFolder.stderr.contains("no such folder"), missingFolder.stderr) @@ -662,7 +688,8 @@ extension CLINotifyProcessIntegrationRegressionTests { let (badPort, badPortLog) = try runVMDev( "vm-dev-bad-port", arguments: ["vm", "dev", "brave-otter", fixture.project.path, "--port", "http"], - home: fixture.home + home: fixture.home, + expectsSocketConnection: false ) { _, _ in nil } XCTAssertNotEqual(badPort.status, 0) XCTAssertTrue(badPort.stderr.contains("--port must be a port number"), badPort.stderr) @@ -671,7 +698,8 @@ extension CLINotifyProcessIntegrationRegressionTests { let (noMachine, noMachineLog) = try runVMDev( "vm-dev-no-machine", arguments: ["vm", "dev"], - home: fixture.home + home: fixture.home, + expectsSocketConnection: false ) { _, _ in nil } XCTAssertNotEqual(noMachine.status, 0) XCTAssertTrue(noMachine.stderr.contains("Usage: cmux vm dev "), noMachine.stderr) diff --git a/cmuxTests/CloudPlacementCoordinatorTests.swift b/cmuxTests/CloudPlacementCoordinatorTests.swift index bf2014acc3f5..705262eb44e8 100644 --- a/cmuxTests/CloudPlacementCoordinatorTests.swift +++ b/cmuxTests/CloudPlacementCoordinatorTests.swift @@ -383,6 +383,11 @@ struct CloudPlacementCoordinatorTests { // graph install; an empty view list means the move must project a // fresh tab rather than move a stale one. catalog.upsert(term, from: provider) + // The daemon's projection reply is fenced by its mutation cursor + // (`CmuxTuiSnapshotParser.placedTab` requires one). When the lane drains it + // reconciles against the installed graph above, which predates the new tab; + // the cursor is what keeps that older graph from clearing the new tab ID. + provider.projectCursor = CloudVMCursor(generation: "g", revision: 21) catalog.record(SurfaceProjection(resource: term.id, workspaceID: viewer, panelID: panel, remoteWorkspaceID: "ws_main", remoteTabID: "tab_gone")) let state = try #require(CmuxTuiSnapshotParser.state(fromSnapshot: stateSnapshot, machine: Self.machine)) catalog.reconcileCloudRemoteState(machine: Self.machine, state: state, observation: .current) diff --git a/cmuxTests/CloudPlacementTestProvider.swift b/cmuxTests/CloudPlacementTestProvider.swift index 9d1935b956ae..5131d1598478 100644 --- a/cmuxTests/CloudPlacementTestProvider.swift +++ b/cmuxTests/CloudPlacementTestProvider.swift @@ -19,6 +19,8 @@ final class CloudPlacementTestProvider: SurfaceProvider, SurfacePlacementSyncing var beforeMaterialization: (() async throws -> Void)? var refreshCount = 0 var moveCursor: CloudVMCursor? + /// The daemon cursor a projection reply carries. The real reply always has one. + var projectCursor: CloudVMCursor? var workspaceRenames: [String] = [] var tabRenames: [String] = [] @@ -62,7 +64,7 @@ final class CloudPlacementTestProvider: SurfaceProvider, SurfacePlacementSyncing func projectTerminal(_ id: SurfaceResourceID, intoRemoteWorkspace remoteWorkspaceID: String) async throws -> SurfaceRemotePlacement { try await beforeMutation?() projected.append((id.key, remoteWorkspaceID)) - return SurfaceRemotePlacement(workspaceID: remoteWorkspaceID, tabID: "tab_projected") + return SurfaceRemotePlacement(workspaceID: remoteWorkspaceID, tabID: "tab_projected", cursor: projectCursor) } func closeRemoteTab(id: String, inRemoteWorkspace remoteWorkspaceID: String) async throws { events.append("close:" + id) diff --git a/cmuxTests/GhosttyConfigTests.swift b/cmuxTests/GhosttyConfigTests.swift index c95d00f2b4ce..367eefb22fd6 100644 --- a/cmuxTests/GhosttyConfigTests.swift +++ b/cmuxTests/GhosttyConfigTests.swift @@ -5396,8 +5396,11 @@ final class ZshShellIntegrationHandoffTests: XCTestCase { _CMUX_TTY_REPORTED=1 _CMUX_PORTS_LAST_RUN=-999 _cmux_precmd - repeat 20; do - [[ -s "\(logPath.path)" ]] && break + # precmd reports the prompt state and kicks the port scan from two + # independent background children. Wait for the kick itself, not + # for whichever line lands first. + repeat 200; do + /usr/bin/grep -q 'surface.ports_kick' "\(logPath.path)" && break sleep 0.05 done cat "\(logPath.path)" @@ -5541,8 +5544,10 @@ final class ZshShellIntegrationHandoffTests: XCTestCase { _CMUX_TTY_REPORTED=1 _CMUX_PORTS_LAST_RUN=-999 _cmux_prompt_command - for _cmux_i in $(seq 1 20); do - [ -s "\(logPath.path)" ] && break + # The prompt hook may send other relay RPCs from separate background + # children. Wait for the kick itself, not for whichever line lands first. + for _cmux_i in $(seq 1 200); do + /usr/bin/grep -q 'surface.ports_kick' "\(logPath.path)" && break sleep 0.05 done cat "\(logPath.path)" diff --git a/cmuxTests/TabManagerSessionSnapshotTests.swift b/cmuxTests/TabManagerSessionSnapshotTests.swift index 0c0065dc9a07..b0fd628e2fc0 100644 --- a/cmuxTests/TabManagerSessionSnapshotTests.swift +++ b/cmuxTests/TabManagerSessionSnapshotTests.swift @@ -2205,27 +2205,31 @@ final class TabManagerSessionSnapshotTests: XCTestCase { defer { catalog.unregister(machine: remote.machine) } let restored = TabManager() - restored.restoreSessionSnapshot(snapshot) - let restoredWorkspace = try XCTUnwrap(restored.tabs.first { $0.customTitle == workspaceTitle }) - let restoredPanelId = try XCTUnwrap(restoredWorkspace.panels.first { $0.value is TerminalPanel }?.key) - - // Until the machine's provider reports the terminal, the pane is a reserved Cloud - // pane (#12675): it has no live projection yet, but its persisted remote identity - // stays staged so the restore is never mistaken for a plain local shell. - XCTAssertNil(catalog.projection(forPanel: restoredPanelId)) - XCTAssertEqual(catalog.projectionIncludingPendingRestore(forPanel: restoredPanelId)?.resource, remote) - XCTAssertEqual(restoredWorkspace.terminalPanel(for: restoredPanelId)?.surface.ioMode, .manualMirror) - catalog.upsert(SurfaceResource(id: remote, title: "root@\(machineId)", detail: "/root", lifecycle: .running, agent: nil, remoteWorkspace: nil, port: nil, url: nil), from: provider) - XCTAssertEqual(catalog.projection(forPanel: restoredPanelId)?.resource, remote, "the restored pane re-links to the remote terminal") - XCTAssertEqual(catalog.projection(forPanel: restoredPanelId)?.workspaceID, restoredWorkspace.id) - - // The projection round-trips through the next save with the live panel id. - let resaved = restored.sessionSnapshot(includeScrollback: false) - XCTAssertEqual( - resaved.workspaces.first { $0.customTitle == workspaceTitle }?.surfaceProjections, - [SurfaceProjectionRecord(panelID: restoredPanelId, resource: remote)] - ) - restored.closeWorkspace(restoredWorkspace, recordHistory: false) + // `SurfaceCatalog.shared` relinks a restored projection only into a workspace the + // app resolves as live (#13196), so the restored window must be registered. + try LiveWorkspaceFixture.withAppRegistration(of: restored) { + restored.restoreSessionSnapshot(snapshot) + let restoredWorkspace = try XCTUnwrap(restored.tabs.first { $0.customTitle == workspaceTitle }) + let restoredPanelId = try XCTUnwrap(restoredWorkspace.panels.first { $0.value is TerminalPanel }?.key) + + // Until the machine's provider reports the terminal, the pane is a reserved Cloud + // pane (#12675): it has no live projection yet, but its persisted remote identity + // stays staged so the restore is never mistaken for a plain local shell. + XCTAssertNil(catalog.projection(forPanel: restoredPanelId)) + XCTAssertEqual(catalog.projectionIncludingPendingRestore(forPanel: restoredPanelId)?.resource, remote) + XCTAssertEqual(restoredWorkspace.terminalPanel(for: restoredPanelId)?.surface.ioMode, .manualMirror) + catalog.upsert(SurfaceResource(id: remote, title: "root@\(machineId)", detail: "/root", lifecycle: .running, agent: nil, remoteWorkspace: nil, port: nil, url: nil), from: provider) + XCTAssertEqual(catalog.projection(forPanel: restoredPanelId)?.resource, remote, "the restored pane re-links to the remote terminal") + XCTAssertEqual(catalog.projection(forPanel: restoredPanelId)?.workspaceID, restoredWorkspace.id) + + // The projection round-trips through the next save with the live panel id. + let resaved = restored.sessionSnapshot(includeScrollback: false) + XCTAssertEqual( + resaved.workspaces.first { $0.customTitle == workspaceTitle }?.surfaceProjections, + [SurfaceProjectionRecord(panelID: restoredPanelId, resource: remote)] + ) + restored.closeWorkspace(restoredWorkspace, recordHistory: false) + } } func testWorkspaceSnapshotWithoutSurfaceProjectionsDecodesAndRestoresLocalOnly() throws { diff --git a/tests/test_vm_scp.py b/tests/test_vm_scp.py index b8017d48094f..4fec69a2c26d 100644 --- a/tests/test_vm_scp.py +++ b/tests/test_vm_scp.py @@ -10,7 +10,6 @@ import json import os import pty -import select from pathlib import Path import shlex import socket @@ -25,6 +24,18 @@ def run(argv, **kwargs): return subprocess.run(argv, text=True, capture_output=True, timeout=45, **kwargs) +def drain_terminal(master, sink): + """Collects a pty's output until every process holding its slave has exited.""" + while True: + try: + chunk = os.read(master, 65536) + except OSError: + return + if not chunk: + return + sink.extend(chunk) + + def main(cli): with tempfile.TemporaryDirectory(prefix="scp-", dir="/tmp") as raw: root = Path(raw) @@ -270,28 +281,46 @@ def push(local, remote, *args): wrong_host_key = fails destination = f"human-{tty}-{fails}" master, slave = pty.openpty() if tty else (None, None) + terminal_output = bytearray() + reader = None try: - result = subprocess.run( + process = subprocess.Popen( [cli, "vm", "push", "test-vm", str(payload), destination], env=env, stdin=subprocess.DEVNULL, stdout=subprocess.PIPE, - stderr=slave if tty else subprocess.PIPE, timeout=30, + stderr=slave if tty else subprocess.PIPE, ) - stderr = result.stderr or b"" - if master is not None: - while select.select([master], [], [], 0)[0]: - stderr += os.read(master, 65536) + if slave is not None: + # The CLI owns the slave now. A terminal's output queue + # holds only about 1 KiB and the host-key failure report is + # longer, so the master must be drained while the CLI runs + # or its stderr writes block until the timeout. + os.close(slave) + slave = None + reader = threading.Thread(target=drain_terminal, args=(master, terminal_output), daemon=True) + reader.start() + try: + stdout, piped_stderr = process.communicate(timeout=30) + except subprocess.TimeoutExpired: + process.kill() + process.communicate() + raise + if reader is not None: + # Returns once every process holding the slave has exited. + reader.join(timeout=10) + assert not reader.is_alive(), "terminal stderr stayed open after the CLI exited" + stderr = piped_stderr if piped_stderr is not None else bytes(terminal_output) finally: if slave is not None: os.close(slave) if master is not None: os.close(master) wrong_host_key = False - assert result.returncode == (1 if fails else 0), stderr + assert process.returncode == (1 if fails else 0), stderr if fails: assert b"Host key verification failed" in stderr, stderr assert b"Cloud diagnostic reference:" in stderr, stderr - assert b"Pushed" not in result.stdout, result.stdout + assert b"Pushed" not in stdout, stdout assert not (guest / destination).exists() else: - assert b"Pushed" in result.stdout and destination.encode() in result.stdout, result.stdout + assert b"Pushed" in stdout and destination.encode() in stdout, stdout assert stderr == b"", stderr assert (guest / destination).read_bytes() == payload.read_bytes() print("PASS SCP human output and errors over pipes and terminals", flush=True)