diff --git a/CHANGELOG.md b/CHANGELOG.md index 244be61b33..1483d2e490 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -130,6 +130,7 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 - Pre-connect script failures sometimes reported without the script's own error message. - Failed MongoDB statements, including writes the server rejected, reported as successful with an empty result. - `tablepro-mcp` crashing when its standard input was non-blocking. +- `tablepro-mcp` using a full CPU core, or crashing, when its standard output or error was non-blocking. - 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/BridgeMain.swift b/TablePro/CLI/BridgeMain.swift index 63292cec06..c17264d2dd 100644 --- a/TablePro/CLI/BridgeMain.swift +++ b/TablePro/CLI/BridgeMain.swift @@ -82,6 +82,6 @@ struct TableProMcpBridge { ) ) guard let data = try? JsonRpcCodec.encodeLine(envelope) else { return } - FileHandle.standardOutput.write(data) + try? DescriptorWrite.allBytes(data, to: FileHandle.standardOutput.fileDescriptor) } } diff --git a/TablePro/CLI/BridgeProxy.swift b/TablePro/CLI/BridgeProxy.swift index 4884d949f1..b9d57fcb3f 100644 --- a/TablePro/CLI/BridgeProxy.swift +++ b/TablePro/CLI/BridgeProxy.swift @@ -142,7 +142,7 @@ actor BridgeProxy { self.upstream = upstream self.discovery = discovery self.logger = logger - self.stdout = BridgeStdout(handle: stdout) + self.stdout = BridgeStdout(handle: stdout, logger: logger) self.hostLines = BridgeStdin.lines(from: stdin, logger: logger) } @@ -564,18 +564,20 @@ enum BridgeStdin { actor BridgeStdout { private let handle: FileHandle + private let logger: any MCPBridgeLogger - init(handle: FileHandle) { + init(handle: FileHandle, logger: any MCPBridgeLogger) { self.handle = handle + self.logger = logger } func write(_ payload: Data) { var line = payload line.append(0x0A) do { - try handle.write(contentsOf: line) + try DescriptorWrite.allBytes(line, to: handle.fileDescriptor) } catch { - FileHandle.standardError.write(Data("[error] stdout write failed: \(error)\n".utf8)) + logger.log(.error, "Writing stdout failed: \(error.localizedDescription)") } } } diff --git a/TablePro/Core/Database/CLIToolVersionProbe.swift b/TablePro/Core/Database/CLIToolVersionProbe.swift index 6cc5055df5..43ad412312 100644 --- a/TablePro/Core/Database/CLIToolVersionProbe.swift +++ b/TablePro/Core/Database/CLIToolVersionProbe.swift @@ -6,22 +6,13 @@ import Foundation import os -/// Asks a command line tool what it is, by running it with `--version`. -/// -/// Synchronous on purpose: `NativeDumpDescriptor.CommandLineTool`'s resolution hooks are -/// synchronous closures that `NativeDumpService` already runs inside a detached task, so an async -/// probe would have to change every one of them. enum CLIToolVersionProbe { private static let logger = Logger(subsystem: "com.TablePro", category: "CLIToolVersionProbe") static let defaultTimeout: TimeInterval = 3 - /// A version banner is one line. Reading past this is a tool doing something other than - /// answering the question, and the answer is taken from what arrived rather than waited for. static let outputCap = 64 * 1_024 - /// Standard output of ` --version`, or nil when the tool cannot run, does not answer in - /// time, or exits non-zero. static func versionOutput(of path: String, timeout: TimeInterval = defaultTimeout) -> String? { let process = Process() process.executableURL = URL(fileURLWithPath: path) diff --git a/TablePro/Core/MCP/Transport/MCPBridgeLogger.swift b/TablePro/Core/MCP/Transport/MCPBridgeLogger.swift index 7238da22ab..d29a308296 100644 --- a/TablePro/Core/MCP/Transport/MCPBridgeLogger.swift +++ b/TablePro/Core/MCP/Transport/MCPBridgeLogger.swift @@ -36,7 +36,15 @@ public struct MCPOSBridgeLogger: MCPBridgeLogger { public struct MCPStderrBridgeLogger: MCPBridgeLogger { private static let lock = NSLock() - public init() {} + private let descriptor: Int32 + + public init() { + self.init(descriptor: FileHandle.standardError.fileDescriptor) + } + + internal init(descriptor: Int32) { + self.descriptor = descriptor + } public func log(_ level: MCPBridgeLogLevel, _ message: String) { let prefix: String @@ -50,7 +58,7 @@ public struct MCPStderrBridgeLogger: MCPBridgeLogger { guard let data = payload.data(using: .utf8) else { return } Self.lock.lock() defer { Self.lock.unlock() } - FileHandle.standardError.write(data) + try? DescriptorWrite.allBytes(data, to: descriptor) } } diff --git a/TablePro/Core/Process/DescriptorWrite.swift b/TablePro/Core/Process/DescriptorWrite.swift new file mode 100644 index 0000000000..1718608b1e --- /dev/null +++ b/TablePro/Core/Process/DescriptorWrite.swift @@ -0,0 +1,42 @@ +// +// DescriptorWrite.swift +// TablePro +// + +import Darwin +import Foundation + +internal enum DescriptorWrite { + static func allBytes(_ bytes: Data, to descriptor: Int32) throws { + try bytes.withUnsafeBytes { buffer in + guard let start = buffer.baseAddress else { return } + var offset = 0 + while offset < buffer.count { + let written = Darwin.write(descriptor, start + offset, buffer.count - offset) + if written >= 0 { + offset += written + continue + } + let failure = errno + if failure == EAGAIN { + try waitForRoom(descriptor) + continue + } + try throwUnlessInterrupted(failure) + } + } + } + + private static func waitForRoom(_ descriptor: Int32) throws { + var request = pollfd(fd: descriptor, events: Int16(POLLOUT), revents: 0) + while poll(&request, 1, -1) == -1 { + try throwUnlessInterrupted(errno) + } + } + + private static func throwUnlessInterrupted(_ failure: Int32) throws { + guard failure == EINTR else { + throw POSIXError(POSIXErrorCode(rawValue: failure) ?? .EIO) + } + } +} diff --git a/TableProTests/Core/MCP/Transport/BridgeStdinTests.swift b/TableProTests/Core/MCP/Transport/BridgeStdinTests.swift index ca0e703984..961ce7df7f 100644 --- a/TableProTests/Core/MCP/Transport/BridgeStdinTests.swift +++ b/TableProTests/Core/MCP/Transport/BridgeStdinTests.swift @@ -10,6 +10,16 @@ import Testing @testable import TablePro struct BridgeStdinTests { + private static func everyLine(of stream: AsyncStream) async -> [Data]? { + await BoundedCall.result { + var lines: [Data] = [] + for await line in stream { + lines.append(line) + } + return lines + } + } + @Test("A non-blocking stdin keeps the session reading until the host closes it", .timeLimit(.minutes(1))) func nonBlockingStdinReadsUntilEndOfFile() async { let pipe = Pipe() @@ -23,10 +33,7 @@ struct BridgeStdinTests { ) let logger = RecordingBridgeLogger() - var lines: [Data] = [] - for await line in BridgeStdin.lines(from: pipe.fileHandleForReading, logger: logger) { - lines.append(line) - } + let lines = await Self.everyLine(of: BridgeStdin.lines(from: pipe.fileHandleForReading, logger: logger)) #expect(lines == [Data("{\"id\":1}".utf8), Data("{\"id\":2}".utf8)]) #expect(logger.entries.isEmpty) @@ -39,12 +46,9 @@ struct BridgeStdinTests { let directory = FileHandle(fileDescriptor: descriptor, closeOnDealloc: true) let logger = RecordingBridgeLogger() - var lines: [Data] = [] - for await line in BridgeStdin.lines(from: directory, logger: logger) { - lines.append(line) - } + let lines = await Self.everyLine(of: BridgeStdin.lines(from: directory, logger: logger)) - #expect(lines.isEmpty) + #expect(lines?.isEmpty == true) #expect(logger.entries.map(\.level) == [.error]) } } diff --git a/TableProTests/Core/MCP/Transport/BridgeStdoutTests.swift b/TableProTests/Core/MCP/Transport/BridgeStdoutTests.swift new file mode 100644 index 0000000000..7cd778d485 --- /dev/null +++ b/TableProTests/Core/MCP/Transport/BridgeStdoutTests.swift @@ -0,0 +1,48 @@ +// +// BridgeStdoutTests.swift +// TableProTests +// + +import Darwin +import Foundation +import Testing + +@testable import TablePro + +struct BridgeStdoutTests { + @Test("A line larger than the pipe reaches a non-blocking stdout whole", .timeLimit(.minutes(1))) + func nonBlockingStdoutGetsTheWholeLine() async throws { + let pipe = Pipe() + defer { withExtendedLifetime(pipe) {} } + let descriptor = pipe.fileHandleForWriting.fileDescriptor + #expect(fcntl(descriptor, F_SETFL, fcntl(descriptor, F_GETFL) | O_NONBLOCK) != -1) + #expect(fcntl(descriptor, F_SETNOSIGPIPE, 1) != -1) + let payload = Data(repeating: UInt8(ascii: "a"), count: 3 * DescriptorRead.pipeCapacity) + let logger = RecordingBridgeLogger() + let stdout = BridgeStdout(handle: pipe.fileHandleForWriting, logger: logger) + + async let drained = BackgroundPipeReader.everything(from: pipe.fileHandleForReading, pausingFirst: 0.2) + let wrote: Void? = await BoundedCall.result { await stdout.write(payload) } + try pipe.fileHandleForWriting.close() + + let received = try #require(await drained) + #expect(wrote != nil) + #expect(received == payload + Data([0x0A])) + #expect(logger.entries.isEmpty) + } + + @Test("A stdout nobody reads says why through the bridge's logger", .timeLimit(.minutes(1))) + func unwritableStdoutIsLogged() async throws { + let pipe = Pipe() + defer { withExtendedLifetime(pipe) {} } + #expect(fcntl(pipe.fileHandleForWriting.fileDescriptor, F_SETNOSIGPIPE, 1) != -1) + try pipe.fileHandleForReading.close() + let logger = RecordingBridgeLogger() + let stdout = BridgeStdout(handle: pipe.fileHandleForWriting, logger: logger) + + let wrote: Void? = await BoundedCall.result { await stdout.write(Data("{}".utf8)) } + + #expect(wrote != nil) + #expect(logger.entries.map(\.level) == [.error]) + } +} diff --git a/TableProTests/Core/MCP/Transport/MCPStderrBridgeLoggerTests.swift b/TableProTests/Core/MCP/Transport/MCPStderrBridgeLoggerTests.swift new file mode 100644 index 0000000000..7e05895f62 --- /dev/null +++ b/TableProTests/Core/MCP/Transport/MCPStderrBridgeLoggerTests.swift @@ -0,0 +1,46 @@ +// +// MCPStderrBridgeLoggerTests.swift +// TableProTests +// + +import Darwin +import Foundation +import Testing + +@testable import TablePro + +struct MCPStderrBridgeLoggerTests { + private static func fill(_ descriptor: Int32) -> Int { + let chunk = [UInt8](repeating: UInt8(ascii: "x"), count: DescriptorRead.pipeCapacity) + var total = 0 + while true { + let written = chunk.withUnsafeBytes { Darwin.write(descriptor, $0.baseAddress, $0.count) } + guard written > 0 else { return total } + total += written + } + } + + @Test("A log line waits for room on a full non-blocking stderr", .timeLimit(.minutes(1))) + func fullNonBlockingStderrGetsTheLine() async throws { + let pipe = Pipe() + defer { withExtendedLifetime(pipe) {} } + let descriptor = pipe.fileHandleForWriting.fileDescriptor + #expect(fcntl(descriptor, F_SETFL, fcntl(descriptor, F_GETFL) | O_NONBLOCK) != -1) + #expect(fcntl(descriptor, F_SETNOSIGPIPE, 1) != -1) + let filled = Self.fill(descriptor) + #expect(filled > 0) + let logger = MCPStderrBridgeLogger(descriptor: descriptor) + + async let drained = BackgroundPipeReader.everything(from: pipe.fileHandleForReading, pausingFirst: 0.2) + let logged: Void? = await BoundedCall.resultOnItsOwnThread { + logger.log(.error, "Upstream stream ended") + } + try pipe.fileHandleForWriting.close() + + let received = try #require(await drained) + let line = Data("[error] Upstream stream ended\n".utf8) + #expect(logged != nil) + #expect(received.count == filled + line.count) + #expect(received.suffix(line.count) == line) + } +} diff --git a/TableProTests/Core/Process/DescriptorReadTests.swift b/TableProTests/Core/Process/DescriptorReadTests.swift index 2d2b78f447..a2db003edb 100644 --- a/TableProTests/Core/Process/DescriptorReadTests.swift +++ b/TableProTests/Core/Process/DescriptorReadTests.swift @@ -29,6 +29,7 @@ struct DescriptorReadTests { let pipe = Pipe() defer { withExtendedLifetime(pipe) {} } try pipe.fileHandleForWriting.write(contentsOf: Data("abcdef".utf8)) + try pipe.fileHandleForWriting.close() let descriptor = pipe.fileHandleForReading.fileDescriptor #expect(try DescriptorRead.availableBytes(from: descriptor, upTo: 4) == Data("abcd".utf8)) @@ -46,15 +47,22 @@ struct DescriptorReadTests { #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 { + @Test( + "An empty non-blocking pipe is a thrown EAGAIN, the read that raised in a dump's stderr callback", + .timeLimit(.minutes(1)) + ) + func emptyNonBlockingPipeThrows() async 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 outcome = await HeldOpenWriter(pipe.fileHandleForWriting).finishedOnItsOwnThread { + Result { try DescriptorRead.availableBytes(from: descriptor) } + } + let read = try #require(outcome) let error = #expect(throws: POSIXError.self) { - try DescriptorRead.availableBytes(from: descriptor) + try read.get() } #expect(error?.code == .EAGAIN) @@ -79,7 +87,7 @@ struct DescriptorReadTests { try pipe.fileHandleForWriting.write(contentsOf: Data("x".utf8)) #expect(DescriptorRead.hasInputWithoutWaiting(descriptor)) - _ = try DescriptorRead.availableBytes(from: descriptor) + #expect(try DescriptorRead.availableBytes(from: descriptor, upTo: 1) == Data("x".utf8)) #expect(!DescriptorRead.hasInputWithoutWaiting(descriptor)) try pipe.fileHandleForWriting.close() @@ -91,15 +99,22 @@ struct DescriptorReadTests { } @Test("On a non-blocking pipe the next read waits for bytes instead of failing with EAGAIN", .timeLimit(.minutes(1))) - func nextBytesWaitsOnANonBlockingPipe() throws { + func nextBytesWaitsOnANonBlockingPipe() async throws { let pipe = Pipe() defer { withExtendedLifetime(pipe) {} } let descriptor = pipe.fileHandleForReading.fileDescriptor #expect(fcntl(descriptor, F_SETFL, fcntl(descriptor, F_GETFL) | O_NONBLOCK) != -1) BackgroundPipeWriter.write([Data("late".utf8)], to: pipe.fileHandleForWriting, pausingBeforeEach: 0.2) - #expect(try DescriptorRead.nextBytes(from: descriptor) == Data("late".utf8)) - #expect(try DescriptorRead.nextBytes(from: descriptor).isEmpty) + let arrived = try #require(await BoundedCall.resultOnItsOwnThread { + Result { try DescriptorRead.nextBytes(from: descriptor) } + }) + #expect(try arrived.get() == Data("late".utf8)) + + let ended = try #require(await BoundedCall.resultOnItsOwnThread { + Result { try DescriptorRead.nextBytes(from: descriptor) } + }) + #expect(try ended.get().isEmpty) } @Test("Buffered bytes stop at the limit and never wait on a writer that is still open", .timeLimit(.minutes(1))) diff --git a/TableProTests/Core/Process/DescriptorWriteTests.swift b/TableProTests/Core/Process/DescriptorWriteTests.swift new file mode 100644 index 0000000000..6959f806b3 --- /dev/null +++ b/TableProTests/Core/Process/DescriptorWriteTests.swift @@ -0,0 +1,67 @@ +// +// DescriptorWriteTests.swift +// TableProTests +// + +import Darwin +import Foundation +import Testing + +@testable import TablePro + +struct DescriptorWriteTests { + private struct Measured: Sendable { + let outcome: Result + let processorTime: Duration + } + + private static func threadProcessorTime() -> Duration { + .nanoseconds(Int64(clock_gettime_nsec_np(CLOCK_THREAD_CPUTIME_ID))) + } + + @Test("A write to a full non-blocking pipe waits for room instead of spinning", .timeLimit(.minutes(1))) + func writeWaitsForRoomWithoutSpinning() async throws { + let pipe = Pipe() + defer { withExtendedLifetime(pipe) {} } + let writer = pipe.fileHandleForWriting + let descriptor = writer.fileDescriptor + #expect(fcntl(descriptor, F_SETFL, fcntl(descriptor, F_GETFL) | O_NONBLOCK) != -1) + #expect(fcntl(descriptor, F_SETNOSIGPIPE, 1) != -1) + let payload = Data((0 ..< 4 * DescriptorRead.pipeCapacity).map { UInt8($0 % 251) }) + + async let drained = BackgroundPipeReader.everything(from: pipe.fileHandleForReading, pausingFirst: 1) + let measured = await BoundedCall.resultOnItsOwnThread { + let started = Self.threadProcessorTime() + let outcome = Result { try DescriptorWrite.allBytes(payload, to: descriptor) } + let processorTime = Self.threadProcessorTime() - started + try? writer.close() + return Measured(outcome: outcome, processorTime: processorTime) + } + + let written = try #require(measured) + let received = try #require(await drained) + #expect(throws: Never.self) { try written.outcome.get() } + #expect(received == payload) + #expect(written.processorTime < .milliseconds(50)) + } + + @Test("A write to a pipe nobody reads is a thrown EPIPE, not a wait", .timeLimit(.minutes(1))) + func writeWithNoReaderThrows() async throws { + let pipe = Pipe() + defer { withExtendedLifetime(pipe) {} } + let descriptor = pipe.fileHandleForWriting.fileDescriptor + #expect(fcntl(descriptor, F_SETFL, fcntl(descriptor, F_GETFL) | O_NONBLOCK) != -1) + #expect(fcntl(descriptor, F_SETNOSIGPIPE, 1) != -1) + try pipe.fileHandleForReading.close() + + let outcome = await BoundedCall.resultOnItsOwnThread { + Result { try DescriptorWrite.allBytes(Data("lost".utf8), to: descriptor) } + } + + let written = try #require(outcome) + let error = #expect(throws: POSIXError.self) { + try written.get() + } + #expect(error?.code == .EPIPE) + } +} diff --git a/TableProTests/Core/Process/PipeReaderTests.swift b/TableProTests/Core/Process/PipeReaderTests.swift index ed838ecb91..7265f46be0 100644 --- a/TableProTests/Core/Process/PipeReaderTests.swift +++ b/TableProTests/Core/Process/PipeReaderTests.swift @@ -73,15 +73,17 @@ struct PipeReaderTests { let dispatched = try #require(handle.readabilityHandler) try pipe.fileHandleForWriting.write(contentsOf: Data("early".utf8)) - let drained: Void? = await HeldOpenWriter(pipe.fileHandleForWriting).finishedOnItsOwnThread { + let writer = HeldOpenWriter(pipe.fileHandleForWriting) + let drained: Void? = await writer.finishedOnItsOwnThread { reader.stop(drainingUpTo: DescriptorRead.pipeCapacity) } try #require(drained != nil) #expect(received.data == Data("early".utf8)) try pipe.fileHandleForWriting.write(contentsOf: Data("late".utf8)) - dispatched(handle) + let returned: Void? = await writer.finishedOnItsOwnThread { dispatched(handle) } + try #require(returned != nil) #expect(received.data == Data("early".utf8)) let descriptor = handle.fileDescriptor #expect(fcntl(descriptor, F_GETFL) & O_NONBLOCK == 0) @@ -101,7 +103,8 @@ struct PipeReaderTests { reader.start { received.append($0) } let dispatched = try #require(handle.readabilityHandler) - dispatched(handle) + let returned: Void? = await BoundedCall.resultOnItsOwnThread { dispatched(handle) } + try #require(returned != nil) #expect(received.data.isEmpty) #expect(handle.readabilityHandler != nil) @@ -138,7 +141,7 @@ struct PipeReaderTests { } @Test("Stopping at end of file takes everything written until the last writer closes", .timeLimit(.minutes(1))) - func stopAtEndOfFileWaitsForTheWriter() throws { + func stopAtEndOfFileWaitsForTheWriter() async throws { let pipe = Pipe() defer { withExtendedLifetime(pipe) {} } let handle = pipe.fileHandleForReading @@ -149,8 +152,9 @@ struct PipeReaderTests { try pipe.fileHandleForWriting.write(contentsOf: Data("before".utf8)) BackgroundPipeWriter.write([Data(" after".utf8)], to: pipe.fileHandleForWriting, pausingBeforeEach: 0.2) - reader.stopAtEndOfFile() + let stopped: Void? = await BoundedCall.resultOnItsOwnThread { reader.stopAtEndOfFile() } + #expect(stopped != nil) #expect(received.data == Data("before after".utf8)) } diff --git a/TableProTests/Database/ProcessNativeDumpRunnerTests.swift b/TableProTests/Database/ProcessNativeDumpRunnerTests.swift index b384767597..7ae77202a8 100644 --- a/TableProTests/Database/ProcessNativeDumpRunnerTests.swift +++ b/TableProTests/Database/ProcessNativeDumpRunnerTests.swift @@ -19,10 +19,6 @@ struct ProcessNativeDumpRunnerTests { ) } - /// A tool that rejects one of its arguments writes its whole complaint and exits at once, and - /// the pipe still held those bytes when the buffer was read. Measured with a harness mirroring - /// this class: 2 of 300 such runs captured nothing, so the sheet said "Process exited with - /// code 7" and named no cause (#3046). @Test("A tool that exits at once still has everything it said") func stderrSurvivesAnImmediateExit() async throws { let message = "/opt/homebrew/bin/mysqldump: unknown variable 'ssl-mode=PREFERRED'" @@ -43,9 +39,6 @@ struct ProcessNativeDumpRunnerTests { #expect(result.stderr.isEmpty) } - /// A tool that leaves a child holding its standard error keeps the pipe open and writable after - /// it has gone. Nothing downstream runs until the drain returns, the temporary credentials file - /// included, so the drain stops at the cap instead of following whatever arrives next. @Test("Output from a surviving child cannot grow the buffer past its cap", .timeLimit(.minutes(1))) func outputAfterExitIsBounded() async throws { let cap = 4_096 @@ -84,22 +77,29 @@ struct ProcessNativeDumpRunnerTests { let survivingWriter = FileHandle(fileDescriptor: survivingDescriptor, closeOnDealloc: true) defer { try? survivingWriter.close() } - let runner = ProcessNativeDumpRunner( - command: command("while [ ! -e '\(gate.path)' ]; do sleep 0.01; done; printf '%s' 'refused' >&2; exit 5"), - stderrPipe: pipe - ) + let script = """ + i=0 + while [ ! -e '\(gate.path)' ] && [ $i -lt 1000 ]; do sleep 0.01; i=$((i+1)); done + [ -e '\(gate.path)' ] || exit 9 + printf '%s' 'refused' >&2 + exit 5 + """ + let runner = ProcessNativeDumpRunner(command: command(script), stderrPipe: pipe) try runner.start() + defer { runner.cancel() } let handle = pipe.fileHandleForReading let dispatched = try #require(handle.readabilityHandler) #expect(FileManager.default.createFile(atPath: gate.path, contents: nil)) - let result = try #require(await HeldOpenWriter(survivingWriter).finished { await runner.result }) + let writer = HeldOpenWriter(survivingWriter) + let result = try #require(await writer.finished { await runner.result }) #expect(result.exitCode == 5) #expect(result.stderr == "refused") try survivingWriter.write(contentsOf: Data("late".utf8)) - dispatched(handle) + let returned: Void? = await writer.finishedOnItsOwnThread { dispatched(handle) } + try #require(returned != nil) let descriptor = handle.fileDescriptor #expect(handle.readabilityHandler == nil) #expect(fcntl(descriptor, F_GETFL) & O_NONBLOCK == 0) diff --git a/TableProTests/Helpers/BackgroundPipeReader.swift b/TableProTests/Helpers/BackgroundPipeReader.swift new file mode 100644 index 0000000000..844f8fe07c --- /dev/null +++ b/TableProTests/Helpers/BackgroundPipeReader.swift @@ -0,0 +1,22 @@ +// +// BackgroundPipeReader.swift +// TableProTests +// + +import Foundation + +@testable import TablePro + +internal enum BackgroundPipeReader { + static func everything(from reader: FileHandle, pausingFirst pause: TimeInterval) async -> Data? { + let descriptor = reader.fileDescriptor + return await BoundedCall.resultOnItsOwnThread { + Thread.sleep(forTimeInterval: pause) + var bytes = Data() + while let chunk = try? DescriptorRead.availableBytes(from: descriptor), !chunk.isEmpty { + bytes.append(chunk) + } + return bytes + } + } +} diff --git a/TableProTests/Helpers/BoundedCall.swift b/TableProTests/Helpers/BoundedCall.swift new file mode 100644 index 0000000000..1e49550f95 --- /dev/null +++ b/TableProTests/Helpers/BoundedCall.swift @@ -0,0 +1,57 @@ +// +// BoundedCall.swift +// TableProTests +// + +import Foundation + +internal enum BoundedCall { + private final class FirstArrival: @unchecked Sendable { + private let lock = NSLock() + private var continuation: CheckedContinuation? + + init(_ continuation: CheckedContinuation) { + self.continuation = continuation + } + + func claim() -> CheckedContinuation? { + lock.withLock { + defer { continuation = nil } + return continuation + } + } + } + + static let deadline = Duration.seconds(10) + + static func result( + onDeadline: @escaping @Sendable () -> Void = {}, + of work: @escaping @Sendable () async -> Value + ) async -> Value? { + await withCheckedContinuation { continuation in + let arrival = FirstArrival(continuation) + let timer = Task { + try? await Task.sleep(for: deadline) + guard let pending = arrival.claim() else { return } + onDeadline() + pending.resume(returning: nil) + } + Task { + let value = await work() + arrival.claim()?.resume(returning: value) + timer.cancel() + } + } + } + + static func resultOnItsOwnThread( + onDeadline: @escaping @Sendable () -> Void = {}, + of work: @escaping @Sendable () -> Value + ) async -> Value? { + await result(onDeadline: onDeadline) { + await withCheckedContinuation { continuation in + Thread.detachNewThread { continuation.resume(returning: work()) } + } + } + } +} diff --git a/TableProTests/Helpers/HeldOpenWriter.swift b/TableProTests/Helpers/HeldOpenWriter.swift index 470a6eaeeb..1b5cd6b3e9 100644 --- a/TableProTests/Helpers/HeldOpenWriter.swift +++ b/TableProTests/Helpers/HeldOpenWriter.swift @@ -6,13 +6,6 @@ import Foundation internal struct HeldOpenWriter: Sendable { - private enum Arrival: Sendable { - case finished(Value) - case deadlinePassed - } - - private static let deadline = Duration.seconds(10) - private let handle: FileHandle init(_ handle: FileHandle) { @@ -20,27 +13,15 @@ internal struct HeldOpenWriter: Sendable { } func finished(_ work: @escaping @Sendable () async -> Value) async -> Value? { - let handle = handle - return await withTaskGroup(of: Arrival.self) { group in - group.addTask { .finished(await work()) } - group.addTask { - try? await Task.sleep(for: Self.deadline) - return .deadlinePassed - } - guard case .finished(let value)? = await group.next() else { - try? handle.close() - return nil - } - group.cancelAll() - return value - } + await BoundedCall.result(onDeadline: closeWriter, of: work) } func finishedOnItsOwnThread(_ work: @escaping @Sendable () -> Value) async -> Value? { - await finished { - await withCheckedContinuation { continuation in - Thread.detachNewThread { continuation.resume(returning: work()) } - } - } + await BoundedCall.resultOnItsOwnThread(onDeadline: closeWriter, of: work) + } + + private var closeWriter: @Sendable () -> Void { + let handle = handle + return { try? handle.close() } } } diff --git a/project.yml b/project.yml index cd67df7c40..5ef6d9d521 100644 --- a/project.yml +++ b/project.yml @@ -347,6 +347,7 @@ targets: - TablePro/Core/MCP/Transport/MCPUpstreamCredentials.swift - TablePro/Core/MCP/Wire - TablePro/Core/Process/DescriptorRead.swift + - TablePro/Core/Process/DescriptorWrite.swift - TablePro/Core/Services/Infrastructure/BackgroundLaunchFlag.swift dependencies: - package: TableProCore