Skip to content
Merged
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
Original file line number Diff line number Diff line change
Expand Up @@ -75,14 +75,13 @@ struct RemoteDaemonRPCClientTimeoutIsolationTests {

@Test("a timed-out PTY attach does not wait for a blocked cancellation write")
func timedOutPTYAttachBoundsCancellationWrite() throws {
let executable = try makeTransport()
let executable = try makeIdleTransport()
defer {
try? FileManager.default.removeItem(
at: URL(fileURLWithPath: executable).deletingLastPathComponent()
)
}

let stalledAttachRead = DispatchSemaphore(value: 0)
let unexpectedTermination = DispatchSemaphore(value: 0)
let client = RemoteDaemonRPCClient(
configuration: configuration(),
Expand All @@ -99,86 +98,39 @@ struct RemoteDaemonRPCClientTimeoutIsolationTests {
}
defer { client.stop() }
client.transportExecutableOverride = executable

try client.start()
_ = try client.attachPTY(
sessionID: "existing-session",
attachmentID: "existing-attachment",
cols: 80,
rows: 24,
command: nil,
requireExisting: true,
queue: .global()
) { event in
if case .data(let data) = event, data == Data("attach-read".utf8) {
stalledAttachRead.signal()
}
}

// Occupy the transport writer before the attach deadline handler
// runs. The writer must already be held when the cancellation is
// queued; holding it in response to a daemon event instead races the
// attach RPC deadline, because nothing bounds that event's round trip
// on a loaded host. The write stays held until the test ends, so the
// cancellation deadline, not a write completion, decides the outcome.
let writeBlockEntered = DispatchSemaphore(value: 0)
let releaseWrite = DispatchSemaphore(value: 0)
let writeBlockQueue = DispatchQueue(label: "com.cmux.tests.remote-daemon.block-cancellation-write")
writeBlockQueue.async {
guard stalledAttachRead.wait(timeout: .now() + 2) == .success else {
releaseWrite.signal()
return
}
client.writeQueue.async {
writeBlockEntered.signal()
releaseWrite.wait()
}
}
// One-shot failure safety keeps a regressed synchronous cancellation
// from stranding the test process; ordinary cleanup is deterministic.
let cleanupFired = DispatchSemaphore(value: 0)
let cleanupTimer = DispatchSource.makeTimerSource(queue: writeBlockQueue)
cleanupTimer.schedule(deadline: .now() + 5)
cleanupTimer.setEventHandler {
cleanupFired.signal()
releaseWrite.signal()
}
cleanupTimer.resume()
defer {
cleanupTimer.cancel()
releaseWrite.signal()
client.writeQueue.async {
writeBlockEntered.signal()
releaseWrite.wait()
}
defer { releaseWrite.signal() }
try #require(writeBlockEntered.wait(timeout: .now() + 5) == .success)

let callFinished = DispatchSemaphore(value: 0)
let callTimedOut = DispatchSemaphore(value: 0)
let unexpectedCallResult = DispatchSemaphore(value: 0)
let callQueue = DispatchQueue(label: "com.cmux.tests.remote-daemon.call-with-blocked-cancellation")
callQueue.async {
defer { callFinished.signal() }
do {
_ = try client.call(
method: "pty.attach",
params: [
"session_id": "stalled-session",
"attachment_id": "stalled-attachment",
"client_attachment_token": "stalled-token",
],
timeout: 1
)
unexpectedCallResult.signal()
} catch {
let nsError = error as NSError
if nsError.domain == "cmux.remote.daemon.rpc", nsError.code == 11 {
callTimedOut.signal()
} else {
unexpectedCallResult.signal()
}
}
// This is the handler a timed-out `pty.attach` invokes before it
// throws the timeout error; the first test covers that wiring.
let cancellationReturned = DispatchSemaphore(value: 0)
DispatchQueue(label: "com.cmux.tests.remote-daemon.timed-out-attach").async {
client.sendPTYAttachCancellation(
requestID: 2,
attachParams: [
"session_id": "stalled-session",
"attachment_id": "stalled-attachment",
"client_attachment_token": "stalled-token",
]
)
cancellationReturned.signal()
}

#expect(writeBlockEntered.wait(timeout: .now() + 2) == .success)
#expect(callFinished.wait(timeout: .now() + 6) == .success)
#expect(callTimedOut.wait(timeout: .now()) == .success)
#expect(unexpectedCallResult.wait(timeout: .now()) == .timedOut)
#expect(cleanupFired.wait(timeout: .now()) == .timedOut)
// The cancellation write is deliberately held until the cleanup timer
// fires. Give the child-process teardown its own budget after that
// release; heavily loaded hosted runners can otherwise report a false
// failure even though the transport exits correctly.
#expect(cancellationReturned.wait(timeout: .now() + 5) == .success)
#expect(unexpectedTermination.wait(timeout: .now() + 10) == .success)
}

Expand Down Expand Up @@ -253,4 +205,25 @@ struct RemoteDaemonRPCClientTimeoutIsolationTests {
chmod(scriptURL.path, 0o755)
return scriptURL.path
}

/// A daemon that answers `hello` and then only drains its input.
private func makeIdleTransport() throws -> String {
let directory = FileManager.default.temporaryDirectory
.appendingPathComponent("cmux-remote-daemon-idle-\(UUID().uuidString)", isDirectory: true)
try FileManager.default.createDirectory(at: directory, withIntermediateDirectories: true)
let scriptURL = directory.appendingPathComponent("fake-ssh-idle")
let script = """
#!/bin/sh
if IFS= read -r line; then
id=$(printf '%s\\n' "$line" | sed -n 's/.*"id":\\([0-9][0-9]*\\).*/\\1/p')
printf '{"id":%s,"ok":true,"result":{"version":"0.64.22","capabilities":["proxy.stream.push"]}}\\n' "$id"
else
exit 1
fi
while IFS= read -r _line; do :; done
"""
try Data(script.utf8).write(to: scriptURL, options: .atomic)
chmod(scriptURL.path, 0o755)
return scriptURL.path
}
}
Loading