From 25b87b35ab440177f9534ba4f3731e4da241522b Mon Sep 17 00:00:00 2001 From: Abdulaziz Albahar <67667005+azooz2003-bit@users.noreply.github.com> Date: Fri, 11 Sep 2026 22:39:25 -0700 Subject: [PATCH 01/17] test: preserve Iroh sessions through application stalls --- .../MobileSyncProtocolTests.swift | 42 ++++----- .../IrxLiveQUICTests.swift | 61 ++++++++++--- .../ControllableResponseTransport.swift | 16 ++++ ...obileCoreRPCNativeWriteLifetimeTests.swift | 86 +++++++++++++++++++ .../NativeObservedResponseTransport.swift | 13 +++ .../IrohConnectionRecoveryOwnerTests.swift | 24 ++---- .../MobileHostConnectionEventLaneTests.swift | 55 ++++++------ .../MobileHostConnectionLifecycleTests.swift | 10 ++- cmuxTests/MobileHostOrderedInputTests.swift | 12 ++- 9 files changed, 235 insertions(+), 84 deletions(-) create mode 100644 Packages/iOS/CmuxMobileRPC/Tests/CmuxMobileRPCTests/MobileCoreRPCNativeWriteLifetimeTests.swift create mode 100644 Packages/iOS/CmuxMobileRPC/Tests/CmuxMobileRPCTests/NativeObservedResponseTransport.swift diff --git a/Packages/Shared/CMUXMobileCore/Tests/CMUXMobileCoreTests/MobileSyncProtocolTests.swift b/Packages/Shared/CMUXMobileCore/Tests/CMUXMobileCoreTests/MobileSyncProtocolTests.swift index 6c97dde42115..d6a08205ee8e 100644 --- a/Packages/Shared/CMUXMobileCore/Tests/CMUXMobileCoreTests/MobileSyncProtocolTests.swift +++ b/Packages/Shared/CMUXMobileCore/Tests/CMUXMobileCoreTests/MobileSyncProtocolTests.swift @@ -199,39 +199,31 @@ import Testing } } -@Test func frameCodecRejectsAZeroLengthFrameFloodAtTheCallerLimit() throws { - let emptyFrame = try MobileSyncFrameCodec.encodeFrame(Data()) +@Test func frameCodecDrainsLargeBatchesWithoutTreatingPacketSizeAsFailure() throws { + let frame = try MobileSyncFrameCodec.encodeFrame(Data("valid payload".utf8)) var buffer = Data() - for _ in 0..<10_000 { - buffer.append(emptyFrame) - } - - do { - _ = try MobileSyncFrameCodec.decodeFrames( - from: &buffer, - maximumDecodedFrameCount: 16 - ) - Issue.record("Expected the decoded frame count limit to fail closed") - } catch let error as MobileSyncFrameCodecError { - #expect(error == .tooManyFrames(16)) - #expect(buffer.count == (10_000 - 16) * emptyFrame.count) + for _ in 0..<1_000 { buffer.append(frame) } + var decoded = 0 + while !buffer.isEmpty { + let frames = try MobileSyncFrameCodec.decodeFrames(from: &buffer, maximumDecodedFrameCount: 16) + #expect(!frames.isEmpty) + #expect(frames.count <= 16) + #expect(frames.allSatisfy { $0 == Data("valid payload".utf8) }) + decoded += frames.count } + #expect(decoded == 1_000) } @Test func frameCodecDefaultFrameCountLimitIsFinite() throws { - let emptyFrame = try MobileSyncFrameCodec.encodeFrame(Data()) + let frame = try MobileSyncFrameCodec.encodeFrame(Data()) var buffer = Data() - for _ in 0...MobileSyncFrameCodec.defaultMaximumDecodedFrameCount { - buffer.append(emptyFrame) - } - - #expect(throws: MobileSyncFrameCodecError.tooManyFrames( - MobileSyncFrameCodec.defaultMaximumDecodedFrameCount - )) { - _ = try MobileSyncFrameCodec.decodeFrames(from: &buffer) - } + for _ in 0...MobileSyncFrameCodec.defaultMaximumDecodedFrameCount { buffer.append(frame) } + let frames = try MobileSyncFrameCodec.decodeFrames(from: &buffer) + #expect(frames.count == MobileSyncFrameCodec.defaultMaximumDecodedFrameCount) + #expect(buffer == frame) } + private func base64URLEncode(_ data: Data) -> String { data.base64EncodedString() .replacingOccurrences(of: "+", with: "-") diff --git a/Packages/Shared/CmuxIrxTransport/Tests/CmuxIrxTransportTests/IrxLiveQUICTests.swift b/Packages/Shared/CmuxIrxTransport/Tests/CmuxIrxTransportTests/IrxLiveQUICTests.swift index ba3394a5eb31..7d456275a910 100644 --- a/Packages/Shared/CmuxIrxTransport/Tests/CmuxIrxTransportTests/IrxLiveQUICTests.swift +++ b/Packages/Shared/CmuxIrxTransport/Tests/CmuxIrxTransportTests/IrxLiveQUICTests.swift @@ -66,6 +66,49 @@ enum IrxLiveTestSupport { @Suite("live QUIC", .serialized) struct IrxLiveQUICTests { + @Test("missing application pongs preserve a usable QUIC connection") + func missingPongPreservesConnection() async throws { + let journal = IrxLiveTestSupport.journal() + let server = try await IrxLiveTestSupport.bindLoopback( + seed: IrxLiveTestSupport.identitySeed(), remoteBiCredit: 1) + let client = try await IrxLiveTestSupport.bindLoopback( + seed: IrxLiveTestSupport.identitySeed(), remoteBiCredit: 0) + let serverTask = Task { () -> (IrxConnection, IrxLaneStream)? in + guard let incoming = await server.acceptNext() else { return nil } + let accepting = try await incoming.accept() + let connection = try await accepting.connect() + let irx = IrxConnection(connection: connection, role: .acceptor, journal: journal) + guard let (_, control, _) = await IrxAdmission.performServer( + connection: irx, + judgment: IrxLiveTestSupport.fixedJudgment(accepting: "good-grant"), + journal: journal + ) else { return nil } + return (irx, control) + } + let connection = try await client.connect( + addr: IrxLiveTestSupport.loopbackAddr(of: server), alpn: IrxProtocol.alpnData) + let irx = IrxConnection(connection: connection, role: .dialer, journal: journal) + let (_, control) = try await IrxAdmission.performClient( + connection: irx, grantJWS: "good-grant", journal: journal) + let (host, hostControl) = try #require(try await serverTask.value) + let death = IrxControlReleaseProbe() + try await irx.startClientKeepalive(interval: .milliseconds(1), deadline: .milliseconds(10)) { + await death.record() + } + // Deliberately exceed both old pong deadlines while the host continues + // acknowledging QUIC packets. Its application sends no pong at all. + try await Task.sleep(for: .milliseconds(100)) + #expect(await !irx.isClosed) + #expect(await death.count == 0) + let payload = Data("still usable".utf8) + try await control.writer.write(payload) + #expect(try await hostControl.reader.readRaw() == payload) + await irx.close(code: .userRequested, origin: .local) + await host.close(code: .userRequested, origin: .local) + try? await client.close() + try? await server.close() + } + @Test("closing a control transport releases its owner without closing the session") func controlTransportReleasesOwner() async throws { let journal = IrxLiveTestSupport.journal() @@ -334,18 +377,14 @@ struct IrxLiveQUICTests { // Foreground recovery must not replace a healthy session merely // because it is older than the historical 15-second threshold. - await engine.foregroundKick(staleAfter: .seconds(15)) - var retained = false - for _ in 0..<20 { - try await Task.sleep(for: .milliseconds(50)) - if let current = await engine.currentSession(), - current.admit.session == first.admit.session - { - retained = true - break - } + await engine.foregroundKick(staleAfter: .zero) + let foregroundDeadline = ContinuousClock.now + .seconds(2) + while journal.counterSnapshot()["foreground-session-retained"] == nil, + ContinuousClock.now < foregroundDeadline { + try await Task.sleep(for: .milliseconds(10)) } - #expect(retained, "foreground recovery replaced a healthy session") + #expect(journal.counterSnapshot()["foreground-session-retained"] == 1) + #expect(await engine.currentSession()?.admit.session == first.admit.session) // Host closes (e.g. shutdown): the engine must redial by itself and // reach ready again without any external trigger. diff --git a/Packages/iOS/CmuxMobileRPC/Tests/CmuxMobileRPCTests/ControllableResponseTransport.swift b/Packages/iOS/CmuxMobileRPC/Tests/CmuxMobileRPCTests/ControllableResponseTransport.swift index 8d0b751d8e9c..58220dcd51c0 100644 --- a/Packages/iOS/CmuxMobileRPC/Tests/CmuxMobileRPCTests/ControllableResponseTransport.swift +++ b/Packages/iOS/CmuxMobileRPC/Tests/CmuxMobileRPCTests/ControllableResponseTransport.swift @@ -95,6 +95,22 @@ actor ControllableResponseTransport: CmxByteTransport { } } + func deliverCoalescedResponse(id: String, eventCount: Int) throws { + let event = try MobileSyncFrameCodec.encodeFrame( + Data(#"{"kind":"event","topic":"workspace.updated","payload":{}}"#.utf8) + ) + var batch = Data() + for _ in 0.. Data { + try MobileCoreRPCClient.requestData(method: "terminal.input", params: ["text": id], id: id) + } +} diff --git a/Packages/iOS/CmuxMobileRPC/Tests/CmuxMobileRPCTests/NativeObservedResponseTransport.swift b/Packages/iOS/CmuxMobileRPC/Tests/CmuxMobileRPCTests/NativeObservedResponseTransport.swift new file mode 100644 index 000000000000..703fff6396b1 --- /dev/null +++ b/Packages/iOS/CmuxMobileRPC/Tests/CmuxMobileRPCTests/NativeObservedResponseTransport.swift @@ -0,0 +1,13 @@ +import CMUXMobileCore +import Foundation + +/// Supplies the same independent connection-state snapshot as the Iroh adapter. +struct NativeObservedResponseTransport: CmxByteTransport, CmxByteTransportLivenessObserving { + let base: ControllableResponseTransport + + func connect() async throws { try await base.connect() } + func send(_ data: Data) async throws { try await base.send(data) } + func receive() async throws -> Data? { try await base.receive() } + func close() async { await base.close() } + func isTransportClosed() async -> Bool { await base.closed() } +} diff --git a/Packages/iOS/CmuxMobileShell/Tests/CmuxMobileShellTests/IrohConnectionRecoveryOwnerTests.swift b/Packages/iOS/CmuxMobileShell/Tests/CmuxMobileShellTests/IrohConnectionRecoveryOwnerTests.swift index 9e8e6df79c00..2c746656e855 100644 --- a/Packages/iOS/CmuxMobileShell/Tests/CmuxMobileShellTests/IrohConnectionRecoveryOwnerTests.swift +++ b/Packages/iOS/CmuxMobileShell/Tests/CmuxMobileShellTests/IrohConnectionRecoveryOwnerTests.swift @@ -398,32 +398,26 @@ extension ReconnectRouteSelectionTests { ) } - @Test func stalledWriteRedialsExactCurrentIrohClientOnce() async throws { + @Test func operationFailurePreservesOpenIrohConnection() async throws { let fixture = try await makeRecoveryOwnerFixture() defer { fixture.release() } - #expect(await fixture.store.reconnectActiveMacIfAvailable(stackUserID: "user-1")) #expect(await fixture.router.waitForCount(of: "mobile.events.subscribe", atLeast: 1)) let client = try #require(fixture.store.remoteClient) let generation = fixture.store.connectionGeneration - fixture.store.handleMacAvailabilityFailureIfCurrent( after: MobileShellConnectionError.transportWriteTimedOut, expectedClient: client, expectedGeneration: generation ) - - #expect(try await pollUntil { - guard let replacement = fixture.store.remoteClient else { return false } - let subscribeCount = await fixture.router.count(of: "mobile.events.subscribe") - return replacement !== client - && fixture.store.connectionState == .connected - && fixture.store.activeRoute?.kind == .iroh - && fixture.store.macConnectionStatus == .connected - && fixture.store.isRecoveringConnection == false - && subscribeCount >= 2 - }) - #expect(fixture.factory.attemptedKinds() == [.iroh, .iroh]) + // Exercise the connection after the failure, so preservation also + // proves that the installed client continues serving real requests. + #expect(await fixture.store.reloadWorkspaceListFromMac()) + #expect(fixture.store.remoteClient === client) + #expect(fixture.store.connectionGeneration == generation) + #expect(fixture.store.connectionState == .connected) + #expect(!fixture.store.isRecoveringConnection) + #expect(fixture.factory.attemptedKinds() == [.iroh]) } @Test func failedIrohRecoveryPreservesTheTypedFailureCategory() async throws { diff --git a/cmuxTests/MobileHostConnectionEventLaneTests.swift b/cmuxTests/MobileHostConnectionEventLaneTests.swift index 39b220fc6149..be2ce2027456 100644 --- a/cmuxTests/MobileHostConnectionEventLaneTests.swift +++ b/cmuxTests/MobileHostConnectionEventLaneTests.swift @@ -183,7 +183,7 @@ extension MobileHostAuthorizationTests { await session.close(reason: "test complete") } - @Test func testIndependentEventBackpressureClosesAtBoundedQueueCapacity() async throws { + @Test func testIndependentEventBackpressurePreservesOrderedEvents() async throws { let control = RecordingMobileHostByteTransport() let independent = TestMobileHostIndependentEventWriter( behavior: .blockAfterProbe @@ -198,10 +198,7 @@ extension MobileHostAuthorizationTests { handleRequest: { _ in .ok([:]) }, onClose: { _ in } ) - // A non-droppable topic (state-sync deltas cannot be re-derived by the - // client) keeps the close-on-overflow contract; recoverable topics like - // terminal.render_grid are shed instead — see - // testStalledRenderGridSubscriberStaysOpenWithBoundedEventQueue. + // Ordered state changes remain queued while the event lane is busy. _ = await session.debugHandleSubscriptionRPCForTesting( MobileHostRPCRequest( id: "subscribe", @@ -232,15 +229,16 @@ extension MobileHostAuthorizationTests { ) } #expect( - !(await session.sendEvent( + await session.sendEvent( topic: "mobile.sync.delta", payload: ["seq": 257] - )) + ) ) - #expect(await control.observedCloseCount() == 1) - #expect(await independent.observedCloseCount() == 1) - #expect(await session.debugQueuedEventCountForTesting() == 0) + #expect(await control.observedCloseCount() == 0) + #expect(await independent.observedCloseCount() == 0) + #expect(await session.debugQueuedEventCountForTesting() == 257) + await session.close(reason: "test complete") } @Test func testIdempotentSubscriptionDoesNotReprobeHealthyIndependentLane() async throws { @@ -333,11 +331,8 @@ extension MobileHostAuthorizationTests { #expect(await transport.observedCloseCount() == 1) } - /// Events that cannot be re-derived by the client (state-sync deltas and - /// other non-refresh topics) must keep the close-on-overflow contract: the - /// host may never silently drop them, so a subscriber that stops draining - /// is torn down at the bounded capacity instead of growing without bound. - @Test func testStalledSubscriberOverflowOnNonRecoverableTopicClosesConnection() async throws { + /// Ordered events must survive congestion without forcing a reconnect. + @Test func testStalledSubscriberPreservesOrderedEventsBeyondSheddingBudget() async throws { let transport = StalledSendMobileHostByteTransport() let session = MobileHostConnection( id: UUID(), @@ -359,9 +354,10 @@ extension MobileHostAuthorizationTests { } } - #expect(await transport.observedCloseCount() == 1) - #expect(admitted <= 258) - #expect(await session.debugQueuedEventCountForTesting() == 0) + #expect(await transport.observedCloseCount() == 0) + #expect(admitted == 300) + #expect(await session.debugQueuedEventCountForTesting() >= 299) + await session.close(reason: "test complete") } /// End-to-end fan-out proof for issue #8842: sustained emission through the @@ -464,7 +460,7 @@ extension MobileHostAuthorizationTests { /// write (TCP zero-window peer) is torn down by the bounded event-send /// stall deadline instead of pinning the connection's queue, tasks, and /// socket forever. - @Test func testEventSendStallDeadlineClosesStalledConnection() async throws { + @Test func testEventSendStallDoesNotCloseConnection() async throws { let transport = StalledSendMobileHostByteTransport() let recorder = MobileHostConnectionCloseRecorder() let connectionID = UUID() @@ -485,12 +481,12 @@ extension MobileHostAuthorizationTests { payload: ["surface_id": "surface-stall-8842", "full": true] ) await transport.waitUntilSendStalled() - for _ in 0..<2_000 { - if !(await recorder.recordedIDs().isEmpty) { break } - try await Task.sleep(nanoseconds: 1_000_000) - } - #expect(await recorder.recordedIDs() == [connectionID]) - #expect(await transport.observedCloseCount() == 1) + // Wait beyond the configured deadline to prove it cannot terminate + // an established session merely because this write has not completed. + try await Task.sleep(for: .milliseconds(30)) + #expect(await recorder.recordedIDs().isEmpty) + #expect(await transport.observedCloseCount() == 0) + await session.close(reason: "test complete") } // MARK: - Bounded event queue admission policy @@ -540,7 +536,7 @@ extension MobileHostAuthorizationTests { #expect(queue.count == 2) } - @Test func testEventQueueOverflowOnNonDroppableTopicRequestsClose() { + @Test func testEventQueuePreservesOrderedEventsBeyondSheddingBudget() { let queue = MobileHostConnectionEventQueue( maximumEventCount: 1, maximumByteCount: 1_000_000 @@ -555,8 +551,11 @@ extension MobileHostAuthorizationTests { topic: "mobile.sync.delta", coalesceKey: nil, isFullRenderGridFrame: false, frame: frame ) - #expect(!overflow.admitted) - #expect(overflow.shouldClose) + #expect(overflow.admitted) + #expect(!overflow.shouldClose) + #expect(queue.count == 2) + #expect(queue.dequeue()?.frame == frame) + #expect(queue.dequeue()?.frame == frame) } @Test func testEventQueueEnforcesByteBudgetBySheddingOldestDroppable() { diff --git a/cmuxTests/MobileHostConnectionLifecycleTests.swift b/cmuxTests/MobileHostConnectionLifecycleTests.swift index 7a7e20618fbb..15817f4bb509 100644 --- a/cmuxTests/MobileHostConnectionLifecycleTests.swift +++ b/cmuxTests/MobileHostConnectionLifecycleTests.swift @@ -613,7 +613,7 @@ extension MobileHostAuthorizationTests { let recordedMethods = await requestRecorder.recordedMethods() #expect(recordedMethods == ["workspace.list"]) } - @Test func testMobileHostConnectionClosesBeforeStartingAnUnboundedRPCBatch() async throws { + @Test func testMobileHostConnectionProcessesLargeBatchWithoutDisconnecting() async throws { let transport = RecordingMobileHostByteTransport() let invocationRecorder = MobileHostAuthorizationInvocationRecorder() let session = MobileHostConnection( @@ -637,8 +637,12 @@ extension MobileHostAuthorizationTests { await session.debugHandleReceiveDataForTesting(batch) - #expect(await transport.observedCloseCount() == 1) - #expect(await invocationRecorder.count() == 0) + let responses = await transport.waitForSentBufferCount( + MobileHostRPCWorkQuota.recommendedMaximumConcurrentRequestCount + 1 + ) + #expect(responses.count == MobileHostRPCWorkQuota.recommendedMaximumConcurrentRequestCount + 1) + #expect(await transport.observedCloseCount() == 0) + await session.close(reason: "test complete") } // MARK: - Advertised mobile host capabilities @Test func testMobileHostAdvertisesWorkspaceActionCapabilities() { diff --git a/cmuxTests/MobileHostOrderedInputTests.swift b/cmuxTests/MobileHostOrderedInputTests.swift index 0f6147290b98..4038c55a946b 100644 --- a/cmuxTests/MobileHostOrderedInputTests.swift +++ b/cmuxTests/MobileHostOrderedInputTests.swift @@ -75,9 +75,13 @@ struct MobileHostOrderedInputTests { try Self.framedBatch([("input-overflow", "terminal.input")]) ) - #expect(await transport.closeCount() == 1) + let responses = await transport.waitForResponseCount(1) + #expect(responses == ["input-overflow"]) + #expect(await transport.errorCode(for: "input-overflow") == "server_busy") + #expect(await transport.closeCount() == 0) #expect(await gate.handledRequestCount() == 1) await gate.releaseFirstInput() + await connection.close(reason: "test complete") } @Test @@ -256,6 +260,7 @@ private actor OrderedInputHandlerGate { private actor OrderedInputRecordingTransport: CmxByteTransport { private var responses: [String] = [] + private var errorCodes: [String: String] = [:] private var responseWaiters: [(Int, CheckedContinuation<[String], Never>)] = [] private var closes = 0 @@ -286,7 +291,9 @@ private actor OrderedInputRecordingTransport: CmxByteTransport { for payload in payloads { let envelope = try JSONSerialization.jsonObject(with: payload) as? [String: Any] - responses.append(envelope?["id"] as? String ?? "") + let id = envelope?["id"] as? String ?? "" + responses.append(id) + errorCodes[id] = (envelope?["error"] as? [String: Any])?["code"] as? String } let ready = responseWaiters.filter { responses.count >= $0.0 } responseWaiters.removeAll { responses.count >= $0.0 } @@ -306,6 +313,7 @@ private actor OrderedInputRecordingTransport: CmxByteTransport { } } + func errorCode(for id: String) -> String? { errorCodes[id] } func responseIDs() -> [String] { responses } func closeCount() -> Int { closes } } From 9d62c242c0dd48d6b371704e37e27470d281ebff Mon Sep 17 00:00:00 2001 From: Abdulaziz Albahar <67667005+azooz2003-bit@users.noreply.github.com> Date: Fri, 11 Sep 2026 22:59:51 -0700 Subject: [PATCH 02/17] fix: let Iroh own established connection lifetime --- .../CMUXMobileCore/MobileSyncProtocol.swift | 12 +- .../CmuxIrxTransport/IrxAdmission.swift | 77 ++++++----- .../CmuxIrxTransport/IrxConnection.swift | 100 ++++++-------- .../CmuxIrxTransport/IrxPeerEngine.swift | 36 ++--- .../CmuxIrxTransport/IrxProtocol.swift | 14 +- .../IrxLiveQUICTests.swift | 12 ++ ...bileCoreRPCSession+IndependentEvents.swift | 14 +- .../CmuxMobileRPC/MobileCoreRPCSession.swift | 38 ++++-- ...ileShellComposite+ConnectionRecovery.swift | 54 +++++++- .../MobileShellComposite.swift | 25 +--- .../IrohConnectionRecoveryOwnerTests.swift | 47 +++++-- .../MobileShellRenderGridLivenessTests.swift | 8 +- .../MobileHostConnectionEventQueue.swift | 51 ++----- Sources/Mobile/MobileHostService.swift | 128 +++++------------- .../MobileHostConnectionEventLaneTests.swift | 12 +- 15 files changed, 293 insertions(+), 335 deletions(-) diff --git a/Packages/Shared/CMUXMobileCore/Sources/CMUXMobileCore/MobileSyncProtocol.swift b/Packages/Shared/CMUXMobileCore/Sources/CMUXMobileCore/MobileSyncProtocol.swift index e23e26d95bd0..63e0f58456b6 100644 --- a/Packages/Shared/CMUXMobileCore/Sources/CMUXMobileCore/MobileSyncProtocol.swift +++ b/Packages/Shared/CMUXMobileCore/Sources/CMUXMobileCore/MobileSyncProtocol.swift @@ -201,6 +201,10 @@ public struct MobileSyncFrameCodec { return frame } + /// Decodes at most one bounded batch, leaving later frames in `buffer`. + /// Callers drain another batch before awaiting more network bytes whenever + /// the returned count reaches `maximumDecodedFrameCount`. A transport read + /// may coalesce any number of valid frames; its boundary is not a message limit. public static func decodeFrames( from buffer: inout Data, maximumFrameByteCount: Int = defaultMaximumFrameByteCount, @@ -222,7 +226,8 @@ public struct MobileSyncFrameCodec { } } - while buffer.count - consumedByteCount >= headerByteCount { + while frames.count < maximumDecodedFrameCount, + buffer.count - consumedByteCount >= headerByteCount { let frameStart = buffer.index( buffer.startIndex, offsetBy: consumedByteCount @@ -241,11 +246,6 @@ public struct MobileSyncFrameCodec { guard buffer.count - consumedByteCount >= headerByteCount + payloadLength else { break } - guard frames.count < maximumDecodedFrameCount else { - throw MobileSyncFrameCodecError.tooManyFrames( - maximumDecodedFrameCount - ) - } let payloadStart = headerEnd let payloadEnd = buffer.index( payloadStart, diff --git a/Packages/Shared/CmuxIrxTransport/Sources/CmuxIrxTransport/IrxAdmission.swift b/Packages/Shared/CmuxIrxTransport/Sources/CmuxIrxTransport/IrxAdmission.swift index 4022e1bbfb13..93ef8346f91d 100644 --- a/Packages/Shared/CmuxIrxTransport/Sources/CmuxIrxTransport/IrxAdmission.swift +++ b/Packages/Shared/CmuxIrxTransport/Sources/CmuxIrxTransport/IrxAdmission.swift @@ -56,46 +56,57 @@ public enum IrxAdmission { grantJWS: String? = nil, journal: IrxJournal ) async throws -> (IrxAdmit, IrxLaneStream) { - let startedAt = DispatchTime.now() - let control = try await connection.openLane(IrxLaneDescriptor(lane: .control)) - try await control.writer.writeControlFrame(IrxHello(grant: grantJWS)) - let admit = try await withIrxDeadline(deadline, onTimeout: { - await connection.close(code: .admissionTimeout, origin: .transport) - }) { - guard let admit = try await control.reader.readControlFrame(IrxAdmit.self) else { - throw IrxConnectionError.closed(await connection.termination()) - } - return admit - } - guard let admit else { - // A stalled QUIC read can outlive the deadline and ignore task - // cancellation. Preserve a close reason already received from the - // peer; otherwise close locally so the read loses its transport - // owner before we inspect the termination reason. - if await connection.closeReason() == nil { + do { + let startedAt = DispatchTime.now() + let control = try await connection.openLane(IrxLaneDescriptor(lane: .control)) + try await control.writer.writeControlFrame(IrxHello(grant: grantJWS)) + let admit = try await withIrxDeadline(deadline, onTimeout: { await connection.close(code: .admissionTimeout, origin: .transport) + }) { + guard let admit = try await control.reader.readControlFrame(IrxAdmit.self) else { + throw IrxConnectionError.closed(await connection.termination()) + } + return admit + } + guard let admit else { + // A stalled QUIC read can outlive the deadline and ignore task + // cancellation. Preserve a close reason already received from the + // peer; otherwise close locally so the read loses its transport + // owner before we inspect the termination reason. + if await connection.closeReason() == nil { + await connection.close(code: .admissionTimeout, origin: .transport) + } + let termination = await connection.termination() + journal.record( + "admission", "denied-or-timeout", + ["code": termination.code] + ) + if let code = IrxCloseCode(rawValue: termination.code) { + throw IrxAdmissionDenied(code: code) + } + throw IrxConnectionError.admissionTimeout } - let termination = await connection.termination() + let elapsedMs = + (DispatchTime.now().uptimeNanoseconds - startedAt.uptimeNanoseconds) / 1_000_000 journal.record( - "admission", "denied-or-timeout", - ["code": termination.code] + "admission", "admitted", + [ + "session": admit.session, + "elapsed_ms": String(elapsedMs), + "path": connection.selectedPathDescription(), + ] ) - if let code = IrxCloseCode(rawValue: termination.code) { + return (admit, control) + } catch { + // A native read/write can throw as soon as the peer's close + // arrives, before the deadline path inspects the reason. Preserve + // the denial code so authorization failure cannot become a retry loop. + if let reason = await connection.closeReason(), + let code = IrxCloseCode.parse(fromRenderedCause: reason) { throw IrxAdmissionDenied(code: code) } - throw IrxConnectionError.admissionTimeout + throw error } - let elapsedMs = - (DispatchTime.now().uptimeNanoseconds - startedAt.uptimeNanoseconds) / 1_000_000 - journal.record( - "admission", "admitted", - [ - "session": admit.session, - "elapsed_ms": String(elapsedMs), - "path": connection.selectedPathDescription(), - ] - ) - return (admit, control) } /// Server half: read the control descriptor + hello off the first stream, diff --git a/Packages/Shared/CmuxIrxTransport/Sources/CmuxIrxTransport/IrxConnection.swift b/Packages/Shared/CmuxIrxTransport/Sources/CmuxIrxTransport/IrxConnection.swift index f9251cc74746..2a7b56aec026 100644 --- a/Packages/Shared/CmuxIrxTransport/Sources/CmuxIrxTransport/IrxConnection.swift +++ b/Packages/Shared/CmuxIrxTransport/Sources/CmuxIrxTransport/IrxConnection.swift @@ -139,8 +139,7 @@ public actor IrxConnection { nonisolated public let remoteEndpointIDHex: String private let connection: Connection /// Instant of the most recent keepalive pong; nil before the first pong. - /// Foreground staleness checks read this to decide zombie-vs-live after - /// a suspension (a QUIC connection can be long dead without isClosed). + /// Diagnostic only; pong age never determines connection lifetime. public private(set) var lastPongAt: ContinuousClock.Instant? private let journal: IrxJournal private var closedFlag = false @@ -172,10 +171,8 @@ public actor IrxConnection { isClosed } - /// Returns whether the connection has demonstrated liveness recently. - /// This catches suspended-app zombie sessions whose native closed flag has - /// not flipped, while allowing a healthy long-lived session to survive a - /// foreground event. + /// Returns whether an application pong was observed recently, for diagnostics. + /// A false result is not evidence that the QUIC connection has closed. public func hasRecentKeepalive(within age: Duration) -> Bool { guard let lastPongAt else { return false } return ContinuousClock.now - lastPongAt <= age @@ -258,42 +255,50 @@ public actor IrxConnection { return writer } - /// Accepts the next bidirectional lane (server side, post-admission). - /// Returns nil once the connection is closed. + /// Accepts the next bidirectional lane. A malformed descriptor retires + /// only that stream; native connection termination ends acceptance. public func acceptLane() async -> IrxLaneStream? { - while true { + while !Task.isCancelled { do { let stream = try await connection.acceptBi() let reader = IrxStreamReader(stream.recv()) let writer = IrxStreamWriter(stream.send()) - guard - let descriptor = try await reader.readControlFrame( - IrxLaneDescriptor.self) - else { - await writer.finish() - continue + do { + if let descriptor = try await reader.readControlFrame(IrxLaneDescriptor.self) { + return IrxLaneStream(descriptor: descriptor, writer: writer, reader: reader) + } + } catch { + // Stream framing failure does not change native connection state. } - return IrxLaneStream(descriptor: descriptor, writer: writer, reader: reader) + await writer.reset(errorCode: 2) + await reader.stop(errorCode: 2) } catch { - closedFlag = true return nil } } + return nil } - /// Accepts the next unidirectional lane (client side: the events lane). + /// Accepts the next usable server event lane without treating a malformed + /// descriptor or stream EOF as complete-connection closure. public func acceptUniLane() async throws -> (IrxLaneDescriptor, IrxStreamReader)? { - do { - let stream = try await connection.acceptUni() - let reader = IrxStreamReader(stream) - guard - let descriptor = try await reader.readControlFrame(IrxLaneDescriptor.self) - else { return nil } - return (descriptor, reader) - } catch { - closedFlag = true - return nil + while !Task.isCancelled { + do { + let stream = try await connection.acceptUni() + let reader = IrxStreamReader(stream) + do { + if let descriptor = try await reader.readControlFrame(IrxLaneDescriptor.self) { + return (descriptor, reader) + } + } catch { + // A later event stream may repair this optional feature. + } + await reader.stop(errorCode: 2) + } catch { + return nil + } } + return nil } /// The selected QUIC path right now, for relay attribution evidence. @@ -305,13 +310,9 @@ public actor IrxConnection { return "\(selected.isRelay ? "relay" : "direct"):\(selected.remoteAddr)" } - /// Continuous client-side keepalive on a dedicated lane: one tiny ping - /// every interval, pong deadline enforced per ping, every exchange - /// journaled with RTT and the selected path (the soak's relay-attribution - /// evidence). A single miss re-pings immediately (journaled as a `miss`, - /// not a death: one transient stall must never sever a healthy session); - /// `IrxProtocol.keepaliveStrikeLimit` consecutive misses close with - /// `keepalive-timeout` and report death so the engine redials at once. + /// Samples application round-trip latency on an optional lane. + /// A missed pong retires only this diagnostic lane. Native Iroh keepalives + /// and connection closure observation own dead-peer detection. public func startClientKeepalive( interval: Duration = IrxProtocol.keepaliveInterval, deadline: Duration = IrxProtocol.keepaliveDeadline, @@ -320,11 +321,8 @@ public actor IrxConnection { guard keepaliveTask == nil else { return } let lane = try await openLane(IrxLaneDescriptor(lane: .keepalive)) keepaliveTask = Task { [journal] in - var strikes = 0 while !Task.isCancelled { - if strikes == 0 { - try? await Task.sleep(for: interval) - } + try? await Task.sleep(for: interval) guard !Task.isCancelled else { return } let seq = await self.nextPingSeq() let sentAt = DispatchTime.now() @@ -357,27 +355,19 @@ public actor IrxConnection { ] ) await self.notePong() - strikes = 0 } catch { guard !Task.isCancelled else { return } - strikes += 1 - if strikes < IrxProtocol.keepaliveStrikeLimit { - journal.record( - "keepalive", "miss", - [ - "seq": String(seq), - "strike": String(strikes), - "path": self.selectedPathDescription(), - ] - ) - continue - } journal.record( - "keepalive", "timeout", + "keepalive", "miss", ["seq": String(seq), "path": self.selectedPathDescription()] ) - await self.close(code: .keepaliveTimeout, origin: .transport) - await onDeath() + await lane.close() + // A stopped stream or slow application pong cannot close + // the other streams. The native termination watcher is + // still running even after this diagnostic loop ends. + if await self.isConnectionClosed() { + await onDeath() + } return } } diff --git a/Packages/Shared/CmuxIrxTransport/Sources/CmuxIrxTransport/IrxPeerEngine.swift b/Packages/Shared/CmuxIrxTransport/Sources/CmuxIrxTransport/IrxPeerEngine.swift index 66fe2420ddda..28aeb0a0e9f8 100644 --- a/Packages/Shared/CmuxIrxTransport/Sources/CmuxIrxTransport/IrxPeerEngine.swift +++ b/Packages/Shared/CmuxIrxTransport/Sources/CmuxIrxTransport/IrxPeerEngine.swift @@ -255,36 +255,20 @@ public actor IrxPeerEngine { Task { _ = try? await self.ensureSession(trigger: trigger) } } - /// Foreground resume: retain a session that recently proved liveness, but - /// replace a native zombie whose closed flag stayed false during - /// suspension. This avoids age-based churn while preserving recovery. + /// Foreground resume retains every connection that native Iroh still owns. + /// The historical age argument is diagnostic only; suspension and pong age + /// cannot invalidate an otherwise usable connection. public func foregroundKick(staleAfter: Duration = .seconds(15)) { Task { if let session = self.currentSessionForKick(), - await !session.connection.isConnectionClosed() - { - let recentlyAlive: Bool - if clockNow() - session.establishedAtMonotonic <= staleAfter { - // A newly admitted session has not necessarily completed - // its first keepalive round yet, but is still within the - // bounded fresh-session grace period. - recentlyAlive = true - } else { - recentlyAlive = await session.connection.hasRecentKeepalive( - within: staleAfter) - } - if recentlyAlive { - self.record( - "foreground-session-retained", - [ - "session": session.admit.session, - "stale_after": String(describing: staleAfter), - ] - ) - return - } + await !session.connection.isConnectionClosed() { + self.record( + "foreground-session-retained", + ["session": session.admit.session, "stale_after": String(describing: staleAfter)] + ) + return } - _ = try? await self.ensureSession(explicit: true, trigger: "foreground") + _ = try? await self.ensureSession(trigger: "foreground") } } diff --git a/Packages/Shared/CmuxIrxTransport/Sources/CmuxIrxTransport/IrxProtocol.swift b/Packages/Shared/CmuxIrxTransport/Sources/CmuxIrxTransport/IrxProtocol.swift index a5980248881d..96564d0bbb3c 100644 --- a/Packages/Shared/CmuxIrxTransport/Sources/CmuxIrxTransport/IrxProtocol.swift +++ b/Packages/Shared/CmuxIrxTransport/Sources/CmuxIrxTransport/IrxProtocol.swift @@ -11,18 +11,12 @@ public enum IrxProtocol { /// Control frames are small (hello/admit/keepalive/descriptors); anything /// larger is a protocol error, never buffered. public static let maximumControlFrameByteCount = 256 * 1024 - /// Keepalive cadence: one tiny ping per interval, pong deadline after - /// which the connection is declared dead. Hard closes (the realistic - /// relay-expiry case) are detected instantly by the termination watcher; - /// the ping loop bounds SILENT path blackholes to interval + deadline, - /// keeping worst-case detection-plus-redial inside single-digit seconds. + /// Application latency sampling cadence. Connection lifetime is owned by + /// Iroh's native keepalives and negotiated connection idle timeout. public static let keepaliveInterval: Duration = .seconds(5) + /// A missed application pong retires only the diagnostic stream. public static let keepaliveDeadline: Duration = .seconds(2) - /// Consecutive pong misses before the connection is declared dead. One - /// transient stall (relay hiccup, brief peer pause) must never sever a - /// healthy session; a re-ping fires immediately after a miss, so real - /// death still detects in ~strikeLimit x deadline. - public static let keepaliveStrikeLimit = 2 + } /// Machine-readable close/denial codes. The code travels in the QUIC diff --git a/Packages/Shared/CmuxIrxTransport/Tests/CmuxIrxTransportTests/IrxLiveQUICTests.swift b/Packages/Shared/CmuxIrxTransport/Tests/CmuxIrxTransportTests/IrxLiveQUICTests.swift index 7d456275a910..7380b19d34e9 100644 --- a/Packages/Shared/CmuxIrxTransport/Tests/CmuxIrxTransportTests/IrxLiveQUICTests.swift +++ b/Packages/Shared/CmuxIrxTransport/Tests/CmuxIrxTransportTests/IrxLiveQUICTests.swift @@ -103,6 +103,18 @@ struct IrxLiveQUICTests { let payload = Data("still usable".utf8) try await control.writer.write(payload) #expect(try await hostControl.reader.readRaw() == payload) + // A malformed descriptor on an optional event lane must also leave + // the connection usable and allow the next valid event lane through. + await irx.raiseRemoteStreamCredit(bi: 0, uni: 2) + let acceptingEvent = Task { try await irx.acceptUniLane() } + let malformed = try await host.underlying.openUni() + try await malformed.writeAll(buf: IrxFrameCodec.encode(IrxPing(seq: 1, pong: false))) + try await malformed.finish() + let events = try await host.openUniLane(IrxLaneDescriptor(lane: .events)) + let acceptedEvent = try await acceptingEvent.value + #expect(acceptedEvent?.0.lane == .events) + #expect(await !irx.isClosed) + await events.finish() await irx.close(code: .userRequested, origin: .local) await host.close(code: .userRequested, origin: .local) try? await client.close() diff --git a/Packages/iOS/CmuxMobileRPC/Sources/CmuxMobileRPC/MobileCoreRPCSession+IndependentEvents.swift b/Packages/iOS/CmuxMobileRPC/Sources/CmuxMobileRPC/MobileCoreRPCSession+IndependentEvents.swift index baff29f16f77..affcef688487 100644 --- a/Packages/iOS/CmuxMobileRPC/Sources/CmuxMobileRPC/MobileCoreRPCSession+IndependentEvents.swift +++ b/Packages/iOS/CmuxMobileRPC/Sources/CmuxMobileRPC/MobileCoreRPCSession+IndependentEvents.swift @@ -96,11 +96,15 @@ extension MobileCoreRPCSession { throw MobileSyncFrameCodecError.frameTooLarge(buffer.count + chunk.count) } buffer.append(chunk) - let frames = try MobileSyncFrameCodec.decodeFrames( - from: &buffer, - maximumDecodedFrameCount: Self.maximumDecodedFrameCountPerRead - ) - for frame in frames { dispatch(frame: frame) } + while !Task.isCancelled, independentEventReader?.id == id { + let frames = try MobileSyncFrameCodec.decodeFrames( + from: &buffer, + maximumDecodedFrameCount: Self.maximumDecodedFrameCountPerRead + ) + for frame in frames { dispatch(frame: frame) } + guard frames.count == Self.maximumDecodedFrameCountPerRead else { break } + await Task.yield() + } } } catch { // The host falls back to control delivery after optional-lane failure. diff --git a/Packages/iOS/CmuxMobileRPC/Sources/CmuxMobileRPC/MobileCoreRPCSession.swift b/Packages/iOS/CmuxMobileRPC/Sources/CmuxMobileRPC/MobileCoreRPCSession.swift index 2ac6cf35428b..165b06dad97a 100644 --- a/Packages/iOS/CmuxMobileRPC/Sources/CmuxMobileRPC/MobileCoreRPCSession.swift +++ b/Packages/iOS/CmuxMobileRPC/Sources/CmuxMobileRPC/MobileCoreRPCSession.swift @@ -1047,22 +1047,20 @@ actor MobileCoreRPCSession { return } buffer.append(chunk) - let frames: [Data] do { - frames = try MobileSyncFrameCodec.decodeFrames( - from: &buffer, - maximumDecodedFrameCount: Self.maximumDecodedFrameCountPerRead - ) + while !Task.isCancelled, installedConnectionID == connectionID { + let frames = try MobileSyncFrameCodec.decodeFrames( + from: &buffer, + maximumDecodedFrameCount: Self.maximumDecodedFrameCountPerRead + ) + for frame in frames { dispatch(frame: frame) } + guard frames.count == Self.maximumDecodedFrameCountPerRead else { break } + await Task.yield() + } } catch { - await tearDownIfInstalled( - connectionID: connectionID, - error: .invalidResponse - ) + await tearDownIfInstalled(connectionID: connectionID, error: .invalidResponse) return } - for frame in frames { - dispatch(frame: frame) - } } } @@ -1235,7 +1233,10 @@ actor MobileCoreRPCSession { /// the wedged transport installed (their timeout cannot recycle a write /// owned by another request ID). private func startQueuedDemandRecovery(requestID: String) { - guard !queuedWriteIDs.isEmpty else { return } + // Native connection observation owns Iroh's lifetime. Queued demand + // has its own request deadline and cannot condemn an unfinished frame. + guard !(transport is any CmxByteTransportLivenessObserving), + !queuedWriteIDs.isEmpty else { return } Task { [self, taskTimeout, cancelledWriteCompletionGraceNanoseconds] in let waitTask = Task { await self.awaitCancelledWriteResolution() @@ -1266,6 +1267,10 @@ actor MobileCoreRPCSession { private func waitForCancelledActiveWriteResolution( deadlineUptimeNanoseconds: UInt64 ) async throws { + // Preserve serialization until the native write completes or fails. + // Cancelling writeAll can leave a frame prefix on the control stream. + // Later requests may queue and expire without cancelling that write. + if transport is any CmxByteTransportLivenessObserving { return } while let write = activeWrite, write.cancelledRequestResolutionTask != nil { try Task.checkCancellation() @@ -1367,7 +1372,12 @@ actor MobileCoreRPCSession { } private func recycleTransportIfActiveWrite(requestID: String) async -> Bool { - guard activeWrite?.requestID == requestID else { return false } + guard let write = activeWrite, write.requestID == requestID else { return false } + if let observing = transport as? any CmxByteTransportLivenessObserving { + guard await observing.isTransportClosed() else { return false } + guard activeWrite?.connectionID == write.connectionID, + activeWrite?.requestID == requestID else { return false } + } activeWrite?.task.cancel() activeWrite?.cancelledRequestResolutionTask?.cancel() activeWrite = nil diff --git a/Packages/iOS/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite+ConnectionRecovery.swift b/Packages/iOS/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite+ConnectionRecovery.swift index 2594ca466d0d..206a484e07f5 100644 --- a/Packages/iOS/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite+ConnectionRecovery.swift +++ b/Packages/iOS/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite+ConnectionRecovery.swift @@ -148,14 +148,45 @@ extension MobileShellComposite { } } - /// A definitive event-stream failure bypasses same-client resubscription. - /// Once the exact session is proven dead, rebuilding its listener only hides - /// the failure behind the transport's reconnect behavior and leaves the - /// shell owner stale. Instead, transition the one lifecycle owner to a fresh - /// authenticated stored-Mac dial. + /// Checks native connection state before promoting a feature failure to + /// recovery. Transports without native observation retain their existing + /// error-driven recovery; Iroh owns its own dead-peer detection. func recoverDeadConnection( trigger: RecoveryTrigger, expectedClient: MobileCoreRPCClient + ) { + guard remoteClient === expectedClient, connectionState == .connected else { return } + if trigger == .eventStreamEnded { + // The RPC listener ends only after its required control session + // has torn down. Detach it synchronously so another producer + // cannot reopen the retired RPC client while recovery starts. + recoverClosedControlSession(trigger: trigger, expectedClient: expectedClient) + return + } + Task { @MainActor [weak self] in + let closed = await expectedClient.isTransportClosed() + guard let self, + self.remoteClient === expectedClient, + self.connectionState == .connected else { return } + if closed == false { + // A failed subscription or request is feature-level evidence. + // Native Iroh closure is observed independently by the RPC session. + if trigger == .subscriptionStartFailed || trigger == .eventStreamEnded { + self.resyncTerminalOutput( + reason: "event_subscription_repair", + restartEventStream: true, + recoversConnectionOnSubscriptionFailure: false + ) + } + return + } + self.recoverClosedControlSession(trigger: trigger, expectedClient: expectedClient) + } + } + + func recoverClosedControlSession( + trigger: RecoveryTrigger, + expectedClient: MobileCoreRPCClient ) { guard remoteClient === expectedClient, connectionState == .connected else { return } guard foregroundRefreshIsActive else { @@ -233,7 +264,6 @@ extension MobileShellComposite { // cannot safely invent a redial route and must remain unavailable. switch trigger { case .liveness, .networkChange: - markMacConnectionReconnecting() resyncTerminalOutput(reason: trigger.description, restartEventStream: true) case .manual, .presencePush, .foreground, .eventStreamEnded, .subscriptionStartFailed, .transportWriteTimedOut, .automaticBackoffExpired, @@ -288,6 +318,18 @@ extension MobileShellComposite { self.applyConnectionRecoveryOwnerState() return } + let transportClosed = await expectedClient.isTransportClosed() + guard !Task.isCancelled, + self.connectionRecoveryOwner.isCurrent(attempt), + self.remoteClient === expectedClient, + self.connectionGeneration == attempt.sourceConnectionGeneration else { return } + if transportClosed == false { + // A slow application response after a path change or + // resume does not invalidate the native connection. + _ = self.completeConnectionRecovery(attempt) + self.applyConnectionRecoveryOwnerState() + return + } if self.lastBackgroundedAt != nil || self.foregroundResumeEpoch != epochAtProbeStart { // The probe spanned a background window: its wall-clock diff --git a/Packages/iOS/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite.swift b/Packages/iOS/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite.swift index c76b1be56d80..aa2052797e8c 100644 --- a/Packages/iOS/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite.swift +++ b/Packages/iOS/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite.swift @@ -1243,7 +1243,6 @@ public final class MobileShellComposite: MobileTerminalOutputSinking { private var renderGridLivenessProbeTask: Task? private var renderGridLivenessProbeID: UUID? private var renderGridLivenessConsecutiveProbeFailures = 0 - private var renderGridLivenessLaneRepairAttempts = 0 var lastTerminalEventAt: Date? @ObservationIgnored var terminalInputAckResubscribeRetryTask: Task? @ObservationIgnored var terminalInputAckResubscribeRetryTaskID: UUID? @@ -11949,8 +11948,8 @@ public final class MobileShellComposite: MobileTerminalOutputSinking { } /// Applies an availability failure only to the connection that produced it. - /// A blocked transport write is definitive and enters the single recovery - /// owner; an ordinary response timeout remains scoped to that one RPC. + /// Request deadlines remain scoped to their operation. A transport failure + /// enters recovery only after checking the native connection state. func handleMacAvailabilityFailureIfCurrent( after error: any Error, expectedClient: MobileCoreRPCClient, @@ -13915,7 +13914,7 @@ public final class MobileShellComposite: MobileTerminalOutputSinking { mobileShellLog.info("terminal event stream ended before subscribe ack, marking unavailable") MobileDebugLog.anchormux("sync.stream_ended before subscribe ack; failed start") diagnosticLog?.record(DiagnosticEvent(.error)) - recoverDeadConnection( + recoverClosedControlSession( trigger: .subscriptionStartFailed, expectedClient: client ) @@ -13988,7 +13987,6 @@ public final class MobileShellComposite: MobileTerminalOutputSinking { renderGridLivenessProbeTask = nil renderGridLivenessProbeID = nil renderGridLivenessConsecutiveProbeFailures = 0 - renderGridLivenessLaneRepairAttempts = 0 } /// Single ownership point for the liveness clock the watchdog reads. @@ -14003,7 +14001,6 @@ public final class MobileShellComposite: MobileTerminalOutputSinking { private func recordTerminalEventStreamLiveness() { lastTerminalEventAt = runtime?.now() ?? Date() renderGridLivenessConsecutiveProbeFailures = 0 - renderGridLivenessLaneRepairAttempts = 0 } #if DEBUG @@ -14130,19 +14127,13 @@ public final class MobileShellComposite: MobileTerminalOutputSinking { // stream stalled. Keep the shared Iroh session, whose other // lanes may still carry terminal input and keepalives, and // restart only the event listener. - self.renderGridLivenessLaneRepairAttempts += 1 - let laneRepairAttempts = self.renderGridLivenessLaneRepairAttempts - let escalate = laneRepairAttempts >= 2 - if escalate { - self.renderGridLivenessLaneRepairAttempts = 0 - } MobileDebugLog.anchormux( - "sync.liveness event_lane_repair transport_alive attempts=\(laneRepairAttempts) escalate=\(escalate) silentMs=\(silentMs)" + "sync.liveness event_lane_repair transport_alive silentMs=\(silentMs)" ) self.resyncTerminalOutput( - reason: escalate ? "liveness_event_lane_escalated" : "liveness_event_lane", + reason: "liveness_event_lane", restartEventStream: true, - recoversConnectionOnSubscriptionFailure: escalate + recoversConnectionOnSubscriptionFailure: false ) return } @@ -14154,9 +14145,7 @@ public final class MobileShellComposite: MobileTerminalOutputSinking { ) self.diagnosticLog?.record(DiagnosticEvent(.livenessResubscribe, ms: UInt32(clamping: silentMs))) mobileShellLog.info("render-grid stream silent for \(silentMs, privacy: .public)ms and subscription probe failed, re-subscribing") - // The bounded probe proved this exact client dead. Hand the session - // to the single recovery owner instead of rebuilding another listener - // on the same stale shell. + // Confirm native closure before handing the session to recovery. self.recoverDeadConnection(trigger: .liveness, expectedClient: client) } } diff --git a/Packages/iOS/CmuxMobileShell/Tests/CmuxMobileShellTests/IrohConnectionRecoveryOwnerTests.swift b/Packages/iOS/CmuxMobileShell/Tests/CmuxMobileShellTests/IrohConnectionRecoveryOwnerTests.swift index 2c746656e855..8a182c0ddaf0 100644 --- a/Packages/iOS/CmuxMobileShell/Tests/CmuxMobileShellTests/IrohConnectionRecoveryOwnerTests.swift +++ b/Packages/iOS/CmuxMobileShell/Tests/CmuxMobileShellTests/IrohConnectionRecoveryOwnerTests.swift @@ -12,7 +12,7 @@ extension ReconnectRouteSelectionTests { defer { fixture.release() } #expect(await fixture.store.reconnectActiveMacIfAvailable(stackUserID: "user-1")) - #expect(await fixture.router.waitForCount(of: "mobile.events.subscribe", atLeast: 1)) + #expect(try await pollUntil { fixture.store.lastSuccessfulTerminalSubscription != nil }) let initialReconnectGeneration = fixture.store.storedMacReconnectGeneration let firstClient = try #require(fixture.store.remoteClient) let first = try #require(fixture.box.get()) @@ -40,7 +40,7 @@ extension ReconnectRouteSelectionTests { await fixture.router.setWorkspaceIDs(["cmux-master", "hevy-cli"]) #expect(await fixture.store.reconnectActiveMacIfAvailable(stackUserID: "user-1")) - #expect(await fixture.router.waitForCount(of: "mobile.events.subscribe", atLeast: 1)) + #expect(try await pollUntil { fixture.store.lastSuccessfulTerminalSubscription != nil }) let firstClient = try #require(fixture.store.remoteClient) let hevyWorkspace = try #require( fixture.store.workspaces.first { $0.rpcWorkspaceID.rawValue == "hevy-cli" } @@ -69,7 +69,7 @@ extension ReconnectRouteSelectionTests { } #expect(await fixture.store.reconnectActiveMacIfAvailable(stackUserID: "user-1")) - #expect(await fixture.router.waitForCount(of: "mobile.events.subscribe", atLeast: 1)) + #expect(try await pollUntil { fixture.store.lastSuccessfulTerminalSubscription != nil }) let firstClient = try #require(fixture.store.remoteClient) fixture.store.recoverDeadConnection( @@ -95,7 +95,7 @@ extension ReconnectRouteSelectionTests { defer { fixture.release() } #expect(await fixture.store.reconnectActiveMacIfAvailable(stackUserID: "user-1")) - #expect(await fixture.router.waitForCount(of: "mobile.events.subscribe", atLeast: 1)) + #expect(try await pollUntil { fixture.store.lastSuccessfulTerminalSubscription != nil }) let firstClient = try #require(fixture.store.remoteClient) fixture.store.suspendForegroundRefresh() fixture.clock.advance(by: 61) @@ -237,7 +237,7 @@ extension ReconnectRouteSelectionTests { } #expect(await fixture.store.reconnectActiveMacIfAvailable(stackUserID: "user-1")) - #expect(await fixture.router.waitForCount(of: "mobile.events.subscribe", atLeast: 1)) + #expect(try await pollUntil { fixture.store.lastSuccessfulTerminalSubscription != nil }) let firstClient = try #require(fixture.store.remoteClient) let backing = try #require(fixture.store.pairedMacStore as? BackingUpPairedMacStore) await backup.blockFutureFetches() @@ -263,7 +263,7 @@ extension ReconnectRouteSelectionTests { defer { fixture.release() } #expect(await fixture.store.reconnectActiveMacIfAvailable(stackUserID: "user-1")) - #expect(await fixture.router.waitForCount(of: "mobile.events.subscribe", atLeast: 1)) + #expect(try await pollUntil { fixture.store.lastSuccessfulTerminalSubscription != nil }) let first = try #require(fixture.box.get()) await first.close() @@ -282,7 +282,7 @@ extension ReconnectRouteSelectionTests { defer { fixture.release() } #expect(await fixture.store.reconnectActiveMacIfAvailable(stackUserID: "user-1")) - #expect(await fixture.router.waitForCount(of: "mobile.events.subscribe", atLeast: 1)) + #expect(try await pollUntil { fixture.store.lastSuccessfulTerminalSubscription != nil }) let firstClient = try #require(fixture.store.remoteClient) // Keep every post-drop dial failing until the first owner reaches its // terminal failed phase. This avoids coupling the assertion to a @@ -329,7 +329,7 @@ extension ReconnectRouteSelectionTests { defer { fixture.release() } #expect(await fixture.store.reconnectActiveMacIfAvailable(stackUserID: "user-1")) - #expect(await fixture.router.waitForCount(of: "mobile.events.subscribe", atLeast: 1)) + #expect(try await pollUntil { fixture.store.lastSuccessfulTerminalSubscription != nil }) let client = try #require(fixture.store.remoteClient) let generation = fixture.store.connectionGeneration @@ -350,7 +350,7 @@ extension ReconnectRouteSelectionTests { defer { fixture.release() } #expect(await fixture.store.reconnectActiveMacIfAvailable(stackUserID: "user-1")) - #expect(await fixture.router.waitForCount(of: "mobile.events.subscribe", atLeast: 1)) + #expect(try await pollUntil { fixture.store.lastSuccessfulTerminalSubscription != nil }) let liveClient = try #require(fixture.store.remoteClient) let liveGeneration = fixture.store.connectionGeneration @@ -368,7 +368,7 @@ extension ReconnectRouteSelectionTests { defer { fixture.release() } #expect(await fixture.store.reconnectActiveMacIfAvailable(stackUserID: "user-1")) - #expect(await fixture.router.waitForCount(of: "mobile.events.subscribe", atLeast: 1)) + #expect(try await pollUntil { fixture.store.lastSuccessfulTerminalSubscription != nil }) let foregroundKey = fixture.store.foregroundMacKey var staleState = try #require(fixture.store.workspacesByMac[foregroundKey]) staleState.status = .unavailable @@ -398,11 +398,32 @@ extension ReconnectRouteSelectionTests { ) } + @Test func failedNetworkRefreshPreservesOpenIrohConnection() async throws { + let fixture = try await makeRecoveryOwnerFixture() + defer { fixture.release() } + #expect(await fixture.store.reconnectActiveMacIfAvailable(stackUserID: "user-1")) + #expect(try await pollUntil { fixture.store.lastSuccessfulTerminalSubscription != nil }) + let client = try #require(fixture.store.remoteClient) + let generation = fixture.store.connectionGeneration + fixture.store.stateSyncAuthorityClientID = nil + let nextList = await fixture.router.count(of: "mobile.workspace.list") + 1 + await fixture.router.failWorkspaceListRequest(number: nextList) + fixture.store.recoverMobileConnection(trigger: .networkChange) + #expect(try await pollUntil { !fixture.store.connectionRecoveryOwner.isActive }) + #expect(fixture.store.remoteClient === client) + #expect(fixture.store.connectionGeneration == generation) + #expect(fixture.store.macConnectionStatus == .connected) + #expect(!fixture.store.isRecoveringConnection) + #expect(await fixture.store.reloadWorkspaceListFromMac()) + #expect(fixture.factory.attemptedKinds() == [.iroh]) + } + @Test func operationFailurePreservesOpenIrohConnection() async throws { let fixture = try await makeRecoveryOwnerFixture() defer { fixture.release() } #expect(await fixture.store.reconnectActiveMacIfAvailable(stackUserID: "user-1")) - #expect(await fixture.router.waitForCount(of: "mobile.events.subscribe", atLeast: 1)) + #expect(try await pollUntil { fixture.store.lastSuccessfulTerminalSubscription != nil }) + #expect(try await pollUntil { !fixture.store.isRecoveringConnection }) let client = try #require(fixture.store.remoteClient) let generation = fixture.store.connectionGeneration fixture.store.handleMacAvailabilityFailureIfCurrent( @@ -425,7 +446,7 @@ extension ReconnectRouteSelectionTests { defer { fixture.release() } #expect(await fixture.store.reconnectActiveMacIfAvailable(stackUserID: "user-1")) - #expect(await fixture.router.waitForCount(of: "mobile.events.subscribe", atLeast: 1)) + #expect(try await pollUntil { fixture.store.lastSuccessfulTerminalSubscription != nil }) await fixture.diagnosticLog.clear() fixture.factory.setConnectFailure(.timedOut) let first = try #require(fixture.box.get()) @@ -486,7 +507,7 @@ extension ReconnectRouteSelectionTests { defer { fixture.release() } #expect(await fixture.store.reconnectActiveMacIfAvailable(stackUserID: "user-1")) - #expect(await fixture.router.waitForCount(of: "mobile.events.subscribe", atLeast: 1)) + #expect(try await pollUntil { fixture.store.lastSuccessfulTerminalSubscription != nil }) await fixture.diagnosticLog.clear() let client = try #require(fixture.store.remoteClient) let generation = fixture.store.connectionGeneration diff --git a/Packages/iOS/CmuxMobileShell/Tests/CmuxMobileShellTests/MobileShellRenderGridLivenessTests.swift b/Packages/iOS/CmuxMobileShell/Tests/CmuxMobileShellTests/MobileShellRenderGridLivenessTests.swift index c1eaf28b97cc..a8093f3fec87 100644 --- a/Packages/iOS/CmuxMobileShell/Tests/CmuxMobileShellTests/MobileShellRenderGridLivenessTests.swift +++ b/Packages/iOS/CmuxMobileShell/Tests/CmuxMobileShellTests/MobileShellRenderGridLivenessTests.swift @@ -576,12 +576,12 @@ import Testing clock.advance(by: 10) store.debugRunRenderGridLivenessCheckForTesting() #expect(await router.waitForCount(of: "mobile.events.probe", atLeast: 1)) - // The second independent probe succeeds and resets the liveness window. - // Reconnecting makes `.connected` prove that its response was applied. + // The idempotent subscription retry succeeds on the same connection. + // The status change proves that its acknowledgement was applied. store.markMacConnectionReconnecting() let followUpSucceeded = try await pollUntil { store.debugRunRenderGridLivenessCheckForTesting() - return await router.count(of: "mobile.events.probe") >= 2 + return await router.count(of: "mobile.events.subscribe") >= 2 && store.macConnectionStatus == .connected } #expect(followUpSucceeded, "the follow-up probe must complete successfully") @@ -623,7 +623,7 @@ import Testing store.debugRunRenderGridLivenessCheckForTesting() let probeCount = await router.count(of: "mobile.events.probe") let subscribeCount = await router.count(of: "mobile.events.subscribe") - return probeCount >= 2 && subscribeCount >= 2 + return probeCount >= 1 && subscribeCount >= 2 } #expect(repaired, "a live transport must restart only the stalled event listener") #expect(store.remoteClient === originalClient) diff --git a/Sources/Mobile/MobileHostConnectionEventQueue.swift b/Sources/Mobile/MobileHostConnectionEventQueue.swift index c6422fb3655e..6577ed140ba4 100644 --- a/Sources/Mobile/MobileHostConnectionEventQueue.swift +++ b/Sources/Mobile/MobileHostConnectionEventQueue.swift @@ -21,10 +21,8 @@ import Foundation /// - `terminal.updated` / `workspace.updated`: level-triggered pings; the newer /// occurrence that forced the shed supersedes the shed one. /// -/// Every other topic keeps the close-on-overflow contract: those payloads -/// cannot be re-derived by the client, so tearing the connection down (the -/// client reconnects and re-syncs from authoritative state) is the only -/// lossless bound. +/// Other topics retain their ordered payloads even beyond the shedding budget. +/// Congestion is not evidence that the connection has closed. enum MobileHostEventTopicPolicy { static let renderGridTopic = "terminal.render_grid" static let simulatorFrameTopic = "simulator.frame" @@ -33,7 +31,7 @@ enum MobileHostEventTopicPolicy { switch topic { case renderGridTopic: // A render-grid event without a surface key cannot be resynced - // per-surface, so it keeps the lossless close-on-overflow path. + // per-surface, so its ordered payload is retained. return coalesceKey != nil case simulatorFrameTopic: // Simulator frames are whole-screen snapshots; a later frame fully @@ -49,13 +47,10 @@ enum MobileHostEventTopicPolicy { /// Outcome of one synchronous admission attempt on a connection's event queue. struct MobileHostEventEnqueueResult: Sendable { - /// The event was appended to the bounded queue. + /// The event was appended to the queue. let admitted: Bool /// The caller must start the (single) drain task for this connection. let startDrain: Bool - /// A non-droppable event overflowed the bounded queue; the caller must - /// close the connection — the lossless-topic contract. - let shouldClose: Bool /// Surfaces whose queued render-grid frames were shed; the caller must ask /// the producer for a full-frame resync of each. let renderGridResyncSurfaceIDs: Set @@ -71,7 +66,6 @@ struct MobileHostEventEnqueueResult: Sendable { static let rejected = MobileHostEventEnqueueResult( admitted: false, startDrain: false, - shouldClose: false, renderGridResyncSurfaceIDs: [], depthAfterEnqueue: nil, shedEventCount: 0, @@ -95,18 +89,10 @@ private struct MobileHostEventShedSummary: Sendable { } } -/// Bounded, synchronously-admitted mailbox between the event fan-out -/// (``MobileHostService/emitEvent(topic:payload:)``) and one connection's -/// drain loop. -/// -/// This is the owner boundary for issue #8842: admission runs *before* any -/// per-connection work is scheduled, on the emitter's thread, against an -/// explicit bound — so the memory pinned per connection is O(capacity) no -/// matter how far emission runs ahead of a slow, paused, or half-dead -/// subscriber. The previous path spawned one unstructured Task per connection -/// per event, each retaining the full payload dictionary until the connection -/// actor reached its own bounded check, which left everything upstream of the -/// bound unbounded. +/// Synchronously admitted mailbox between event fan-out and a single drain. +/// Refresh events have a shedding budget; ordered events are retained until +/// delivery. Admission happens before task creation, so producers never create +/// a separate task retaining each event while the network is slow. final class MobileHostConnectionEventQueue: @unchecked Sendable { struct QueuedEvent: Sendable { let topic: String @@ -174,7 +160,7 @@ final class MobileHostConnectionEventQueue: @unchecked Sendable { return subscribedTopics.contains(topic) } - /// Synchronous bounded admission. Safe to call from any thread; never + /// Synchronous admission with refresh-event shedding. Safe on any thread; never /// blocks on the network, the connection actor, or the runtime. func enqueue( topic: String, @@ -214,7 +200,6 @@ final class MobileHostConnectionEventQueue: @unchecked Sendable { return MobileHostEventEnqueueResult( admitted: false, startDrain: false, - shouldClose: false, renderGridResyncSurfaceIDs: resyncSurfaceIDs, depthAfterEnqueue: nil, shedEventCount: shedSummary.eventCount, @@ -222,20 +207,8 @@ final class MobileHostConnectionEventQueue: @unchecked Sendable { simulatorFrameShedPanelIDs: shedSummary.simulatorFramePanelIDs ) } - guard hasRoomLocked(for: frame) else { - guard MobileHostEventTopicPolicy.isDroppable(topic: topic, coalesceKey: coalesceKey) else { - lock.unlock() - return MobileHostEventEnqueueResult( - admitted: false, - startDrain: false, - shouldClose: true, - renderGridResyncSurfaceIDs: resyncSurfaceIDs, - depthAfterEnqueue: nil, - shedEventCount: shedSummary.eventCount, - shedByteCount: shedSummary.byteCount, - simulatorFrameShedPanelIDs: shedSummary.simulatorFramePanelIDs - ) - } + if !hasRoomLocked(for: frame), + MobileHostEventTopicPolicy.isDroppable(topic: topic, coalesceKey: coalesceKey) { if isRenderGrid, let coalesceKey { if poisonedRenderGridSurfaceIDs.insert(coalesceKey).inserted { resyncSurfaceIDs.insert(coalesceKey) @@ -249,7 +222,6 @@ final class MobileHostConnectionEventQueue: @unchecked Sendable { return MobileHostEventEnqueueResult( admitted: false, startDrain: false, - shouldClose: false, renderGridResyncSurfaceIDs: resyncSurfaceIDs, depthAfterEnqueue: nil, shedEventCount: shedSummary.eventCount, @@ -279,7 +251,6 @@ final class MobileHostConnectionEventQueue: @unchecked Sendable { return MobileHostEventEnqueueResult( admitted: true, startDrain: startDrain, - shouldClose: false, renderGridResyncSurfaceIDs: resyncSurfaceIDs, depthAfterEnqueue: depthAfterEnqueue, shedEventCount: shedSummary.eventCount, diff --git a/Sources/Mobile/MobileHostService.swift b/Sources/Mobile/MobileHostService.swift index 0108ac2f9969..43ba7f71b7df 100644 --- a/Sources/Mobile/MobileHostService.swift +++ b/Sources/Mobile/MobileHostService.swift @@ -703,17 +703,7 @@ final class MobileHostService { if result.startDrain { Task { await connection.drainQueuedEvents() } } - if result.shouldClose { - Task { - await connection.close( - reason: "event queue exceeded bounded capacity", - exit: CmxIrohAdmittedConnectionExit( - lifecycle: .controlWriteFailed, - failure: .sendQueueOverflow - ) - ) - } - } + } if !resyncSurfaceIDs.isEmpty { MobileTerminalRenderObserver.requestRenderGridFullResync( @@ -2051,12 +2041,6 @@ actor MobileHostConnection { private static let maximumReceiveBufferByteCount = MobileSyncFrameCodec.defaultMaximumFrameByteCount + MobileSyncFrameCodec.headerByteCount private static let defaultFirstFrameTimeoutNanoseconds: UInt64 = 15 * 1_000_000_000 fileprivate static let defaultIdleTimeoutNanoseconds: UInt64 = 30 * 1_000_000_000 - /// Bounded deadline for one control-lane event write. A peer that accepted - /// the connection but stopped reading (TCP zero-window, QUIC flow-control - /// stall) would otherwise pin the drain — and with it this connection's - /// queue, transport, and tasks — indefinitely (issue #8842). - private static let defaultEventSendStallTimeoutNanoseconds: UInt64 = 30 * 1_000_000_000 - private struct EventSubscription: Sendable { let topics: Set let transport: MobileHostEventTransport @@ -2104,15 +2088,10 @@ actor MobileHostConnection { private let onClose: @Sendable (UUID) async -> Void private let requestSimulatorFrameReplay: @Sendable (UUID, Set) async -> Void private let responseWorkQuota = MobileHostRPCWorkQuota() - /// Bounded pre-write mailbox with synchronous admission from the event + /// Pre-write mailbox with synchronous admission from the event /// fan-out. Nonisolated so ``MobileHostService/emitEvent(topic:payload:)`` /// admits events without scheduling any per-event actor work. nonisolated let eventQueue: MobileHostConnectionEventQueue - private let eventSendStallTimeoutNanoseconds: UInt64 - /// Invalidates the pending event-send stall deadline: bumped when a send - /// starts and again when it settles, so a deadline armed for send N can - /// never close the connection after N completed. - private var eventSendGeneration: UInt64 = 0 private var receiveBuffer = Data() private var firstFrameTimeoutTask: Task? private var idleTimeoutTask: Task? @@ -2145,7 +2124,6 @@ actor MobileHostConnection { eventQueue: MobileHostConnectionEventQueue = MobileHostConnectionEventQueue(), firstFrameTimeoutNanoseconds: UInt64 = MobileHostConnection.defaultFirstFrameTimeoutNanoseconds, idleTimeoutNanoseconds: UInt64 = MobileHostConnection.defaultIdleTimeoutNanoseconds, - eventSendStallTimeoutNanoseconds: UInt64 = MobileHostConnection.defaultEventSendStallTimeoutNanoseconds, independentEventWriter: (any MobileHostIndependentEventWriting)? = nil, authorizeRequest: @escaping @Sendable (MobileHostRPCRequest) async -> MobileHostRPCResult?, onAuthorizedRequest: @escaping @Sendable (MobileHostRPCRequest) async -> Void, @@ -2161,7 +2139,6 @@ actor MobileHostConnection { self.independentEventWriter = independentEventWriter self.firstFrameTimeoutNanoseconds = firstFrameTimeoutNanoseconds self.idleTimeoutNanoseconds = idleTimeoutNanoseconds - self.eventSendStallTimeoutNanoseconds = eventSendStallTimeoutNanoseconds self.authorizeRequest = authorizeRequest self.onAuthorizedRequest = onAuthorizedRequest self.onUsableSession = onUsableSession @@ -2177,7 +2154,6 @@ actor MobileHostConnection { eventQueue: MobileHostConnectionEventQueue = MobileHostConnectionEventQueue(), firstFrameTimeoutNanoseconds: UInt64 = MobileHostConnection.defaultFirstFrameTimeoutNanoseconds, idleTimeoutNanoseconds: UInt64 = MobileHostConnection.defaultIdleTimeoutNanoseconds, - eventSendStallTimeoutNanoseconds: UInt64 = MobileHostConnection.defaultEventSendStallTimeoutNanoseconds, independentEventWriter: (any MobileHostIndependentEventWriting)? = nil, authorizeRequest: @escaping @Sendable (MobileHostRPCRequest) async -> MobileHostRPCResult?, onAuthorizedRequest: @escaping @Sendable (MobileHostRPCRequest) async -> Void, @@ -2192,7 +2168,6 @@ actor MobileHostConnection { self.independentEventWriter = independentEventWriter self.firstFrameTimeoutNanoseconds = firstFrameTimeoutNanoseconds self.idleTimeoutNanoseconds = idleTimeoutNanoseconds - self.eventSendStallTimeoutNanoseconds = eventSendStallTimeoutNanoseconds self.authorizeRequest = authorizeRequest self.onAuthorizedRequest = onAuthorizedRequest self.onUsableSession = onUsableSession @@ -2325,30 +2300,32 @@ actor MobileHostConnection { } receiveBuffer.append(data) do { - let frames = try MobileSyncFrameCodec.decodeFrames( - from: &receiveBuffer, - maximumDecodedFrameCount: responseWorkQuota - .maximumConcurrentRequestCount - ) - if !frames.isEmpty { - didDecodeFirstFrame = true - firstFrameTimeoutTask?.cancel() - firstFrameTimeoutTask = nil - } - for frame in frames { - guard !isClosed else { - return + let batchLimit = responseWorkQuota.maximumConcurrentRequestCount + while !isClosed, !Task.isCancelled { + let frames = try MobileSyncFrameCodec.decodeFrames( + from: &receiveBuffer, + maximumDecodedFrameCount: batchLimit + ) + if !frames.isEmpty { + didDecodeFirstFrame = true + firstFrameTimeoutTask?.cancel() + firstFrameTimeoutTask = nil } - guard startResponseTask(for: frame) else { - await close( - reason: "rpc work capacity exceeded", - exit: CmxIrohAdmittedConnectionExit( - lifecycle: .controlReadFailed, - failure: .protocolViolation - ) - ) - return + for frame in frames { + guard !isClosed else { return } + if !startResponseTask(for: frame) { + // Work pressure fails this request explicitly; it + // does not invalidate the authenticated connection. + let request = try? MobileHostRPCEnvelope.decodeRequest(frame).get() + guard await sendResponse(MobileHostRPCEnvelope.error( + id: request?.id, + code: "server_busy", + message: "Too many requests are pending" + )) else { return } + } } + guard frames.count == batchLimit else { break } + await Task.yield() } guard !isClosed else { return @@ -2943,26 +2920,12 @@ actor MobileHostConnection { if result.startDrain { Task { await self.drainQueuedEvents() } } - if result.shouldClose { - // The bounded queue fills when the control stream stops draining - // (e.g. the peer's network path died mid-write) while terminal - // events keep arriving. The peer violated nothing; field host - // rings (2026-07-23 WiFi path flap) showed this close mislabeled - // protocolViolation seconds after admission. - await close( - reason: "event queue exceeded bounded capacity", - exit: CmxIrohAdmittedConnectionExit( - lifecycle: .controlWriteFailed, - failure: .sendQueueOverflow - ) - ) - } return result.admitted } /// Synchronous bounded admission from the fan-out path. Never blocks and /// never schedules per-event work; the caller acts on the returned - /// outcome (drain start, overflow close, render-grid resync). + /// outcome (drain start, refresh shedding, render-grid resync). nonisolated func enqueueEventFrame( _ frame: Data, topic: String, @@ -3020,7 +2983,7 @@ actor MobileHostConnection { /// (enforced by the queue's drain claim), pulling from the bounded queue /// and writing to the negotiated lane. Exits when the queue is empty, the /// connection closes, lane negotiation pauses delivery, or a delivery - /// fails or stalls (which closes the connection). + /// fails (which closes the unusable control session). func drainQueuedEvents() async { while true { if isClosed || independentEventNegotiationInProgress { @@ -3099,38 +3062,11 @@ actor MobileHostConnection { return await sendEventControlFrame(event.frame) } - /// Writes one event frame on the control lane under the bounded stall - /// deadline. On a stall the connection is closed — `transport.close()` - /// resolves the pending write — converting a half-dead subscriber into - /// deterministic teardown instead of a forever-pinned drain. + /// Writes one serialized frame until the transport completes or fails. + /// An application deadline cannot cancel writeAll safely: it may already + /// have sent a prefix. Native transport failure still ends the drain. private func sendEventControlFrame(_ frame: Data) async -> Bool { - guard !isClosed else { return false } - let timeoutNanoseconds = eventSendStallTimeoutNanoseconds - guard timeoutNanoseconds > 0 else { - return await sendControlFrame(frame) - } - eventSendGeneration &+= 1 - let generation = eventSendGeneration - let deadlineTask = Task { [weak self] in - try? await Task.sleep(nanoseconds: timeoutNanoseconds) - guard !Task.isCancelled else { return } - await self?.closeIfEventSendStillInFlight(generation: generation) - } - let delivered = await sendControlFrame(frame) - eventSendGeneration &+= 1 - deadlineTask.cancel() - return delivered && !isClosed - } - - private func closeIfEventSendStillInFlight(generation: UInt64) async { - guard eventSendGeneration == generation, !isClosed else { return } - await close( - reason: "event send stalled past the bounded deadline", - exit: CmxIrohAdmittedConnectionExit( - lifecycle: .controlWriteFailed, - failure: .timedOut - ) - ) + await sendControlFrame(frame) } private func downgradeIndependentSubscriptionsToControl() { diff --git a/cmuxTests/MobileHostConnectionEventLaneTests.swift b/cmuxTests/MobileHostConnectionEventLaneTests.swift index be2ce2027456..ff0b31b3ec67 100644 --- a/cmuxTests/MobileHostConnectionEventLaneTests.swift +++ b/cmuxTests/MobileHostConnectionEventLaneTests.swift @@ -456,10 +456,7 @@ extension MobileHostAuthorizationTests { )) } - /// A subscriber whose transport accepted a frame but never completes the - /// write (TCP zero-window peer) is torn down by the bounded event-send - /// stall deadline instead of pinning the connection's queue, tasks, and - /// socket forever. + /// A stalled event write does not end the connection or interrupt framing. @Test func testEventSendStallDoesNotCloseConnection() async throws { let transport = StalledSendMobileHostByteTransport() let recorder = MobileHostConnectionCloseRecorder() @@ -467,7 +464,6 @@ extension MobileHostAuthorizationTests { let session = MobileHostConnection( id: connectionID, transport: transport, - eventSendStallTimeoutNanoseconds: 5_000_000, authorizeRequest: { _ in nil }, onAuthorizedRequest: { _ in }, handleRequest: { _ in .ok([:]) }, @@ -481,8 +477,8 @@ extension MobileHostAuthorizationTests { payload: ["surface_id": "surface-stall-8842", "full": true] ) await transport.waitUntilSendStalled() - // Wait beyond the configured deadline to prove it cannot terminate - // an established session merely because this write has not completed. + // Exercise an unresolved send across suspension before verifying + // that the connection has not been closed on the application's behalf. try await Task.sleep(for: .milliseconds(30)) #expect(await recorder.recordedIDs().isEmpty) #expect(await transport.observedCloseCount() == 0) @@ -514,7 +510,6 @@ extension MobileHostAuthorizationTests { isFullRenderGridFrame: false, frame: frame ) #expect(!overflow.admitted) - #expect(!overflow.shouldClose) #expect(overflow.renderGridResyncSurfaceIDs == ["s1"]) #expect(queue.count == 0) // While poisoned, deltas stay refused even though there is room. @@ -552,7 +547,6 @@ extension MobileHostAuthorizationTests { isFullRenderGridFrame: false, frame: frame ) #expect(overflow.admitted) - #expect(!overflow.shouldClose) #expect(queue.count == 2) #expect(queue.dequeue()?.frame == frame) #expect(queue.dequeue()?.frame == frame) From bcbece487baf83ebe2904335de11c263817dea9e Mon Sep 17 00:00:00 2001 From: Abdulaziz Albahar <67667005+azooz2003-bit@users.noreply.github.com> Date: Fri, 11 Sep 2026 23:04:38 -0700 Subject: [PATCH 03/17] fix: clear cloud panel failures through the workspace API --- Sources/Workspace+PanelLifecycle.swift | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/Sources/Workspace+PanelLifecycle.swift b/Sources/Workspace+PanelLifecycle.swift index f04dc1569856..011ac2bc09c5 100644 --- a/Sources/Workspace+PanelLifecycle.swift +++ b/Sources/Workspace+PanelLifecycle.swift @@ -445,7 +445,7 @@ extension Workspace { preservesTerminalForTransfer: Bool = false, preservesRemoteTerminalTracking: Bool = false ) -> WorkspaceRemoteConfiguration? { - cloudMaterializationFailures.removeValue(forKey: panelId) + clearCloudMaterializationFailure(surfaceID: panelId) appLinkHandoffCoordinator.cancel(sourcePanelID: panelId) if publishSurfaceClosedEvent { publishCmuxSurfaceClosed(panelId, paneId: paneId, panel: panel, origin: origin) From 5f0830f36bf55a97e0ef57eff1d8320a2c3a5f7d Mon Sep 17 00:00:00 2001 From: Abdulaziz Albahar <67667005+azooz2003-bit@users.noreply.github.com> Date: Fri, 11 Sep 2026 23:15:09 -0700 Subject: [PATCH 04/17] test: cover coalesced tail of maximum-size mobile frame --- .../ControllableResponseTransport.swift | 8 ++++++ ...obileCoreRPCNativeWriteLifetimeTests.swift | 25 +++++++++++++++++++ .../MobileHostConnectionLifecycleTests.swift | 23 +++++++++++++++++ 3 files changed, 56 insertions(+) diff --git a/Packages/iOS/CmuxMobileRPC/Tests/CmuxMobileRPCTests/ControllableResponseTransport.swift b/Packages/iOS/CmuxMobileRPC/Tests/CmuxMobileRPCTests/ControllableResponseTransport.swift index 58220dcd51c0..7901b4dad44e 100644 --- a/Packages/iOS/CmuxMobileRPC/Tests/CmuxMobileRPCTests/ControllableResponseTransport.swift +++ b/Packages/iOS/CmuxMobileRPC/Tests/CmuxMobileRPCTests/ControllableResponseTransport.swift @@ -116,4 +116,12 @@ actor ControllableResponseTransport: CmxByteTransport { receiveWaiters = [] for waiter in waiters { waiter.resume(returning: nil) } } + + func deliverRawChunk(_ chunk: Data) { + if !receiveWaiters.isEmpty { + receiveWaiters.removeFirst().resume(returning: chunk) + } else { + queuedFrames.append(chunk) + } + } } diff --git a/Packages/iOS/CmuxMobileRPC/Tests/CmuxMobileRPCTests/MobileCoreRPCNativeWriteLifetimeTests.swift b/Packages/iOS/CmuxMobileRPC/Tests/CmuxMobileRPCTests/MobileCoreRPCNativeWriteLifetimeTests.swift index 52af70770bcc..a82f5d94f649 100644 --- a/Packages/iOS/CmuxMobileRPC/Tests/CmuxMobileRPCTests/MobileCoreRPCNativeWriteLifetimeTests.swift +++ b/Packages/iOS/CmuxMobileRPC/Tests/CmuxMobileRPCTests/MobileCoreRPCNativeWriteLifetimeTests.swift @@ -80,6 +80,31 @@ struct MobileCoreRPCNativeWriteLifetimeTests { await session.tearDown(error: .connectionClosed) } + @Test func coalescedFrameTailRespectsPerFrameLimit() async throws { + let base = ControllableResponseTransport(closeEndsReceive: true) + let session = MobileCoreRPCSession(makeTransport: { base }) + let request = Task { + try await session.send( + payload: try Self.request("after-large-frame"), + requestID: "after-large-frame", + deadlineUptimeNanoseconds: DispatchTime.now().uptimeNanoseconds + 5_000_000_000 + ) + } + await base.waitUntilSent(count: 1) + var payload = Data(#"{"kind":"event","topic":"workspace.updated","payload":{}}"#.utf8) + payload.append(Data(repeating: 0x20, count: MobileSyncFrameCodec.defaultMaximumFrameByteCount - payload.count)) + let frame = try MobileSyncFrameCodec.encodeFrame(payload) + await base.deliverRawChunk(Data(frame.dropLast())) + var tail = Data(frame.suffix(1)) + tail.append(try MobileSyncFrameCodec.encodeFrame( + Data(#"{"id":"after-large-frame","ok":true,"result":{}}"#.utf8) + )) + await base.deliverRawChunk(tail) + _ = try await request.value + #expect(await !base.closed()) + await session.tearDown(error: .connectionClosed) + } + private static func request(_ id: String) throws -> Data { try MobileCoreRPCClient.requestData(method: "terminal.input", params: ["text": id], id: id) } diff --git a/cmuxTests/MobileHostConnectionLifecycleTests.swift b/cmuxTests/MobileHostConnectionLifecycleTests.swift index 15817f4bb509..4b5434f4bd87 100644 --- a/cmuxTests/MobileHostConnectionLifecycleTests.swift +++ b/cmuxTests/MobileHostConnectionLifecycleTests.swift @@ -644,6 +644,29 @@ extension MobileHostAuthorizationTests { #expect(await transport.observedCloseCount() == 0) await session.close(reason: "test complete") } + @Test func testMobileHostAcceptsCoalescedTailOfMaximumSizeFrame() async throws { + let transport = RecordingMobileHostByteTransport() + let session = MobileHostConnection( + id: UUID(), transport: transport, + authorizeRequest: { _ in nil }, + onAuthorizedRequest: { _ in }, + handleRequest: { _ in .ok([:]) }, + onClose: { _ in } + ) + var payload = Data(#"{"id":"large","method":"workspace.list","params":{}}"#.utf8) + payload.append(Data(repeating: 0x20, count: MobileSyncFrameCodec.defaultMaximumFrameByteCount - payload.count)) + let frame = try MobileSyncFrameCodec.encodeFrame(payload) + await session.debugHandleReceiveDataForTesting(Data(frame.dropLast())) + var tail = Data(frame.suffix(1)) + tail.append(try MobileSyncFrameCodec.encodeFrame( + Data(#"{"id":"following","method":"workspace.list","params":{}}"#.utf8) + )) + await session.debugHandleReceiveDataForTesting(tail) + try #require(await transport.observedCloseCount() == 0) + #expect(await transport.waitForSentBufferCount(2).count == 2) + await session.close(reason: "test complete") + } + // MARK: - Advertised mobile host capabilities @Test func testMobileHostAdvertisesWorkspaceActionCapabilities() { let capabilities = MobileHostService.mobileHostCapabilities From 73b78067a7c8c5c06999306fb4bf305776c3c42b Mon Sep 17 00:00:00 2001 From: Abdulaziz Albahar <67667005+azooz2003-bit@users.noreply.github.com> Date: Fri, 11 Sep 2026 23:16:53 -0700 Subject: [PATCH 05/17] fix: apply mobile size limits per frame rather than per read --- ...bileCoreRPCSession+IndependentEvents.swift | 5 ++--- .../CmuxMobileRPC/MobileCoreRPCSession.swift | 12 ++--------- Sources/Mobile/MobileHostService.swift | 20 ++----------------- 3 files changed, 6 insertions(+), 31 deletions(-) diff --git a/Packages/iOS/CmuxMobileRPC/Sources/CmuxMobileRPC/MobileCoreRPCSession+IndependentEvents.swift b/Packages/iOS/CmuxMobileRPC/Sources/CmuxMobileRPC/MobileCoreRPCSession+IndependentEvents.swift index affcef688487..cc92e2fa1c0e 100644 --- a/Packages/iOS/CmuxMobileRPC/Sources/CmuxMobileRPC/MobileCoreRPCSession+IndependentEvents.swift +++ b/Packages/iOS/CmuxMobileRPC/Sources/CmuxMobileRPC/MobileCoreRPCSession+IndependentEvents.swift @@ -92,9 +92,8 @@ extension MobileCoreRPCSession { for try await chunk in stream { try Task.checkCancellation() guard !chunk.isEmpty else { continue } - guard chunk.count <= Self.maximumReceiveBufferByteCount - buffer.count else { - throw MobileSyncFrameCodecError.frameTooLarge(buffer.count + chunk.count) - } + // The codec checks each frame's size, independent of how the + // transport groups adjacent frames into chunks. buffer.append(chunk) while !Task.isCancelled, independentEventReader?.id == id { let frames = try MobileSyncFrameCodec.decodeFrames( diff --git a/Packages/iOS/CmuxMobileRPC/Sources/CmuxMobileRPC/MobileCoreRPCSession.swift b/Packages/iOS/CmuxMobileRPC/Sources/CmuxMobileRPC/MobileCoreRPCSession.swift index 165b06dad97a..20a69ea63caa 100644 --- a/Packages/iOS/CmuxMobileRPC/Sources/CmuxMobileRPC/MobileCoreRPCSession.swift +++ b/Packages/iOS/CmuxMobileRPC/Sources/CmuxMobileRPC/MobileCoreRPCSession.swift @@ -30,9 +30,6 @@ actor MobileCoreRPCSession { static let defaultAbandonedConnectCleanupTimeoutNanoseconds: UInt64 = 1_000_000_000 static let defaultLateAbandonedConnectCloseTimeoutNanoseconds: UInt64 = 5_000_000_000 static let defaultCancelledWriteCompletionGraceNanoseconds: UInt64 = 250_000_000 - static let maximumReceiveBufferByteCount = - MobileSyncFrameCodec.defaultMaximumFrameByteCount - + MobileSyncFrameCodec.headerByteCount static let maximumDecodedFrameCountPerRead = 256 struct EventSubscription { @@ -1039,13 +1036,8 @@ actor MobileCoreRPCSession { installedConnectionID == connectionID else { return } - guard chunk.count <= Self.maximumReceiveBufferByteCount - buffer.count else { - await tearDownIfInstalled( - connectionID: connectionID, - error: .invalidResponse - ) - return - } + // Enforce size per decoded frame. A chunk can finish one valid + // maximum-size frame and also contain bytes from the next frame. buffer.append(chunk) do { while !Task.isCancelled, installedConnectionID == connectionID { diff --git a/Sources/Mobile/MobileHostService.swift b/Sources/Mobile/MobileHostService.swift index 43ba7f71b7df..286ea0e2e22c 100644 --- a/Sources/Mobile/MobileHostService.swift +++ b/Sources/Mobile/MobileHostService.swift @@ -2038,7 +2038,6 @@ extension MobileHostService { #endif actor MobileHostConnection { - private static let maximumReceiveBufferByteCount = MobileSyncFrameCodec.defaultMaximumFrameByteCount + MobileSyncFrameCodec.headerByteCount private static let defaultFirstFrameTimeoutNanoseconds: UInt64 = 15 * 1_000_000_000 fileprivate static let defaultIdleTimeoutNanoseconds: UInt64 = 30 * 1_000_000_000 private struct EventSubscription: Sendable { @@ -2281,23 +2280,8 @@ actor MobileHostConnection { if !data.isEmpty { idleTimeoutTask?.cancel() idleTimeoutTask = nil - guard receiveBuffer.count + data.count <= Self.maximumReceiveBufferByteCount else { - _ = await sendResponse( - MobileHostRPCEnvelope.error( - id: nil, - code: "frame_decode_error", - message: "Invalid frame" - ) - ) - await close( - reason: "receive buffer exceeded frame limit", - exit: CmxIrohAdmittedConnectionExit( - lifecycle: .controlReadFailed, - failure: .protocolViolation - ) - ) - return - } + // Message limits belong to individual frames. A receive chunk may + // contain the tail of a maximum-size frame followed by another. receiveBuffer.append(data) do { let batchLimit = responseWorkQuota.maximumConcurrentRequestCount From 499f8f4985cee9e1e47e0fa46ef9d1f5303e78fb Mon Sep 17 00:00:00 2001 From: Abdulaziz Albahar <67667005+azooz2003-bit@users.noreply.github.com> Date: Fri, 11 Sep 2026 23:44:48 -0700 Subject: [PATCH 06/17] test: derive mobile lane quota expectation --- cmuxTests/MobileHostIrohAdmissionTests.swift | 6 +++++- 1 file changed, 5 insertions(+), 1 deletion(-) diff --git a/cmuxTests/MobileHostIrohAdmissionTests.swift b/cmuxTests/MobileHostIrohAdmissionTests.swift index c8cd8265d6f5..cafc18db0122 100644 --- a/cmuxTests/MobileHostIrohAdmissionTests.swift +++ b/cmuxTests/MobileHostIrohAdmissionTests.swift @@ -722,7 +722,11 @@ extension MobileHostAuthorizationTests { @Test func testIrohApplicationLaneQuotasReserveArtifactCapacity() { #expect(MobileHostIrohApplicationLaneRouter.maximumConcurrentTerminalLaneCount == 4) #expect(MobileHostIrohApplicationLaneRouter.maximumConcurrentArtifactLaneCount == 1) - #expect(MobileHostIrohApplicationLaneRouter.maximumConcurrentLaneCount == 5) + #expect( + MobileHostIrohApplicationLaneRouter.maximumConcurrentLaneCount + == MobileHostIrohApplicationLaneRouter.maximumConcurrentTerminalLaneCount + + MobileHostIrohApplicationLaneRouter.maximumConcurrentArtifactLaneCount + ) var quota = MobileHostIrohApplicationLaneQuota() let terminalIDs = (0..<5).map { _ in UUID() } From 50baaaad0a49c870072f3f522073299a23153ac1 Mon Sep 17 00:00:00 2001 From: Abdulaziz Albahar <67667005+azooz2003-bit@users.noreply.github.com> Date: Sat, 12 Sep 2026 19:07:13 -0700 Subject: [PATCH 07/17] test: include simulator streams in mobile lane quota --- cmuxTests/MobileHostIrohAdmissionTests.swift | 17 +++++++++++++++++ 1 file changed, 17 insertions(+) diff --git a/cmuxTests/MobileHostIrohAdmissionTests.swift b/cmuxTests/MobileHostIrohAdmissionTests.swift index cafc18db0122..e3e1f06b151a 100644 --- a/cmuxTests/MobileHostIrohAdmissionTests.swift +++ b/cmuxTests/MobileHostIrohAdmissionTests.swift @@ -722,10 +722,12 @@ extension MobileHostAuthorizationTests { @Test func testIrohApplicationLaneQuotasReserveArtifactCapacity() { #expect(MobileHostIrohApplicationLaneRouter.maximumConcurrentTerminalLaneCount == 4) #expect(MobileHostIrohApplicationLaneRouter.maximumConcurrentArtifactLaneCount == 1) + #expect(MobileHostIrohApplicationLaneRouter.maximumConcurrentSimulatorStreamLaneCount == 2) #expect( MobileHostIrohApplicationLaneRouter.maximumConcurrentLaneCount == MobileHostIrohApplicationLaneRouter.maximumConcurrentTerminalLaneCount + MobileHostIrohApplicationLaneRouter.maximumConcurrentArtifactLaneCount + + MobileHostIrohApplicationLaneRouter.maximumConcurrentSimulatorStreamLaneCount ) var quota = MobileHostIrohApplicationLaneQuota() @@ -741,8 +743,18 @@ extension MobileHostAuthorizationTests { #expect(didReserveArtifact) let didReserveSecondArtifact = quota.reserve(UUID(), laneClass: .artifact) #expect(!didReserveSecondArtifact) + let simulatorStreamIDs = (0..<3).map { _ in UUID() } + for id in simulatorStreamIDs.prefix(2) { + let didReserve = quota.reserve(id, laneClass: .simulatorStream) + #expect(didReserve) + } + let didReserveThirdSimulatorStream = quota.reserve( + simulatorStreamIDs[2], laneClass: .simulatorStream + ) + #expect(!didReserveThirdSimulatorStream) #expect(quota.terminalCount == 4) #expect(quota.artifactCount == 1) + #expect(quota.simulatorStreamCount == 2) quota.release(terminalIDs[0]) let didReuseTerminalCredit = quota.reserve(terminalIDs[4], laneClass: .terminal) @@ -750,6 +762,11 @@ extension MobileHostAuthorizationTests { quota.release(artifactID) let didReuseArtifactCredit = quota.reserve(UUID(), laneClass: .artifact) #expect(didReuseArtifactCredit) + quota.release(simulatorStreamIDs[0]) + let didReuseSimulatorStreamCredit = quota.reserve( + simulatorStreamIDs[2], laneClass: .simulatorStream + ) + #expect(didReuseSimulatorStreamCredit) } private func irohPeer( From ebe4c47af8b4e169b083051afcd3b89781ff382e Mon Sep 17 00:00:00 2001 From: Abdulaziz Albahar <67667005+azooz2003-bit@users.noreply.github.com> Date: Sun, 13 Sep 2026 22:21:11 -0700 Subject: [PATCH 08/17] test: keep mobile host connections alive while idle --- .../MobileHostConnectionEventLaneTests.swift | 34 ++++++----------- .../MobileHostConnectionLifecycleTests.swift | 37 ++----------------- 2 files changed, 15 insertions(+), 56 deletions(-) diff --git a/cmuxTests/MobileHostConnectionEventLaneTests.swift b/cmuxTests/MobileHostConnectionEventLaneTests.swift index ff0b31b3ec67..8d81b8dd1b17 100644 --- a/cmuxTests/MobileHostConnectionEventLaneTests.swift +++ b/cmuxTests/MobileHostConnectionEventLaneTests.swift @@ -43,7 +43,7 @@ extension MobileHostAuthorizationTests { let finalRecordedIDs = await recorder.recordedIDs() #expect(finalRecordedIDs == [connectionID]) } - @Test func testMobileHostConnectionClosesWhenIdleAfterFirstFrame() async throws { + @Test func testMobileHostConnectionStaysOpenWhenIdleAfterFirstFrame() async throws { let connectionID = UUID() let recorder = MobileHostConnectionCloseRecorder() let connection = NWConnection( @@ -54,7 +54,6 @@ extension MobileHostAuthorizationTests { let session = MobileHostConnection( id: connectionID, connection: connection, - idleTimeoutNanoseconds: 1_000_000, authorizeRequest: { _ in nil }, onAuthorizedRequest: { _ in }, handleRequest: { _ in .ok([:]) }, @@ -62,16 +61,13 @@ extension MobileHostAuthorizationTests { await recorder.record(id) } ) - await session.debugStartIdleTimeoutAfterFrameForTesting() - for _ in 0..<100 { - let recordedIDs = await recorder.recordedIDs() - if !recordedIDs.isEmpty { - break - } - try await Task.sleep(nanoseconds: 1_000_000) - } - let finalRecordedIDs = await recorder.recordedIDs() - #expect(finalRecordedIDs == [connectionID]) + let frame = try MobileSyncFrameCodec.encodeFrame( + Data(#"{"id":"status","method":"mobile.host.status","params":{}}"#.utf8) + ) + await session.debugHandleReceiveDataForTesting(frame) + try await Task.sleep(nanoseconds: 25_000_000) + #expect(await recorder.recordedIDs().isEmpty) + await session.close(reason: "test cleanup") } @Test func testMobileHostConnectionKeepsSubscribedEventStreamPastIdleTimeout() async throws { let connectionID = UUID() @@ -84,7 +80,6 @@ extension MobileHostAuthorizationTests { let session = MobileHostConnection( id: connectionID, connection: connection, - idleTimeoutNanoseconds: 1_000_000, authorizeRequest: { _ in nil }, onAuthorizedRequest: { _ in }, handleRequest: { _ in .ok([:]) }, @@ -93,7 +88,6 @@ extension MobileHostAuthorizationTests { } ) await session.subscribe(streamID: "events", topics: ["terminal.updated"]) - await session.debugStartIdleTimeoutAfterFrameForTesting() // An active subscription suppresses the idle-after-frame timeout: the // arm path early-returns without scheduling any close. Awaiting an // actor-isolated round-trip on the connection guarantees the arm call @@ -104,15 +98,9 @@ extension MobileHostAuthorizationTests { let subscribedCloseIDs = await recorder.recordedIDs() #expect(subscribedCloseIDs.isEmpty) _ = await session.unsubscribe(streamID: "events") - for _ in 0..<100 { - let recordedIDs = await recorder.recordedIDs() - if !recordedIDs.isEmpty { - break - } - try await Task.sleep(nanoseconds: 1_000_000) - } - let finalRecordedIDs = await recorder.recordedIDs() - #expect(finalRecordedIDs == [connectionID]) + try await Task.sleep(nanoseconds: 25_000_000) + #expect(await recorder.recordedIDs().isEmpty) + await session.close(reason: "test cleanup") } @Test func testDeadIndependentEventLaneFallsBackCurrentAndFutureEventsToControl() async throws { diff --git a/cmuxTests/MobileHostConnectionLifecycleTests.swift b/cmuxTests/MobileHostConnectionLifecycleTests.swift index 4b5434f4bd87..308993ffe2d0 100644 --- a/cmuxTests/MobileHostConnectionLifecycleTests.swift +++ b/cmuxTests/MobileHostConnectionLifecycleTests.swift @@ -189,7 +189,7 @@ extension MobileHostAuthorizationTests { service.debugResetMobileLifecycleStateForTesting() } - @Test func testIrohTransportCanDisableTheControlIdleTimeout() async throws { + @Test func testMobileHostTransportStaysOpenWhenIdleAfterAdmission() async throws { let service = MobileHostService.shared service.debugResetMobileLifecycleStateForTesting() let registry = MobileHostConnectionRegistry.shared @@ -200,33 +200,12 @@ extension MobileHostAuthorizationTests { service.debugResetMobileLifecycleStateForTesting() } - let expiringTransport = ScriptedMobileHostByteTransport() let authorization = try irohAdmissionContext() - let expiringTask = Task { - await MobileHostService.acceptTransport( - expiringTransport, - authorization: authorization, - idleTimeoutNanoseconds: 1_000_000, - isCurrent: { true } - ) - } - await waitForMobileHostConnectionCount(1) - try await expiringTransport.enqueue(Self.mobileHostStatusFrame(id: "expiring")) - _ = await expiringTransport.waitForSentBufferCount(1) - await expiringTransport.waitForCloseCount(1) - #expect( - await expiringTask.value == CmxIrohAdmittedConnectionExit( - lifecycle: .controlReadFailed, - failure: .timedOut - ) - ) - let persistentTransport = ScriptedMobileHostByteTransport() let persistentTask = Task { await MobileHostService.acceptTransport( persistentTransport, authorization: authorization, - idleTimeoutNanoseconds: 0, isCurrent: { true } ) } @@ -525,7 +504,6 @@ extension MobileHostAuthorizationTests { let session = MobileHostConnection( id: connectionID, connection: socket.connection, - idleTimeoutNanoseconds: 1_000_000, authorizeRequest: { _ in .failure(MobileHostRPCError(code: "unauthorized", message: "no")) }, @@ -540,16 +518,9 @@ extension MobileHostAuthorizationTests { ) await session.debugHandleReceiveDataForTesting(frame) try await Task.sleep(nanoseconds: 25_000_000) - await session.debugStartIdleTimeoutAfterFrameForTesting() - for _ in 0..<100 { - let recordedIDs = await recorder.recordedIDs() - if !recordedIDs.isEmpty { - break - } - try await Task.sleep(nanoseconds: 1_000_000) - } - let finalRecordedIDs = await recorder.recordedIDs() - #expect(finalRecordedIDs == [connectionID]) + try await Task.sleep(nanoseconds: 25_000_000) + #expect(await recorder.recordedIDs().isEmpty) + await session.close(reason: "test cleanup") } @Test func testMobileHostConnectionStopsBatchedFrameProcessingAfterClose() async throws { let connectionID = UUID() From c0b792e28097b4ddc3b7df015fe552a2dce0b9d7 Mon Sep 17 00:00:00 2001 From: Abdulaziz Albahar <67667005+azooz2003-bit@users.noreply.github.com> Date: Sun, 13 Sep 2026 22:34:03 -0700 Subject: [PATCH 09/17] fix: let native transport own idle session lifetime --- ...ileShellComposite+ConnectionRecovery.swift | 37 ++++++---- .../MobileShellComposite.swift | 5 +- .../IrohConnectionRecoveryOwnerTests.swift | 3 + .../MobileHostIrohRuntime+Activation.swift | 1 - .../MobileHostIrxLegacyDialectServer.swift | 1 - Sources/Mobile/MobileHostIrxRuntime.swift | 4 -- Sources/Mobile/MobileHostService.swift | 69 ------------------- .../MobileHostConnectionEventLaneTests.swift | 26 ++----- cmuxTests/MobileHostIrohAdmissionTests.swift | 2 - 9 files changed, 36 insertions(+), 112 deletions(-) diff --git a/Packages/iOS/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite+ConnectionRecovery.swift b/Packages/iOS/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite+ConnectionRecovery.swift index 206a484e07f5..8d7c862ddb6a 100644 --- a/Packages/iOS/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite+ConnectionRecovery.swift +++ b/Packages/iOS/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite+ConnectionRecovery.swift @@ -355,14 +355,19 @@ extension MobileShellComposite { self.connectionRecoveryOwner.transitionToRedialing(attempt) else { return } if let expectedClient { guard self.remoteClient === expectedClient else { return } - // Detach the stale shell synchronously on the main actor - // before awaiting its transport teardown. This cancels every - // tracked producer and makes untracked producers fail their - // identity guard, so they cannot reopen the old endpoint - // while the fresh stored-Mac dial starts. - self.connectionState = .disconnected - self.macConnectionStatus = .unavailable - self.clearRemoteConnectionContext() + // Retire the stale client before the replacement dial, but + // keep the foreground identity and last-known workspace + // presentation while the replacement is being validated. + // A path handoff is an implementation detail until the new + // session fails; clearing connectionState here creates a + // false disconnected flash and deactivates every terminal + // lane even when replacement succeeds immediately. + self.connectionGeneration = UUID() + self.cancelRemoteOperationTasks() + self.rawTerminalInputBuffer.clear() + self.terminalInputRPCPipeline.clear() + self.resumeRawTerminalInputDrainWaiters() + await self.releaseRemoteClientForReplacement() self.applyConnectionRecoveryOwnerState() MobileDebugLog.anchormux( "connection.recovery waiting for physical transport drain " @@ -390,11 +395,6 @@ extension MobileShellComposite { + "attempt=\(attempt.id.uuidString)" ) } - if self.connectionState == .connected { - self.connectionState = .disconnected - self.macConnectionStatus = .unavailable - self.clearRemoteConnectionContext() - } self.applyConnectionRecoveryOwnerState() // Recovery uses authenticated local Iroh state first. A stuck @@ -414,6 +414,12 @@ extension MobileShellComposite { outcome: reconnectOutcome, connectionGeneration: self.connectionGeneration ) else { return } + if !reconnectOutcome.didConnect, + self.connectionState == .connected { + self.connectionState = .disconnected + self.macConnectionStatus = .unavailable + self.clearRemoteConnectionContext() + } self.applyConnectionRecoveryOwnerState() } onCancel: { MobileDebugLog.anchormux( @@ -447,15 +453,16 @@ extension MobileShellComposite { try? await ContinuousClock().sleep(for: .seconds(seconds)) guard let self, !Task.isCancelled else { return } guard self.connectionRecoveryOwner.isCurrent(attempt), - self.connectionRecoveryOwner.isRedialingOrValidating, - self.connectionState != .connected else { return } + self.connectionRecoveryOwner.isRedialingOrValidating else { return } MobileDebugLog.anchormux( "connection.recovery attempt deadline forced failure " + "trigger=\(attempt.trigger) attempt=\(attempt.id.uuidString)" ) guard self.connectionRecoveryOwner.failReplacement() != nil else { return } self.recordConnectionRecoveryFailed(attempt, failure: .timedOut) + self.connectionState = .disconnected self.macConnectionStatus = .unavailable + self.clearRemoteConnectionContext() self.applyConnectionRecoveryOwnerState() self.armAutomaticReconnectRetryAfterFailedAttempt( failure: .timedOut, diff --git a/Packages/iOS/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite.swift b/Packages/iOS/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite.swift index aa2052797e8c..b94e47794c53 100644 --- a/Packages/iOS/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite.swift +++ b/Packages/iOS/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite.swift @@ -11915,7 +11915,10 @@ public final class MobileShellComposite: MobileTerminalOutputSinking { } func markMacConnectionReconnecting() { - guard connectionState == .connected, remoteClient != nil else { + // Recovery retires the old RPC client before dialing its replacement, + // while the logical session remains presented as connected. Keep the + // reconnecting status visible during that ownership gap. + guard connectionState == .connected else { macConnectionStatus = .unavailable return } diff --git a/Packages/iOS/CmuxMobileShell/Tests/CmuxMobileShellTests/IrohConnectionRecoveryOwnerTests.swift b/Packages/iOS/CmuxMobileShell/Tests/CmuxMobileShellTests/IrohConnectionRecoveryOwnerTests.swift index 8a182c0ddaf0..e7a4c19b4543 100644 --- a/Packages/iOS/CmuxMobileShell/Tests/CmuxMobileShellTests/IrohConnectionRecoveryOwnerTests.swift +++ b/Packages/iOS/CmuxMobileShell/Tests/CmuxMobileShellTests/IrohConnectionRecoveryOwnerTests.swift @@ -78,6 +78,9 @@ extension ReconnectRouteSelectionTests { ) #expect(await closeGate.waitUntilCloseStarted()) + #expect(fixture.store.connectionState == .connected) + #expect(fixture.store.macConnectionStatus == .reconnecting) + #expect(fixture.store.isRecoveringConnection) #expect(fixture.factory.attemptedKinds() == [.iroh]) await closeGate.release() diff --git a/Sources/Mobile/MobileHostIrohRuntime+Activation.swift b/Sources/Mobile/MobileHostIrohRuntime+Activation.swift index 43599e1e2741..274e77e29c3d 100644 --- a/Sources/Mobile/MobileHostIrohRuntime+Activation.swift +++ b/Sources/Mobile/MobileHostIrohRuntime+Activation.swift @@ -345,7 +345,6 @@ extension MobileHostIrohRuntime { authorization: .irohAdmission(session.peer), artifactTransfers: artifactTransfers, independentEventWriter: eventWriter, - idleTimeoutNanoseconds: 0, promoteUsableSession: { await session.markUsable() }, diff --git a/Sources/Mobile/MobileHostIrxLegacyDialectServer.swift b/Sources/Mobile/MobileHostIrxLegacyDialectServer.swift index 3f6ec1e0646f..9858ba6749fa 100644 --- a/Sources/Mobile/MobileHostIrxLegacyDialectServer.swift +++ b/Sources/Mobile/MobileHostIrxLegacyDialectServer.swift @@ -113,7 +113,6 @@ enum MobileHostIrxLegacyDialectServer { authorization: .irohAdmission(admitted.peer), artifactTransfers: artifactTransfers, independentEventWriter: eventWriter, - idleTimeoutNanoseconds: 0, promoteUsableSession: { await admitted.markUsable() }, isCurrent: isCurrent ) diff --git a/Sources/Mobile/MobileHostIrxRuntime.swift b/Sources/Mobile/MobileHostIrxRuntime.swift index dffa07a20f2b..3bb86e575368 100644 --- a/Sources/Mobile/MobileHostIrxRuntime.swift +++ b/Sources/Mobile/MobileHostIrxRuntime.swift @@ -941,10 +941,6 @@ final class MobileHostIrxRuntime { authorization: .irohAdmission(admittedPeer), artifactTransfers: artifactRegistry, independentEventWriter: eventWriter, - // The bounded Iroh peer pool stays alive via transport keepalives. - // Control-idle timeout is for unowned legacy TCP connections and - // must not tear down a healthy multi-lane QUIC session. - idleTimeoutNanoseconds: 0, isCurrent: { [weak self] in let runtime = self return await MainActor.run { runtime?.generationToken == token } diff --git a/Sources/Mobile/MobileHostService.swift b/Sources/Mobile/MobileHostService.swift index 286ea0e2e22c..51219e6eb775 100644 --- a/Sources/Mobile/MobileHostService.swift +++ b/Sources/Mobile/MobileHostService.swift @@ -1367,7 +1367,6 @@ final class MobileHostService { authorization: MobileHostConnectionAuthorizationContext, artifactTransfers: MobileHostIrohArtifactTransferRegistry? = nil, independentEventWriter: (any MobileHostIndependentEventWriting)? = nil, - idleTimeoutNanoseconds: UInt64? = nil, promoteUsableSession: @escaping @Sendable () async -> Bool = { true }, remoteControlDisabledByPolicy: @escaping @Sendable () -> Bool = { MobileRemoteControlPolicy.isDisabled @@ -1398,8 +1397,6 @@ final class MobileHostService { let session = MobileHostConnection( id: id, transport: transport, - idleTimeoutNanoseconds: idleTimeoutNanoseconds - ?? MobileHostConnection.defaultIdleTimeoutNanoseconds, independentEventWriter: independentEventWriter, authorizeRequest: { request in await Self.connectionAuthorizationError( @@ -2039,7 +2036,6 @@ extension MobileHostService { actor MobileHostConnection { private static let defaultFirstFrameTimeoutNanoseconds: UInt64 = 15 * 1_000_000_000 - fileprivate static let defaultIdleTimeoutNanoseconds: UInt64 = 30 * 1_000_000_000 private struct EventSubscription: Sendable { let topics: Set let transport: MobileHostEventTransport @@ -2079,7 +2075,6 @@ actor MobileHostConnection { private let writer: MobileHostSerializedTransportWriter private let independentEventWriter: (any MobileHostIndependentEventWriting)? private let firstFrameTimeoutNanoseconds: UInt64 - private let idleTimeoutNanoseconds: UInt64 private let authorizeRequest: @Sendable (MobileHostRPCRequest) async -> MobileHostRPCResult? private let onAuthorizedRequest: @Sendable (MobileHostRPCRequest) async -> Void private let onUsableSession: @Sendable () async -> Bool @@ -2093,7 +2088,6 @@ actor MobileHostConnection { nonisolated let eventQueue: MobileHostConnectionEventQueue private var receiveBuffer = Data() private var firstFrameTimeoutTask: Task? - private var idleTimeoutTask: Task? private var responseTasks: [UUID: ResponseTask] = [:] /// PTY-writing requests are ordered PER SURFACE: ordering is only a /// property of one terminal, and a connection-wide FIFO would let one @@ -2122,7 +2116,6 @@ actor MobileHostConnection { connection: NWConnection, eventQueue: MobileHostConnectionEventQueue = MobileHostConnectionEventQueue(), firstFrameTimeoutNanoseconds: UInt64 = MobileHostConnection.defaultFirstFrameTimeoutNanoseconds, - idleTimeoutNanoseconds: UInt64 = MobileHostConnection.defaultIdleTimeoutNanoseconds, independentEventWriter: (any MobileHostIndependentEventWriting)? = nil, authorizeRequest: @escaping @Sendable (MobileHostRPCRequest) async -> MobileHostRPCResult?, onAuthorizedRequest: @escaping @Sendable (MobileHostRPCRequest) async -> Void, @@ -2137,7 +2130,6 @@ actor MobileHostConnection { self.writer = MobileHostSerializedTransportWriter(transport: transport) self.independentEventWriter = independentEventWriter self.firstFrameTimeoutNanoseconds = firstFrameTimeoutNanoseconds - self.idleTimeoutNanoseconds = idleTimeoutNanoseconds self.authorizeRequest = authorizeRequest self.onAuthorizedRequest = onAuthorizedRequest self.onUsableSession = onUsableSession @@ -2152,7 +2144,6 @@ actor MobileHostConnection { transport: any CmxByteTransport, eventQueue: MobileHostConnectionEventQueue = MobileHostConnectionEventQueue(), firstFrameTimeoutNanoseconds: UInt64 = MobileHostConnection.defaultFirstFrameTimeoutNanoseconds, - idleTimeoutNanoseconds: UInt64 = MobileHostConnection.defaultIdleTimeoutNanoseconds, independentEventWriter: (any MobileHostIndependentEventWriting)? = nil, authorizeRequest: @escaping @Sendable (MobileHostRPCRequest) async -> MobileHostRPCResult?, onAuthorizedRequest: @escaping @Sendable (MobileHostRPCRequest) async -> Void, @@ -2166,7 +2157,6 @@ actor MobileHostConnection { self.writer = MobileHostSerializedTransportWriter(transport: transport) self.independentEventWriter = independentEventWriter self.firstFrameTimeoutNanoseconds = firstFrameTimeoutNanoseconds - self.idleTimeoutNanoseconds = idleTimeoutNanoseconds self.authorizeRequest = authorizeRequest self.onAuthorizedRequest = onAuthorizedRequest self.onUsableSession = onUsableSession @@ -2243,8 +2233,6 @@ actor MobileHostConnection { self.exit = exit firstFrameTimeoutTask?.cancel() firstFrameTimeoutTask = nil - idleTimeoutTask?.cancel() - idleTimeoutTask = nil receiveTask?.cancel() receiveTask = nil // Rejects all future admissions and releases every queued payload; the @@ -2278,8 +2266,6 @@ actor MobileHostConnection { private func handleReceive(data: Data) async { if !data.isEmpty { - idleTimeoutTask?.cancel() - idleTimeoutTask = nil // Message limits belong to individual frames. A receive chunk may // contain the tail of a maximum-size frame followed by another. receiveBuffer.append(data) @@ -2314,7 +2300,6 @@ actor MobileHostConnection { guard !isClosed else { return } - startIdleTimeout() } catch { _ = await sendResponse( MobileHostRPCEnvelope.error( @@ -2410,9 +2395,6 @@ actor MobileHostConnection { startOrderedRequestWorkerIfNeeded(surfaceKey: surfaceKey) } else { orderedRequestQueuesBySurfaceKey[surfaceKey] = nil - if !hasActiveResponseWork { - startIdleTimeout() - } } } @@ -2434,18 +2416,8 @@ actor MobileHostConnection { ) } - private var hasActiveResponseWork: Bool { - !responseTasks.isEmpty - || !orderedRequestWorkerTasksBySurfaceKey.isEmpty - || !orderedRequestRunningFrameByteCountsBySurfaceKey.isEmpty - || orderedRequestQueuesBySurfaceKey.values.contains { !$0.isEmpty } - } - private func finishResponseTask(_ taskID: UUID) { responseTasks[taskID] = nil - if !hasActiveResponseWork { - startIdleTimeout() - } } private func startFirstFrameTimeout() { @@ -2475,37 +2447,6 @@ actor MobileHostConnection { ) } - private func startIdleTimeout() { - guard idleTimeoutNanoseconds > 0, - didDecodeFirstFrame, - !isClosed, - subscriptions.isEmpty, - !hasActiveResponseWork else { - return - } - idleTimeoutTask?.cancel() - let timeoutNanoseconds = idleTimeoutNanoseconds - idleTimeoutTask = Task { [weak self] in - do { - try await Task.sleep(nanoseconds: timeoutNanoseconds) - await self?.closeIfIdleAfterFrame() - } catch {} - } - } - - private func closeIfIdleAfterFrame() async { - guard didDecodeFirstFrame, subscriptions.isEmpty, !hasActiveResponseWork else { - return - } - await close( - reason: "idle after frame timed out", - exit: CmxIrohAdmittedConnectionExit( - lifecycle: .controlReadFailed, - failure: .timedOut - ) - ) - } - private func respond( to decodedRequest: Result ) async { @@ -2814,8 +2755,6 @@ actor MobileHostConnection { previousTopics: previousTopics, nextTopics: topics ) - idleTimeoutTask?.cancel() - idleTimeoutTask = nil if currentSubscribedTopics().contains(MobileHostEventTopicPolicy.simulatorFrameTopic) { await dispatchPendingSimulatorFrameReplay() } @@ -2841,9 +2780,6 @@ actor MobileHostConnection { }) { await resetIndependentEventWriter() } - if subscriptions.isEmpty { - startIdleTimeout() - } return removed } @@ -3132,11 +3068,6 @@ extension MobileHostConnection { startFirstFrameTimeout() } - func debugStartIdleTimeoutAfterFrameForTesting() { - didDecodeFirstFrame = true - startIdleTimeout() - } - func debugHandleReceiveDataForTesting(_ data: Data) async { await handleReceive(data: data) } diff --git a/cmuxTests/MobileHostConnectionEventLaneTests.swift b/cmuxTests/MobileHostConnectionEventLaneTests.swift index 8d81b8dd1b17..1684310f2f54 100644 --- a/cmuxTests/MobileHostConnectionEventLaneTests.swift +++ b/cmuxTests/MobileHostConnectionEventLaneTests.swift @@ -16,14 +16,10 @@ extension MobileHostAuthorizationTests { @Test func testMobileHostConnectionClosesWhenFirstFrameTimesOut() async throws { let connectionID = UUID() let recorder = MobileHostConnectionCloseRecorder() - let connection = NWConnection( - host: NWEndpoint.Host("127.0.0.1"), - port: NWEndpoint.Port(rawValue: 9)!, - using: .tcp - ) + let transport = RecordingMobileHostByteTransport() let session = MobileHostConnection( id: connectionID, - connection: connection, + transport: transport, firstFrameTimeoutNanoseconds: 1_000_000, authorizeRequest: { _ in nil }, onAuthorizedRequest: { _ in }, @@ -46,14 +42,10 @@ extension MobileHostAuthorizationTests { @Test func testMobileHostConnectionStaysOpenWhenIdleAfterFirstFrame() async throws { let connectionID = UUID() let recorder = MobileHostConnectionCloseRecorder() - let connection = NWConnection( - host: NWEndpoint.Host("127.0.0.1"), - port: NWEndpoint.Port(rawValue: 9)!, - using: .tcp - ) + let transport = RecordingMobileHostByteTransport() let session = MobileHostConnection( id: connectionID, - connection: connection, + transport: transport, authorizeRequest: { _ in nil }, onAuthorizedRequest: { _ in }, handleRequest: { _ in .ok([:]) }, @@ -69,7 +61,7 @@ extension MobileHostAuthorizationTests { #expect(await recorder.recordedIDs().isEmpty) await session.close(reason: "test cleanup") } - @Test func testMobileHostConnectionKeepsSubscribedEventStreamPastIdleTimeout() async throws { + @Test func testMobileHostConnectionKeepsSubscribedEventStreamIdle() async throws { let connectionID = UUID() let recorder = MobileHostConnectionCloseRecorder() let connection = NWConnection( @@ -88,12 +80,8 @@ extension MobileHostAuthorizationTests { } ) await session.subscribe(streamID: "events", topics: ["terminal.updated"]) - // An active subscription suppresses the idle-after-frame timeout: the - // arm path early-returns without scheduling any close. Awaiting an - // actor-isolated round-trip on the connection guarantees the arm call - // was fully processed and that the connection is still alive and - // subscribed, so the recorder reflects the final state with no - // wall-clock window to race. + // Subscriptions remain usable without requiring synthetic traffic. The + // host has no application-level idle deadline after admission. #expect(await session.isSubscribed(to: "terminal.updated")) let subscribedCloseIDs = await recorder.recordedIDs() #expect(subscribedCloseIDs.isEmpty) diff --git a/cmuxTests/MobileHostIrohAdmissionTests.swift b/cmuxTests/MobileHostIrohAdmissionTests.swift index e3e1f06b151a..b680852f9f48 100644 --- a/cmuxTests/MobileHostIrohAdmissionTests.swift +++ b/cmuxTests/MobileHostIrohAdmissionTests.swift @@ -231,7 +231,6 @@ struct IrohTailscaleVersionSkewMacGateTests { id: UUID(), transport: transport, firstFrameTimeoutNanoseconds: 0, - idleTimeoutNanoseconds: 0, authorizeRequest: { request in await MobileHostService.connectionAuthorizationError( for: request, @@ -292,7 +291,6 @@ struct IrohTailscaleVersionSkewMacGateTests { id: UUID(), transport: transport, firstFrameTimeoutNanoseconds: 0, - idleTimeoutNanoseconds: 0, authorizeRequest: { request in await MobileHostService.connectionAuthorizationError( for: request, From 39c35acbb7d15b5046c64065a2cc8d7808e36dac Mon Sep 17 00:00:00 2001 From: Abdulaziz Albahar <67667005+azooz2003-bit@users.noreply.github.com> Date: Mon, 14 Sep 2026 11:20:56 -0700 Subject: [PATCH 10/17] test: distinguish retired transport from recovery presentation --- .../IrohConnectionRecoveryOwnerTests.swift | 5 ++++- 1 file changed, 4 insertions(+), 1 deletion(-) diff --git a/Packages/iOS/CmuxMobileShell/Tests/CmuxMobileShellTests/IrohConnectionRecoveryOwnerTests.swift b/Packages/iOS/CmuxMobileShell/Tests/CmuxMobileShellTests/IrohConnectionRecoveryOwnerTests.swift index e7a4c19b4543..e39a09c5b35b 100644 --- a/Packages/iOS/CmuxMobileShell/Tests/CmuxMobileShellTests/IrohConnectionRecoveryOwnerTests.swift +++ b/Packages/iOS/CmuxMobileShell/Tests/CmuxMobileShellTests/IrohConnectionRecoveryOwnerTests.swift @@ -78,7 +78,10 @@ extension ReconnectRouteSelectionTests { ) #expect(await closeGate.waitUntilCloseStarted()) - #expect(fixture.store.connectionState == .connected) + #expect(fixture.store.connectionState == .disconnected) + #expect(fixture.store.remoteClient == nil) + #expect(fixture.store.foregroundMacDeviceID == "test-mac") + #expect(!fixture.store.workspaces.isEmpty) #expect(fixture.store.macConnectionStatus == .reconnecting) #expect(fixture.store.isRecoveringConnection) #expect(fixture.factory.attemptedKinds() == [.iroh]) From 4edfea487ab02c111f01d593474e3c6008eb09b0 Mon Sep 17 00:00:00 2001 From: Abdulaziz Albahar <67667005+azooz2003-bit@users.noreply.github.com> Date: Mon, 14 Sep 2026 11:21:18 -0700 Subject: [PATCH 11/17] fix: retain recovery presentation while retiring dead transports --- ...ileShellComposite+ConnectionRecovery.swift | 20 +++---------- .../MobileShellComposite.swift | 28 ++++++++++++++++--- 2 files changed, 28 insertions(+), 20 deletions(-) diff --git a/Packages/iOS/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite+ConnectionRecovery.swift b/Packages/iOS/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite+ConnectionRecovery.swift index f613fc9fcb37..41cca1d5b9c0 100644 --- a/Packages/iOS/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite+ConnectionRecovery.swift +++ b/Packages/iOS/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite+ConnectionRecovery.swift @@ -355,19 +355,7 @@ extension MobileShellComposite { self.connectionRecoveryOwner.transitionToRedialing(attempt) else { return } if let expectedClient { guard self.remoteClient === expectedClient else { return } - // Retire the stale client before the replacement dial, but - // keep the foreground identity and last-known workspace - // presentation while the replacement is being validated. - // A path handoff is an implementation detail until the new - // session fails; clearing connectionState here creates a - // false disconnected flash and deactivates every terminal - // lane even when replacement succeeds immediately. - self.connectionGeneration = UUID() - self.cancelRemoteOperationTasks() - self.rawTerminalInputBuffer.clear() - self.terminalInputRPCPipeline.clear() - self.resumeRawTerminalInputDrainWaiters() - await self.releaseRemoteClientForReplacement() + self.retireRemoteClientForConnectionRecovery() self.applyConnectionRecoveryOwnerState() MobileDebugLog.anchormux( "connection.recovery waiting for physical transport drain " @@ -414,8 +402,7 @@ extension MobileShellComposite { outcome: reconnectOutcome, connectionGeneration: self.connectionGeneration ) else { return } - if !reconnectOutcome.didConnect, - self.connectionState == .connected { + if !reconnectOutcome.didConnect { self.connectionState = .disconnected self.macConnectionStatus = .unavailable self.clearRemoteConnectionContext() @@ -604,10 +591,11 @@ extension MobileShellComposite { case .redialing, .validatingReplacement: isRecoveringConnection = true connectionRecoveryFailed = false - if connectionState == .connected { markMacConnectionReconnecting() } + markMacConnectionReconnecting() case .failed: isRecoveringConnection = false connectionRecoveryFailed = true + if connectionState != .connected { macConnectionStatus = .unavailable } } } diff --git a/Packages/iOS/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite.swift b/Packages/iOS/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite.swift index 586d6986a1d9..8961e1af4494 100644 --- a/Packages/iOS/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite.swift +++ b/Packages/iOS/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite.swift @@ -11075,6 +11075,26 @@ public final class MobileShellComposite: MobileTerminalOutputSinking { await replaceRemoteClientAwaitingTeardownRegistration(with: nil) } + /// Retire the dead transport synchronously while retaining the Mac identity + /// and workspace rows. Recovery owns the visible status; connectionState + /// still records real transport availability so streams restart on adoption. + func retireRemoteClientForConnectionRecovery() { + connectionGeneration = UUID() + connectionAttemptGeneration = UUID() + cancelRemoteOperationTasks() + macConnectionStatus = .reconnecting + connectionState = .disconnected + rawTerminalInputBuffer.clear() + terminalInputRPCPipeline.clear() + resumeRawTerminalInputDrainWaiters() + if let focused = focusedForegroundConnection, + focused.client === remoteClient { + removeControlCapability(ifMatching: focused) + removeFocusedConnection(ifMatching: focused) + } + replaceRemoteClient(with: nil) + } + /// Retire the current pre-authentication candidate before a newer connect /// competes for the same physical route. The registry lease transfers to /// teardown before the replacement reaches transport admission. @@ -11969,10 +11989,10 @@ public final class MobileShellComposite: MobileTerminalOutputSinking { } func markMacConnectionReconnecting() { - // Recovery retires the old RPC client before dialing its replacement, - // while the logical session remains presented as connected. Keep the - // reconnecting status visible during that ownership gap. - guard connectionState == .connected else { + // An active replacement owns the presentation while its transport is + // absent. Probes on a usable connection never enter this phase. + guard connectionState == .connected + || connectionRecoveryOwner.isRedialingOrValidating else { macConnectionStatus = .unavailable return } From e7ca3cb36b8e3c79179a07831f4dc902fccda307 Mon Sep 17 00:00:00 2001 From: Abdulaziz Albahar <67667005+azooz2003-bit@users.noreply.github.com> Date: Tue, 15 Sep 2026 10:46:25 -0700 Subject: [PATCH 12/17] test: cover live transport foreground recovery --- .../MobileShellForegroundResumeTests.swift | 51 +++++++++++++++++-- 1 file changed, 46 insertions(+), 5 deletions(-) diff --git a/Packages/iOS/CmuxMobileShell/Tests/CmuxMobileShellTests/MobileShellForegroundResumeTests.swift b/Packages/iOS/CmuxMobileShell/Tests/CmuxMobileShellTests/MobileShellForegroundResumeTests.swift index 30bf050c8297..528442a06e5f 100644 --- a/Packages/iOS/CmuxMobileShell/Tests/CmuxMobileShellTests/MobileShellForegroundResumeTests.swift +++ b/Packages/iOS/CmuxMobileShell/Tests/CmuxMobileShellTests/MobileShellForegroundResumeTests.swift @@ -133,7 +133,7 @@ struct MobileShellForegroundConnectionRecoveryTests { } @MainActor -@Test func failedForegroundProbeStillSurfacesRedialRecoveryState() async throws { +@Test func failedForegroundProbeWithLiveTransportDoesNotRedial() async throws { let router = LivenessHostRouter() let box = TransportBox() let clock = TestClock() @@ -152,14 +152,55 @@ struct MobileShellForegroundConnectionRecoveryTests { store.resumeForegroundRefresh() #expect(try await pollUntil { - if case .redialing = store.connectionRecoveryOwner.phase { - return store.isRecoveringConnection - } - return false + store.connectionRecoveryOwner.phase == .idle + && !store.isRecoveringConnection + && store.connectionState == .connected }) await router.releaseAllHeld() } +@MainActor +@Test func foregroundProbeTimeoutWithLiveTransportRepairsMountedTerminal() async throws { + let router = LivenessHostRouter() + let box = TransportBox() + let clock = TestClock() + let (store, directory) = try await makeForegroundRecoveryStore( + router: router, + box: box, + clock: clock, + probeTimeoutNanoseconds: 50_000_000 + ) + defer { + Task { await router.releaseAllHeld() } + try? FileManager.default.removeItem(at: directory) + } + let collector = OutputCollector() + collector.mount(store: store, surfaceID: "live-terminal") + await router.waitForCount(of: "mobile.terminal.replay", atLeast: 1) + try await waitForReplayResponsesServed( + 1, + router: router, + "the cold replay response must settle before testing a foreground probe timeout" + ) + let replayCount = await router.count(of: "mobile.terminal.replay") + let originalTransport = try #require(box.get()) + + store.suspendForegroundRefresh() + clock.advance(by: 31) + await router.holdNextWorkspaceListRequests() + store.resumeForegroundRefresh() + #expect(await router.waitForCount(of: "mobile.sync.fetch", atLeast: 2)) + + #expect(await router.waitForCount( + of: "mobile.terminal.replay", + atLeast: replayCount + 1, + timeoutNanoseconds: 1_000_000_000 + )) + #expect(store.connectionState == .connected) + #expect(box.get() === originalTransport) + collector.unmount() +} + @MainActor @Test func suspendingForegroundRefreshCancelsInFlightProbeWithoutRedial() async throws { let router = LivenessHostRouter() From 7205608265f4bf5a0f97559c90f34e28d0629a4e Mon Sep 17 00:00:00 2001 From: Abdulaziz Albahar <67667005+azooz2003-bit@users.noreply.github.com> Date: Tue, 15 Sep 2026 10:46:25 -0700 Subject: [PATCH 13/17] fix: retain live Iroh sessions during foreground recovery --- .../DeviceRegistryService.swift | 22 ++++++++++------ ...ileShellComposite+ConnectionRecovery.swift | 25 ++++++++++++------- 2 files changed, 31 insertions(+), 16 deletions(-) diff --git a/Packages/iOS/CmuxMobileShell/Sources/CmuxMobileShell/DeviceRegistryService.swift b/Packages/iOS/CmuxMobileShell/Sources/CmuxMobileShell/DeviceRegistryService.swift index 06c6120abeb0..d54a8823eca8 100644 --- a/Packages/iOS/CmuxMobileShell/Sources/CmuxMobileShell/DeviceRegistryService.swift +++ b/Packages/iOS/CmuxMobileShell/Sources/CmuxMobileShell/DeviceRegistryService.swift @@ -813,14 +813,22 @@ public extension MobileIOSAppNamespace { deviceWitness: String? = nil, evidence: any SameDeviceEvidenceProbing = IrohEndpointIdentityEvidenceProbe() ) -> String? { - DeviceRegistryService.durableDeviceID( - store: KeychainDeviceIdentityStore( - service: keychainService( - base: "com.cmuxterm.deviceRegistry.iosDeviceID.v1" - ), - accessGroup: keychainAccessGroup, - legacyService: "com.cmuxterm.deviceRegistry.iosDeviceID.v1" + #if targetEnvironment(simulator) + let store: any DeviceIdentityStoring = SimulatorDeviceIdentityStore( + defaults: defaults, + seededDeviceID: ProcessInfo.processInfo.environment["CMUX_SIMULATOR_DEVICE_ID"] + ) + #else + let store: any DeviceIdentityStoring = KeychainDeviceIdentityStore( + service: keychainService( + base: "com.cmuxterm.deviceRegistry.iosDeviceID.v1" ), + accessGroup: keychainAccessGroup, + legacyService: "com.cmuxterm.deviceRegistry.iosDeviceID.v1" + ) + #endif + return DeviceRegistryService.durableDeviceID( + store: store, defaults: defaults, deviceWitness: deviceWitness, evidence: evidence diff --git a/Packages/iOS/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite+ConnectionRecovery.swift b/Packages/iOS/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite+ConnectionRecovery.swift index 41cca1d5b9c0..d209c01eb9a6 100644 --- a/Packages/iOS/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite+ConnectionRecovery.swift +++ b/Packages/iOS/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite+ConnectionRecovery.swift @@ -149,20 +149,15 @@ extension MobileShellComposite { } /// Checks native connection state before promoting a feature failure to - /// recovery. Transports without native observation retain their existing - /// error-driven recovery; Iroh owns its own dead-peer detection. + /// recovery. An event stream can end while Iroh has already replaced the + /// underlying path, so the stream ending alone does not justify retiring + /// the RPC client. Transports without native observation retain their + /// existing error-driven recovery; Iroh owns its own dead-peer detection. func recoverDeadConnection( trigger: RecoveryTrigger, expectedClient: MobileCoreRPCClient ) { guard remoteClient === expectedClient, connectionState == .connected else { return } - if trigger == .eventStreamEnded { - // The RPC listener ends only after its required control session - // has torn down. Detach it synchronously so another producer - // cannot reopen the retired RPC client while recovery starts. - recoverClosedControlSession(trigger: trigger, expectedClient: expectedClient) - return - } Task { @MainActor [weak self] in let closed = await expectedClient.isTransportClosed() guard let self, @@ -327,6 +322,18 @@ extension MobileShellComposite { // A slow application response after a path change or // resume does not invalidate the native connection. _ = self.completeConnectionRecovery(attempt) + self.markMacConnectionHealthy() + // The probe timed out, so the control request does + // not prove that the terminal event stream survived + // the background transition. Keep the native session + // and repair only the feature lane that can leave a + // mounted Ghostty surface blank. + if resyncAfterHealthy { + self.resyncTerminalOutput( + reason: "connectionRecovery.\(trigger).transportAlive", + restartEventStream: true + ) + } self.applyConnectionRecoveryOwnerState() return } From 6f64d4b2dc2f277e7f58ca9aebf57ee34166b5f8 Mon Sep 17 00:00:00 2001 From: austinpower1258 Date: Sun, 13 Sep 2026 21:49:48 -0700 Subject: [PATCH 14/17] test: avoid AppKit keyWindow collision in Cloud drag fixture --- cmuxTests/CloudTreeNativeDragOwnershipTests.swift | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/cmuxTests/CloudTreeNativeDragOwnershipTests.swift b/cmuxTests/CloudTreeNativeDragOwnershipTests.swift index 5ff6711c7237..4b28d30007dc 100644 --- a/cmuxTests/CloudTreeNativeDragOwnershipTests.swift +++ b/cmuxTests/CloudTreeNativeDragOwnershipTests.swift @@ -12,10 +12,10 @@ import Testing @Suite("Cloud tree native drag ownership", .serialized) struct CloudTreeNativeDragOwnershipTests { private final class HoverWindow: NSWindow { - var keyWindow = false + var simulatedKeyWindow = false var pointerOnScreen = NSPoint.zero - override var isKeyWindow: Bool { keyWindow } + override var isKeyWindow: Bool { simulatedKeyWindow } override var mouseLocationOutsideOfEventStream: NSPoint { pointerOnScreen } } @@ -317,7 +317,7 @@ struct CloudTreeNativeDragOwnershipTests { NotificationCenter.default.post(name: NSWindow.didResignKeyNotification, object: window) #expect(buttons.alphaValue == 0) - window.keyWindow = true + window.simulatedKeyWindow = true NotificationCenter.default.post(name: NSWindow.didBecomeKeyNotification, object: window) #expect(buttons.alphaValue == 1) _ = window From a18ca1410d06dde40c9775b22cab123ad8b3eabf Mon Sep 17 00:00:00 2001 From: austinpower1258 Date: Sun, 13 Sep 2026 22:10:43 -0700 Subject: [PATCH 15/17] test: convert Cloud hover fixture coordinates as a point --- cmuxTests/CloudTreeNativeDragOwnershipTests.swift | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/cmuxTests/CloudTreeNativeDragOwnershipTests.swift b/cmuxTests/CloudTreeNativeDragOwnershipTests.swift index 4b28d30007dc..572776fc2751 100644 --- a/cmuxTests/CloudTreeNativeDragOwnershipTests.swift +++ b/cmuxTests/CloudTreeNativeDragOwnershipTests.swift @@ -313,7 +313,7 @@ struct CloudTreeNativeDragOwnershipTests { let cell = try #require(outline.view(atColumn: 0, row: 0, makeIfNecessary: true) as? CloudTreeCellView) let buttons = try #require(cell.subviews.last) let rowPoint = NSPoint(x: outline.rect(ofRow: 0).midX, y: outline.rect(ofRow: 0).midY) - window.pointerOnScreen = window.convertToScreen(outline.convert(rowPoint, to: nil)) + window.pointerOnScreen = window.convertPoint(toScreen: outline.convert(rowPoint, to: nil)) NotificationCenter.default.post(name: NSWindow.didResignKeyNotification, object: window) #expect(buttons.alphaValue == 0) From 844a43a0a62d14d4accbb45947f6812c0c32021f Mon Sep 17 00:00:00 2001 From: Abdulaziz Albahar <67667005+azooz2003-bit@users.noreply.github.com> Date: Tue, 15 Sep 2026 10:51:22 -0700 Subject: [PATCH 16/17] fix: disable first-frame timeout for admitted Iroh --- Sources/Mobile/MobileHostService.swift | 10 ++++++++++ 1 file changed, 10 insertions(+) diff --git a/Sources/Mobile/MobileHostService.swift b/Sources/Mobile/MobileHostService.swift index 51219e6eb775..c337a12f69db 100644 --- a/Sources/Mobile/MobileHostService.swift +++ b/Sources/Mobile/MobileHostService.swift @@ -1394,9 +1394,19 @@ final class MobileHostService { } let id = UUID() + let firstFrameTimeoutNanoseconds: UInt64 = switch authorization { + case .irohAdmission: + // Iroh owns admission and native connection liveness. A delayed + // first control frame is valid while the admitted session is + // settling, so an application timer must not retire it. + 0 + case .stackBearer: + MobileHostConnection.defaultFirstFrameTimeoutNanoseconds + } let session = MobileHostConnection( id: id, transport: transport, + firstFrameTimeoutNanoseconds: firstFrameTimeoutNanoseconds, independentEventWriter: independentEventWriter, authorizeRequest: { request in await Self.connectionAuthorizationError( From 20535b6d1522964400b68ec499b829d33b769fa2 Mon Sep 17 00:00:00 2001 From: Abdulaziz Albahar <67667005+azooz2003-bit@users.noreply.github.com> Date: Tue, 15 Sep 2026 11:16:53 -0700 Subject: [PATCH 17/17] fix: keep closed control transports bound to their session --- .../IrxControlByteTransport.swift | 13 ++- .../IrxLiveQUICTests.swift | 108 ++++++++++++++++++ 2 files changed, 120 insertions(+), 1 deletion(-) diff --git a/Packages/Shared/CmuxIrxTransport/Sources/CmuxIrxTransport/IrxControlByteTransport.swift b/Packages/Shared/CmuxIrxTransport/Sources/CmuxIrxTransport/IrxControlByteTransport.swift index 7187294a4a54..4f80abcaec6e 100644 --- a/Packages/Shared/CmuxIrxTransport/Sources/CmuxIrxTransport/IrxControlByteTransport.swift +++ b/Packages/Shared/CmuxIrxTransport/Sources/CmuxIrxTransport/IrxControlByteTransport.swift @@ -11,6 +11,9 @@ public import Foundation /// session, so the transport retires that exact session before closing the /// complete QUIC connection. Session ownership belongs to /// ``IrxPeerEngine``. +/// +/// Each transport binds to at most one admitted session. Native closure makes +/// that transport terminal; the RPC owner creates a new transport to recover. public actor IrxControlByteTransport: CmxByteTransport { /// Factory for an admitted connection and its control lane. public typealias Establish = @Sendable () async throws -> (IrxConnection, IrxLaneStream) @@ -113,7 +116,15 @@ public actor IrxControlByteTransport: CmxByteTransport { private func establishedPair() async throws -> (IrxConnection, IrxLaneStream) { guard !isClosed else { throw IrxConnectionError.closed(nil) } - if let pair, await !pair.0.isConnectionClosed() { + if let pair { + let connectionIsClosed = await pair.0.isConnectionClosed() + guard !isClosed else { throw IrxConnectionError.closed(nil) } + if connectionIsClosed { + // Reads and writes may still be unwinding on this pair. Keep + // their eventual close tied to this RPC generation's session. + await close() + throw IrxConnectionError.closed(nil) + } return pair } if let connectInFlight { diff --git a/Packages/Shared/CmuxIrxTransport/Tests/CmuxIrxTransportTests/IrxLiveQUICTests.swift b/Packages/Shared/CmuxIrxTransport/Tests/CmuxIrxTransportTests/IrxLiveQUICTests.swift index 430bf5a3f586..d4f245fb17bb 100644 --- a/Packages/Shared/CmuxIrxTransport/Tests/CmuxIrxTransportTests/IrxLiveQUICTests.swift +++ b/Packages/Shared/CmuxIrxTransport/Tests/CmuxIrxTransportTests/IrxLiveQUICTests.swift @@ -201,6 +201,114 @@ struct IrxLiveQUICTests { try? await client.close() } + @Test( + "native closure cannot rebind an existing control transport", + .timeLimit(.minutes(1)), + arguments: [false, true] + ) + func nativeClosureCannotRebindControlTransport(sendAfterClosure: Bool) async throws { + let journal = IrxLiveTestSupport.journal() + let server = try await IrxLiveTestSupport.bindLoopback( + seed: IrxLiveTestSupport.identitySeed(), remoteBiCredit: 1) + let client = try await IrxLiveTestSupport.bindLoopback( + seed: IrxLiveTestSupport.identitySeed(), remoteBiCredit: 0) + defer { + Task { + try? await server.close() + try? await client.close() + } + } + let serverTask = Task { () throws -> [(IrxConnection, IrxLaneStream)] in + var pairs: [(IrxConnection, IrxLaneStream)] = [] + for _ in 0..<2 { + guard let incoming = await server.acceptNext() else { break } + let native = try await incoming.accept().connect() + let connection = IrxConnection( + connection: native, role: .acceptor, journal: journal) + guard let (_, control, _) = await IrxAdmission.performServer( + connection: connection, + judgment: IrxLiveTestSupport.fixedJudgment(accepting: "good-grant"), + journal: journal + ) else { break } + pairs.append((connection, control)) + } + return pairs + } + var clientPairs: [(IrxConnection, IrxLaneStream)] = [] + for _ in 0..<2 { + let native = try await client.connect( + addr: IrxLiveTestSupport.loopbackAddr(of: server), alpn: IrxProtocol.alpnData) + let connection = IrxConnection(connection: native, role: .dialer, journal: journal) + let (_, control) = try await IrxAdmission.performClient( + connection: connection, grantJWS: "good-grant", journal: journal) + clientPairs.append((connection, control)) + } + let serverPairs = try await serverTask.value + try #require(serverPairs.count == 2) + let first = clientPairs[0] + let replacement = clientPairs[1] + let establishments = AsyncStream.makeStream() + let releaseProbe = IrxControlReleaseProbe() + let transport = IrxControlByteTransport( + closeCode: .explicitRedial, + establish: { + establishments.continuation.yield(()) + return await first.0.isConnectionClosed() ? replacement : first + }, + onClose: { connection, code, retiresConnection in + #expect(connection.underlying.stableId() == first.0.underlying.stableId()) + await releaseProbe.record(closeCode: code, retiresConnection: retiresConnection) + } + ) + try await transport.connect() + try await transport.connect() + await serverPairs[0].0.close(code: .hostShutdown, origin: .local) + // Synchronize with Iroh's native closure without receiving on the old + // control lane, which would independently mark the transport closed. + _ = await first.0.underlying.closed() + #expect(await transport.isTransportClosed()) + + do { + if sendAfterClosure { + try await transport.send(Data("old RPC generation".utf8)) + } else { + try await transport.connect() + } + Issue.record("closed control transport adopted a replacement session") + } catch let error as IrxConnectionError { + guard case .closed = error else { + Issue.record("unexpected connection error after native closure: \(error)") + return + } + } catch { + Issue.record("unexpected error after native closure: \(error)") + } + + // A delayed close from this RPC generation must only release its + // original claim; it must never reach the replacement connection. + await transport.close() + #expect(await releaseProbe.count == 1) + #expect(await releaseProbe.retiresConnections == [false]) + establishments.continuation.finish() + var establishmentCount = 0 + for await _ in establishments.stream { establishmentCount += 1 } + #expect(establishmentCount == 1) + #expect(await !replacement.0.isConnectionClosed()) + + let replacementTransport = IrxControlByteTransport( + connection: replacement.0, control: replacement.1, closeCode: .explicitRedial) + try await replacementTransport.connect() + let message = Data("new RPC generation".utf8) + let received = try await withIrxDeadline(.seconds(1), onTimeout: { + await serverPairs[1].1.reader.stop() + }) { + try await replacementTransport.send(message) + return try await serverPairs[1].1.reader.readRaw() + } + #expect(received == message) + await replacementTransport.close() + } + @Test("cancelling a control read retires the owner locally") func cancelledControlReadRetiresLocally() async throws { let journal = IrxLiveTestSupport.journal()