diff --git a/Sources/JSONRPCStdio/ProcessTransport.swift b/Sources/JSONRPCStdio/ProcessTransport.swift index c3be9d1..fe3e156 100644 --- a/Sources/JSONRPCStdio/ProcessTransport.swift +++ b/Sources/JSONRPCStdio/ProcessTransport.swift @@ -20,6 +20,12 @@ public struct ProcessExit: Sendable { } } +/// How long ``ProcessTransport/waitForExit()`` waits for a captured stderr to reach EOF +/// once the child has exited. EOF normally lands immediately, but it is not guaranteed to +/// arrive at all — a grandchild that inherited stderr holds the write end open — so the +/// wait is bounded rather than indefinite. +private let stderrDrainGrace: DispatchTimeInterval = .milliseconds(250) + /// Errors specific to launching the `Foundation.Process` transport. public enum ProcessTransportError: Error, LocalizedError { case launchFailed(String) @@ -46,22 +52,43 @@ public final class ProcessTransport: JSONRPCMessageTran private let process = Process() private let stdinPipe = Pipe() private let stdoutPipe = Pipe() + private let stderrTail: StderrTail? private let writeLock = NSLock() private let stateLock = NSLock() private var isClosed = false private var exitResult: ProcessExit? private var exitWaiters: [CheckedContinuation] = [] + /// Set while a captured stderr is still being drained. The child can exit with bytes + /// left queued in the pipe, so exit alone does not mean the tail is complete — + /// `waitForExit()` waits for both, and a caller can read the final diagnostic + /// straight after it returns. + private var awaitingStderrEOF = false + /// The read end of a captured stderr pipe, kept so ``close()`` can cancel its + /// readability source: a handler left installed outlives the transport, and on Linux + /// keeps the whole process from exiting. + private var stderrReadHandle: FileHandle? /// The child's process identifier (pid), valid once launched. public var processIdentifier: Int32 { process.processIdentifier } + /// The tail of the child's stderr kept under ``StderrDisposition/capture(maxBytes:)``, + /// as text. Empty for any other disposition, and for a child that wrote nothing. + public func capturedStandardError() -> String { + stderrTail?.text ?? "" + } + public init(launch: ProcessLaunch, framing: Framing) throws { self.framing = framing + if case .capture(let maxBytes) = launch.stderr { + self.stderrTail = StderrTail(maxBytes: maxBytes) + } else { + self.stderrTail = nil + } process.executableURL = Self.resolveExecutable(launch.executable) process.arguments = launch.arguments process.standardInput = stdinPipe process.standardOutput = stdoutPipe - process.standardError = launch.inheritStderr ? FileHandle.standardError : nil + attachStandardError(launch.stderr, to: process) if let env = launch.environment { process.environment = env } @@ -74,10 +101,21 @@ public final class ProcessTransport: JSONRPCMessageTran let result = ProcessExit(code: proc.terminationStatus, reason: proc.terminationReason) self.stateLock.lock() self.exitResult = result - let waiters = self.exitWaiters - self.exitWaiters = [] + // Bytes can still be queued in a captured stderr pipe: hold the waiters + // until its EOF arrives so the tail they read is the child's last word. + let awaitingDrain = self.awaitingStderrEOF + let waiters = awaitingDrain ? [] : self.exitWaiters + if !awaitingDrain { self.exitWaiters = [] } self.stateLock.unlock() for waiter in waiters { waiter.resume(returning: result) } + guard awaitingDrain else { return } + // Wait briefly for the tail to complete, then release regardless. A tail + // missing its last bytes is a far better failure than a caller that never + // wakes — which is exactly what an indefinite wait produced on Linux, where + // a child that writes once and exits never delivered the empty read. + DispatchQueue.global().asyncAfter(deadline: .now() + stderrDrainGrace) { [weak self] in + self?.finishStderrDrain() + } } do { @@ -87,11 +125,57 @@ public final class ProcessTransport: JSONRPCMessageTran } } + /// Point the child's stderr at whatever the disposition asks for. + /// + /// Under `.capture` the pipe is read continuously rather than at exit: one nobody + /// reads fills up and stalls the child. Only the tail is kept, and EOF on it is what + /// tells ``waitForExit()`` the tail is final. + private func attachStandardError(_ disposition: StderrDisposition, to process: Process) { + switch disposition { + case .inherit: + process.standardError = FileHandle.standardError + case .discard: + process.standardError = nil + case .capture: + let pipe = Pipe() + process.standardError = pipe + let tail = stderrTail + awaitingStderrEOF = true + stderrReadHandle = pipe.fileHandleForReading + pipe.fileHandleForReading.readabilityHandler = { [weak self] handle in + let chunk = handle.availableData + guard chunk.isEmpty else { + tail?.append(chunk) + return + } + handle.readabilityHandler = nil + self?.finishStderrDrain() + } + } + } + + /// Called on stderr EOF: the tail is complete, so an exit that already happened can + /// now be reported. + private func finishStderrDrain() { + stateLock.lock() + awaitingStderrEOF = false + let result = exitResult + let waiters = result == nil ? [] : exitWaiters + if result != nil { exitWaiters = [] } + stateLock.unlock() + if let result { for waiter in waiters { waiter.resume(returning: result) } } + } + /// Suspends until the child exits, returning its termination status. + /// + /// Under ``StderrDisposition/capture(maxBytes:)`` this also waits for the stderr + /// pipe to reach EOF, so ``capturedStandardError()`` read straight afterwards + /// includes whatever the child said last — which is usually the part that explains + /// the exit. public func waitForExit() async -> ProcessExit { await withCheckedContinuation { continuation in stateLock.lock() - if let result = exitResult { + if let result = exitResult, !awaitingStderrEOF { stateLock.unlock() continuation.resume(returning: result) } else { @@ -143,6 +227,14 @@ public final class ProcessTransport: JSONRPCMessageTran stateLock.unlock() try? stdinPipe.fileHandleForWriting.close() + // Leaving a readability handler installed keeps its dispatch source — and on + // Linux the whole process — alive after the transport is done with. + if let handle = stderrReadHandle { + handle.readabilityHandler = nil + try? handle.close() + stderrReadHandle = nil + } + finishStderrDrain() if process.isRunning { process.terminate() } diff --git a/Sources/JSONRPCSubprocess/StdioMessageTransport.swift b/Sources/JSONRPCSubprocess/StdioMessageTransport.swift index 5d086b9..b9f42b5 100644 --- a/Sources/JSONRPCSubprocess/StdioMessageTransport.swift +++ b/Sources/JSONRPCSubprocess/StdioMessageTransport.swift @@ -37,6 +37,7 @@ public final class StdioTransport: JSONRPCMessageTransp private let outbound: AsyncStream.Continuation private let inbound: AsyncThrowingStream private let runTask: Task + private let stderrTail: StderrTail? public init(endpoint: StdioEndpoint, framing: Framing) { let (outboundStream, outboundContinuation) = AsyncStream.makeStream() @@ -46,15 +47,32 @@ public final class StdioTransport: JSONRPCMessageTransp switch endpoint { case .childProcess(let launch): + let tail: StderrTail? + if case .capture(let maxBytes) = launch.stderr { + tail = StderrTail(maxBytes: maxBytes) + } else { + tail = nil + } + self.stderrTail = tail self.runTask = Self.runChild( - launch: launch, framing: framing, + launch: launch, framing: framing, stderrTail: tail, outbound: outboundStream, inbound: inboundContinuation) case .currentProcess: + self.stderrTail = nil self.runTask = Self.runCurrentProcess( framing: framing, outbound: outboundStream, inbound: inboundContinuation) } } + /// The tail of the child's stderr kept under ``StderrDisposition/capture(maxBytes:)``, + /// as text. Empty for any other disposition, and for a child that wrote nothing. + /// + /// Read it after a failure: whatever the child said on its way out is usually what + /// explains the exit. + public func capturedStandardError() -> String { + stderrTail?.text ?? "" + } + public func send(_ message: JSONRPCMessage) throws { guard case .enqueued = outbound.yield(message) else { throw JSONRPCPeerError.closed @@ -74,6 +92,7 @@ public final class StdioTransport: JSONRPCMessageTransp private static func runChild( launch: ProcessLaunch, framing: Framing, + stderrTail: StderrTail?, outbound: AsyncStream, inbound: AsyncThrowingStream.Continuation ) -> Task { @@ -82,7 +101,6 @@ public final class StdioTransport: JSONRPCMessageTransp : .name(launch.executable) let arguments = Arguments(launch.arguments) let workingDirectory = launch.workingDirectory.map { FilePath($0) } - let inheritStderr = launch.inheritStderr // Honor a caller-supplied environment as a full replacement — matching the // `Foundation.Process` transport's `process.environment = launch.environment` // (e.g. an ACP/MCP client injecting auth vars into the agent it spawns). `nil` @@ -100,7 +118,8 @@ public final class StdioTransport: JSONRPCMessageTransp // The server's stderr (its logs) either passes through to ours or is // discarded — distinct output types, so the `run` call is branched; // the I/O pump is shared. - if inheritStderr { + switch launch.stderr { + case .inherit: _ = try await run( executable, arguments: arguments, environment: environment, workingDirectory: workingDirectory, @@ -108,7 +127,7 @@ public final class StdioTransport: JSONRPCMessageTransp ) { execution in try await pump(execution, framing: framing, outbound: outbound, inbound: inbound) } - } else { + case .discard: _ = try await run( executable, arguments: arguments, environment: environment, workingDirectory: workingDirectory, @@ -116,6 +135,16 @@ public final class StdioTransport: JSONRPCMessageTransp ) { execution in try await pump(execution, framing: framing, outbound: outbound, inbound: inbound) } + case .capture: + _ = try await run( + executable, arguments: arguments, environment: environment, + workingDirectory: workingDirectory, + input: .inputWriter, output: .sequence, error: .sequence + ) { execution in + try await pumpCapturingStderr( + execution, framing: framing, tail: stderrTail, + outbound: outbound, inbound: inbound) + } } inbound.finish() } catch is CancellationError { @@ -128,6 +157,31 @@ public final class StdioTransport: JSONRPCMessageTransp } } + /// ``pump`` plus a stderr drain, for ``StderrDisposition/capture(maxBytes:)``. + /// + /// The drain runs *alongside* the message pump rather than after it: a stderr pipe + /// nobody reads fills up and stalls — or kills — the child. Everything is read; + /// only the tail is kept. + private static func pumpCapturingStderr( + _ execution: Execution, + framing: Framing, + tail: StderrTail?, + outbound: AsyncStream, + inbound: AsyncThrowingStream.Continuation + ) async throws { + try await withThrowingTaskGroup(of: Void.self) { group in + group.addTask { + for try await buffer in execution.standardError { + tail?.append(Data(buffer.withUnsafeBytes { Array($0) })) + } + } + group.addTask { + try await pump(execution, framing: framing, outbound: outbound, inbound: inbound) + } + try await group.waitForAll() + } + } + /// Bridge the child's stdio to the two streams: stdout → framing → `inbound`; /// `outbound` → framing → stdin. Generic over the error output (passthrough vs /// discarded) since the body never touches stderr. diff --git a/Sources/JSONRPCWire/ProcessLaunch.swift b/Sources/JSONRPCWire/ProcessLaunch.swift index 8ebed1f..05e01ba 100644 --- a/Sources/JSONRPCWire/ProcessLaunch.swift +++ b/Sources/JSONRPCWire/ProcessLaunch.swift @@ -1,3 +1,5 @@ +import Foundation + /// How to launch a child process that speaks JSON-RPC over stdio. /// /// A generic, transport-agnostic launch descriptor shared by every stdio @@ -10,22 +12,93 @@ public struct ProcessLaunch: Sendable { public var arguments: [String] public var environment: [String: String]? public var workingDirectory: String? - /// When true the child's stderr (its own logs) passes through to this process's - /// stderr; when false it is discarded. Either way it stays off the JSON-RPC - /// stdout stream. - public var inheritStderr: Bool + /// What becomes of the child's stderr — its own logs, never part of the JSON-RPC + /// stdout stream either way. + public var stderr: StderrDisposition + + /// The two-way view this type had before ``StderrDisposition`` existed: reads as + /// `true` only for ``StderrDisposition/inherit``, and writing it selects `inherit` + /// or `discard`. + public var inheritStderr: Bool { + get { stderr == .inherit } + set { stderr = newValue ? .inherit : .discard } + } public init( executable: String, arguments: [String] = [], environment: [String: String]? = nil, workingDirectory: String? = nil, - inheritStderr: Bool = false + stderr: StderrDisposition = .discard ) { self.executable = executable self.arguments = arguments self.environment = environment self.workingDirectory = workingDirectory - self.inheritStderr = inheritStderr + self.stderr = stderr + } + + /// Convenience for the pre-``StderrDisposition`` spelling. + public init( + executable: String, + arguments: [String] = [], + environment: [String: String]? = nil, + workingDirectory: String? = nil, + inheritStderr: Bool + ) { + self.init( + executable: executable, arguments: arguments, environment: environment, + workingDirectory: workingDirectory, stderr: inheritStderr ? .inherit : .discard) + } +} + +/// What a stdio transport does with its child's stderr. +public enum StderrDisposition: Sendable, Equatable { + /// Passes through to this process's stderr — the child's logs appear in ours. + case inherit + /// Dropped. + case discard + /// Kept, but only the last `maxBytes`, and still drained: a child whose stderr is + /// never read can block on a full pipe or die of `EPIPE`, so the bytes past the + /// limit are read and thrown away rather than left unread. + /// + /// For diagnosing a child that dies: whatever it said on the way out is what + /// explains the exit, and the tail is where that lives. Read it back with the + /// transport's `capturedStandardError()`. + case capture(maxBytes: Int) +} + +/// The last `maxBytes` of a stream, kept across concurrent appends. +/// +/// Internal to the package's transports; exposed to callers as a `String` through +/// their `capturedStandardError()`. +package final class StderrTail: @unchecked Sendable { + private let lock = NSLock() + private var bytes = Data() + private let maxBytes: Int + + package init(maxBytes: Int) { + self.maxBytes = max(0, maxBytes) + } + + package func append(_ chunk: Data) { + guard maxBytes > 0, !chunk.isEmpty else { return } + lock.lock() + defer { lock.unlock() } + bytes.append(chunk) + if bytes.count > maxBytes { bytes.removeFirst(bytes.count - maxBytes) } + } + + /// The tail as text. A cut may land mid-character, so decoding is lossy rather + /// than failing — a diagnostic string is worth more than nothing. + package var text: String { + lock.lock() + defer { lock.unlock() } + // Deliberately lossy: keeping a *tail* means the first bytes may be half a + // character, and a diagnostic string with one replacement character is worth + // more than no diagnostic at all — which is what the failable initializer + // would give here. + // swiftlint:disable:next optional_data_string_conversion + return String(decoding: bytes, as: UTF8.self) } } diff --git a/Tests/JSONRPCStdioTests/JSONRPCStdioTests.swift b/Tests/JSONRPCStdioTests/JSONRPCStdioTests.swift index 664d945..9b8321e 100644 --- a/Tests/JSONRPCStdioTests/JSONRPCStdioTests.swift +++ b/Tests/JSONRPCStdioTests/JSONRPCStdioTests.swift @@ -40,4 +40,61 @@ func processTransportChildReceivesCustomEnvironment() async throws { #expect(received?.method == "hello-env") transport.close() } + +// MARK: - Captured stderr + +/// The `Foundation.Process` transport keeps the same bounded tail as the subprocess +/// one: a child that dies has usually explained itself on stderr. +/// `waitForExit()` then read — no polling. A child can exit with stderr still queued, +/// so the wait covers the drain as well; otherwise the tail read here could miss the +/// very line that explains the exit. +@Test(.timeLimit(.minutes(1))) +func processTransportCapturesTheStderrTail() async throws { + let transport = try ProcessTransport( + launch: ProcessLaunch( + executable: "/bin/sh", arguments: ["-c", "echo 'boom: no such model' >&2; exit 2"], + stderr: .capture(maxBytes: 4096)), + framing: LineFraming()) + + let exit = await transport.waitForExit() + #expect(exit.code == 2) + #expect(transport.capturedStandardError().contains("no such model")) + transport.close() +} + +/// A child that writes a backlog and exits immediately: the final line is the one worth +/// having, and it must survive the exit. +/// +/// This pins the contract rather than reproducing the race. `terminationHandler` and the +/// readability drain have no specified ordering, but on macOS the drain wins here — the +/// test passes without the coordination too, even at ~800 KB of backlog. It is kept +/// because the guarantee is what callers rely on, not because it reproduces the failure. +@Test(.timeLimit(.minutes(1))) +func theFinalStderrLineSurvivesAnImmediateExit() async throws { + let script = "for i in $(seq 1 2000); do echo 'chatter chatter chatter chatter' >&2; done;" + + " echo 'FINAL: the reason' >&2; exit 7" + let transport = try ProcessTransport( + launch: ProcessLaunch( + executable: "/bin/sh", arguments: ["-c", script], stderr: .capture(maxBytes: 512)), + framing: LineFraming()) + + let exit = await transport.waitForExit() + #expect(exit.code == 7) + #expect(transport.capturedStandardError().contains("FINAL: the reason")) + transport.close() +} + +@Test(.timeLimit(.minutes(1))) +func processTransportCapturesNothingByDefault() async throws { + let transport = try ProcessTransport( + launch: ProcessLaunch(executable: "/bin/sh", arguments: ["-c", "echo noise >&2; cat -u"]), + framing: LineFraming()) + try transport.send(.request(id: 1, method: "ping", params: nil)) + var inbound = transport.makeInboundStream().makeAsyncIterator() + _ = try await inbound.next() + + #expect(transport.capturedStandardError().isEmpty) + transport.close() +} + #endif diff --git a/Tests/JSONRPCSubprocessTests/JSONRPCSubprocessTests.swift b/Tests/JSONRPCSubprocessTests/JSONRPCSubprocessTests.swift index 4152322..d3a6d41 100644 --- a/Tests/JSONRPCSubprocessTests/JSONRPCSubprocessTests.swift +++ b/Tests/JSONRPCSubprocessTests/JSONRPCSubprocessTests.swift @@ -39,4 +39,84 @@ func childProcessReceivesCustomEnvironment() async throws { #expect(received?.method == "hello-env") transport.close() } + +// MARK: - Captured stderr + +/// A child that dies has usually said why on stderr. `.capture` keeps the tail of it +/// so the caller can quote that instead of reporting a bare exit. +@Test(.timeLimit(.minutes(1))) +func capturedStderrExplainsAChildThatDies() async throws { + let script = "echo 'agent failed: missing credentials' >&2; exit 1" + let transport = StdioTransport( + endpoint: .childProcess(ProcessLaunch( + executable: "/bin/sh", arguments: ["-c", script], + stderr: .capture(maxBytes: 4096))), + framing: LineFraming()) + var inbound = transport.makeInboundStream().makeAsyncIterator() + // The child writes no JSON-RPC, so the stream ends; the stderr tail is the story. + _ = try? await inbound.next() + for _ in 0 ..< 50 where transport.capturedStandardError().isEmpty { + try await Task.sleep(nanoseconds: 20_000_000) + } + + #expect(transport.capturedStandardError().contains("missing credentials")) + transport.close() +} + +/// Only the tail is kept — the point is a bounded buffer, not a transcript. +@Test(.timeLimit(.minutes(1))) +func capturedStderrKeepsOnlyTheTail() async throws { + // 400 lines of padding, then the line that matters. + let script = "for i in $(seq 1 400); do echo 'pad pad pad pad pad' >&2; done;" + + " echo 'LAST LINE' >&2; exit 3" + let transport = StdioTransport( + endpoint: .childProcess(ProcessLaunch( + executable: "/bin/sh", arguments: ["-c", script], + stderr: .capture(maxBytes: 256))), + framing: LineFraming()) + var inbound = transport.makeInboundStream().makeAsyncIterator() + _ = try? await inbound.next() + for _ in 0 ..< 50 where !transport.capturedStandardError().contains("LAST LINE") { + try await Task.sleep(nanoseconds: 20_000_000) + } + + let captured = transport.capturedStandardError() + #expect(captured.contains("LAST LINE")) + #expect(captured.utf8.count <= 256) + transport.close() +} + +/// The default disposition captures nothing, so nothing changes for callers that +/// never ask for it. +@Test(.timeLimit(.minutes(1))) +func stderrIsNotCapturedByDefault() async throws { + let transport = StdioTransport( + endpoint: .childProcess(ProcessLaunch( + executable: "/bin/sh", arguments: ["-c", "echo noise >&2; cat -u"])), + framing: LineFraming()) + try transport.send(.request(id: 1, method: "ping", params: nil)) + var inbound = transport.makeInboundStream().makeAsyncIterator() + _ = try await inbound.next() + + #expect(transport.capturedStandardError().isEmpty) + transport.close() +} + +/// A child whose stderr is never drained can stall on a full pipe; `.capture` keeps +/// reading past the limit, so a chatty child still gets through its work. +@Test(.timeLimit(.minutes(1))) +func aChattyChildIsNotStalledByTheLimit() async throws { + let script = "for i in $(seq 1 2000); do echo 'noisy diagnostic line' >&2; done;" + + #" printf '{"jsonrpc":"2.0","method":"survived","params":null}\n'"# + let transport = StdioTransport( + endpoint: .childProcess(ProcessLaunch( + executable: "/bin/sh", arguments: ["-c", script], + stderr: .capture(maxBytes: 64))), + framing: LineFraming()) + var inbound = transport.makeInboundStream().makeAsyncIterator() + let received = try await inbound.next() + + #expect(received?.method == "survived") + transport.close() +} #endif