Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
16 changes: 14 additions & 2 deletions CLI/CMUXCLI+SSHStartupScripts.swift
Original file line number Diff line number Diff line change
Expand Up @@ -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 {",
Expand Down Expand Up @@ -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 {",
Expand Down
40 changes: 34 additions & 6 deletions cmuxTests/CLIVMDevTests.swift
Original file line number Diff line number Diff line change
Expand Up @@ -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()
Expand All @@ -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)
}
Expand Down Expand Up @@ -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)
Expand All @@ -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)
Expand All @@ -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)
Expand All @@ -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 <machine>"), noMachine.stderr)
Expand Down
5 changes: 5 additions & 0 deletions cmuxTests/CloudPlacementCoordinatorTests.swift
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down
4 changes: 3 additions & 1 deletion cmuxTests/CloudPlacementTestProvider.swift
Original file line number Diff line number Diff line change
Expand Up @@ -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] = []

Expand Down Expand Up @@ -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)
Expand Down
13 changes: 9 additions & 4 deletions cmuxTests/GhosttyConfigTests.swift
Original file line number Diff line number Diff line change
Expand Up @@ -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)"
Expand Down Expand Up @@ -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)"
Expand Down
46 changes: 25 additions & 21 deletions cmuxTests/TabManagerSessionSnapshotTests.swift
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down
49 changes: 39 additions & 10 deletions tests/test_vm_scp.py
Original file line number Diff line number Diff line change
Expand Up @@ -10,7 +10,6 @@
import json
import os
import pty
import select
from pathlib import Path
import shlex
import socket
Expand All @@ -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)
Expand Down Expand Up @@ -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)
Expand Down
Loading