From 442149787a3befcb08457aeaf890af8c297c37f2 Mon Sep 17 00:00:00 2001 From: Ngo Quoc Dat Date: Sat, 26 Sep 2026 00:22:50 +0700 Subject: [PATCH] fix(connections): read every subprocess pipe through one reader that cannot raise or read after it stops --- CHANGELOG.md | 1 + TablePro/CLI/BridgeProxy.swift | 13 +- .../Database/ProcessNativeDumpRunner.swift | 63 +----- TablePro/Core/LSP/LSPTransport.swift | 15 +- TablePro/Core/Process/DescriptorRead.swift | 51 +++++ TablePro/Core/Process/PipeReader.swift | 94 +++++++++ .../Process/SupervisedProcessRunner.swift | 46 ++--- .../Infrastructure/PreConnectHookRunner.swift | 10 +- .../MCPBridgeIntegrationTests.swift | 11 +- .../Core/Process/DescriptorReadTests.swift | 94 +++++++++ .../Core/Process/PipeReaderTests.swift | 194 ++++++++++++++++++ .../PreConnectHookRunnerTests.swift | 45 ++++ project.yml | 1 + 13 files changed, 530 insertions(+), 108 deletions(-) create mode 100644 TablePro/Core/Process/DescriptorRead.swift create mode 100644 TablePro/Core/Process/PipeReader.swift create mode 100644 TableProTests/Core/Process/DescriptorReadTests.swift create mode 100644 TableProTests/Core/Process/PipeReaderTests.swift create mode 100644 TableProTests/Core/Services/Infrastructure/PreConnectHookRunnerTests.swift diff --git a/CHANGELOG.md b/CHANGELOG.md index 641b5d02a1..74aca1a0b3 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -127,6 +127,7 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ### Fixed +- Pre-connect script failures sometimes reported without the script's own error message. - Server connections piling up while browsing many databases or schemas, and staying open after a failed connect. (#3103) - Variables declared in a SQL Server script lost after its first statement. (#3078) - Later SQL Server result sets shown under the first one's columns, or crashing the app. diff --git a/TablePro/CLI/BridgeProxy.swift b/TablePro/CLI/BridgeProxy.swift index 731cc2d7bc..0c5d5cb056 100644 --- a/TablePro/CLI/BridgeProxy.swift +++ b/TablePro/CLI/BridgeProxy.swift @@ -143,7 +143,7 @@ actor BridgeProxy { self.discovery = discovery self.logger = logger self.stdout = BridgeStdout(handle: stdout) - self.hostLines = BridgeStdin.lines(from: stdin) + self.hostLines = BridgeStdin.lines(from: stdin, logger: logger) } func run() async { @@ -524,12 +524,19 @@ enum BridgeJson { } enum BridgeStdin { - static func lines(from handle: FileHandle) -> AsyncStream { + static func lines(from handle: FileHandle, logger: any MCPBridgeLogger) -> AsyncStream { AsyncStream { continuation in let reader = Task.detached(priority: .userInitiated) { + let descriptor = handle.fileDescriptor var buffer = Data() while !Task.isCancelled { - let chunk = handle.availableData + let chunk: Data + do { + chunk = try DescriptorRead.availableBytes(from: descriptor) + } catch { + logger.log(.error, "Reading stdin failed: \(error.localizedDescription)") + break + } if chunk.isEmpty { break } buffer.append(chunk) while let newline = buffer.firstIndex(of: 0x0A) { diff --git a/TablePro/Core/Database/ProcessNativeDumpRunner.swift b/TablePro/Core/Database/ProcessNativeDumpRunner.swift index 6bba657d66..3f212ddbfb 100644 --- a/TablePro/Core/Database/ProcessNativeDumpRunner.swift +++ b/TablePro/Core/Database/ProcessNativeDumpRunner.swift @@ -10,10 +10,9 @@ final class ProcessNativeDumpRunner: NativeDumpRunner, @unchecked Sendable { private let command: NativeDumpCommand private let process = Process() private let stderrPipe = Pipe() + private let stderrReader: PipeReader private let stateLock = NSLock() - /// Held across the read as well as the append, so a chunk can never be taken out of the pipe by - /// one reader and still be missing from the buffer when another snapshots it. Separate from - /// `stateLock`, which `cancel()` takes and which must never wait on a pipe. + /// Separate from `stateLock`, which `cancel()` takes and which must never wait on a pipe. private let stderrLock = NSLock() private var stderrBuffer = Data() private var wasCancelled = false @@ -24,6 +23,7 @@ final class ProcessNativeDumpRunner: NativeDumpRunner, @unchecked Sendable { init(command: NativeDumpCommand) { self.command = command + stderrReader = PipeReader(stderrPipe.fileHandleForReading) } func start() throws { @@ -37,18 +37,13 @@ final class ProcessNativeDumpRunner: NativeDumpRunner, @unchecked Sendable { try attachRedirection(for: command) - stderrPipe.fileHandleForReading.readabilityHandler = { [weak self] handle in - guard let self else { return } - self.stderrLock.lock() - let chunk = handle.availableData - self.append(chunk, cap: stderrCap) - self.stderrLock.unlock() + stderrReader.start { [weak self] chunk in + self?.append(chunk, cap: stderrCap) } process.terminationHandler = { [weak self] proc in guard let self else { return } - self.stderrPipe.fileHandleForReading.readabilityHandler = nil - self.drainStderr(cap: stderrCap) + self.stderrReader.stop(drainingUpTo: stderrCap) self.releaseRedirection() self.stderrLock.lock() @@ -82,51 +77,9 @@ final class ProcessNativeDumpRunner: NativeDumpRunner, @unchecked Sendable { } } - /// Takes whatever the pipe still holds once the process has gone. - /// - /// The readability source and the process reaper run on independent queues, so bytes written - /// just before the child exits can still be in the pipe when `terminationHandler` reads the - /// buffer, and a tool that exits on its first argument writes everything it has to say in that - /// window. Measured with a harness mirroring this class against a child that writes 66 bytes - /// and exits immediately: 2 of 300 runs captured nothing at all, and with this drain 0 of 300 - /// did. What the user saw instead was "Process exited with code 7" and no message. - /// Takes whatever the pipe still holds once the process has gone. - /// - /// The readability source and the process reaper run on independent queues, so a tool that - /// exits on its first argument can have written everything it has to say and still be waiting - /// to be read. Measured with a harness mirroring this class against a child that writes 66 - /// bytes and exits at once: between 1 and 6 of every 300 runs captured nothing at all, the rate - /// rising with load, and 0 of 1,200 with the drain and the lock above. What the sheet showed - /// instead was the exit code alone. - /// - /// The child has exited, so everything it wrote is already in the pipe's buffer and a - /// non-blocking read takes all of it. `readDataToEndOfFile` would take it too and then wait for - /// every writer to close, which a grandchild that inherited this end would never do, hanging - /// the termination handler and with it the run. The read stops at the cap for the same reason: - /// a grandchild still writing would otherwise keep the loop fed for as long as it cared to, and - /// nothing downstream, including the credentials file's removal, happens until it returns. - private func drainStderr(cap: Int) { - let descriptor = stderrPipe.fileHandleForReading.fileDescriptor - let flags = fcntl(descriptor, F_GETFL) - guard flags != -1, fcntl(descriptor, F_SETFL, flags | O_NONBLOCK) != -1 else { return } - - var buffer = [UInt8](repeating: 0, count: 4_096) - var taken = 0 - stderrLock.lock() - while taken < cap { - let received = buffer.withUnsafeMutableBytes { raw in - read(descriptor, raw.baseAddress, raw.count) - } - guard received > 0 else { break } - taken += received - append(Data(buffer[0 ..< received]), cap: cap) - } - stderrLock.unlock() - } - - /// Call with `stderrLock` held. private func append(_ chunk: Data, cap: Int) { - guard !chunk.isEmpty else { return } + stderrLock.lock() + defer { stderrLock.unlock() } stderrBuffer.append(chunk) if stderrBuffer.count > cap { stderrBuffer = Data(stderrBuffer.suffix(cap)) diff --git a/TablePro/Core/LSP/LSPTransport.swift b/TablePro/Core/LSP/LSPTransport.swift index b06970fd3a..1ec3b109d1 100644 --- a/TablePro/Core/LSP/LSPTransport.swift +++ b/TablePro/Core/LSP/LSPTransport.swift @@ -40,6 +40,7 @@ actor LSPTransport { private var stdinPipe: Pipe? private var stdoutPipe: Pipe? private var stderrPipe: Pipe? + private var stderrReader: PipeReader? private var nextRequestID: Int = 1 private var pendingRequests: [Int: CheckedContinuation] = [:] private var notificationHandlers: [String: @Sendable (Data) -> Void] = [:] @@ -81,13 +82,12 @@ actor LSPTransport { } } - // Drain stderr to prevent pipe buffer from filling - stderr.fileHandleForReading.readabilityHandler = { handle in - let data = handle.availableData - if !data.isEmpty, let text = String(data: data, encoding: .utf8) { - Self.logger.debug("LSP stderr: \(text)") - } + let stderrReader = PipeReader(stderr.fileHandleForReading) + stderrReader.start { data in + guard let text = String(data: data, encoding: .utf8) else { return } + Self.logger.debug("LSP stderr: \(text)") } + self.stderrReader = stderrReader try proc.run() @@ -115,8 +115,9 @@ actor LSPTransport { if let stdoutHandle = stdoutPipe?.fileHandleForReading { try? stdoutHandle.close() } + stderrReader?.stop() + stderrReader = nil if let stderrHandle = stderrPipe?.fileHandleForReading { - stderrHandle.readabilityHandler = nil try? stderrHandle.close() } diff --git a/TablePro/Core/Process/DescriptorRead.swift b/TablePro/Core/Process/DescriptorRead.swift new file mode 100644 index 0000000000..b05cf41fba --- /dev/null +++ b/TablePro/Core/Process/DescriptorRead.swift @@ -0,0 +1,51 @@ +// +// DescriptorRead.swift +// TablePro +// + +import Darwin +import Foundation + +/// Reads a descriptor with `read(2)` and reports a failed read as a thrown `POSIXError`. +/// +/// `FileHandle.availableData` raises `NSFileHandleOperationException` on any failed read, and an +/// Objective-C exception unwinding through Swift ends the process. `FileHandle.read(upToCount:)` +/// throws instead, but on a pipe it waits for the whole count or for every writer to close, and on +/// a non-blocking pipe it throws away the bytes it read before `EAGAIN`, so it cannot return +/// whatever has arrived so far. +enum DescriptorRead { + /// The most a macOS pipe holds, so a writer that has exited cannot have left more than this. + static let pipeCapacity = 65_536 + + /// The bytes one read returns, at most `limit`, empty at end of file. + static func availableBytes(from descriptor: Int32, upTo limit: Int = pipeCapacity) throws -> Data { + var bytes = Data(count: limit) + let received = try bytes.withUnsafeMutableBytes { buffer in + try read(descriptor, into: buffer) + } + bytes.count = received + return bytes + } + + /// Whether a read would return at once, with bytes or at end of file, instead of waiting on a + /// writer. Asks `poll(2)` rather than setting `O_NONBLOCK`, which would change the descriptor for + /// every other reader of it too. + static func hasInputWithoutWaiting(_ descriptor: Int32) -> Bool { + var request = pollfd(fd: descriptor, events: Int16(POLLIN), revents: 0) + guard poll(&request, 1, 0) > 0 else { return false } + return request.revents & Int16(POLLIN | POLLHUP) != 0 + } + + private static func read(_ descriptor: Int32, into buffer: UnsafeMutableRawBufferPointer) throws -> Int { + while true { + let received = Darwin.read(descriptor, buffer.baseAddress, buffer.count) + if received >= 0 { + return received + } + let failure = errno + guard failure == EINTR else { + throw POSIXError(POSIXErrorCode(rawValue: failure) ?? .EIO) + } + } + } +} diff --git a/TablePro/Core/Process/PipeReader.swift b/TablePro/Core/Process/PipeReader.swift new file mode 100644 index 0000000000..82ccc69549 --- /dev/null +++ b/TablePro/Core/Process/PipeReader.swift @@ -0,0 +1,94 @@ +// +// PipeReader.swift +// TablePro +// + +import Foundation +import os + +/// Hands a pipe's output to a consumer as it arrives, one chunk at a time and in order. +/// +/// Clearing `readabilityHandler` does not recall a callback Foundation has already dispatched, so +/// that callback can run its read after the owner has drained the pipe or closed it. Every read +/// here happens under one lock and only while a consumer is installed, and both ways of stopping +/// remove the consumer under that lock, so once either returns nothing is reading the descriptor +/// and nothing will. The handle can be closed straight after. +final class PipeReader: @unchecked Sendable { + private static let logger = Logger(subsystem: "com.TablePro", category: "PipeReader") + + private let handle: FileHandle + private let descriptor: Int32 + private let lock = NSLock() + private var consumer: (@Sendable (Data) -> Void)? + + init(_ handle: FileHandle) { + self.handle = handle + descriptor = handle.fileDescriptor + } + + func start(delivering consumer: @escaping @Sendable (Data) -> Void) { + lock.withLock { self.consumer = consumer } + handle.readabilityHandler = { _ in + self.readArrivedBytes() + } + } + + /// Stops reading, first handing the consumer what the pipe already holds, up to `limit` bytes. + /// It never waits on a writer, because a helper that inherited the write end can hold it open + /// for as long as it lives. + func stop(drainingUpTo limit: Int = 0) { + stopReading { consumer in + var taken = 0 + while taken < limit, DescriptorRead.hasInputWithoutWaiting(descriptor) { + let wanted = min(DescriptorRead.pipeCapacity, limit - taken) + guard let chunk = try? DescriptorRead.availableBytes(from: descriptor, upTo: wanted), + !chunk.isEmpty else { return } + taken += chunk.count + consumer(chunk) + } + } + } + + /// Stops reading once every writer has closed the pipe, handing the consumer everything written + /// before that. A helper that inherited the write end keeps this waiting for as long as it lives. + func stopAtEndOfFile() { + stopReading { consumer in + while let chunk = try? DescriptorRead.availableBytes(from: descriptor), !chunk.isEmpty { + consumer(chunk) + } + } + } + + private func stopReading(thenDrainInto drain: (@Sendable (Data) -> Void) -> Void) { + handle.readabilityHandler = nil + lock.lock() + defer { lock.unlock() } + guard let consumer else { return } + self.consumer = nil + drain(consumer) + } + + private func readArrivedBytes() { + lock.lock() + defer { lock.unlock() } + guard let consumer else { return } + do { + let chunk = try DescriptorRead.availableBytes(from: descriptor) + guard !chunk.isEmpty else { + endReading() + return + } + consumer(chunk) + } catch POSIXError.EAGAIN { + return + } catch { + Self.logger.error("Reading a pipe failed: \(error.publicLogShape, privacy: .public)") + endReading() + } + } + + private func endReading() { + consumer = nil + handle.readabilityHandler = nil + } +} diff --git a/TablePro/Core/Process/SupervisedProcessRunner.swift b/TablePro/Core/Process/SupervisedProcessRunner.swift index c172ebacb9..aa27414c26 100644 --- a/TablePro/Core/Process/SupervisedProcessRunner.swift +++ b/TablePro/Core/Process/SupervisedProcessRunner.swift @@ -23,15 +23,10 @@ final class ProcessSupervisedRunner: SupervisedProcessRunner, @unchecked Sendabl private let process = Process() private let stdoutPipe = Pipe() private let stderrPipe = Pipe() + private let stdoutReader: PipeReader + private let stderrReader: PipeReader private let stateLock = NSLock() - /// Held for the whole of a stderr read, from taking the bytes off the pipe to handing the - /// lines to the stream, and for the whole of `finish`. The readability callback runs on - /// Foundation's queue and the termination handler on another thread, so without it a chunk - /// read just before the process exited could still be on its way to the stream when - /// `finish` drained an already empty pipe and closed the stream under it. - private let ingestLock = NSLock() - private static let forcedTerminationGrace = Duration.seconds(2) private var partialLine = "" @@ -46,6 +41,8 @@ final class ProcessSupervisedRunner: SupervisedProcessRunner, @unchecked Sendabl var continuation: AsyncStream.Continuation! stderrLines = AsyncStream(bufferingPolicy: .bufferingNewest(100)) { continuation = $0 } stderrContinuation = continuation + stdoutReader = PipeReader(stdoutPipe.fileHandleForReading) + stderrReader = PipeReader(stderrPipe.fileHandleForReading) } var processIdentifier: Int32? { @@ -61,17 +58,9 @@ final class ProcessSupervisedRunner: SupervisedProcessRunner, @unchecked Sendabl process.standardOutput = stdoutPipe process.standardError = stderrPipe - stdoutPipe.fileHandleForReading.readabilityHandler = { handle in - _ = handle.availableData - } - - stderrPipe.fileHandleForReading.readabilityHandler = { [weak self] handle in - guard let self else { return } - self.ingestLock.lock() - defer { self.ingestLock.unlock() } - let chunk = handle.availableData - guard !chunk.isEmpty else { return } - self.ingestStderr(chunk) + stdoutReader.start { _ in } + stderrReader.start { [weak self] chunk in + self?.ingestStderr(chunk) } process.terminationHandler = { [weak self] proc in @@ -155,21 +144,14 @@ final class ProcessSupervisedRunner: SupervisedProcessRunner, @unchecked Sendabl } } + /// The termination handler can run before the pipe delivers its last readability callback, so + /// what is still buffered is drained on the way out: a dropped final line is how a process that + /// announced itself ready right before exiting reads as one that never did. A callback that was + /// already dispatched has finished its delivery by the time the reader stops, so no line can + /// reach the stream after it has been closed. private func finish(exitCode: Int32) { - ingestLock.lock() - defer { ingestLock.unlock() } - stdoutPipe.fileHandleForReading.readabilityHandler = nil - stderrPipe.fileHandleForReading.readabilityHandler = nil - - /// The termination handler can run before the pipe delivers its last readability - /// callback, and clearing the handler above cancels that callback outright. Whatever is - /// still buffered is drained here, because a dropped final line is how a process that - /// announced itself ready right before exiting reads as one that never did. A callback - /// already dispatched waits on the lock and then finds the pipe at end of file, so it - /// cannot hand the stream a line after it has been closed. - if let remaining = try? stderrPipe.fileHandleForReading.readToEnd(), !remaining.isEmpty { - ingestStderr(remaining) - } + stdoutReader.stop() + stderrReader.stopAtEndOfFile() stateLock.lock() guard terminationResult == nil else { diff --git a/TablePro/Core/Services/Infrastructure/PreConnectHookRunner.swift b/TablePro/Core/Services/Infrastructure/PreConnectHookRunner.swift index 4dbc92fed8..1730f4d3d6 100644 --- a/TablePro/Core/Services/Infrastructure/PreConnectHookRunner.swift +++ b/TablePro/Core/Services/Infrastructure/PreConnectHookRunner.swift @@ -59,11 +59,9 @@ enum PreConnectHookRunner { // the pipe buffer fills and the child blocks on write — deadlocking // with waitUntilExit() on the parent side. let stderrCollector = StderrCollector() - stderrPipe.fileHandleForReading.readabilityHandler = { handle in - let chunk = handle.availableData - if !chunk.isEmpty { - stderrCollector.append(chunk) - } + let stderrReader = PipeReader(stderrPipe.fileHandleForReading) + stderrReader.start { chunk in + stderrCollector.append(chunk) } try process.run() @@ -78,7 +76,7 @@ enum PreConnectHookRunner { process.waitUntilExit() timeoutTask.cancel() - stderrPipe.fileHandleForReading.readabilityHandler = nil + stderrReader.stop(drainingUpTo: DescriptorRead.pipeCapacity) let stderr = stderrCollector.result if process.terminationReason == .uncaughtSignal { diff --git a/TableProTests/Core/MCP/Integration/MCPBridgeIntegrationTests.swift b/TableProTests/Core/MCP/Integration/MCPBridgeIntegrationTests.swift index 311040790e..6bfea553b9 100644 --- a/TableProTests/Core/MCP/Integration/MCPBridgeIntegrationTests.swift +++ b/TableProTests/Core/MCP/Integration/MCPBridgeIntegrationTests.swift @@ -1,8 +1,9 @@ import Foundation import TableProPluginKit -@testable import TablePro import XCTest +@testable import TablePro + final class MCPBridgeIntegrationTests: XCTestCase { private var upstream: MockHttpServer! @@ -371,6 +372,7 @@ final class BridgeHarness: @unchecked Sendable { private let hostToBridge = Pipe() private let bridgeToHost = Pipe() private let output = BridgeOutputBuffer() + private let outputReader: PipeReader private let proxy: BridgeProxy private let runTask: Task @@ -383,9 +385,8 @@ final class BridgeHarness: @unchecked Sendable { stdout: bridgeToHost.fileHandleForWriting ) let buffer = output - bridgeToHost.fileHandleForReading.readabilityHandler = { handle in - let chunk = handle.availableData - guard !chunk.isEmpty else { return } + outputReader = PipeReader(bridgeToHost.fileHandleForReading) + outputReader.start { chunk in buffer.append(chunk) } let runner = proxy @@ -414,7 +415,7 @@ final class BridgeHarness: @unchecked Sendable { } func shutdown() { - bridgeToHost.fileHandleForReading.readabilityHandler = nil + outputReader.stop() try? hostToBridge.fileHandleForWriting.close() runTask.cancel() try? bridgeToHost.fileHandleForWriting.close() diff --git a/TableProTests/Core/Process/DescriptorReadTests.swift b/TableProTests/Core/Process/DescriptorReadTests.swift new file mode 100644 index 0000000000..2999df0962 --- /dev/null +++ b/TableProTests/Core/Process/DescriptorReadTests.swift @@ -0,0 +1,94 @@ +// +// DescriptorReadTests.swift +// TableProTests +// + +import Darwin +import Foundation +import Testing + +@testable import TablePro + +struct DescriptorReadTests { + @Test("A read returns what has arrived without waiting for the rest") + func returnsWhatHasArrived() throws { + let pipe = Pipe() + defer { withExtendedLifetime(pipe) {} } + try pipe.fileHandleForWriting.write(contentsOf: Data("first".utf8)) + + let bytes = try DescriptorRead.availableBytes(from: pipe.fileHandleForReading.fileDescriptor) + + #expect(bytes == Data("first".utf8)) + } + + @Test("A read takes no more than its limit") + func stopsAtTheLimit() throws { + let pipe = Pipe() + defer { withExtendedLifetime(pipe) {} } + try pipe.fileHandleForWriting.write(contentsOf: Data("abcdef".utf8)) + let descriptor = pipe.fileHandleForReading.fileDescriptor + + #expect(try DescriptorRead.availableBytes(from: descriptor, upTo: 4) == Data("abcd".utf8)) + #expect(try DescriptorRead.availableBytes(from: descriptor) == Data("ef".utf8)) + } + + @Test("End of file reads as no bytes") + func endOfFileIsEmpty() throws { + let pipe = Pipe() + defer { withExtendedLifetime(pipe) {} } + try pipe.fileHandleForWriting.close() + + let bytes = try DescriptorRead.availableBytes(from: pipe.fileHandleForReading.fileDescriptor) + + #expect(bytes.isEmpty) + } + + @Test("An empty non-blocking pipe is a thrown EAGAIN, the read that raised in a dump's stderr callback") + func emptyNonBlockingPipeThrows() throws { + let pipe = Pipe() + defer { withExtendedLifetime(pipe) {} } + let descriptor = pipe.fileHandleForReading.fileDescriptor + #expect(fcntl(descriptor, F_SETFL, fcntl(descriptor, F_GETFL) | O_NONBLOCK) != -1) + + let error = #expect(throws: POSIXError.self) { + try DescriptorRead.availableBytes(from: descriptor) + } + + #expect(error?.code == .EAGAIN) + } + + @Test("A descriptor that is not open is a thrown EBADF") + func invalidDescriptorThrows() { + let error = #expect(throws: POSIXError.self) { + try DescriptorRead.availableBytes(from: -1) + } + + #expect(error?.code == .EBADF) + } + + /// End of file comes once every holder of the write end has let go, and a child this process is + /// spawning at that moment briefly holds it too. Measured: with four threads launching + /// `/usr/bin/true`, 14 of 20,000 closes were not yet end of file at once, and 0 of 20,000 + /// without them; beside the other process suites here, 9 of 200 were not, and all 200 were + /// 200ms later. So the last check waits for it rather than expecting it on the spot. + @Test("Only a pipe holding bytes or at end of file reads without waiting", .timeLimit(.minutes(1))) + func readinessFollowsThePipe() async throws { + let pipe = Pipe() + defer { withExtendedLifetime(pipe) {} } + let descriptor = pipe.fileHandleForReading.fileDescriptor + #expect(!DescriptorRead.hasInputWithoutWaiting(descriptor)) + + try pipe.fileHandleForWriting.write(contentsOf: Data("x".utf8)) + #expect(DescriptorRead.hasInputWithoutWaiting(descriptor)) + + _ = try DescriptorRead.availableBytes(from: descriptor) + #expect(!DescriptorRead.hasInputWithoutWaiting(descriptor)) + + try pipe.fileHandleForWriting.close() + let deadline = ContinuousClock.now + .seconds(10) + while !DescriptorRead.hasInputWithoutWaiting(descriptor), ContinuousClock.now < deadline { + try await Task.sleep(for: .milliseconds(5)) + } + #expect(DescriptorRead.hasInputWithoutWaiting(descriptor)) + } +} diff --git a/TableProTests/Core/Process/PipeReaderTests.swift b/TableProTests/Core/Process/PipeReaderTests.swift new file mode 100644 index 0000000000..4b7173b47d --- /dev/null +++ b/TableProTests/Core/Process/PipeReaderTests.swift @@ -0,0 +1,194 @@ +// +// PipeReaderTests.swift +// TableProTests +// + +import Darwin +import Foundation +import Testing + +@testable import TablePro + +struct PipeReaderTests { + private final class Received: @unchecked Sendable { + private let lock = NSLock() + private var bytes = Data() + + func append(_ chunk: Data) { + lock.withLock { bytes.append(chunk) } + } + + var data: Data { + lock.withLock { bytes } + } + } + + private final class Flag: @unchecked Sendable { + private let lock = NSLock() + private var raised = false + + func raise() { + lock.withLock { raised = true } + } + + var isRaised: Bool { + lock.withLock { raised } + } + } + + private struct TimedOut: Error {} + + private func waitUntil(_ condition: () -> Bool) async throws { + let deadline = ContinuousClock.now + .seconds(10) + while !condition() { + guard ContinuousClock.now < deadline else { throw TimedOut() } + try await Task.sleep(for: .milliseconds(5)) + } + } + + @Test("Everything written arrives in order and reading ends by itself at end of file", .timeLimit(.minutes(1))) + func deliversInOrderUntilEndOfFile() async throws { + let pipe = Pipe() + defer { withExtendedLifetime(pipe) {} } + let handle = pipe.fileHandleForReading + let received = Received() + let reader = PipeReader(handle) + reader.start { received.append($0) } + + let payload = Data((0 ..< 3 * DescriptorRead.pipeCapacity).map { UInt8($0 % 251) }) + let writer = pipe.fileHandleForWriting + Thread.detachNewThread { + try? writer.write(contentsOf: payload) + try? writer.close() + } + + try await waitUntil { received.data.count == payload.count && handle.readabilityHandler == nil } + #expect(received.data == payload) + } + + /// The order that crashed a dump: the process exits, the termination handler drains the pipe + /// while a readability callback Foundation had already dispatched waits, and the callback then + /// reads a pipe that is empty with its writer still open. + @Test("A callback dispatched before a drain reads nothing when it runs after it", .timeLimit(.minutes(1))) + func dispatchedCallbackAfterDrainReadsNothing() throws { + let pipe = Pipe() + defer { withExtendedLifetime(pipe) {} } + let handle = pipe.fileHandleForReading + let received = Received() + let reader = PipeReader(handle) + reader.start { received.append($0) } + let dispatched = try #require(handle.readabilityHandler) + + try pipe.fileHandleForWriting.write(contentsOf: Data("early".utf8)) + reader.stop(drainingUpTo: DescriptorRead.pipeCapacity) + #expect(received.data == Data("early".utf8)) + + dispatched(handle) + try pipe.fileHandleForWriting.write(contentsOf: Data("late".utf8)) + dispatched(handle) + + #expect(received.data == Data("early".utf8)) + let descriptor = handle.fileDescriptor + #expect(fcntl(descriptor, F_GETFL) & O_NONBLOCK == 0) + #expect(try DescriptorRead.availableBytes(from: descriptor) == Data("late".utf8)) + } + + @Test("A callback that finds nothing to read keeps the reader going", .timeLimit(.minutes(1))) + func callbackWithNothingToReadIsHarmless() async throws { + let pipe = Pipe() + defer { withExtendedLifetime(pipe) {} } + let handle = pipe.fileHandleForReading + let descriptor = handle.fileDescriptor + #expect(fcntl(descriptor, F_SETFL, fcntl(descriptor, F_GETFL) | O_NONBLOCK) != -1) + let received = Received() + let reader = PipeReader(handle) + reader.start { received.append($0) } + let dispatched = try #require(handle.readabilityHandler) + + dispatched(handle) + #expect(received.data.isEmpty) + #expect(handle.readabilityHandler != nil) + + try pipe.fileHandleForWriting.write(contentsOf: Data("after".utf8)) + try await waitUntil { received.data == Data("after".utf8) } + reader.stop() + } + + @Test("A drain takes what the pipe holds up to its limit and leaves the descriptor blocking") + func drainStopsAtTheLimitWithoutWaitingOnTheWriter() throws { + let pipe = Pipe() + defer { withExtendedLifetime(pipe) {} } + let handle = pipe.fileHandleForReading + let received = Received() + let reader = PipeReader(handle) + reader.start { received.append($0) } + handle.readabilityHandler = nil + + let buffered = Data((0 ..< 10_000).map { UInt8($0 % 251) }) + try pipe.fileHandleForWriting.write(contentsOf: buffered) + reader.stop(drainingUpTo: 4_096) + + #expect(received.data == buffered.prefix(4_096)) + let descriptor = handle.fileDescriptor + #expect(fcntl(descriptor, F_GETFL) & O_NONBLOCK == 0) + #expect(try DescriptorRead.availableBytes(from: descriptor) == buffered.dropFirst(4_096)) + } + + @Test("Stopping at end of file takes everything written until the last writer closes", .timeLimit(.minutes(1))) + func stopAtEndOfFileWaitsForTheWriter() throws { + let pipe = Pipe() + defer { withExtendedLifetime(pipe) {} } + let handle = pipe.fileHandleForReading + let received = Received() + let reader = PipeReader(handle) + reader.start { received.append($0) } + handle.readabilityHandler = nil + + let writer = pipe.fileHandleForWriting + try writer.write(contentsOf: Data("before".utf8)) + Thread.detachNewThread { + Thread.sleep(forTimeInterval: 0.2) + try? writer.write(contentsOf: Data(" after".utf8)) + try? writer.close() + } + reader.stopAtEndOfFile() + + #expect(received.data == Data("before after".utf8)) + } + + @Test("Stop waits for a delivery in progress, and nothing is delivered once it returns", .timeLimit(.minutes(1))) + func stopWaitsForDeliveryInProgress() async throws { + let pipe = Pipe() + defer { withExtendedLifetime(pipe) {} } + let handle = pipe.fileHandleForReading + let received = Received() + let entered = Flag() + let release = DispatchSemaphore(value: 0) + let reader = PipeReader(handle) + reader.start { chunk in + entered.raise() + release.wait() + received.append(chunk) + } + + try pipe.fileHandleForWriting.write(contentsOf: Data("first".utf8)) + try await waitUntil { entered.isRaised } + + let stopped = Flag() + Thread.detachNewThread { + reader.stop(drainingUpTo: DescriptorRead.pipeCapacity) + stopped.raise() + } + try await waitUntil { handle.readabilityHandler == nil } + try await Task.sleep(for: .milliseconds(50)) + #expect(!stopped.isRaised) + + release.signal() + try await waitUntil { stopped.isRaised } + #expect(received.data == Data("first".utf8)) + + try pipe.fileHandleForWriting.write(contentsOf: Data("second".utf8)) + try await Task.sleep(for: .milliseconds(50)) + #expect(received.data == Data("first".utf8)) + } +} diff --git a/TableProTests/Core/Services/Infrastructure/PreConnectHookRunnerTests.swift b/TableProTests/Core/Services/Infrastructure/PreConnectHookRunnerTests.swift new file mode 100644 index 0000000000..216544241e --- /dev/null +++ b/TableProTests/Core/Services/Infrastructure/PreConnectHookRunnerTests.swift @@ -0,0 +1,45 @@ +// +// PreConnectHookRunnerTests.swift +// TableProTests +// + +import Foundation +import Testing + +@testable import TablePro + +struct PreConnectHookRunnerTests { + private static let message = "vault login failed: token expired" + + private static func failureMessage() async -> String? { + do { + try await PreConnectHookRunner.run(script: "printf '%s' '\(Self.message)' >&2; exit 3") + return nil + } catch PreConnectHookRunner.HookError.scriptFailed(let exitCode, let stderr) { + return exitCode == 3 ? stderr : nil + } catch { + return nil + } + } + + /// A script that writes its reason and exits at once can still have it in the pipe when the + /// runner reads what it collected. Measured with sixteen scripts at a time: 352 of 16,000 + /// failures came back without their message before the pipe was drained on the way out, and 0 + /// of 16,000 after. + @Test("A failing script's message reaches the error every time", .timeLimit(.minutes(1))) + func failureCarriesTheScriptsMessage() async { + for _ in 0 ..< 4 { + let messages = await withTaskGroup(of: String?.self) { group in + for _ in 0 ..< 16 { + group.addTask { await Self.failureMessage() } + } + var collected: [String?] = [] + for await message in group { + collected.append(message) + } + return collected + } + #expect(messages.allSatisfy { $0 == Self.message }) + } + } +} diff --git a/project.yml b/project.yml index f9d0849a32..65b08ce77b 100644 --- a/project.yml +++ b/project.yml @@ -346,6 +346,7 @@ targets: - TablePro/Core/MCP/Transport/MCPStreamableHttpClientTransport.swift - TablePro/Core/MCP/Transport/MCPUpstreamCredentials.swift - TablePro/Core/MCP/Wire + - TablePro/Core/Process/DescriptorRead.swift - TablePro/Core/Services/Infrastructure/BackgroundLaunchFlag.swift dependencies: - package: TableProCore