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
1 change: 1 addition & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
13 changes: 10 additions & 3 deletions TablePro/CLI/BridgeProxy.swift
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down Expand Up @@ -524,12 +524,19 @@ enum BridgeJson {
}

enum BridgeStdin {
static func lines(from handle: FileHandle) -> AsyncStream<Data> {
static func lines(from handle: FileHandle, logger: any MCPBridgeLogger) -> AsyncStream<Data> {
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) {
Expand Down
63 changes: 8 additions & 55 deletions TablePro/Core/Database/ProcessNativeDumpRunner.swift
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -24,6 +23,7 @@ final class ProcessNativeDumpRunner: NativeDumpRunner, @unchecked Sendable {

init(command: NativeDumpCommand) {
self.command = command
stderrReader = PipeReader(stderrPipe.fileHandleForReading)
}

func start() throws {
Expand All @@ -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()
Expand Down Expand Up @@ -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))
Expand Down
15 changes: 8 additions & 7 deletions TablePro/Core/LSP/LSPTransport.swift
Original file line number Diff line number Diff line change
Expand Up @@ -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<Data, Error>] = [:]
private var notificationHandlers: [String: @Sendable (Data) -> Void] = [:]
Expand Down Expand Up @@ -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()

Expand Down Expand Up @@ -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()
}

Expand Down
51 changes: 51 additions & 0 deletions TablePro/Core/Process/DescriptorRead.swift
Original file line number Diff line number Diff line change
@@ -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)
}
}
}
}
94 changes: 94 additions & 0 deletions TablePro/Core/Process/PipeReader.swift
Original file line number Diff line number Diff line change
@@ -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
}
}
46 changes: 14 additions & 32 deletions TablePro/Core/Process/SupervisedProcessRunner.swift
Original file line number Diff line number Diff line change
Expand Up @@ -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 = ""
Expand All @@ -46,6 +41,8 @@ final class ProcessSupervisedRunner: SupervisedProcessRunner, @unchecked Sendabl
var continuation: AsyncStream<String>.Continuation!
stderrLines = AsyncStream<String>(bufferingPolicy: .bufferingNewest(100)) { continuation = $0 }
stderrContinuation = continuation
stdoutReader = PipeReader(stdoutPipe.fileHandleForReading)
stderrReader = PipeReader(stderrPipe.fileHandleForReading)
}

var processIdentifier: Int32? {
Expand All @@ -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
Expand Down Expand Up @@ -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 {
Expand Down
Loading
Loading