fix(macos): backport native pipe drain ownership (#155268)

* fix(macos): backport native pipe drain ownership

Backport #136214 to extended-stable/2026.8.33.

* test(macos): avoid blocking stderr drain fixture
This commit is contained in:
Dallin Romney 2026-09-21 17:18:12 -07:00 • committed by GitHub
parent 550b052515
commit 7a0701150c
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
16 changed files with 992 additions and 546 deletions

View file

@ -1,4 +1,3 @@
import Darwin
import Foundation
import Observation
import OSLog
@ -37,10 +36,7 @@ final class CookieSyncManager: NSObject {
@ObservationIgnored private var reconcileTask: Task<Void, Never>?
@ObservationIgnored private var retryTask: Task<Void, Never>?
@ObservationIgnored private var process: Process?
@ObservationIgnored private var stdoutPipe: Pipe?
@ObservationIgnored private var stderrPipe: Pipe?
@ObservationIgnored private var stdoutSource: DispatchSourceRead?
@ObservationIgnored private var stderrSource: DispatchSourceRead?
@ObservationIgnored private var readers: [PipeReadStream] = []
@ObservationIgnored private var startupWatchdog: DispatchSourceTimer?
@ObservationIgnored private var stdoutBuffer = Data()
@ObservationIgnored private var processGeneration: UUID?
@ -208,6 +204,10 @@ final class CookieSyncManager: NSObject {
let process = Process()
let stdoutPipe = Pipe()
let stderrPipe = Pipe()
defer {
try? stdoutPipe.fileHandleForReading.close()
try? stderrPipe.fileHandleForReading.close()
}
let generation = UUID()
let arguments = [
"browser",
@ -235,25 +235,33 @@ final class CookieSyncManager: NSObject {
environment["OPENCLAW_GATEWAY_PASSWORD"] = password
}
process.environment = environment
process.terminationHandler = { [weak self] process in
let terminationStatus = process.terminationStatus
Task { @MainActor [weak self] in
self?.childTerminated(generation: generation, status: terminationStatus)
}
}
self.process = process
self.stdoutPipe = stdoutPipe
self.stderrPipe = stderrPipe
self.processGeneration = generation
self.runningIntent = intent
self.stdoutBuffer.removeAll(keepingCapacity: true)
do {
// Consume on the reader's serial executor so termination cannot
// overtake status chunks waiting for a separate actor hop.
let readers = try [
PipeReadStream(handle: stdoutPipe.fileHandleForReading, queue: .main, onData: { [weak self] data in
MainActor.assumeIsolated { self?.consumeStdout(data, generation: generation) }
}),
PipeReadStream(handle: stderrPipe.fileHandleForReading, queue: .main, onData: { [weak self] data in
MainActor.assumeIsolated { self?.consumeStderr(data, generation: generation) }
}),
]
self.readers = readers
process.terminationHandler = { [weak self] process in
let terminationStatus = process.terminationStatus
Task { @MainActor [weak self] in
for reader in readers {
await reader.finish()
}
self?.childTerminated(generation: generation, status: terminationStatus)
}
}
try process.run()
try? stdoutPipe.fileHandleForWriting.close()
try? stderrPipe.fileHandleForWriting.close()
self.installReadSources(stdoutPipe: stdoutPipe, stderrPipe: stderrPipe, generation: generation)
self.installStartupWatchdog(generation: generation)
self.state = .running
self.logger.info("cookie sync started for \(intent.domains.count, privacy: .public) domain(s)")
@ -264,30 +272,6 @@ final class CookieSyncManager: NSObject {
}
}
private func installReadSources(stdoutPipe: Pipe, stderrPipe: Pipe, generation: UUID) {
let stdoutDescriptor = stdoutPipe.fileHandleForReading.fileDescriptor
let stdoutSource = DispatchSource.makeReadSource(fileDescriptor: stdoutDescriptor, queue: self.queue)
stdoutSource.setEventHandler { [weak self] in
let data = Self.readAvailable(fileDescriptor: stdoutDescriptor, byteCount: stdoutSource.data)
Task { @MainActor [weak self] in
self?.consumeStdout(data, generation: generation)
}
}
self.stdoutSource = stdoutSource
stdoutSource.resume()
let stderrDescriptor = stderrPipe.fileHandleForReading.fileDescriptor
let stderrSource = DispatchSource.makeReadSource(fileDescriptor: stderrDescriptor, queue: self.queue)
stderrSource.setEventHandler { [weak self] in
let data = Self.readAvailable(fileDescriptor: stderrDescriptor, byteCount: stderrSource.data)
Task { @MainActor [weak self] in
self?.consumeStderr(data, generation: generation)
}
}
self.stderrSource = stderrSource
stderrSource.resume()
}
private func installStartupWatchdog(generation: UUID) {
let timer = DispatchSource.makeTimerSource(queue: self.queue)
timer.schedule(deadline: .now() + 5)
@ -308,11 +292,6 @@ final class CookieSyncManager: NSObject {
private func consumeStdout(_ data: Data, generation: UUID) {
guard self.processGeneration == generation else { return }
guard !data.isEmpty else {
self.stdoutSource?.cancel()
self.stdoutSource = nil
return
}
self.stdoutBuffer.append(data)
if self.stdoutBuffer.count > 64 * 1024 {
self.stdoutBuffer.removeFirst(self.stdoutBuffer.count - 64 * 1024)
@ -332,11 +311,6 @@ final class CookieSyncManager: NSObject {
private func consumeStderr(_ data: Data, generation: UUID) {
guard self.processGeneration == generation else { return }
guard !data.isEmpty else {
self.stderrSource?.cancel()
self.stderrSource = nil
return
}
// Preserve lossy decoding so malformed CLI bytes do not hide useful diagnostics.
// swiftlint:disable:next optional_data_string_conversion
let message = String(decoding: data, as: UTF8.self)
@ -378,35 +352,16 @@ final class CookieSyncManager: NSObject {
private func stopChild(nextState: State) {
self.startupWatchdog?.cancel()
self.startupWatchdog = nil
self.stdoutSource?.cancel()
self.stdoutSource = nil
self.stderrSource?.cancel()
self.stderrSource = nil
self.readers.forEach { $0.close() }
self.readers.removeAll()
let process = self.process
self.process = nil
self.processGeneration = nil
self.runningIntent = nil
try? self.stdoutPipe?.fileHandleForReading.close()
try? self.stdoutPipe?.fileHandleForWriting.close()
try? self.stderrPipe?.fileHandleForReading.close()
try? self.stderrPipe?.fileHandleForWriting.close()
self.stdoutPipe = nil
self.stderrPipe = nil
self.stdoutBuffer.removeAll(keepingCapacity: false)
if process?.isRunning == true {
process?.terminate()
}
self.state = nextState
}
private nonisolated static func readAvailable(fileDescriptor: Int32, byteCount: UInt) -> Data {
let count = max(1, min(Int(byteCount), 64 * 1024))
var data = Data(count: count)
let bytesRead = data.withUnsafeMutableBytes { buffer in
Darwin.read(fileDescriptor, buffer.baseAddress, count)
}
guard bytesRead > 0 else { return Data() }
data.removeSubrange(bytesRead..<data.count)
return data
}
}

View file

@ -15,155 +15,6 @@ struct CuaDriverProcessLaunch: Sendable {
let environment: [String: String]
}
enum CuaDriverStderrEvent: Equatable, Sendable {
case notice(String)
case error(String)
}
final class CuaDriverStderrRelay: @unchecked Sendable {
static let managedModeNotice =
"""
CUA embedded driver running in managed unrestricted mode; \
OpenClaw command arming and pairing are the authorization boundary.
"""
private static let dangerBannerPrefix = "DANGER: Cua Driver is running in unrestricted mode"
private static let maximumBufferedBytes = 32 * 1024
private static let readChunkBytes = 4 * 1024
let pipe = Pipe()
private let lock = NSLock()
private let emit: @Sendable (CuaDriverStderrEvent) -> Void
private var buffer = Data()
private var started = false
private var stopped = false
private var emittedManagedModeNotice = false
init(emit: @escaping @Sendable (CuaDriverStderrEvent) -> Void) {
self.emit = emit
}
func startReading() {
let shouldStart = self.lock.withLock {
guard !self.started, !self.stopped else { return false }
self.started = true
return true
}
guard shouldStart else { return }
self.pipe.fileHandleForReading.readabilityHandler = { [weak self] handle in
guard let self else { return }
let data = handle.readSafely(upToCount: Self.readChunkBytes)
guard !data.isEmpty else {
self.stop()
return
}
self.consume(data)
}
}
func reportManagedMode() {
let shouldEmit = self.lock.withLock {
guard !self.stopped, !self.emittedManagedModeNotice else { return false }
self.emittedManagedModeNotice = true
return true
}
if shouldEmit {
self.emit(.notice(Self.managedModeNotice))
}
}
func stop() {
let tail = self.lock.withLock { () -> Data? in
guard !self.stopped else { return nil }
self.stopped = true
defer { self.buffer.removeAll(keepingCapacity: false) }
return self.buffer.isEmpty ? nil : self.buffer
}
self.pipe.fileHandleForReading.readabilityHandler = nil
try? self.pipe.fileHandleForReading.close()
try? self.pipe.fileHandleForWriting.close()
if let tail {
self.forward(tail)
}
}
private func consume(_ data: Data) {
let lines = self.lock.withLock { () -> [Data] in
guard !self.stopped else { return [] }
self.buffer.append(data)
if self.buffer.count > Self.maximumBufferedBytes {
self.buffer = Data(self.buffer.suffix(Self.maximumBufferedBytes))
}
var lines: [Data] = []
while let newline = self.buffer.firstIndex(of: 0x0A) {
lines.append(Data(self.buffer[..<newline]))
self.buffer.removeSubrange(...newline)
}
return lines
}
lines.forEach(self.forward)
}
private func forward(_ data: Data) {
let line = (String(bytes: data, encoding: .utf8) ?? "")
.trimmingCharacters(in: .whitespacesAndNewlines)
guard !line.isEmpty, !line.hasPrefix(Self.dangerBannerPrefix) else { return }
self.emit(.error(line))
}
}
@MainActor
protocol CuaDriverProcessControlling: AnyObject {
var isRunning: Bool { get }
/// Spawned daemon pid. OpenClaw records this itself because `serve` ignores
/// `--pid-file` and writes only the machine-global default path.
var processIdentifier: pid_t { get }
func closeLiveness()
func terminate()
func forceKill()
}
@MainActor
private final class FoundationCuaDriverProcess: CuaDriverProcessControlling {
let process: Process
private let livenessPipe: Pipe
private let stderrRelay: CuaDriverStderrRelay
init(process: Process, livenessPipe: Pipe, stderrRelay: CuaDriverStderrRelay) {
self.process = process
self.livenessPipe = livenessPipe
self.stderrRelay = stderrRelay
}
deinit {
try? self.livenessPipe.fileHandleForWriting.close()
self.stderrRelay.stop()
}
var isRunning: Bool {
self.process.isRunning
}
var processIdentifier: pid_t {
self.process.processIdentifier
}
func closeLiveness() {
try? self.livenessPipe.fileHandleForWriting.close()
}
func terminate() {
guard self.process.isRunning else { return }
self.process.terminate()
}
func forceKill() {
guard self.process.isRunning else { return }
_ = Darwin.kill(self.process.processIdentifier, SIGKILL)
}
}
struct CuaDriverSocketDirectory: Equatable, Sendable {
let url: URL
let socketPath: String
@ -583,11 +434,16 @@ final class CuaDriverHostCoordinator {
process.standardOutput = FileHandle.nullDevice
process.standardError = stderrRelay.pipe
process.terminationHandler = { terminated in
stderrRelay.stop()
onTermination(terminated.terminationStatus)
let status = terminated.terminationStatus
// Readiness and shutdown can release the process wrapper before
// this callback finishes; the handler owns stderr through its drain.
Task {
await stderrRelay.finishReading()
onTermination(status)
}
}
stderrRelay.startReading()
do {
try stderrRelay.startReading()
try process.run()
} catch {
stderrRelay.stop()
@ -596,8 +452,7 @@ final class CuaDriverHostCoordinator {
stderrRelay.reportManagedMode()
return FoundationCuaDriverProcess(
process: process,
livenessPipe: livenessPipe,
stderrRelay: stderrRelay)
livenessPipe: livenessPipe)
}
private static func waitUntilStopped(

View file

@ -0,0 +1,148 @@
import Darwin
import Foundation
@MainActor
protocol CuaDriverProcessControlling: AnyObject {
var isRunning: Bool { get }
/// Spawned daemon pid. OpenClaw records this itself because `serve` ignores
/// `--pid-file` and writes only the machine-global default path.
var processIdentifier: pid_t { get }
func closeLiveness()
func terminate()
func forceKill()
}
@MainActor
final class FoundationCuaDriverProcess: CuaDriverProcessControlling {
let process: Process
private let livenessPipe: Pipe
init(process: Process, livenessPipe: Pipe) {
self.process = process
self.livenessPipe = livenessPipe
}
deinit {
try? self.livenessPipe.fileHandleForWriting.close()
}
var isRunning: Bool {
self.process.isRunning
}
var processIdentifier: pid_t {
self.process.processIdentifier
}
func closeLiveness() {
try? self.livenessPipe.fileHandleForWriting.close()
}
func terminate() {
guard self.process.isRunning else { return }
self.process.terminate()
}
func forceKill() {
guard self.process.isRunning else { return }
_ = Darwin.kill(self.process.processIdentifier, SIGKILL)
}
}
enum CuaDriverStderrEvent: Equatable, Sendable {
case notice(String)
case error(String)
}
final class CuaDriverStderrRelay: @unchecked Sendable {
static let managedModeNotice =
"""
CUA embedded driver running in managed unrestricted mode; \
OpenClaw command arming and pairing are the authorization boundary.
"""
private static let dangerBannerPrefix = "DANGER: Cua Driver is running in unrestricted mode"
private static let maximumBufferedBytes = 32 * 1024
private static let readChunkBytes = 4 * 1024
let pipe = Pipe()
private let lock = NSLock()
private let emit: @Sendable (CuaDriverStderrEvent) -> Void
private var buffer = Data()
private var reader: PipeReadStream?
private var stopped = false
private var emittedManagedModeNotice = false
init(emit: @escaping @Sendable (CuaDriverStderrEvent) -> Void) {
self.emit = emit
}
func startReading() throws {
try self.lock.withLock {
guard self.reader == nil, !self.stopped else { return }
self.reader = try PipeReadStream(
handle: self.pipe.fileHandleForReading,
maximumChunkBytes: Self.readChunkBytes,
onData: { [weak self] in self?.consume($0) },
onClose: { [weak self] in self?.stop() })
}
}
func reportManagedMode() {
let shouldEmit = self.lock.withLock {
guard !self.stopped, !self.emittedManagedModeNotice else { return false }
self.emittedManagedModeNotice = true
return true
}
if shouldEmit {
self.emit(.notice(Self.managedModeNotice))
}
}
func finishReading() async {
let reader = self.lock.withLock { self.reader }
await reader?.finish()
}
func stop() {
let tail = self.lock.withLock { () -> Data? in
guard !self.stopped else { return nil }
self.stopped = true
defer { self.buffer.removeAll(keepingCapacity: false) }
return self.buffer.isEmpty ? nil : self.buffer
}
self.reader?.close()
try? self.pipe.fileHandleForReading.close()
try? self.pipe.fileHandleForWriting.close()
if let tail {
self.forward(tail)
}
}
private func consume(_ data: Data) {
let lines = self.lock.withLock { () -> [Data] in
guard !self.stopped else { return [] }
self.buffer.append(data)
if self.buffer.count > Self.maximumBufferedBytes {
self.buffer = Data(self.buffer.suffix(Self.maximumBufferedBytes))
}
var lines: [Data] = []
while let newline = self.buffer.firstIndex(of: 0x0A) {
lines.append(Data(self.buffer[..<newline]))
self.buffer.removeSubrange(...newline)
}
return lines
}
lines.forEach(self.forward)
}
private func forward(_ data: Data) {
// The byte cap can split UTF-8; keep the remaining diagnostic.
// swiftlint:disable:next optional_data_string_conversion
let line = String(decoding: data, as: UTF8.self)
.trimmingCharacters(in: .whitespacesAndNewlines)
guard !line.isEmpty, !line.hasPrefix(Self.dangerBannerPrefix) else { return }
self.emit(.error(line))
}
}

View file

@ -0,0 +1,12 @@
import Darwin
import Foundation
extension FileHandle {
/// Marks a pipe/socket write end so a vanished reader fails the write with a
/// thrown EPIPE instead of raising SIGPIPE, which kills the whole process.
/// Required on every write end whose reader is another process that can exit.
@discardableResult
func disableSIGPIPE() -> Bool {
fcntl(self.fileDescriptor, F_SETNOSIGPIPE, 1) != -1
}
}

View file

@ -1,37 +0,0 @@
import Darwin
import Foundation
extension FileHandle {
/// Marks a pipe/socket write end so a vanished reader fails the write with a
/// thrown EPIPE instead of raising SIGPIPE, which kills the whole process.
/// Required on every write end whose reader is another process that can exit.
@discardableResult
func disableSIGPIPE() -> Bool {
fcntl(self.fileDescriptor, F_SETNOSIGPIPE, 1) != -1
}
/// Reads until EOF using the throwing FileHandle API and returns empty `Data` on failure.
///
/// Important: Avoid legacy, non-throwing FileHandle read APIs (e.g. `readDataToEndOfFile()` and
/// `availableData`). They can raise Objective-C exceptions when the handle is closed/invalid, which
/// will abort the process.
func readToEndSafely() -> Data {
do {
return try self.readToEnd() ?? Data()
} catch {
return Data()
}
}
/// Reads up to `count` bytes using the throwing FileHandle API and returns empty `Data` on failure/EOF.
///
/// Important: Use this instead of `availableData` in callbacks like `readabilityHandler` to avoid
/// Objective-C exceptions terminating the process.
func readSafely(upToCount count: Int) -> Data {
do {
return try self.read(upToCount: count) ?? Data()
} catch {
return Data()
}
}
}

View file

@ -101,17 +101,13 @@ final class MacNodeHostWorker: MacNodeHostWorking, @unchecked Sendable {
private var process: ManagedProcess?
private var processCleanupTask: Task<Void, Never>?
private var stdinPipe: Pipe?
private var stdoutPipe: Pipe?
private var stderrPipe: Pipe?
private var stdoutSource: DispatchSourceRead?
private var stderrSource: DispatchSourceRead?
private var readers: [PipeReadStream] = []
private var processGeneration: UUID?
private var launchedWorker: MacNodeHostWorkerLaunch?
private var stdoutBuffer = Data()
// Bounded head of worker stderr. CLI startup failures print their cause
// first; without this the operator-visible error is just "exited(1)".
private var stderrHead = ""
private static let maxStderrHeadLength = 700
private var stderrCapture = PipeTextCapture(characterLimit: 700, retention: .head)
private var manifest: MacNodeHostManifest?
private var route: GatewayNodeSessionRoute?
private var routeAuthorityGeneration: UInt64 = 0
@ -373,10 +369,38 @@ final class MacNodeHostWorker: MacNodeHostWorking, @unchecked Sendable {
let stdinPipe = Pipe()
let stdoutPipe = Pipe()
let stderrPipe = Pipe()
defer {
try? stdoutPipe.fileHandleForReading.close()
try? stderrPipe.fileHandleForReading.close()
}
guard stdinPipe.fileHandleForWriting.disableSIGPIPE() else {
self.finishStartLocked(.failure(WorkerError.unavailable(reason: "could not protect worker input pipe")))
return
}
let processGeneration = UUID()
let stderrCapture = self.stderrCapture
let consumeStderr: @Sendable (Data, Bool) -> Void = { [weak self] data, atEOF in
guard let self, self.processGeneration == processGeneration, self.processCleanupTask == nil else { return }
let message = stderrCapture.append(data, atEOF: atEOF)
if !message.isEmpty { self.logger.error("node-host worker stderr: \(message, privacy: .private)") }
}
do {
self.readers = try [
PipeReadStream(handle: stdoutPipe.fileHandleForReading, queue: self.queue, onData: { [weak self] data in
guard let self, self.processGeneration == processGeneration,
self.processCleanupTask == nil else { return }
self.consumeStdoutLocked(data)
}),
PipeReadStream(
handle: stderrPipe.fileHandleForReading,
queue: self.queue,
onData: { consumeStderr($0, false) },
onClose: { consumeStderr(Data(), true) }),
]
} catch {
self.finishStartLocked(.failure(WorkerError.unavailable(reason: "could not read worker output")))
return
}
var environment = ProcessInfo.processInfo.environment.filter { key, _ in
!CuaDriverWorkerEnvironment.inheritedFamilyPrefixes.contains { key.hasPrefix($0) }
}
@ -387,9 +411,6 @@ final class MacNodeHostWorker: MacNodeHostWorking, @unchecked Sendable {
environment["OPENCLAW_NODE_EXEC_FALLBACK"] = "0"
self.launchedWorker = launch
self.stdinPipe = stdinPipe
self.stdoutPipe = stdoutPipe
self.stderrPipe = stderrPipe
let processGeneration = UUID()
self.processGeneration = processGeneration
let timer = DispatchSource.makeTimerSource(queue: self.queue)
@ -431,58 +452,14 @@ final class MacNodeHostWorker: MacNodeHostWorking, @unchecked Sendable {
generation: UUID)
{
guard self.processGeneration == generation, self.processCleanupTask == nil else { return }
guard started,
let process = self.process,
let stdoutPipe = self.stdoutPipe,
let stderrPipe = self.stderrPipe
else {
guard started, let process else {
self.stopLocked(reason: "worker launch failed")
return
}
let stdoutSource = DispatchSource.makeReadSource(
fileDescriptor: stdoutPipe.fileHandleForReading.fileDescriptor,
queue: self.queue)
stdoutSource.setEventHandler { [weak self] in
guard let self, self.processGeneration == generation else { return }
let data = Self.readAvailable(
fileDescriptor: stdoutPipe.fileHandleForReading.fileDescriptor,
byteCount: stdoutSource.data)
if data.isEmpty {
self.stdoutSource?.cancel()
} else {
self.consumeStdoutLocked(data)
}
}
self.stdoutSource = stdoutSource
stdoutSource.resume()
let stderrSource = DispatchSource.makeReadSource(
fileDescriptor: stderrPipe.fileHandleForReading.fileDescriptor,
queue: self.queue)
stderrSource.setEventHandler { [weak self] in
guard let self, self.processGeneration == generation else { return }
let data = Self.readAvailable(
fileDescriptor: stderrPipe.fileHandleForReading.fileDescriptor,
byteCount: stderrSource.data)
guard !data.isEmpty else {
self.stderrSource?.cancel()
return
}
if let message = String(data: data, encoding: .utf8)?
.trimmingCharacters(in: .whitespacesAndNewlines),
!message.isEmpty
{
self.logger.error("node-host worker stderr: \(message, privacy: .private)")
if self.stderrHead.count < Self.maxStderrHeadLength {
self.stderrHead.append(self.stderrHead.isEmpty ? message : "\n" + message)
self.stderrHead = String(self.stderrHead.prefix(Self.maxStderrHeadLength))
}
}
}
self.stderrSource = stderrSource
stderrSource.resume()
Task { [weak self, completionTask = process.completionTask] in
let status = await completionTask.value
// Retire the route before draining queued worker messages; unlike
// diagnostic-only pipes, stdout can request privileged operations.
self?.queue.async { [weak self] in
guard let self,
self.processGeneration == generation,
@ -716,8 +693,8 @@ final class MacNodeHostWorker: MacNodeHostWorking, @unchecked Sendable {
// A worker that dies before its ready manifest still needs its stderr
// surfaced: the raw exit status alone cannot explain a CLI bootstrap
// refusal (missing runtime, incompatible state database, bad install).
let diagnostic = self.stderrHead.nonEmpty
self.stderrHead = ""
let diagnostic = self.stderrCapture.snapshot().nonEmpty
self.stderrCapture = PipeTextCapture(characterLimit: 700, retention: .head)
self.startTimer?.cancel()
self.startTimer = nil
self.launchedWorker = nil
@ -728,6 +705,8 @@ final class MacNodeHostWorker: MacNodeHostWorking, @unchecked Sendable {
self.finishStartLocked(.failure(WorkerError.unavailable(reason: reason, diagnostic: diagnostic)))
}
if let processCleanupTask = self.processCleanupTask { return processCleanupTask }
let readers = self.readers
self.readers.removeAll()
let pending = self.invokeContinuations
self.invokeContinuations.removeAll()
self.pendingInvokeControls.removeAll()
@ -745,25 +724,23 @@ final class MacNodeHostWorker: MacNodeHostWorking, @unchecked Sendable {
return nil
}
let cleanupTask = Task { [weak self] in
// Keep draining through TERM cleanup: closing the pipes early can
// interrupt the child's shutdown handler with SIGPIPE.
await process.terminate()
readers.forEach { $0.close() }
for reader in readers {
await reader.finish()
}
await withCheckedContinuation { continuation in
guard let self else {
continuation.resume()
return
}
self.queue.async {
self.stdoutSource?.cancel()
self.stdoutSource = nil
self.stderrSource?.cancel()
self.stderrSource = nil
try? self.stdinPipe?.fileHandleForWriting.close()
try? self.stdoutPipe?.fileHandleForReading.close()
try? self.stderrPipe?.fileHandleForReading.close()
self.process = nil
self.processCleanupTask = nil
self.stdinPipe = nil
self.stdoutPipe = nil
self.stderrPipe = nil
self.processGeneration = nil
continuation.resume()
}
@ -796,15 +773,4 @@ final class MacNodeHostWorker: MacNodeHostWorking, @unchecked Sendable {
guard JSONSerialization.isValidJSONObject(object) else { return nil }
return try? JSONSerialization.data(withJSONObject: object)
}
private static func readAvailable(fileDescriptor: Int32, byteCount: UInt) -> Data {
let count = max(1, min(Int(byteCount), 64 * 1024))
var data = Data(count: count)
let bytesRead = data.withUnsafeMutableBytes { buffer in
Darwin.read(fileDescriptor, buffer.baseAddress, count)
}
guard bytesRead > 0 else { return Data() }
data.removeSubrange(bytesRead..<data.count)
return data
}
}

View file

@ -0,0 +1,92 @@
import Darwin
import Foundation
final class PipeReadStream: @unchecked Sendable {
// Darwin FIONREAD is _IOR('f', 127, int); Swift cannot import that C macro.
private static let bytesAvailableRequest: UInt = 0x4004_667F
private let source: DispatchSourceRead
private let queue: DispatchQueue
private let maximumChunkBytes: Int
private let onData: @Sendable (Data) -> Void
private let completion: Task<Void, Never>
init(
handle: FileHandle,
maximumChunkBytes: Int = 64 * 1024,
queue: DispatchQueue = DispatchQueue(label: "ai.openclaw.pipe.read"),
onData: @escaping @Sendable (Data) -> Void,
onClose: @escaping @Sendable () -> Void = {}) throws
{
let descriptor = fcntl(handle.fileDescriptor, F_DUPFD_CLOEXEC, 0)
guard descriptor >= 0 else {
throw NSError(domain: NSPOSIXErrorDomain, code: Int(errno))
}
let flags = fcntl(descriptor, F_GETFL)
guard flags >= 0, fcntl(descriptor, F_SETFL, flags | O_NONBLOCK) == 0 else {
let error = NSError(domain: NSPOSIXErrorDomain, code: Int(errno))
_ = Darwin.close(descriptor)
throw error
}
let (closed, finished) = AsyncStream<Void>.makeStream()
// A cancelled waiter must still join cleanup without cancelling other waiters.
self.completion = Task { for await _ in closed {} }
self.queue = queue
self.maximumChunkBytes = maximumChunkBytes
self.onData = onData
self.source = DispatchSource.makeReadSource(fileDescriptor: descriptor, queue: queue)
// Only dispatch cleanup closes the duplicate: cancellation may overlap an
// active callback, and the original FileHandle can already be closed.
self.source.setCancelHandler {
_ = Darwin.close(descriptor)
onClose()
finished.finish()
}
self.source.setEventHandler { [weak self] in
guard let self else { return }
_ = self.read(maximumBytes: max(1, Int(self.source.data)))
}
self.source.resume()
}
deinit { self.close() }
func close() {
self.source.cancel()
}
func finish() async {
self.queue.async {
defer { self.close() }
guard !self.source.isCancelled else { return }
// Drain the exit-time snapshot, not EOF: inherited writers can stay
// open or keep producing after the owned child has terminated.
var available: Int32 = 0
guard ioctl(Int32(self.source.handle), Self.bytesAvailableRequest, &available) == 0 else { return }
while available > 0 {
let count = self.read(maximumBytes: Int(available))
guard count > 0 else { return }
available -= Int32(count)
}
}
await self.completion.value
}
private func read(maximumBytes: Int) -> Int {
guard !self.source.isCancelled else { return 0 }
let count = min(maximumBytes, self.maximumChunkBytes)
var data = Data(count: count)
var bytesRead: Int
repeat {
bytesRead = data.withUnsafeMutableBytes {
Darwin.read(Int32(self.source.handle), $0.baseAddress, count)
}
} while bytesRead < 0 && errno == EINTR && !self.source.isCancelled
guard bytesRead > 0 else {
if bytesRead == 0 || errno != EAGAIN { self.close() }
return 0
}
data.count = bytesRead
self.onData(data)
return bytesRead
}
}

View file

@ -0,0 +1,53 @@
import Foundation
final class PipeTextCapture: @unchecked Sendable {
enum Retention {
case head, tail
}
private let characterLimit: Int
private let retention: Retention
private let lock = NSLock()
private var text = ""
private var pending = Data()
init(characterLimit: Int, retention: Retention) {
self.characterLimit = characterLimit
self.retention = retention
}
func append(_ chunk: Data, atEOF: Bool = false) -> String {
self.lock.withLock {
self.pending.append(chunk)
var end = self.pending.endIndex
// UTF-8's 2–4-byte leaders can leave at most three bytes unfinished.
// Keep only that suffix; the standard decoder repairs malformed text.
if !atEOF,
let start = self.pending.indices.suffix(4).last(where: { !UTF8.isContinuation(self.pending[$0]) }),
(0xC2...0xF4).contains(self.pending[start]),
end - start < (~self.pending[start]).leadingZeroBitCount
{
end = start
}
// swiftlint:disable:next optional_data_string_conversion
let message = String(decoding: self.pending[..<end], as: UTF8.self)
.trimmingCharacters(in: .whitespacesAndNewlines)
self.pending.removeSubrange(..<end)
self.text = self.retaining(message)
return message
}
}
func snapshot() -> String {
self.lock.withLock { self.text.trimmingCharacters(in: .whitespacesAndNewlines) }
}
private func retaining(_ addition: String) -> String {
// Head retention is final once full. Separate records so later combining
// marks cannot extend a retained Character across read callbacks.
guard !addition.isEmpty,
self.retention != .head || self.text.count < self.characterLimit else { return self.text }
let text = self.text.isEmpty ? addition : self.text + "\n" + addition
return String(self.retention == .head ? text.prefix(self.characterLimit) : text.suffix(self.characterLimit))
}
}

View file

@ -28,51 +28,25 @@ final class RemotePortTunnel: @unchecked Sendable {
let processIdentifier: pid_t
private let process: ManagedProcess
private let stderrHandle: FileHandle?
private let stderrReader: PipeReadStream
private let guardianReceipt: PortGuardian.Record
private final class StderrCapture: @unchecked Sendable {
private let lock = NSLock()
private var text = ""
private let limit = 4096
func append(_ chunk: String) {
let trimmed = chunk.trimmingCharacters(in: .whitespacesAndNewlines)
guard !trimmed.isEmpty else { return }
self.lock.lock()
defer { self.lock.unlock() }
if !self.text.isEmpty {
self.text += "\n"
}
self.text += trimmed
if self.text.count > self.limit {
self.text = String(self.text.suffix(self.limit))
}
}
func snapshot() -> String {
self.lock.lock()
defer { self.lock.unlock() }
return self.text.trimmingCharacters(in: .whitespacesAndNewlines)
}
}
private init(
process: ManagedProcess,
processIdentifier: pid_t,
localPort: UInt16?,
stderrHandle: FileHandle?,
stderrReader: PipeReadStream,
guardianReceipt: PortGuardian.Record)
{
self.process = process
self.processIdentifier = processIdentifier
self.localPort = localPort
self.stderrHandle = stderrHandle
self.stderrReader = stderrReader
self.guardianReceipt = guardianReceipt
}
deinit {
Self.cleanupStderr(self.stderrHandle)
self.stderrReader.close()
let receipt = self.guardianReceipt
// deinit cannot wait. Leave the receipt durable until a later sweep proves
// the child exited; deleting it after TERM alone can orphan a resistant SSH.
@ -82,7 +56,7 @@ final class RemotePortTunnel: @unchecked Sendable {
func terminate() async {
await self.process.terminate()
Self.cleanupStderr(self.stderrHandle)
await self.stderrReader.finish()
// Finish retiring this receipt before a replacement spawn reserves the ledger.
await PortGuardian.shared.removeRecord(self.guardianReceipt)
}
@ -139,30 +113,25 @@ final class RemotePortTunnel: @unchecked Sendable {
let pipe = Pipe()
let stderrHandle = pipe.fileHandleForReading
let stderrWriter = pipe.fileHandleForWriting
let stderrCapture = StderrCapture()
let stderrCapture = PipeTextCapture(characterLimit: 4096, retention: .tail)
// Consume stderr so ssh cannot block if it logs.
stderrHandle.readabilityHandler = { handle in
let data = handle.readSafely(upToCount: 64 * 1024)
guard !data.isEmpty else {
// EOF (or read failure): stop monitoring to avoid spinning on a closed pipe.
Self.cleanupStderr(handle)
return
}
guard let line = String(data: data, encoding: .utf8)?
.trimmingCharacters(in: .whitespacesAndNewlines),
!line.isEmpty
else { return }
stderrCapture.append(line)
defer { try? stderrHandle.close() }
let consumeStderr: @Sendable (Data, Bool) -> Void = { data, atEOF in
let line = stderrCapture.append(data, atEOF: atEOF)
guard !line.isEmpty else { return }
Self.logger.error("ssh tunnel stderr: \(line, privacy: .public)")
}
let stderrReader = try PipeReadStream(
handle: stderrHandle,
onData: { consumeStderr($0, false) },
onClose: { consumeStderr(Data(), true) })
let spawnPreparation: PortGuardian.SpawnPreparation
do {
// Legacy reconciliation can inspect many live processes. Complete it
// before spawn so a crash during migration cannot orphan this SSH child.
spawnPreparation = try await PortGuardian.shared.prepareForTunnelSpawn()
} catch {
Self.cleanupStderr(stderrHandle)
stderrReader.close()
throw NSError(
domain: "RemotePortTunnel",
code: 5,
@ -195,7 +164,7 @@ final class RemotePortTunnel: @unchecked Sendable {
// child before releasing its reservation or closing inherited handles.
await process.terminate(gracefully: false)
await PortGuardian.shared.cancelTunnelSpawn(spawnPreparation)
Self.cleanupStderr(stderrHandle)
stderrReader.close()
throw error
}
@ -214,7 +183,7 @@ final class RemotePortTunnel: @unchecked Sendable {
// Keep the reservation exclusive until this exact child is reaped.
// Only then may another operation migrate or open the ledger.
await PortGuardian.shared.cancelTunnelSpawn(spawnPreparation)
Self.cleanupStderr(stderrHandle)
stderrReader.close()
throw NSError(
domain: "RemotePortTunnel",
code: 5,
@ -229,11 +198,11 @@ final class RemotePortTunnel: @unchecked Sendable {
process: process,
processIdentifier: processIdentifier,
localPort: localPort,
stderrHandle: stderrHandle,
stderrReader: stderrReader,
stderrCapture: stderrCapture)
} catch {
await process.terminate()
Self.cleanupStderr(stderrHandle)
stderrReader.close()
await PortGuardian.shared.removeRecord(receipt)
throw error
}
@ -242,7 +211,7 @@ final class RemotePortTunnel: @unchecked Sendable {
process: process,
processIdentifier: processIdentifier,
localPort: localPort,
stderrHandle: stderrHandle,
stderrReader: stderrReader,
guardianReceipt: receipt)
}
@ -250,13 +219,16 @@ final class RemotePortTunnel: @unchecked Sendable {
process: ManagedProcess,
processIdentifier: pid_t,
localPort: UInt16,
stderrHandle: FileHandle,
stderrCapture: StderrCapture) async throws
stderrReader: PipeReadStream,
stderrCapture: PipeTextCapture) async throws
{
let deadline = Date().addingTimeInterval(6)
repeat {
if !process.isRunning {
let stderr = Self.drainStderr(stderrHandle, captured: stderrCapture.snapshot())
// The reader owns the entire pipe; wait for its final bytes instead
// of starting a competing read after the child exits.
await stderrReader.finish()
let stderr = stderrCapture.snapshot()
let msg = stderr.isEmpty ? "ssh tunnel exited before listening" : "ssh tunnel failed: \(stderr)"
throw NSError(domain: "RemotePortTunnel", code: 4, userInfo: [NSLocalizedDescriptionKey: msg])
}
@ -452,39 +424,6 @@ final class RemotePortTunnel: @unchecked Sendable {
}
#endif
private static func cleanupStderr(_ handle: FileHandle?) {
guard let handle else { return }
Self.cleanupStderr(handle)
}
private static func cleanupStderr(_ handle: FileHandle) {
if handle.readabilityHandler != nil {
handle.readabilityHandler = nil
}
try? handle.close()
}
private static func drainStderr(_ handle: FileHandle, captured: String) -> String {
handle.readabilityHandler = nil
defer { try? handle.close() }
do {
let data = try handle.readToEnd() ?? Data()
let remaining = String(data: data, encoding: .utf8)?
.trimmingCharacters(in: .whitespacesAndNewlines) ?? ""
if captured.isEmpty {
return remaining
}
if remaining.isEmpty {
return captured
}
return captured + "\n" + remaining
} catch {
self.logger.debug("Failed to drain ssh stderr: \(error, privacy: .public)")
return captured
}
}
#if SWIFT_PACKAGE
static func _testPortIsFree(_ port: UInt16) -> Bool {
self.portIsFree(port)
@ -502,9 +441,5 @@ final class RemotePortTunnel: @unchecked Sendable {
self.sshOptions(localPort: localPort, remotePort: remotePort, hostKeyPolicy: hostKeyPolicy)
}
static func _testDrainStderr(_ handle: FileHandle) -> String {
self.drainStderr(handle, captured: "")
}
#endif
}

View file

@ -616,7 +616,7 @@ private enum MLXTTSTransportError: Error {
private actor ProcessMLXTTSTransport: MLXTTSTransport {
private let process: ManagedProcess
private let input: FileHandle
private let output: FileHandle
private let output: PipeReadStream
private let chunkContinuation: AsyncStream<Data>.Continuation
private let chunks: MLXChunkIterator
private var decoder = MLXTTSFrameDecoder()
@ -626,7 +626,7 @@ private actor ProcessMLXTTSTransport: MLXTTSTransport {
private init(
process: ManagedProcess,
input: FileHandle,
output: FileHandle,
output: PipeReadStream,
chunks: AsyncStream<Data>,
chunkContinuation: AsyncStream<Data>.Continuation)
{
@ -646,19 +646,11 @@ private actor ProcessMLXTTSTransport: MLXTTSTransport {
// send() to its stdin raises SIGPIPE and kills the app.
inputPipe.fileHandleForWriting.disableSIGPIPE()
let output = outputPipe.fileHandleForReading
let (stream, continuation) = AsyncStream<Data>.makeStream()
output.readabilityHandler = { handle in
// Throwing read wrapper; availableData can raise ObjC exceptions on
// closed/invalid handles and abort the process (FileHandle+SafeRead).
let data = handle.readSafely(upToCount: 64 * 1024)
if data.isEmpty {
handle.readabilityHandler = nil
continuation.finish()
} else {
continuation.yield(data)
}
}
let output = try PipeReadStream(
handle: outputPipe.fileHandleForReading,
onData: { continuation.yield($0) },
onClose: { continuation.finish() })
let configuration = Subprocess.Configuration(
executable: .path(.init(invocation.executableURL.path)),
@ -674,6 +666,7 @@ private actor ProcessMLXTTSTransport: MLXTTSTransport {
error: .currentStandardError,
closeAfterSpawn: [
inputPipe.fileHandleForReading,
outputPipe.fileHandleForReading,
outputPipe.fileHandleForWriting,
])
do {
@ -681,8 +674,9 @@ private actor ProcessMLXTTSTransport: MLXTTSTransport {
} catch {
// The detached launch can still spawn; reap it before closing inherited pipes.
await process.terminate(gracefully: false)
output.readabilityHandler = nil
output.close()
continuation.finish()
await output.finish()
throw error
}
@ -714,11 +708,11 @@ private actor ProcessMLXTTSTransport: MLXTTSTransport {
func close() async {
guard !self.isClosed else { return }
self.isClosed = true
self.output.readabilityHandler = nil
self.output.close()
self.chunkContinuation.finish()
self.input.closeFile()
self.output.closeFile()
await self.process.terminate()
await self.output.finish()
}
}

View file

@ -425,35 +425,6 @@ struct CuaDriverHostCoordinatorTests {
#expect(launch.environment["PATH"] == "/usr/bin:/bin")
}
@Test func `stderr relay replaces the raw danger banner with one managed mode notice`() async throws {
let probe = CuaDriverStderrProbe()
let relay = CuaDriverStderrRelay { probe.append($0) }
relay.startReading()
relay.reportManagedMode()
relay.reportManagedMode()
let driverOutput = """
DANGER: Cua Driver is running in unrestricted mode. Runtime approval prompts are disabled.
driver diagnostic
"""
// The relay's readability handler calls stop() on any empty read, which
// closes the pipe's read end; without suppression a racing stop turns
// this write into a harness-killing SIGPIPE.
try TestProcessSupport.suppressSIGPIPE(relay.pipe.fileHandleForWriting)
try relay.pipe.fileHandleForWriting.write(contentsOf: Data(driverOutput.utf8))
try relay.pipe.fileHandleForWriting.close()
for _ in 0..<1000 where probe.events.count < 2 {
await Task.yield()
}
relay.stop()
#expect(probe.events == [
.notice(CuaDriverStderrRelay.managedModeNotice),
.error("driver diagnostic"),
])
}
@Test func `unexpected exits retry with a bounded budget while advertising unavailable`() async throws {
let delays = CuaRestartDelayProbe()
let root = try ExecApprovalsSocketTestSupport.makeRoot()
@ -692,19 +663,6 @@ struct CuaDriverHostCoordinatorTests {
}
}
private final class CuaDriverStderrProbe: @unchecked Sendable {
private let lock = NSLock()
private var captured: [CuaDriverStderrEvent] = []
var events: [CuaDriverStderrEvent] {
self.lock.withLock { self.captured }
}
func append(_ event: CuaDriverStderrEvent) {
self.lock.withLock { self.captured.append(event) }
}
}
@MainActor
private final class CuaPermissionSnapshotProbe {
var value: [Capability: CapabilityAuthorizationStatus] = [

View file

@ -0,0 +1,268 @@
import Darwin
import Foundation
import Testing
@testable import OpenClaw
struct CuaDriverStderrRelayTests {
@Test(arguments: [false, true])
@MainActor
func `releasing the exited process wrapper preserves its pending stderr drain`(retainWriter: Bool) async throws {
let file = FileManager.default.temporaryDirectory.appendingPathComponent(UUID().uuidString)
try Data(("gate\n" + String(repeating: "x", count: 4090) + "\n").utf8).write(to: file)
defer { try? FileManager.default.removeItem(at: file) }
let probe = CuaDriverStderrProbe()
let entered = DispatchSemaphore(value: 0)
let release = DispatchSemaphore(value: 0)
let exited = DispatchSemaphore(value: 0)
let notified = DispatchSemaphore(value: 0)
let relay = CuaDriverStderrRelay { event in
probe.append(event)
if event == .error("gate") {
entered.signal()
release.wait()
}
}
let descriptor = retainWriter ? fcntl(relay.pipe.fileHandleForWriting.fileDescriptor, F_DUPFD_CLOEXEC, 0) : -1
try #require(!retainWriter || descriptor >= 0)
let writer = retainWriter ? FileHandle(fileDescriptor: descriptor, closeOnDealloc: true) : nil
let input = Pipe()
let process = Process()
process.executableURL = URL(fileURLWithPath: "/bin/sh")
process.arguments = [
"-c", "/bin/cat \"$1\" >&2; IFS= read -r release; printf 'final driver diagnostic' >&2",
"driver", file.path,
]
process.standardInput = input
process.standardOutput = FileHandle.nullDevice
process.standardError = relay.pipe
process.terminationHandler = { _ in
exited.signal()
Task {
await relay.finishReading()
notified.signal()
}
}
try relay.startReading()
try process.run()
var wrapper: FoundationCuaDriverProcess? = FoundationCuaDriverProcess(
process: process, livenessPipe: input)
weak var weakWrapper = wrapper
defer {
release.signal()
relay.stop()
try? writer?.close()
if process.isRunning { process.terminate() }
process.waitUntilExit()
}
try #require(await Self.waitForSignal(entered))
try input.fileHandleForWriting.write(contentsOf: Data("exit\n".utf8))
try #require(await Self.waitForSignal(exited))
#expect(!process.isRunning)
#expect(await !Self.waitForSignal(notified, timeout: 0))
wrapper = nil
#expect(weakWrapper == nil)
release.signal()
#expect(await Self.waitForSignal(notified))
#expect(process.terminationStatus == 0)
#expect(probe.events.contains(.error("final driver diagnostic")))
}
private static func waitForSignal(_ signal: DispatchSemaphore, timeout: TimeInterval = 2) async -> Bool {
await withCheckedContinuation { continuation in
DispatchQueue.global().async {
continuation.resume(returning: signal.wait(timeout: .now() + timeout) == .success)
}
}
}
@Test func `stderr relay filters the banner and forwards diagnostics before the driver exits`() throws {
let probe = CuaDriverStderrProbe()
let relay = CuaDriverStderrRelay { probe.append($0) }
try relay.startReading()
relay.reportManagedMode()
relay.reportManagedMode()
let input = Pipe()
let process = Process()
process.executableURL = URL(fileURLWithPath: "/bin/sh")
process.arguments = [
"-c",
"""
printf '%s\\n' \\
'DANGER: Cua Driver is running in unrestricted mode. Runtime approval prompts are disabled.' \\
'driver diagnostic' >&2
IFS= read -r response
""",
]
process.standardInput = input
process.standardOutput = FileHandle.nullDevice
process.standardError = relay.pipe
try process.run()
defer {
try? input.fileHandleForWriting.close()
if process.isRunning {
process.terminate()
}
process.waitUntilExit()
relay.stop()
}
let delivered = probe.diagnostic.wait(timeout: .now() + 2) == .success
#expect(process.isRunning)
#expect(delivered, "The driver is waiting for input; stderr must not wait for its exit")
try input.fileHandleForWriting.write(contentsOf: Data("done\n".utf8))
process.waitUntilExit()
relay.stop()
#expect(process.terminationStatus == 0)
#expect(probe.events == [
.notice(CuaDriverStderrRelay.managedModeNotice),
.error("driver diagnostic"),
])
}
@Test(arguments: [false, true])
func `truncated UTF8 retains the remaining driver diagnostic`(endsWithNewline: Bool) async throws {
let diagnostic = "final driver diagnostic"
let ending = diagnostic + (endsWithNewline ? "\n" : "")
// The 32 KiB byte tail starts at the continuation byte of é.
let payload = Data((String(repeating: "p", count: 4095) + "é" +
String(repeating: "x", count: 32767 - ending.utf8.count) + ending).utf8)
let file = FileManager.default.temporaryDirectory.appendingPathComponent(UUID().uuidString)
try payload.write(to: file)
defer { try? FileManager.default.removeItem(at: file) }
let probe = CuaDriverStderrProbe()
let relay = CuaDriverStderrRelay { probe.append($0) }
try relay.startReading()
let process = Process()
process.executableURL = URL(fileURLWithPath: "/bin/sh")
process.arguments = ["-c", "exec /bin/cat \"$1\" >&2", "driver", file.path]
process.standardOutput = FileHandle.nullDevice
process.standardError = relay.pipe
try process.run()
defer {
relay.stop()
if process.isRunning { process.terminate() }
process.waitUntilExit()
}
process.waitUntilExit()
await relay.finishReading()
#expect(process.terminationStatus == 0)
#expect(probe.events.count == 1)
guard case let .error(message) = probe.events.first else {
Issue.record("Truncating a UTF-8 scalar discarded the entire diagnostic")
return
}
#expect(message.hasSuffix(diagnostic))
#expect(message.count <= 32768)
}
@Test func `EOF forwards an unterminated diagnostic once`() throws {
let probe = CuaDriverStderrProbe()
let relay = CuaDriverStderrRelay { probe.append($0) }
defer { relay.stop() }
try relay.startReading()
#expect(relay.pipe.fileHandleForWriting.disableSIGPIPE())
try relay.pipe.fileHandleForWriting.write(contentsOf: Data("final diagnostic".utf8))
try relay.pipe.fileHandleForWriting.close()
#expect(probe.diagnostic.wait(timeout: .now() + 2) == .success)
relay.stop()
relay.stop()
#expect(probe.events == [.error("final diagnostic")])
}
@Test(arguments: [false, true])
func `natural exit drains queued stderr but explicit stop cancels it`(naturalExit: Bool) async throws {
let firstChunk = FileManager.default.temporaryDirectory.appendingPathComponent(UUID().uuidString)
try Data(("gate\n" + String(repeating: "x", count: 4090) + "\n").utf8).write(to: firstChunk)
defer { try? FileManager.default.removeItem(at: firstChunk) }
let probe = CuaDriverStderrProbe()
let entered = DispatchSemaphore(value: 0)
let release = DispatchSemaphore(value: 0)
let exited = DispatchSemaphore(value: 0)
let notified = DispatchSemaphore(value: 0)
let closed = DispatchSemaphore(value: 0)
let relay = CuaDriverStderrRelay { event in
probe.append(event)
if event == .error("gate") {
entered.signal()
release.wait()
}
}
let descriptor = fcntl(relay.pipe.fileHandleForWriting.fileDescriptor, F_DUPFD_CLOEXEC, 0)
try #require(descriptor >= 0)
let retainedWriter = FileHandle(fileDescriptor: descriptor, closeOnDealloc: true)
let input = Pipe()
let process = Process()
process.executableURL = URL(fileURLWithPath: "/bin/sh")
process.arguments = [
"-c", "/bin/cat \"$1\" >&2; IFS= read -r release; printf 'final driver diagnostic' >&2",
"worker", firstChunk.path,
]
process.standardInput = input
process.standardOutput = FileHandle.nullDevice
process.standardError = relay.pipe
process.terminationHandler = { _ in
exited.signal()
Task {
if naturalExit {
await relay.finishReading()
} else {
relay.stop()
}
notified.signal()
}
}
try relay.startReading()
try process.run()
defer {
release.signal()
relay.stop()
try? retainedWriter.close()
try? input.fileHandleForWriting.close()
if process.isRunning {
process.terminate()
}
process.waitUntilExit()
}
try #require(await Self.waitForSignal(entered))
try input.fileHandleForWriting.write(contentsOf: Data("exit\n".utf8))
try #require(await Self.waitForSignal(exited))
if naturalExit {
#expect(await !Self.waitForSignal(notified, timeout: 0))
} else {
#expect(await Self.waitForSignal(notified))
}
release.signal()
if naturalExit {
#expect(await Self.waitForSignal(notified))
}
Task { await relay.finishReading()
closed.signal()
}
#expect(await Self.waitForSignal(closed))
#expect(process.terminationStatus == 0)
#expect(probe.events.contains(.error("final driver diagnostic")) == naturalExit)
}
}
private final class CuaDriverStderrProbe: @unchecked Sendable {
let diagnostic = DispatchSemaphore(value: 0)
private let lock = NSLock()
private var captured: [CuaDriverStderrEvent] = []
var events: [CuaDriverStderrEvent] {
self.lock.withLock { self.captured }
}
func append(_ event: CuaDriverStderrEvent) {
self.lock.withLock { self.captured.append(event) }
if case .error = event {
self.diagnostic.signal()
}
}
}

View file

@ -1,47 +0,0 @@
import Foundation
import Testing
@testable import OpenClaw
struct FileHandleSafeReadTests {
@Test func `read to end safely returns empty for closed handle`() {
let pipe = Pipe()
let handle = pipe.fileHandleForReading
try? handle.close()
let data = handle.readToEndSafely()
#expect(data.isEmpty)
}
@Test func `read safely up to count returns empty for closed handle`() {
let pipe = Pipe()
let handle = pipe.fileHandleForReading
try? handle.close()
let data = handle.readSafely(upToCount: 16)
#expect(data.isEmpty)
}
@Test func `read to end safely reads pipe contents`() {
let pipe = Pipe()
let writeHandle = pipe.fileHandleForWriting
writeHandle.write(Data("hello".utf8))
try? writeHandle.close()
let data = pipe.fileHandleForReading.readToEndSafely()
#expect(String(data: data, encoding: .utf8) == "hello")
}
@Test func `read safely up to count reads incrementally`() {
let pipe = Pipe()
let writeHandle = pipe.fileHandleForWriting
writeHandle.write(Data("hello world".utf8))
try? writeHandle.close()
let readHandle = pipe.fileHandleForReading
let first = readHandle.readSafely(upToCount: 5)
let second = readHandle.readSafely(upToCount: 32)
#expect(String(data: first, encoding: .utf8) == "hello")
#expect(String(data: second, encoding: .utf8) == " world")
}
}

View file

@ -0,0 +1,162 @@
import Darwin
import Foundation
import Testing
@testable import OpenClaw
struct PipeReadStreamTests {
@Test func `cancelled finish still joins reader cleanup`() throws {
let pipe = Pipe()
let probe = PipeReadProbe()
let entered = DispatchSemaphore(value: 0)
let release = DispatchSemaphore(value: 0)
let started = DispatchSemaphore(value: 0)
let joined = DispatchSemaphore(value: 0)
let reader = try PipeReadStream(
handle: pipe.fileHandleForReading,
onData: { data in
probe.append(data)
entered.signal()
release.wait()
},
onClose: { probe.finish() })
defer {
release.signal()
reader.close()
try? pipe.fileHandleForWriting.close()
}
try pipe.fileHandleForReading.close()
try pipe.fileHandleForWriting.write(contentsOf: Data("first".utf8))
try #require(entered.wait(timeout: .now() + 2) == .success)
let closing = Task {
started.signal()
await reader.finish()
joined.signal()
}
closing.cancel()
try #require(started.wait(timeout: .now() + 2) == .success)
#expect(joined.wait(timeout: .now() + 0.1) == .timedOut)
release.signal()
#expect(joined.wait(timeout: .now() + 2) == .success)
#expect(probe.finishCount == 1)
}
@Test func `reader drains large child output in bounded chunks before closing`() throws {
let file = FileManager.default.temporaryDirectory.appendingPathComponent(UUID().uuidString)
let payload = Data((0..<200_003).map { UInt8(truncatingIfNeeded: $0) })
try payload.write(to: file)
defer { try? FileManager.default.removeItem(at: file) }
let pipe = Pipe()
let probe = PipeReadProbe()
let reader = try PipeReadStream(
handle: pipe.fileHandleForReading,
maximumChunkBytes: 4096,
onData: { probe.append($0) },
onClose: { probe.finish() })
let child = Process()
child.executableURL = URL(fileURLWithPath: "/bin/cat")
child.arguments = [file.path]
child.standardOutput = pipe
child.standardError = FileHandle.nullDevice
try child.run()
defer {
reader.close()
try? pipe.fileHandleForReading.close()
if child.isRunning {
child.terminate()
}
child.waitUntilExit()
}
let closed = probe.finished.wait(timeout: .now() + 3) == .success
#expect(closed)
if !closed {
child.terminate()
}
child.waitUntilExit()
reader.close()
#expect(child.terminationStatus == 0)
#expect(probe.contents == payload)
#expect(probe.chunkSizes.count > 1)
#expect(probe.chunkSizes.allSatisfy { $0 > 0 && $0 <= 4096 })
#expect(probe.finishCount == 1)
}
@Test(arguments: [false, true])
func `reader survives original handle closure and releases its own descriptor`(dropOwner: Bool) throws {
let pipe = Pipe()
let probe = PipeReadProbe()
var reader: PipeReadStream? = try PipeReadStream(
handle: pipe.fileHandleForReading,
onData: { probe.append($0) },
onClose: { probe.finish() })
defer {
reader?.close()
try? pipe.fileHandleForWriting.close()
}
#expect(pipe.fileHandleForWriting.disableSIGPIPE())
try pipe.fileHandleForReading.close()
try pipe.fileHandleForWriting.write(contentsOf: Data("still streaming".utf8))
#expect(probe.received.wait(timeout: .now() + 2) == .success)
#expect(probe.contents == Data("still streaming".utf8))
let deadline = DispatchTime.now() + 2
if dropOwner {
reader = nil
} else {
reader?.close()
reader?.close()
}
let closed = probe.finished.wait(timeout: deadline) == .success
#expect(closed)
var byte: UInt8 = 0
var written: Int
var writeError: Int32
// Concurrent child startup briefly inherits even CLOEXEC descriptors.
// EPIPE observes all process references, so share the cleanup deadline.
repeat {
written = Darwin.write(pipe.fileHandleForWriting.fileDescriptor, &byte, 1)
writeError = errno
if written == -1 { break }
Thread.sleep(forTimeInterval: 0.001)
} while DispatchTime.now() < deadline
#expect(written == -1)
#expect(writeError == EPIPE)
#expect(probe.finishCount == 1)
}
}
private final class PipeReadProbe: @unchecked Sendable {
let received = DispatchSemaphore(value: 0)
let finished = DispatchSemaphore(value: 0)
private let lock = NSLock()
private var data = Data()
private var sizes: [Int] = []
private var closes = 0
var contents: Data {
self.lock.withLock { self.data }
}
var chunkSizes: [Int] {
self.lock.withLock { self.sizes }
}
var finishCount: Int {
self.lock.withLock { self.closes }
}
func append(_ chunk: Data) {
self.lock.withLock {
self.data.append(chunk)
self.sizes.append(chunk.count)
}
self.received.signal()
}
func finish() {
self.lock.withLock { self.closes += 1 }
self.finished.signal()
}
}

View file

@ -0,0 +1,141 @@
import Foundation
import Testing
@testable import OpenClaw
struct PipeTextCaptureTests {
@Test(arguments: [false, true])
func `split UTF8 retains the final SSH diagnostic`(endsWithNewline: Bool) {
let capture = PipeTextCapture(characterLimit: 4096, retention: .tail)
let diagnostic = "Permission denied (publickey)."
let output = Data((String(repeating: "x", count: 65535) + "é\n" + diagnostic +
(endsWithNewline ? "\n" : "")).utf8)
let firstChunk = capture.append(Data(output.prefix(65536)))
#expect(firstChunk == String(repeating: "x", count: 65535))
let completeLines = capture.append(Data(output.dropFirst(65536)))
let finalLine = capture.append(Data(), atEOF: true)
let loggedLines = [completeLines, finalLine].filter { !$0.isEmpty }.joined(separator: "\n")
#expect(loggedLines.contains("é\n" + diagnostic))
#expect(capture.snapshot().hasSuffix("é\n" + diagnostic))
#expect(!capture.snapshot().contains("<EFBFBD>"))
#expect(capture.snapshot().count <= 4096)
#expect(capture.append(Data(), atEOF: true).isEmpty)
}
@Test func `live diagnostic chunks retain a bounded tail without repeating at EOF`() {
let capture = PipeTextCapture(characterLimit: 4096, retention: .tail)
let tail = String(repeating: "é", count: 4090) + "failure"
#expect(capture.append(Data("earlier line\n".utf8)) == "earlier line")
let padding = String(repeating: "x", count: 64 * 1024)
for _ in 0..<4 {
#expect(capture.append(Data(padding.utf8)) == padding)
}
#expect(capture.append(Data(tail.utf8)) == tail)
#expect(capture.snapshot() == String(tail.suffix(4096)))
#expect(capture.append(Data(), atEOF: true).isEmpty)
#expect(capture.snapshot() == String(tail.suffix(4096)))
}
@Test func `worker retains the first 700 normalized characters across split UTF8 and long output`() {
let capture = PipeTextCapture(characterLimit: 700, retention: .head)
let head = "refused:\n" + String(repeating: "é", count: 691)
let output = Data((head + String(repeating: "x", count: 64 * 1024)).utf8)
#expect(capture.append(Data(output.prefix(10))) == "refused:")
#expect(capture.append(Data(output.dropFirst(10))).hasPrefix("é"))
#expect(capture.snapshot() == head)
#expect(capture.append(Data(repeating: 0x78, count: 64 * 1024)).count == 64 * 1024)
#expect(capture.snapshot() == head)
#expect(capture.append(Data("\nlast error\n".utf8)).hasSuffix("last error"))
#expect(capture.snapshot() == head)
}
@Test func `completed worker head cannot grow through later combining marks`() {
let capture = PipeTextCapture(characterLimit: 700, retention: .head)
let head = String(repeating: "x", count: 700)
_ = capture.append(Data((head + "\n").utf8))
for _ in 0..<256 {
#expect(capture.append(Data("\u{0301}\n".utf8)) == "\u{0301}")
}
#expect(capture.snapshot() == head)
#expect(capture.snapshot().utf8.count == head.utf8.count)
}
@Test(arguments: [false, true], [false, true])
func `whitespace cannot displace a diagnostic`(retainHead: Bool, separateChunks: Bool) {
let limit = retainHead ? 700 : 4096
let capture = PipeTextCapture(characterLimit: limit, retention: retainHead ? .head : .tail)
let diagnostic = "Permission denied (publickey)."
let padding = String(repeating: " \n", count: limit)
let records = retainHead ? [padding, diagnostic] : [diagnostic + "\n", padding]
for chunk in separateChunks ? records : [records.joined()] {
_ = capture.append(Data(chunk.utf8))
}
#expect(capture.snapshot() == diagnostic)
_ = capture.append(Data(), atEOF: true)
#expect(capture.snapshot() == diagnostic)
}
@Test(arguments: [false, true])
func `diagnostic snapshots trim the retained boundary`(retainHead: Bool) {
let limit = retainHead ? 700 : 4096
let capture = PipeTextCapture(characterLimit: limit, retention: retainHead ? .head : .tail)
let diagnostic = String(repeating: "x", count: limit - 1)
let output = retainHead ? diagnostic + "\nnext\n" : "earlier\n" + diagnostic + "\n"
_ = capture.append(Data(output.utf8))
#expect(capture.snapshot() == diagnostic)
}
@Test(arguments: [false, true], [" ", "\u{2003}"])
func `unterminated padding cannot evict diagnostics`(retainHead: Bool, whitespace: String) {
let capture = PipeTextCapture(characterLimit: retainHead ? 700 : 4096, retention: retainHead ? .head : .tail)
let diagnostic = "Permission denied (publickey)."
let padding = String(repeating: whitespace, count: 64 * 1024)
let output = Data((retainHead ? padding + diagnostic : diagnostic + padding).utf8)
for offset in stride(from: 0, to: output.count, by: 65536) {
_ = capture.append(output.subdata(in: offset..<min(offset + 65536, output.count)))
}
#expect(capture.snapshot() == diagnostic)
_ = capture.append(Data("next diagnostic".utf8))
#expect(capture.snapshot() == diagnostic + "\nnext diagnostic")
_ = capture.append(Data(), atEOF: true)
#expect(capture.snapshot() == diagnostic + "\nnext diagnostic")
}
@Test(arguments: ["é", "€", "💡"])
func `split scalars emit as soon as their bytes are complete`(scalar: String) {
let bytes = Data(scalar.utf8)
for split in 1..<bytes.count {
let capture = PipeTextCapture(characterLimit: 4096, retention: .tail)
#expect(capture.append(Data(bytes.prefix(split))).isEmpty)
#expect(capture.snapshot().isEmpty)
#expect(capture.append(Data(bytes.dropFirst(split))) == scalar)
#expect(capture.snapshot() == scalar)
#expect(capture.append(Data(), atEOF: true).isEmpty)
}
}
@Test(arguments: ["é", "€", "💡"])
func `incomplete scalars flush once at EOF`(scalar: String) {
let bytes = Data(scalar.utf8)
for split in 1..<bytes.count {
let capture = PipeTextCapture(characterLimit: 4096, retention: .tail)
let prefix = Data(bytes.prefix(split))
#expect(capture.append(prefix).isEmpty)
let finalText = String(decoding: prefix, as: UTF8.self)
#expect(capture.append(Data(), atEOF: true) == finalText)
#expect(capture.snapshot() == finalText)
#expect(capture.append(Data(), atEOF: true).isEmpty)
}
}
@Test(arguments: [Data([0x80]), Data([0xC0, 0xAF]), Data([0xED, 0xA0, 0x80]), Data([0xFF, 0x61])])
func `malformed complete bytes use the standard replacement behavior`(bytes: Data) {
let capture = PipeTextCapture(characterLimit: 4096, retention: .tail)
let expected = String(decoding: bytes, as: UTF8.self)
#expect(capture.append(bytes) == expected)
#expect(capture.snapshot() == expected)
#expect(capture.append(Data(), atEOF: true).isEmpty)
}
}

View file

@ -28,15 +28,6 @@ struct RemotePortTunnelTests {
#expect(!options.contains { $0.hasPrefix("UpdateHostKeys=") })
}
@Test func `drain stderr does not crash when handle closed`() {
let pipe = Pipe()
let handle = pipe.fileHandleForReading
try? handle.close()
let drained = RemotePortTunnel._testDrainStderr(handle)
#expect(drained.isEmpty)
}
@Test func `port is free detects I pv4 listener`() {
var fd = socket(AF_INET, SOCK_STREAM, 0)
#expect(fd >= 0)