Skip to content
Closed
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
Original file line number Diff line number Diff line change
Expand Up @@ -201,6 +201,9 @@ public enum DiagnosticSessionLifecycleKind: Int, Sendable, Codable, CaseIterable
case runtimeReconfigured = 9
/// A caller explicitly invalidated one exact peer session.
case explicitlyInvalidated = 10
/// The pool evicted a session that reported no usable network path for a
/// full bounded grace window while its closure callback never fired.
case allPathsClosed = 11
}

/// Which component produced a diagnostic report.
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -64,11 +64,23 @@ actor CmxIrohClientSessionPool {
private var controlOwners: [SessionKey: ControlOwner] = [:]
private var controlWaiters: [SessionKey: [ControlWaiter]] = [:]
private var selectedPathContinuations: [UUID: AsyncStream<Void>.Continuation] = [:]
private var allPathsClosedEvictions: [
SessionKey: (sessionID: UUID, task: Task<Void, Never>)
] = [:]

/// Never-hit safety bound for a dial that ignores cancellation, retained so
/// one wedged dial can never become a permanent connect outage.
static var retiredDialSettleWaitLimitSeconds: TimeInterval { 10 }

/// Grace a pooled session gets with no usable path before the pool
/// declares it dead. The iroh boundary normally reports closure itself
/// (dead-path failover, QUIC idle timeout), so this is the level-triggered
/// backstop for a session whose closure callback never fires; the
/// 2026-07-29 outage session lost both paths and stayed pooled for 73+s
/// (https://github.com/manaflow-ai/cmux/issues/9178). Sized above a normal
/// path migration, below the point where recovery visibly stalls.
static var allPathsClosedEvictionGraceSeconds: TimeInterval { 15 }

init(
supervisor: CmxIrohEndpointSupervisor,
contextProvider: any CmxIrohClientContextProvider,
Expand Down Expand Up @@ -202,9 +214,10 @@ actor CmxIrohClientSessionPool {
}
let pathObservationTask = Task { [weak self] in
let changes = await connected.observedSelectedPathChanges()
for await _ in changes {
for await observed in changes {
guard !Task.isCancelled else { return }
await self?.publishSelectedPathChange(
await self?.handleObservedSelectedPathChange(
observed,
key: key,
sessionID: sessionID
)
Expand Down Expand Up @@ -345,6 +358,10 @@ actor CmxIrohClientSessionPool {
for key in Array(connectionTasks.keys) {
retirePendingConnection(for: key)
}
for pending in allPathsClosedEvictions.values {
pending.task.cancel()
}
allPathsClosedEvictions.removeAll(keepingCapacity: false)
let closing = sessions
let closingOwners = controlOwners
sessions.removeAll(keepingCapacity: false)
Expand Down Expand Up @@ -397,9 +414,76 @@ actor CmxIrohClientSessionPool {
return controlWaiters[key]?.count ?? 0
}

/// Reacts to one session's selected-path evidence. `.unavailable` arms a
/// bounded eviction so a session whose closure callback never fires cannot
/// sit in the pool as a corpse; any usable path disarms it.
private func handleObservedSelectedPathChange(
_ observed: CmxIrohObservedConnectionPath,
key: SessionKey,
sessionID: UUID
) {
if observed == .unavailable {
armAllPathsClosedEviction(for: key, sessionID: sessionID)
} else {
cancelAllPathsClosedEviction(for: key, sessionID: sessionID)
}
publishSelectedPathChange(key: key, sessionID: sessionID)
}

private func armAllPathsClosedEviction(for key: SessionKey, sessionID: UUID) {
guard sessions[key]?.id == sessionID else { return }
if let pending = allPathsClosedEvictions[key], pending.sessionID == sessionID {
return
}
allPathsClosedEvictions[key]?.task.cancel()
let clock = clock
let deadline = clock.now()
.addingTimeInterval(Self.allPathsClosedEvictionGraceSeconds)
let task = Task { [weak self] in
// Injected-clock sleep bounds the grace window; the task is
// cancelled when a usable path returns or the session leaves the
// pool, so a recovered connection never pays for this deadline.
try? await clock.sleep(until: deadline)
guard !Task.isCancelled else { return }
await self?.evictIfPathsStillClosed(for: key, sessionID: sessionID)
}
allPathsClosedEvictions[key] = (sessionID: sessionID, task: task)
}

private func cancelAllPathsClosedEviction(for key: SessionKey, sessionID: UUID) {
guard let pending = allPathsClosedEvictions[key],
pending.sessionID == sessionID else { return }
pending.task.cancel()
allPathsClosedEvictions[key] = nil
}

private func clearAllPathsClosedEviction(for key: SessionKey) {
allPathsClosedEvictions[key]?.task.cancel()
allPathsClosedEvictions[key] = nil
}

private func evictIfPathsStillClosed(for key: SessionKey, sessionID: UUID) async {
if let pending = allPathsClosedEvictions[key], pending.sessionID == sessionID {
allPathsClosedEvictions[key] = nil
}
guard let pooled = sessions[key], pooled.id == sessionID else { return }
// Level-triggered: re-read live path state at the deadline instead of
// trusting the event that armed this eviction.
guard await pooled.session.observedSelectedPath() == .unavailable else {
return
}
await invalidateSession(
for: key,
matching: sessionID,
reason: .allPathsClosed,
failure: .noRoute
)
}

private func sessionDidClose(key: SessionKey, sessionID: UUID) async {
guard let pooled = sessions[key], pooled.id == sessionID else { return }
let owner = controlOwners[key]
clearAllPathsClosedEviction(for: key)
sessions[key] = nil
sessionOrder.removeAll { $0 == key }
pooled.pathObservationTask.cancel()
Expand Down Expand Up @@ -442,6 +526,7 @@ actor CmxIrohClientSessionPool {
if let expectedID, sessions[key]?.id != expectedID { return }
let currentOwner = controlOwners[key]
let owner = releasesControlOwner ? currentOwner : nil
clearAllPathsClosedEviction(for: key)
retirePendingConnection(for: key)
let pooled = sessions.removeValue(forKey: key)
sessionOrder.removeAll { $0 == key }
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,155 @@
import CMUXMobileCore
import Foundation
import Testing
@testable import CmuxIrohTransport

/// Completes eviction grace windows immediately so the all-paths-closed
/// deadline fires as soon as it is armed.
private struct ImmediateGraceClock: CmxIrohRelayClock {
private let date = Date(timeIntervalSince1970: 1_800_000_000)

func now() -> Date { date }
func sleep(until _: Date) async throws {}
}

/// Holds every grace-window sleep until the test releases it, so path
/// recovery can be observed disarming the eviction deterministically.
private actor HoldingGraceClock: CmxIrohRelayClock {
private let date = Date(timeIntervalSince1970: 1_800_000_000)
private var waiters: [CheckedContinuation<Void, Never>] = []
private var sleepCount = 0

nonisolated func now() -> Date { date }

func sleep(until _: Date) async throws {
sleepCount += 1
await withCheckedContinuation { waiters.append($0) }
}

func release() {
let resumable = waiters
waiters = []
for waiter in resumable { waiter.resume() }
}

func observedSleepCount() -> Int { sleepCount }

func waitForSleeper() async -> Bool {
for _ in 0..<400 {
if sleepCount > 0 { return true }
try? await Task.sleep(nanoseconds: 5_000_000)
}
return sleepCount > 0
}
}

@Suite
struct CmxIrohClientSessionPoolPathEvictionTests {
@Test
func sessionWithNoUsablePathIsEvictedAndAttributedAfterGrace() async throws {
let fixture = try PoolFixture()
let connection = TestIrohConnection(
remoteIdentity: fixture.remoteIdentity,
bidirectionalStreams: [fixture.controlStream()],
selectedPath: .privateNetwork
)
let endpoint = TestDialingIrohEndpoint(
localIdentity: fixture.localIdentity,
dialResults: [.connection(connection)]
)
let diagnosticLog = DiagnosticLog()
let pool = try await fixture.pool(
endpoint: endpoint,
generation: 1,
diagnosticLog: diagnosticLog,
clock: ImmediateGraceClock()
)

let transport = try CmxIrohByteTransportFactory(sessionPool: pool)
.makeTransport(for: fixture.request)
try await transport.connect()
#expect(await pool.selectedObservedPath() == .privateNetwork)

await connection.setObservedSelectedPath(.unavailable)

let evicted = await pollUntilTrue {
let pathUnavailable = await pool.selectedObservedPath() == .unavailable
let closed = await connection.observedCloseCallCount() > 0
return pathUnavailable && closed
}
#expect(evicted)

let events = await drainedEvents(from: diagnosticLog)
let closure = events.last { $0.code == .sessionClosed }
#expect(closure != nil)
#expect(closure?.diagnosticFailureKind == .noRoute)
let lifecycleRemoval = events.last { event in
event.code == .transportSessionLifecycle
&& event.diagnosticSessionLifecycleKind == .allPathsClosed
}
#expect(lifecycleRemoval != nil)
}

@Test
func pathRecoveryWithinGraceDisarmsTheEviction() async throws {
let fixture = try PoolFixture()
let connection = TestIrohConnection(
remoteIdentity: fixture.remoteIdentity,
bidirectionalStreams: [fixture.controlStream()],
selectedPath: .privateNetwork
)
let endpoint = TestDialingIrohEndpoint(
localIdentity: fixture.localIdentity,
dialResults: [.connection(connection)]
)
let clock = HoldingGraceClock()
let pool = try await fixture.pool(
endpoint: endpoint,
generation: 1,
clock: clock
)

let transport = try CmxIrohByteTransportFactory(sessionPool: pool)
.makeTransport(for: fixture.request)
try await transport.connect()

await connection.setObservedSelectedPath(.unavailable)
#expect(await clock.waitForSleeper())

// The path recovers while the grace window is still pending; releasing
// the window afterwards must not evict the healthy session.
await connection.setObservedSelectedPath(.relay(url: "https://relay.example"))
let recovered = await pollUntilTrue {
await pool.selectedObservedPath() == .relay(url: "https://relay.example")
}
#expect(recovered)
await clock.release()

// Give a mistaken eviction every chance to land before asserting.
await Task.yield()
try await Task.sleep(nanoseconds: 50_000_000)
#expect(await connection.observedCloseCallCount() == 0)
#expect(await pool.selectedObservedPath() == .relay(url: "https://relay.example"))
}
Comment on lines +128 to +133

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🩺 Stability & Availability | 🔵 Trivial | ⚡ Quick win

Fixed real-time sleep before a negative assertion.

The fixed Task.sleep(nanoseconds: 50_000_000) before asserting observedCloseCallCount() == 0 is a wall-clock-dependent, single fixed-duration wait immediately preceding a correctness assertion. If a mistaken eviction landed slightly later than 50ms (e.g. under CI load), this passes without ever exercising the failure path; conversely it could occasionally flake. Since HoldingGraceClock already tracks state, prefer a deterministic check (e.g. assert no eviction task was re-armed after release()) or a bounded poll loop instead of a single fixed wait.

♻️ Example bounded-poll alternative
-        // Give a mistaken eviction every chance to land before asserting.
-        await Task.yield()
-        try await Task.sleep(nanoseconds: 50_000_000)
-        `#expect`(await connection.observedCloseCallCount() == 0)
+        // Poll for a bounded window to give a mistaken eviction every chance
+        // to land, without depending on a single fixed-duration wait.
+        for _ in 0..<10 {
+            `#expect`(await connection.observedCloseCallCount() == 0)
+            try? await Task.sleep(nanoseconds: 5_000_000)
+        }
         `#expect`(await pool.selectedObservedPath() == .relay(url: "https://relay.example"))

As per coding guidelines, {cmuxTests,cmuxUITests,ios/cmuxUITests,Packages/**/Tests,tests,tests_v2,web/tests,webviews/test}/**: "Do not use fixed sleeps, measured wall-clock assertions, or hard absolute latency ceilings in correctness tests." Note this conflicts with the **/*Tests.swift carve-out ("Test-only synchronization or sleeps are allowed") that also matches this filename; flagging per the more specific, directly-applicable test-timing rule.

📝 Committable suggestion

‼️ IMPORTANT
Carefully review the code before committing. Ensure that it accurately replaces the highlighted code, contains no missing lines, and has no issues with indentation. Thoroughly test & benchmark the code to ensure it meets the requirements.

Suggested change
// Give a mistaken eviction every chance to land before asserting.
await Task.yield()
try await Task.sleep(nanoseconds: 50_000_000)
#expect(await connection.observedCloseCallCount() == 0)
#expect(await pool.selectedObservedPath() == .relay(url: "https://relay.example"))
}
// Poll for a bounded window to give a mistaken eviction every chance
// to land, without depending on a single fixed-duration wait.
for _ in 0..<10 {
`#expect`(await connection.observedCloseCallCount() == 0)
try? await Task.sleep(nanoseconds: 5_000_000)
}
`#expect`(await pool.selectedObservedPath() == .relay(url: "https://relay.example"))
}
🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

In
`@Packages/Shared/CmuxIrohTransport/Tests/CmuxIrohTransportTests/CmxIrohClientSessionPoolPathEvictionTests.swift`
around lines 128 - 133, Replace the fixed 50ms Task.sleep in the test assertion
around selectedObservedPath with deterministic synchronization using
HoldingGraceClock state, preferably verifying that release() did not re-arm an
eviction task. If state inspection is unavailable, use a bounded polling loop
that waits for the relevant condition without relying on a single wall-clock
delay, while preserving the observedCloseCallCount() == 0 assertion.

Source: Coding guidelines

}

private func pollUntilTrue(
_ condition: @escaping () async -> Bool
) async -> Bool {
for _ in 0..<400 {
if await condition() { return true }
try? await Task.sleep(nanoseconds: 5_000_000)
}
return await condition()
}

private func drainedEvents(from log: DiagnosticLog) async -> [DiagnosticEvent] {
for _ in 0..<400 {
let report = await log.snapshot()
if report.events.contains(where: { $0.code == .sessionClosed }) {
return report.events
}
try? await Task.sleep(nanoseconds: 5_000_000)
}
return await log.snapshot().events
}
Original file line number Diff line number Diff line change
Expand Up @@ -919,7 +919,7 @@ private func waitForSelectedPathChangeCount(
return false
}

private struct PoolFixture {
struct PoolFixture {
let localIdentity: CmxIrohPeerIdentity
let remoteIdentity: CmxIrohPeerIdentity
let request: CmxByteTransportRequest
Expand Down Expand Up @@ -951,7 +951,8 @@ private struct PoolFixture {
endpoint: any CmxIrohEndpoint,
generation: UInt64,
contextProvider: (any CmxIrohClientContextProvider)? = nil,
diagnosticLog: DiagnosticLog? = nil
diagnosticLog: DiagnosticLog? = nil,
clock: (any CmxIrohRelayClock)? = nil
) async throws -> CmxIrohClientSessionPool {
let configuration = try CmxIrohEndpointConfiguration(
secretKey: CmxIrohSecretKey(bytes: Data(repeating: 7, count: 32)),
Expand All @@ -969,7 +970,8 @@ private struct PoolFixture {
contextProvider: contextProvider
?? TestIrohClientContextProvider(context: context),
protocolConfiguration: .testApplicationLanes,
diagnosticLog: diagnosticLog
diagnosticLog: diagnosticLog,
clock: clock ?? CmxIrohSystemRelayClock()
)
await pool.activate(runtimeGeneration: generation)
return pool
Expand Down