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
8 changes: 8 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -7,8 +7,16 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0

## [Unreleased]

### Changed

- PostgreSQL, Redshift and CockroachDB keep one connection per database in the connections strip, and send TCP keepalives so idle ones stay open.

### Fixed

- Switching between two databases of a PostgreSQL connection reconnecting each time and dropping the open transaction and temp tables.
- A tab on a database the connection had switched away from running on a shared connection that closed after 10 minutes idle.
- Import and Copy To committing a transaction left open on the target connection.
- Health check reconnecting a PostgreSQL session that was sitting in a failed transaction.
- Each switch between connections in the connections strip reloading that connection's schema.
- Opening a connection fetching every column of its schema twice.
- Unsaved grid edits not kept with their tab after jumping to a tab of a background connection.
Expand Down
32 changes: 32 additions & 0 deletions Plugins/PostgreSQLDriverPlugin/LibPQConnectionLoss.swift
Original file line number Diff line number Diff line change
Expand Up @@ -51,6 +51,38 @@ enum LibPQServerMessage {
}
}

/// How a health check asks whether a session is still there without disturbing what the user has
/// open in it.
internal enum LibPQSessionCheck {
static let statement = "SELECT 1"

/// Never inside a transaction block. In an aborted one the server refuses every statement with
/// `25P02` while the session is fine, and in an open one a statement can take a repeatable-read
/// snapshot early and resets `idle_in_transaction_session_timeout`. Reading the socket, which
/// the caller does first, still sees a server that closed the session.
static func sendsStatement(in state: LibPQTransactionState) -> Bool {
switch state {
case .inTransaction, .inError:
return false
case .idle, .active, .unknown:
return true
}
}

/// The SQLSTATE of a server that refused the check statement, or nil when the failure is not
/// the server's answer. libpq's own failures carry no SQLSTATE and a lost session arrives as
/// `LibPQConnectionLostError`. A FATAL carries a SQLSTATE too, so the caller also requires
/// `PQstatus` to read `CONNECTION_OK` afterwards: `PQexec` reads past a FATAL to the closed
/// socket and turns it `CONNECTION_BAD`.
static func refusalState(of error: Error) -> String? {
guard let error = error as? LibPQPluginError,
let sqlState = error.sqlState,
!sqlState.isEmpty
else { return nil }
return sqlState
}
}

