diff --git a/apps/macos/Sources/OpenClaw/CookieSyncManager.swift b/apps/macos/Sources/OpenClaw/CookieSyncManager.swift index fe59a29a13f5..b5449b297072 100644 --- a/apps/macos/Sources/OpenClaw/CookieSyncManager.swift +++ b/apps/macos/Sources/OpenClaw/CookieSyncManager.swift @@ -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? @ObservationIgnored private var retryTask: Task? @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.. 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[.. 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[.. Bool { + fcntl(self.fileDescriptor, F_SETNOSIGPIPE, 1) != -1 + } +} diff --git a/apps/macos/Sources/OpenClaw/FileHandle+SafeRead.swift b/apps/macos/Sources/OpenClaw/FileHandle+SafeRead.swift deleted file mode 100644 index 50e5e18ff3f2..000000000000 --- a/apps/macos/Sources/OpenClaw/FileHandle+SafeRead.swift +++ /dev/null @@ -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() - } - } -} diff --git a/apps/macos/Sources/OpenClaw/NodeMode/MacNodeHostWorker.swift b/apps/macos/Sources/OpenClaw/NodeMode/MacNodeHostWorker.swift index 11a5699dbed4..e2e6c83a50e8 100644 --- a/apps/macos/Sources/OpenClaw/NodeMode/MacNodeHostWorker.swift +++ b/apps/macos/Sources/OpenClaw/NodeMode/MacNodeHostWorker.swift @@ -101,17 +101,13 @@ final class MacNodeHostWorker: MacNodeHostWorking, @unchecked Sendable { private var process: ManagedProcess? private var processCleanupTask: Task? 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.. Void + private let completion: Task + + 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.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 + } +} diff --git a/apps/macos/Sources/OpenClaw/PipeTextCapture.swift b/apps/macos/Sources/OpenClaw/PipeTextCapture.swift new file mode 100644 index 000000000000..7024d2f32978 --- /dev/null +++ b/apps/macos/Sources/OpenClaw/PipeTextCapture.swift @@ -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[.. 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)) + } +} diff --git a/apps/macos/Sources/OpenClaw/RemotePortTunnel.swift b/apps/macos/Sources/OpenClaw/RemotePortTunnel.swift index 4f2035bf7773..a15574fe07c6 100644 --- a/apps/macos/Sources/OpenClaw/RemotePortTunnel.swift +++ b/apps/macos/Sources/OpenClaw/RemotePortTunnel.swift @@ -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 } diff --git a/apps/macos/Sources/OpenClaw/TalkMLXSpeechSynthesizer.swift b/apps/macos/Sources/OpenClaw/TalkMLXSpeechSynthesizer.swift index bace8977154e..6f715866ece1 100644 --- a/apps/macos/Sources/OpenClaw/TalkMLXSpeechSynthesizer.swift +++ b/apps/macos/Sources/OpenClaw/TalkMLXSpeechSynthesizer.swift @@ -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.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, chunkContinuation: AsyncStream.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.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() } } diff --git a/apps/macos/Tests/OpenClawIPCTests/CuaDriverHostCoordinatorTests.swift b/apps/macos/Tests/OpenClawIPCTests/CuaDriverHostCoordinatorTests.swift index 6a60ba9fc4a5..28f5b02e08ce 100644 --- a/apps/macos/Tests/OpenClawIPCTests/CuaDriverHostCoordinatorTests.swift +++ b/apps/macos/Tests/OpenClawIPCTests/CuaDriverHostCoordinatorTests.swift @@ -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] = [ diff --git a/apps/macos/Tests/OpenClawIPCTests/CuaDriverStderrRelayTests.swift b/apps/macos/Tests/OpenClawIPCTests/CuaDriverStderrRelayTests.swift new file mode 100644 index 000000000000..f7c6309cce16 --- /dev/null +++ b/apps/macos/Tests/OpenClawIPCTests/CuaDriverStderrRelayTests.swift @@ -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() + } + } +} diff --git a/apps/macos/Tests/OpenClawIPCTests/FileHandleSafeReadTests.swift b/apps/macos/Tests/OpenClawIPCTests/FileHandleSafeReadTests.swift deleted file mode 100644 index 5fb2e1c86ded..000000000000 --- a/apps/macos/Tests/OpenClawIPCTests/FileHandleSafeReadTests.swift +++ /dev/null @@ -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") - } -} diff --git a/apps/macos/Tests/OpenClawIPCTests/PipeReadStreamTests.swift b/apps/macos/Tests/OpenClawIPCTests/PipeReadStreamTests.swift new file mode 100644 index 000000000000..f3d9048d058b --- /dev/null +++ b/apps/macos/Tests/OpenClawIPCTests/PipeReadStreamTests.swift @@ -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() + } +} diff --git a/apps/macos/Tests/OpenClawIPCTests/PipeTextCaptureTests.swift b/apps/macos/Tests/OpenClawIPCTests/PipeTextCaptureTests.swift new file mode 100644 index 000000000000..88ec1b4ba193 --- /dev/null +++ b/apps/macos/Tests/OpenClawIPCTests/PipeTextCaptureTests.swift @@ -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("�")) + #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..