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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
100 changes: 96 additions & 4 deletions Sources/JSONRPCStdio/ProcessTransport.swift
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand All @@ -46,22 +52,43 @@ public final class ProcessTransport<Framing: MessageFraming>: 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<ProcessExit, Never>] = []
/// 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
}
Expand All @@ -74,10 +101,21 @@ public final class ProcessTransport<Framing: MessageFraming>: 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 {
Expand All @@ -87,11 +125,57 @@ public final class ProcessTransport<Framing: MessageFraming>: 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 {
Expand Down Expand Up @@ -143,6 +227,14 @@ public final class ProcessTransport<Framing: MessageFraming>: 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()
}
Expand Down
62 changes: 58 additions & 4 deletions Sources/JSONRPCSubprocess/StdioMessageTransport.swift
Original file line number Diff line number Diff line change
Expand Up @@ -37,6 +37,7 @@ public final class StdioTransport<Framing: MessageFraming>: JSONRPCMessageTransp
private let outbound: AsyncStream<JSONRPCMessage>.Continuation
private let inbound: AsyncThrowingStream<JSONRPCMessage, any Error>
private let runTask: Task<Void, Never>
private let stderrTail: StderrTail?

public init(endpoint: StdioEndpoint, framing: Framing) {
let (outboundStream, outboundContinuation) = AsyncStream<JSONRPCMessage>.makeStream()
Expand All @@ -46,15 +47,32 @@ public final class StdioTransport<Framing: MessageFraming>: 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
Expand All @@ -74,6 +92,7 @@ public final class StdioTransport<Framing: MessageFraming>: JSONRPCMessageTransp
private static func runChild(
launch: ProcessLaunch,
framing: Framing,
stderrTail: StderrTail?,
outbound: AsyncStream<JSONRPCMessage>,
inbound: AsyncThrowingStream<JSONRPCMessage, any Error>.Continuation
) -> Task<Void, Never> {
Expand All @@ -82,7 +101,6 @@ public final class StdioTransport<Framing: MessageFraming>: 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`
Expand All @@ -100,22 +118,33 @@ public final class StdioTransport<Framing: MessageFraming>: 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,
input: .inputWriter, output: .sequence, error: .currentStandardError
) { 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,
input: .inputWriter, output: .sequence, error: .discarded
) { 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 {
Expand All @@ -128,6 +157,31 @@ public final class StdioTransport<Framing: MessageFraming>: 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<CustomWriteInput, SequenceOutput, SequenceOutput>,
framing: Framing,
tail: StderrTail?,
outbound: AsyncStream<JSONRPCMessage>,
inbound: AsyncThrowingStream<JSONRPCMessage, any Error>.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.
Expand Down
85 changes: 79 additions & 6 deletions Sources/JSONRPCWire/ProcessLaunch.swift
Original file line number Diff line number Diff line change
@@ -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
Expand All @@ -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)
}
}
Loading
Loading