enum LibPQConnectionLoss: Sendable, Equatable {
case beforeSending(transactionMayBeOpen: Bool)
case afterSending
Expand Down
13 changes: 13 additions & 0 deletions Plugins/PostgreSQLDriverPlugin/LibPQConnectionString.swift
Original file line number Diff line number Diff line change
Expand Up @@ -74,6 +74,18 @@ internal enum LibPQConnectionString {
static let sessionApplicationName = "TablePro"
static let metadataApplicationName = "TablePro Metadata"

/// Left to the kernel, macOS sends the first keepalive after 7,200 s of idle, while an AWS
/// Network Load Balancer drops an idle flow after 350 s and an Azure Load Balancer silently after
/// 4 minutes. These keep a parked session's flow open and let the kernel declare a dead peer
/// about 90 s after the last traffic. `tcp_user_timeout` is left out: macOS has no
/// `TCP_USER_TIMEOUT`, so libpq accepts it and changes nothing.
static let keepaliveParameters: [(String, String)] = [
("keepalives", "1"),
("keepalives_idle", "60"),
("keepalives_interval", "10"),
("keepalives_count", "3")
]

static func build(
host: String,
port: Int,
Expand All @@ -100,6 +112,7 @@ internal enum LibPQConnectionString {
if let connectTimeoutSeconds {
parameters.append(("connect_timeout", String(max(connectTimeoutSeconds, 1))))
}
parameters.append(contentsOf: keepaliveParameters)

parameters.append(("sslmode", LibPQSSLMapping.sslmode(for: sslConfig.mode)))
if sslConfig.verifiesCertificate, !sslConfig.caCertificatePath.isEmpty {
Expand Down
2 changes: 1 addition & 1 deletion Plugins/PostgreSQLDriverPlugin/LibPQDriverCore.swift
Original file line number Diff line number Diff line change
Expand Up @@ -164,7 +164,7 @@ final class LibPQDriverCore: @unchecked Sendable {
guard let pqConn = libpqConnection else {
throw LibPQPluginError.notConnected
}
_ = try await pqConn.executeQuery("SELECT 1")
try await pqConn.ping()
}

// MARK: - Query Execution
Expand Down
24 changes: 24 additions & 0 deletions Plugins/PostgreSQLDriverPlugin/LibPQPluginConnection.swift
Original file line number Diff line number Diff line change
Expand Up @@ -658,6 +658,30 @@ final class LibPQPluginConnection: @unchecked Sendable {
return Self.transactionState(PQtransactionStatus(conn))
}

/// libpq has no round trip that sends no statement, so the socket is read first: that sees a
/// server that closed the session, and inside a transaction block it is the whole check
/// (`LibPQSessionCheck.sendsStatement`). A statement the server refuses still proves the
/// backend answered, so it does not fail the ping.
func ping() async throws {
try await pluginDispatchAsync(on: queue) { [self] in
guard !isShuttingDown, let conn = connectionHandle else { throw LibPQPluginError.notConnected }
if let ended = sessionEndedBeforeSending(conn) { throw ended }
guard LibPQSessionCheck.sendsStatement(in: transactionStateOnQueue()) else { return }
do {
if let deadline = activeConnectDeadline {
_ = try executeConnectQuerySync(LibPQSessionCheck.statement, deadline: deadline)
} else {
_ = try executeQuerySync(LibPQSessionCheck.statement)
}
} catch {
guard PQstatus(conn) == CONNECTION_OK,
let sqlState = LibPQSessionCheck.refusalState(of: error)
else { throw error }
Self.logger.info("Ping refused with SQLSTATE \(sqlState, privacy: .public) on a live session")
}
}
}

func boundedQuery(_ query: String, rowCap: Int) async throws -> LibPQPluginQueryResult {
let queryToRun = String(query)
let cap = max(rowCap, 1)
Expand Down
84 changes: 55 additions & 29 deletions TablePro/Core/Concurrency/SessionDriverGate.swift
Original file line number Diff line number Diff line change
Expand Up @@ -5,12 +5,16 @@

import Foundation

/// Serialises access to a connection's single shared driver.
/// Serialises access to a connection's shared session drivers.
///
/// The driver carries one mutable position (its current database and schema), so an
/// operation has to move it before it runs. Without ordering, two windows interleave
/// their moves and each runs against the other's database.
///
/// A connection that keeps one driver per database (see `SessionLanes`) takes one turn per
/// database: each of those drivers sits on its own database for good, so work on two databases
/// never has to wait for the other, while two operations on one database still take turns.
///
/// The body runs inline in the caller's own task rather than in a detached one, so
/// cancellation still reaches the work.
///
Expand All @@ -19,97 +23,119 @@ import Foundation
/// release must not free or hand off a turn a later session has taken since.
@MainActor
final class SessionDriverGate {
/// One turn per connection, or per database of a connection that keeps a driver per database.
struct Key: Hashable {
let connectionId: UUID
let database: String?
}

private struct Waiter {
let ticket: UUID
let continuation: CheckedContinuation<Void, Error>
}

private var owners: [UUID: UUID] = [:]
private var waiters: [UUID: [Waiter]] = [:]
private var owners: [Key: UUID] = [:]
private var waiters: [Key: [Waiter]] = [:]

func withExclusiveAccess<T>(
_ connectionId: UUID,
_ body: () async throws -> T
) async throws -> T {
let ticket = try await acquire(connectionId)
defer { release(connectionId, ticket: ticket) }
try await withExclusiveAccess(Key(connectionId: connectionId, database: nil), body)
}

func withExclusiveAccess<T>(
_ key: Key,
_ body: () async throws -> T
) async throws -> T {
let ticket = try await acquire(key)
defer { release(key, ticket: ticket) }
return try await body()
}

/// Whether a turn is running on `key`, which for a connection's own driver is proof that it is
/// answering without asking it again.
func isHeld(_ key: Key) -> Bool {
owners[key] != nil
}

#if DEBUG
/// How many callers are queued behind the holder, so a test can wait for one to reach the
/// gate instead of guessing how many scheduler turns that takes.
internal func waiterCount(for connectionId: UUID) -> Int {
waiters[connectionId]?.count ?? 0
waiters.filter { $0.key.connectionId == connectionId }.values.reduce(0) { $0 + $1.count }
}
#endif

/// Releases a connection that is going away, failing everyone still queued for it.
/// Releases a connection that is going away, failing everyone still queued for any of its turns.
func drain(connectionId: UUID) {
owners.removeValue(forKey: connectionId)
let pending = waiters.removeValue(forKey: connectionId) ?? []
for waiter in pending {
waiter.continuation.resume(throwing: CancellationError())
let keys = Set(owners.keys).union(waiters.keys).filter { $0.connectionId == connectionId }
for key in keys {
owners.removeValue(forKey: key)
let pending = waiters.removeValue(forKey: key) ?? []
for waiter in pending {
waiter.continuation.resume(throwing: CancellationError())
}
}
}

private func acquire(_ connectionId: UUID) async throws -> UUID {
private func acquire(_ key: Key) async throws -> UUID {
let ticket = UUID()
guard owners[connectionId] != nil else {
owners[connectionId] = ticket
guard owners[key] != nil else {
owners[key] = ticket
return ticket
}
try await withTaskCancellationHandler(
operation: { try await enqueue(ticket: ticket, connectionId: connectionId) },
operation: { try await enqueue(ticket: ticket, key: key) },
onCancel: { [weak self] in
Task { @MainActor in
self?.failWaiter(ticket: ticket, connectionId: connectionId)
self?.failWaiter(ticket: ticket, key: key)
}
}
)
/// A hand-off resumes this caller before it runs, so a drain can land in between, and the
/// turn it was handed ended with that drain.
guard owners[connectionId] == ticket else {
guard owners[key] == ticket else {
throw CancellationError()
}
return ticket
}

private func enqueue(ticket: UUID, connectionId: UUID) async throws {
private func enqueue(ticket: UUID, key: Key) async throws {
try await withCheckedThrowingContinuation { (continuation: CheckedContinuation<Void, Error>) in
guard !Task.isCancelled else {
continuation.resume(throwing: CancellationError())
return
}
waiters[connectionId, default: []].append(
waiters[key, default: []].append(
Waiter(ticket: ticket, continuation: continuation)
)
}
}

/// Removes the ticket before resuming it, so a cancellation racing a hand-off
/// can only ever find one of them.
private func failWaiter(ticket: UUID, connectionId: UUID) {
guard var pending = waiters[connectionId],
private func failWaiter(ticket: UUID, key: Key) {
guard var pending = waiters[key],
let index = pending.firstIndex(where: { $0.ticket == ticket })
else {
return
}
let waiter = pending.remove(at: index)
waiters[connectionId] = pending.isEmpty ? nil : pending
waiters[key] = pending.isEmpty ? nil : pending
waiter.continuation.resume(throwing: CancellationError())
}

private func release(_ connectionId: UUID, ticket: UUID) {
guard owners[connectionId] == ticket else { return }
guard var pending = waiters[connectionId], !pending.isEmpty else {
owners.removeValue(forKey: connectionId)
waiters.removeValue(forKey: connectionId)
private func release(_ key: Key, ticket: UUID) {
guard owners[key] == ticket else { return }
guard var pending = waiters[key], !pending.isEmpty else {
owners.removeValue(forKey: key)
waiters.removeValue(forKey: key)
return
}
let next = pending.removeFirst()
waiters[connectionId] = pending.isEmpty ? nil : pending
owners[connectionId] = next.ticket
waiters[key] = pending.isEmpty ? nil : pending
owners[key] = next.ticket
next.continuation.resume()
}
}
34 changes: 29 additions & 5 deletions TablePro/Core/Database/DatabaseManager+Health.swift
Original file line number Diff line number Diff line change
Expand Up @@ -58,6 +58,7 @@ extension DatabaseManager {
Self.logger.debug("Ping skipped — no active driver for \(connectionId)")
return false
}
await self.notePinged(mainDriver, for: connectionId)
do {
try await mainDriver.ping()
await self.markSessionVerified(connectionId)
Expand All @@ -69,7 +70,10 @@ extension DatabaseManager {
},
reconnectHandler: { [weak self] in
guard let self else { return .abort }
return await self.performHealthMonitorReconnect(connectionId: connectionId)
return await self.performHealthMonitorReconnect(
connectionId: connectionId,
failedDriver: await self.pingedDriver(for: connectionId)
)
},
onStateChanged: { [weak self] id, state in
guard let self else { return }
Expand Down Expand Up @@ -112,19 +116,33 @@ extension DatabaseManager {
/// is not a teardown, and clearing the cache here leaves the sidebar and autocomplete empty
/// with nothing scheduled to refill them. Success publishes `databaseDidConnect` so the same
/// listeners that reload after a first connect or a manual reconnect run here too.
internal func performHealthMonitorReconnect(connectionId: UUID) async -> ConnectionHealthMonitor.ReconnectOutcome {
internal func performHealthMonitorReconnect(
connectionId: UUID,
failedDriver: ObjectIdentifier? = nil
) async -> ConnectionHealthMonitor.ReconnectOutcome {
guard let session = activeSessions[connectionId] else { return .abort }
/// The check that failed was made on one driver, and a database switch may have parked it and
/// put another one in its place since. Reconnecting now would disconnect the connection the
/// user just moved onto, with its transaction; the parked one is checked again before use.
if let failedDriver, let current = session.driver, ObjectIdentifier(current) != failedDriver {
return .success
}
/// The driver this attempt is replacing. Every give-up below is fenced on it, because a
/// reconnect blocked inside a C call cannot be cancelled and completes late: without the
/// fence, a losing attempt would report a connection unreachable that a later one restored.
let attemptedDriver = session.driver
await SchemaService.shared.prepareForReload(connectionId: connectionId)
await DatabaseTreeMetadataService.shared.handleReconnect(connectionId: connectionId)
/// A connection that stopped answering has most likely taken its pooled connections with it,
/// and a rebuilt tunnel moves every one of them to a new port, so pooled work waits for the
/// replacement rather than dialing what is being torn down.
/// replacement rather than dialing what is being torn down. Begun before anything suspends,
/// so a database switch cannot promote another connection in the meantime.
MetadataConnectionPool.shared.beginTransportReplacement(connectionId: connectionId)
defer { MetadataConnectionPool.shared.endTransportReplacement(connectionId: connectionId) }
await SchemaService.shared.prepareForReload(connectionId: connectionId)
await DatabaseTreeMetadataService.shared.handleReconnect(connectionId: connectionId)
/// Asked again after the suspensions above: a switch that landed before the replacement
/// began installed another database's connection, which must not be disconnected for a
/// failure it never had.
guard activeSessions[connectionId]?.driver === attemptedDriver else { return .success }

do {
guard let result = try await trackOperation(sessionId: connectionId, operation: {
Expand Down Expand Up @@ -255,6 +273,9 @@ extension DatabaseManager {
// Rebuild the tunnel if needed; otherwise reuse effective connection
let connectionForDriver: DatabaseConnection
if session.connection.activeTunnelKind != nil {
/// Rebuilding the tunnel moves it to a new local port, which strands every other
/// database's connection on the old one.
await sessionLanes.closeAllNotingLostTransactions(for: session.connection.id)
connectionForDriver = try await buildEffectiveConnection(
for: session.connection,
deadline: deadline
Expand Down Expand Up @@ -377,6 +398,9 @@ extension DatabaseManager {
let replacesPooledTransport = session.connection.activeTunnelKind != nil || session.liveness != .live
if replacesPooledTransport {
MetadataConnectionPool.shared.beginTransportReplacement(connectionId: sessionId)
/// The other databases' connections were dialed through the same transport, so they go
/// with it and reopen on their next use.
await sessionLanes.closeAllNotingLostTransactions(for: sessionId)
}
defer {
if replacesPooledTransport {
Expand Down
Loading
Loading