From 650d28f1bd794c7c10f3f90e076df085b7915a8c Mon Sep 17 00:00:00 2001 From: lawrencecchen <54008264+lawrencecchen@users.noreply.github.com> Date: Wed, 10 Jun 2026 18:24:01 -0700 Subject: [PATCH 01/10] iOS render-grid liveness: failing tests for the watchdog false-fire The Release-sim bisect (2026-06-10) caught the liveness watchdog firing "render-grid stream silent for 10499ms, re-subscribing" every ~10.5s forever, plus "subscribe failed reason=start: requestTimedOut", while the Mac kept delivering events. Each false fire makes the Mac re-replay the full grid (constant repaint, battery and bandwidth waste). Two reproduced defects, expressed through the real RPC transport and the real consumer path against a scripted host: 1. renderGridEventsArrivingDuringStartSubscribeAreConsumed: events the transport delivers while the mobile.events.subscribe ack is still in flight pile up unconsumed in the subscription stream buffer, because the listener task awaits the ack before starting the for-await consumer loop that stamps the liveness clock. A healthy establishing stream therefore looks silent to the watchdog, and the watchdog's resync then cancels its own in-flight subscribe, which surfaces as requestTimedOut. 2. watchdogDoesNotResubscribeHealthyIdleStream: a healthy idle terminal legitimately emits zero events (the Mac dedupes render-grid frames by row signature and stateSeq), so wall-clock silence alone cannot distinguish idle from dead. The watchdog tears down and full-replays a perfectly healthy subscription every threshold window. watchdogStillResubscribesGenuinelyDeadStream pins the watchdog's original purpose (the ~85s silent-death hang) so the fix cannot regress it. Includes a DEBUG-only seam to run one watchdog evaluation deterministically with the injected clock. Co-Authored-By: Claude Fable 5 --- .../MobileShellComposite.swift | 11 + .../MobileShellRenderGridLivenessTests.swift | 492 ++++++++++++++++++ 2 files changed, 503 insertions(+) create mode 100644 Packages/CmuxMobileShell/Tests/CmuxMobileShellTests/MobileShellRenderGridLivenessTests.swift diff --git a/Packages/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite.swift b/Packages/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite.swift index c5518439f75b..0ac901944b6d 100644 --- a/Packages/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite.swift +++ b/Packages/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite.swift @@ -3884,6 +3884,17 @@ public final class MobileShellComposite: MobileTerminalOutputSinking { renderGridLivenessListenerID = nil } + #if DEBUG + /// Test-only: run one liveness evaluation for the currently armed watchdog + /// generation, exactly as a `DispatchSourceTimer` tick would. Lets package + /// tests drive the silence check deterministically against an injected + /// clock instead of waiting on the wall-clock tick cadence. + func debugRunRenderGridLivenessCheckForTesting() { + guard let listenerID = renderGridLivenessListenerID else { return } + checkRenderGridLiveness(listenerID: listenerID) + } + #endif + /// One watchdog tick on the main actor: if the subscription generation still /// matches, the store is connected, and the stream has been silent past the /// threshold, tear down + re-subscribe + replay via the existing resync path. diff --git a/Packages/CmuxMobileShell/Tests/CmuxMobileShellTests/MobileShellRenderGridLivenessTests.swift b/Packages/CmuxMobileShell/Tests/CmuxMobileShellTests/MobileShellRenderGridLivenessTests.swift new file mode 100644 index 000000000000..d4038727f6e8 --- /dev/null +++ b/Packages/CmuxMobileShell/Tests/CmuxMobileShellTests/MobileShellRenderGridLivenessTests.swift @@ -0,0 +1,492 @@ +import CMUXMobileCore +import CmuxMobileRPC +import Foundation +import Testing +@testable import CmuxMobileShell + +// Regression coverage for the render-grid liveness watchdog false-fire +// (Release-sim bisect, 2026-06-10): the phone logged "render-grid stream +// silent for 10499ms, re-subscribing" every ~10.5s plus "subscribe failed +// reason=start: requestTimedOut" while the Mac demonstrably kept the +// connection healthy. Two defects combined: +// +// 1. The liveness clock was stamped only inside the listener's `for await` +// consumer loop, which did not start until the `mobile.events.subscribe` +// ack round-trip completed. Events yielded into the subscription stream +// during that window were buffered invisibly, so the watchdog read a +// healthy establishing stream as silence (and its resync then CANCELLED +// the in-flight subscribe, which surfaces as `requestTimedOut`). +// 2. A healthy idle terminal legitimately pushes no events at all (the Mac +// dedupes render-grid emits by row signature + stateSeq), so wall-clock +// silence alone can never distinguish "idle" from "dead". The watchdog +// needs a bounded host probe before it may declare death. + +// MARK: - Injected clock + +private final class TestClock: @unchecked Sendable { + private let lock = NSLock() + private var current: Date + + init() { + current = Date() + } + + var now: Date { + lock.withLock { current } + } + + func advance(by interval: TimeInterval) { + lock.withLock { current = current.addingTimeInterval(interval) } + } +} + +// MARK: - Runtime double + +private struct LivenessTestRuntime: MobileSyncRuntime { + var transportFactory: any CmxByteTransportFactory + var stackAccessTokenProvider: @Sendable () async throws -> String = { "test-stack-token" } + var stackAccessTokenForceRefresher: @Sendable () async throws -> String = { "test-stack-token" } + var rpcRequestTimeoutNanoseconds: UInt64 = 30 * 1_000_000_000 + var now: @Sendable () -> Date + var supportedRouteKinds: [CmxAttachTransportKind] = [.debugLoopback] + var pairingRequestTimeoutNanoseconds: UInt64 = 30 * 1_000_000_000 + var supportsServerPushEvents: Bool = true + /// Bounded deadline for the watchdog's host liveness probe. Short here so + /// the dead-stream test does not wait the production default. + var livenessProbeTimeoutNanoseconds: UInt64 = 200_000_000 +} + +// MARK: - Scripted host (router + transport) + +/// Scripts the Mac side of the persistent RPC connection: answers the +/// connect-time `workspace.list`, the `mobile.host.status` capability and +/// probe requests, `mobile.events.subscribe`, and replay/viewport calls. +/// Individual requests can be held unresolved to model an ack that has not +/// arrived yet (establishment window) or a host that stopped answering +/// (dead stream). +private actor LivenessHostRouter { + struct RecordedRequest: Sendable { + var method: String? + var topics: [String]? + } + + private var recorded: [RecordedRequest] = [] + private var hostStatusRequestCount = 0 + private var heldHostStatusRequestNumbers: Set = [] + private var holdSubscribe = false + private var heldContinuations: [CheckedContinuation] = [] + + func record(method: String?, topics: [String]?) { + recorded.append(RecordedRequest(method: method, topics: topics)) + } + + func count(of method: String) -> Int { + recorded.filter { $0.method == method }.count + } + + /// Hold every `mobile.events.subscribe` response until released. + func setHoldSubscribe(_ hold: Bool) { + holdSubscribe = hold + } + + /// Hold the Nth `mobile.host.status` request (1-based) forever, modeling + /// a host that stopped answering on a half-dead transport. + func holdHostStatusRequest(number: Int) { + heldHostStatusRequestNumbers.insert(number) + } + + /// Resume every held request so parked continuations do not leak past the + /// end of the test. + func releaseAllHeld() { + holdSubscribe = false + heldHostStatusRequestNumbers = [] + let continuations = heldContinuations + heldContinuations = [] + for continuation in continuations { + continuation.resume() + } + } + + func response(method: String?, id: String?) async -> Data? { + switch method { + case "workspace.list", "mobile.workspace.list": + return try? Self.resultFrame(id: id, result: [ + "workspaces": [ + [ + "id": "live-workspace", + "title": "Live Workspace", + "current_directory": "/Users/test/project", + "is_selected": true, + "terminals": [ + [ + "id": "live-terminal", + "title": "Terminal", + "current_directory": "/Users/test/project", + "is_ready": true, + "is_focused": true, + ], + ], + ], + ], + ]) + case "mobile.host.status": + hostStatusRequestCount += 1 + if heldHostStatusRequestNumbers.contains(hostStatusRequestCount) { + await park() + return nil + } + return try? Self.resultFrame(id: id, result: [ + "terminal_fidelity": "render_grid", + "capabilities": ["events.v1", "terminal.render_grid.v1", "terminal.replay.v1"], + ]) + case "mobile.events.subscribe": + if holdSubscribe { + await park() + return nil + } + return try? Self.resultFrame(id: id, result: [ + "stream_id": "test-stream", + "topics": ["workspace.updated", "terminal.render_grid"], + ]) + case "mobile.events.unsubscribe", "mobile.terminal.replay", "mobile.terminal.viewport": + return try? Self.resultFrame(id: id, result: [:]) + default: + return try? Self.errorFrame(id: id, message: "Unexpected method \(method ?? "nil")") + } + } + + private func park() async { + await withCheckedContinuation { continuation in + heldContinuations.append(continuation) + } + } + + private static func resultFrame(id: String?, result: [String: Any]) throws -> Data { + let envelope: [String: Any] = [ + "id": id ?? UUID().uuidString, + "ok": true, + "result": result, + ] + return try MobileSyncFrameCodec.encodeFrame(JSONSerialization.data(withJSONObject: envelope)) + } + + private static func errorFrame(id: String?, message: String) throws -> Data { + let envelope: [String: Any] = [ + "id": id ?? UUID().uuidString, + "ok": false, + "error": ["message": message], + ] + return try MobileSyncFrameCodec.encodeFrame(JSONSerialization.data(withJSONObject: envelope)) + } +} + +/// Holds the live transport instance so the test can push unsolicited +/// server-side event frames through the same receive path production uses. +private final class TransportBox: @unchecked Sendable { + private let lock = NSLock() + private var transport: LivenessTransport? + + func set(_ transport: LivenessTransport) { + lock.withLock { self.transport = transport } + } + + func get() -> LivenessTransport? { + lock.withLock { transport } + } +} + +private struct LivenessTransportFactory: CmxByteTransportFactory { + let router: LivenessHostRouter + let box: TransportBox + + func makeTransport(for route: CmxAttachRoute) throws -> any CmxByteTransport { + let transport = LivenessTransport(router: router) + box.set(transport) + return transport + } +} + +private actor LivenessTransport: CmxByteTransport { + private let router: LivenessHostRouter + private var pendingFrames: [Data] = [] + private var receiveWaiters: [CheckedContinuation] = [] + private var isClosed = false + + init(router: LivenessHostRouter) { + self.router = router + } + + func connect() async throws {} + + func receive() async throws -> Data? { + if !pendingFrames.isEmpty { + return pendingFrames.removeFirst() + } + if isClosed { + return nil + } + return await withCheckedContinuation { continuation in + receiveWaiters.append(continuation) + } + } + + func send(_ data: Data) async throws { + var buffer = data + let payloads = try MobileSyncFrameCodec.decodeFrames(from: &buffer) + for payload in payloads { + let parsed = (try? JSONSerialization.jsonObject(with: payload)) as? [String: Any] + let method = parsed?["method"] as? String + let id = parsed?["id"] as? String + let topics = (parsed?["params"] as? [String: Any])?["topics"] as? [String] + await router.record(method: method, topics: topics) + // Answer each request concurrently so one held response cannot + // head-of-line block later RPCs, matching the Mac host's + // per-frame response tasks. + Task { [router, weak self] in + guard let response = await router.response(method: method, id: id) else { + return + } + await self?.deliver(response) + } + } + } + + func close() async { + isClosed = true + let waiters = receiveWaiters + receiveWaiters = [] + for waiter in waiters { + waiter.resume(returning: nil) + } + } + + /// Deliver a frame to the client's read loop. Also used by tests to push + /// unsolicited server-side event envelopes. + func deliver(_ frame: Data) { + if receiveWaiters.isEmpty { + pendingFrames.append(frame) + return + } + let waiter = receiveWaiters.removeFirst() + waiter.resume(returning: frame) + } +} + +// MARK: - Test helpers + +@MainActor +private final class OutputCollector { + private(set) var lines: [String] = [] + private var task: Task? + + func mount(store: MobileShellComposite, surfaceID: String) { + task = Task { @MainActor [weak self] in + for await data in store.terminalOutputStream(surfaceID: surfaceID) { + self?.lines.append(String(decoding: data, as: UTF8.self)) + } + } + } + + func unmount() { + task?.cancel() + task = nil + } +} + +private func makeTicket(clock: TestClock) throws -> CmxAttachTicket { + let route = try CmxAttachRoute( + id: "debug_loopback", + kind: .debugLoopback, + endpoint: .hostPort(host: "127.0.0.1", port: 56584) + ) + return try CmxAttachTicket( + workspaceID: "live-workspace", + terminalID: "live-terminal", + macDeviceID: "test-mac", + macDisplayName: "Test Mac", + routes: [route], + expiresAt: clock.now.addingTimeInterval(3600) + ) +} + +private func attachURL(for ticket: CmxAttachTicket) throws -> String { + let encoder = JSONEncoder() + encoder.dateEncodingStrategy = .iso8601 + let payload = try encoder.encode(ticket) + .base64EncodedString() + .replacingOccurrences(of: "+", with: "-") + .replacingOccurrences(of: "/", with: "_") + .replacingOccurrences(of: "=", with: "") + return "cmux-ios://attach?v=\(ticket.version)&payload=\(payload)" +} + +private func renderGridEventFrame(surfaceID: String, seq: UInt64, text: String) throws -> Data { + let frame = try MobileTerminalRenderGridFrame.fromPlainRows( + surfaceID: surfaceID, + stateSeq: seq, + columns: 16, + rows: 4, + text: text + ) + let envelope: [String: Any] = [ + "kind": "event", + "topic": "terminal.render_grid", + "payload": try frame.jsonObject(), + ] + return try MobileSyncFrameCodec.encodeFrame(JSONSerialization.data(withJSONObject: envelope)) +} + +/// Poll until `condition` is true, bounded at `attempts` x 10ms. Returns the +/// final value so tests can assert both presence and (bounded) absence. +@MainActor +private func pollUntil( + attempts: Int = 300, + _ condition: @MainActor () async -> Bool +) async throws -> Bool { + for _ in 0.. MobileShellComposite { + let runtime = LivenessTestRuntime( + transportFactory: LivenessTransportFactory(router: router, box: box), + now: { clock.now }, + livenessProbeTimeoutNanoseconds: probeTimeoutNanoseconds + ) + let store = MobileShellComposite.preview(runtime: runtime) + store.signIn() + let ticket = try makeTicket(clock: clock) + let connected = await store.connectPairingURL(try attachURL(for: ticket)) + #expect(connected, "scripted connect must succeed") + return store +} + +// MARK: - Tests + +/// The decoupling found by the bisect: events that the transport delivers +/// while the `mobile.events.subscribe` ack is still in flight must reach the +/// real consumer (and therefore the liveness clock), not pile up unconsumed +/// in the subscription stream's buffer behind the ack await. +@MainActor +@Test func renderGridEventsArrivingDuringStartSubscribeAreConsumed() async throws { + let clock = TestClock() + let router = LivenessHostRouter() + let box = TransportBox() + await router.setHoldSubscribe(true) + let store = try await makeConnectedStore(router: router, box: box, clock: clock) + defer { + Task { await router.releaseAllHeld() } + } + + // The listener has sent its start subscribe; the ack is parked. + let sawSubscribe = try await pollUntil { await router.count(of: "mobile.events.subscribe") >= 1 } + #expect(sawSubscribe, "listener must request the server-side subscription") + + let collector = OutputCollector() + collector.mount(store: store, surfaceID: "live-terminal") + let sawReplay = try await pollUntil { await router.count(of: "mobile.terminal.replay") >= 1 } + #expect(sawReplay, "mounting a sink must arm the cold-attach replay") + + // The Mac pushes a live render-grid event while the subscribe ack is + // still pending (the server-side subscription from a previous generation + // keeps pushing across re-subscribes; the ack is an enable handshake, + // not a delivery precondition). + let event = try renderGridEventFrame(surfaceID: "live-terminal", seq: 5, text: "live") + let transport = try #require(box.get()) + await transport.deliver(event) + + let delivered = try await pollUntil { collector.lines.isEmpty == false } + #expect( + delivered, + "render-grid events must be consumed while the start-subscribe ack is in flight; buffering them unconsumed is what made a healthy stream look silent to the liveness watchdog" + ) + #expect(collector.lines.first?.contains("live") == true) + + await router.releaseAllHeld() + collector.unmount() +} + +/// A healthy idle stream produces zero events (the Mac dedupes unchanged +/// frames), so silence alone must not tear the subscription down. The +/// watchdog must verify with a host probe and stay quiet when the host +/// answers. Without this, the phone re-subscribed and full-grid re-replayed +/// every ~10.5s forever on any idle terminal. +@MainActor +@Test func watchdogDoesNotResubscribeHealthyIdleStream() async throws { + let clock = TestClock() + let router = LivenessHostRouter() + let box = TransportBox() + let store = try await makeConnectedStore(router: router, box: box, clock: clock) + + let sawSubscribe = try await pollUntil { await router.count(of: "mobile.events.subscribe") >= 1 } + #expect(sawSubscribe, "listener must establish the push subscription") + + // Idle past the silence threshold: no events at all, host healthy. + clock.advance(by: 10) + store.debugRunRenderGridLivenessCheckForTesting() + + // The watchdog may probe the host (mobile.host.status request number 2; + // number 1 was the listener's transport-capability resolve), and the + // router answers it. What it must NOT do is restart the subscription. + let resubscribed = try await pollUntil(attempts: 60) { + await router.count(of: "mobile.events.subscribe") >= 2 + } + #expect( + resubscribed == false, + "the watchdog must not re-subscribe a healthy idle stream; the host answered, so silence only means the terminal had nothing to say" + ) + + // The probe outcome must reset the silence window: an immediate second + // evaluation stays quiet too. + store.debugRunRenderGridLivenessCheckForTesting() + let resubscribedAfterRecheck = try await pollUntil(attempts: 30) { + await router.count(of: "mobile.events.subscribe") >= 2 + } + #expect(resubscribedAfterRecheck == false) + let replayCount = await router.count(of: "mobile.terminal.replay") + #expect(replayCount == 0, "no sink is mounted and the stream is healthy, so no replay traffic may be generated") +} + +/// The watchdog's original purpose (the ~85s silent-death hang) must keep +/// working: silence past the threshold plus a host that stops answering the +/// probe must still tear down and re-subscribe. +@MainActor +@Test func watchdogStillResubscribesGenuinelyDeadStream() async throws { + let clock = TestClock() + let router = LivenessHostRouter() + let box = TransportBox() + let store = try await makeConnectedStore(router: router, box: box, clock: clock) + defer { + Task { await router.releaseAllHeld() } + } + + let sawSubscribe = try await pollUntil { await router.count(of: "mobile.events.subscribe") >= 1 } + #expect(sawSubscribe, "listener must establish the push subscription") + + // The host stops answering the next host.status (the watchdog's probe on + // current code's resync path resolves capabilities first, so holding + // request number 2 models a dead host in both worlds). + await router.holdHostStatusRequest(number: 2) + clock.advance(by: 10) + store.debugRunRenderGridLivenessCheckForTesting() + + let resubscribed = try await pollUntil(attempts: 600) { + await router.count(of: "mobile.events.subscribe") >= 2 + } + #expect( + resubscribed, + "a stream that is silent past the threshold AND whose host stops answering must still be torn down and re-subscribed" + ) + await router.releaseAllHeld() +} From 06d9cbdf7adf80c0a1d03fa8d672ccaa116be39c Mon Sep 17 00:00:00 2001 From: lawrencecchen <54008264+lawrencecchen@users.noreply.github.com> Date: Wed, 10 Jun 2026 18:34:04 -0700 Subject: [PATCH 02/10] iOS render-grid liveness: probe before teardown, consume during subscribe ack Fixes the watchdog false-fire loop the 2026-06-10 Release-sim bisect caught: "render-grid stream silent for 10499ms, re-subscribing" every ~10.5s forever on a healthy connection, plus "subscribe failed reason=start: requestTimedOut", with the Mac re-replaying the full grid each cycle (constant repaint, battery and bandwidth waste). Root cause, two halves: 1. Silence is not death. A healthy idle terminal pushes zero events (MobileTerminalRenderObserver dedupes unchanged frames by row signature and stateSeq, and cursor blink changes neither), so the watchdog's wall-clock silence check fired on every idle stream once per threshold window, forever. The watchdog had no positive liveness signal at all. 2. The liveness clock was stamped only inside the listener's for-await loop, which did not start until resolveTerminalOutputTransport and the mobile.events.subscribe ack both completed. Events delivered during that establishment window were buffered unconsumed in the subscription stream, invisible to the clock; each watchdog fire then cancelled its own in-flight start subscribe, which MobileCoreRPCSession.cancelPendingRequest surfaces as requestTimedOut (the real wire timeout is 30s, so the 10.5s cadence of that log was the watchdog, not the network), and the dying generation marked the connection unavailable underneath the fresh one. Fix, structurally: - One liveness ownership point: recordTerminalEventStreamLiveness() is stamped by every consumed envelope, by a successful host probe, and (as the generation reset) when a watchdog generation is armed; the watchdog reads the same record, and generation identity (listenerID) guards every transition. - The silence threshold is now a suspicion, not a verdict: a crossing runs a bounded auth-exempt mobile.host.status probe over the same transport the events ride on. Probe answered = stream healthy but quiet, stamp and stay quiet. Probe failed = re-check silence (events may have resumed mid-probe), then run the existing teardown + re-subscribe + replay recovery. The ~85s silent-death case still heals: dead transport fails the probe at livenessProbeTimeoutNanoseconds (default 3s, runtime-injectable). - The start subscribe ack now runs concurrently with consumption (beginTerminalEventSubscriptionStart): the ack is a server-side enable handshake, not a delivery precondition, so the consumer loop starts immediately and the clock stays coupled to actual event arrival. Success/failure is acted on only while the generation is current, so a superseded ack can no longer mark the connection unavailable, and a cancelled ack logs as cancelled instead of a fake wire timeout. Co-Authored-By: Claude Fable 5 --- .../CmuxMobileRPC/MobileSyncRuntime.swift | 14 ++ .../MobileShellComposite.swift | 218 +++++++++++++++--- 2 files changed, 199 insertions(+), 33 deletions(-) diff --git a/Packages/CmuxMobileRPC/Sources/CmuxMobileRPC/MobileSyncRuntime.swift b/Packages/CmuxMobileRPC/Sources/CmuxMobileRPC/MobileSyncRuntime.swift index 0044b9775541..7c14dbefc097 100644 --- a/Packages/CmuxMobileRPC/Sources/CmuxMobileRPC/MobileSyncRuntime.swift +++ b/Packages/CmuxMobileRPC/Sources/CmuxMobileRPC/MobileSyncRuntime.swift @@ -30,4 +30,18 @@ public protocol MobileSyncRuntime: Sendable { /// skips background subscribe/poll so scripted-transport tests do not /// consume responses intended for foreground methods. var supportsServerPushEvents: Bool { get } + /// Bounded deadline, in nanoseconds, for the render-grid liveness + /// watchdog's host probe (`mobile.host.status`). A healthy idle terminal + /// legitimately pushes no events, so the watchdog verifies prolonged + /// silence with this probe before declaring the stream dead; the deadline + /// bounds how long a dead transport can stall that verdict. + var livenessProbeTimeoutNanoseconds: UInt64 { get } +} + +public extension MobileSyncRuntime { + /// Default probe deadline: generous against a momentarily loaded Mac (the + /// status request is auth-exempt and answered without side effects), while + /// keeping dead-stream recovery within a few seconds of the silence + /// threshold instead of the full ``rpcRequestTimeoutNanoseconds``. + var livenessProbeTimeoutNanoseconds: UInt64 { 3_000_000_000 } } diff --git a/Packages/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite.swift b/Packages/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite.swift index 0ac901944b6d..e166e94e0ba9 100644 --- a/Packages/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite.swift +++ b/Packages/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite.swift @@ -63,11 +63,13 @@ public final class MobileShellComposite: MobileTerminalOutputSinking { private static let terminalOutputCapabilityTimeoutNanoseconds: UInt64 = 750_000_000 /// How long the render-grid stream may stay silent (no event of any topic) - /// before the liveness watchdog assumes the push subscription is dead and - /// forces a re-subscribe + replay. Picked at the low end of the acceptable - /// 8-12s window so a wedged stream recovers in a few seconds instead of the - /// transport's ~85s timeout, while staying well above any normal inter-event - /// gap on a busy shell. + /// before the liveness watchdog suspects the push subscription is dead and + /// runs a bounded host probe; only a failed probe forces the + /// re-subscribe + replay (silence alone is the normal state of an idle + /// terminal). Picked at the low end of the acceptable 8-12s window so a + /// wedged stream recovers in a few seconds instead of the transport's ~85s + /// timeout, while staying well above any normal inter-event gap on a busy + /// shell. private static let renderGridLivenessSilenceThreshold: TimeInterval = 9 /// Cadence of the liveness watchdog tick. It only reads a timestamp and /// compares against the threshold, so a short interval is cheap; it does not @@ -447,6 +449,13 @@ public final class MobileShellComposite: MobileTerminalOutputSinking { /// Tail of the serialized paired-Mac store write chain; see /// ``performSerializedPairedMacWrite(ifStillCurrent:_:)``. private var pairedMacWriteChain: Task? + /// The in-flight `mobile.events.subscribe` (reason `start`) ack for the + /// current listener generation. It runs concurrently with the consumer + /// loop (the ack is a server-side enable handshake, not a delivery + /// precondition: a prior generation's server subscription keeps pushing + /// across re-subscribes) so events arriving during the round-trip are + /// consumed, not buffered invisibly behind the await. + private var terminalSubscriptionStartTask: Task? // Liveness watchdog for the render-grid push subscription. The `for await` // listener loop blocks indefinitely if the underlying connection half-dies // (network blip, Mac stops pushing, background/foreground cycle): the @@ -455,9 +464,16 @@ public final class MobileShellComposite: MobileTerminalOutputSinking { // of render-grid deltas. The transport's own timeout (~85s) is far too slow. // A `DispatchSourceTimer` ticks independently of the (potentially wedged) // stream and compares "now" against the last received event to detect - // prolonged silence, then tears down + re-subscribes + replays. + // prolonged silence. Silence alone is NOT death: a healthy idle terminal + // pushes nothing (the Mac dedupes unchanged render-grid frames), so a + // silence-threshold crossing first runs a bounded `mobile.host.status` + // probe and only tears down + re-subscribes + replays when the host fails + // to answer it. private var renderGridLivenessTimer: (any DispatchSourceTimer)? private var renderGridLivenessListenerID: UUID? + /// The in-flight liveness probe spawned by a silence-threshold crossing. + /// Single-flight: ticks while a probe is pending are no-ops. + private var renderGridLivenessProbeTask: Task? private var lastTerminalEventAt: Date? private var terminalSubscriptionRefreshTask: Task? private var createWorkspaceTask: Task? @@ -597,7 +613,9 @@ public final class MobileShellComposite: MobileTerminalOutputSinking { isolated deinit { networkPathObservationTask?.cancel() terminalEventListenerTask?.cancel() + terminalSubscriptionStartTask?.cancel() renderGridLivenessTimer?.cancel() + renderGridLivenessProbeTask?.cancel() terminalSubscriptionRefreshTask?.cancel() createWorkspaceTask?.cancel() createTerminalTask?.cancel() @@ -3651,6 +3669,14 @@ public final class MobileShellComposite: MobileTerminalOutputSinking { do { responseData = try await client.sendRequest(requestData) } catch { + if Task.isCancelled { + // A superseding generation (resync, disconnect) cancelled this + // request; the session layer surfaces that cancellation as + // `requestTimedOut`. Not a host failure: stay quiet so the log + // does not report a self-inflicted cancel as a wire timeout. + mobileShellLog.info("subscribe cancelled reason=\(reason, privacy: .public)") + return false + } mobileShellLog.error("subscribe failed reason=\(reason, privacy: .public): \(String(describing: error), privacy: .private)") // Event-stream (re)subscribe is the view-only/foreground-resume path. // A definitive auth failure here (RPC layer already tried a @@ -3773,19 +3799,21 @@ public final class MobileShellComposite: MobileTerminalOutputSinking { let outputTransport = await self?.resolveTerminalOutputTransport(client: client) ?? .rawBytes let topics = outputTransport.eventTopics let stream = await client.subscribe(to: Set(topics)) - let subscribed = await self?.requestTerminalEventSubscription( + // Kick off the server-side enable handshake CONCURRENTLY with + // consumption. The old structure awaited the ack here, which + // parked the consumer loop while events from a still-active prior + // server subscription piled up unconsumed in `stream`'s buffer; + // the liveness watchdog (stamped only at consumption) then read a + // healthy establishing stream as silence and false-fired, and its + // resync cancelled this very ack (surfacing a bogus + // `requestTimedOut`). Consuming from the start keeps the liveness + // clock coupled to actual event arrival. + self?.beginTerminalEventSubscriptionStart( client: client, - reason: "start", - topics: topics - ) ?? false - guard subscribed else { - MobileDebugLog.anchormux("sync.subscribe_failed reason=start") - self?.diagnosticLog?.record(DiagnosticEvent(.error)) - self?.markMacConnectionUnavailable() - return - } - self?.markMacConnectionHealthy() - MobileDebugLog.anchormux("sync.subscribe_ok topics=\(topics.count) transport=\(outputTransport)") + listenerID: listenerID, + topics: topics, + transport: outputTransport + ) // Keep the listener alive without keeping the shell store alive. for await event in stream { guard !Task.isCancelled else { return } @@ -3793,7 +3821,7 @@ public final class MobileShellComposite: MobileTerminalOutputSinking { guard self.remoteClient === client, self.connectionState == .connected else { return } // Any yielded envelope proves the transport is still pushing, so // it resets the liveness window (not just render_grid events). - self.lastTerminalEventAt = self.runtime?.now() ?? Date() + self.recordTerminalEventStreamLiveness() self.markMacConnectionHealthy() if event.topic == "workspace.updated" { self.scheduleWorkspaceListRefreshFromEvent() @@ -3811,6 +3839,43 @@ public final class MobileShellComposite: MobileTerminalOutputSinking { } } + /// Run the `mobile.events.subscribe` (reason `start`) handshake for one + /// listener generation, concurrently with that generation's consumer loop. + /// + /// Success and failure are only acted on while the generation is still + /// current: a superseded or cancelled handshake exits silently so a stale + /// generation can never mark the connection unavailable underneath a + /// fresh, healthy one (the bisected false-fire loop did exactly that via + /// its self-cancelled ack). + private func beginTerminalEventSubscriptionStart( + client: MobileCoreRPCClient, + listenerID: UUID, + topics: [String], + transport: TerminalOutputTransport + ) { + guard terminalEventListenerID == listenerID else { return } + terminalSubscriptionStartTask?.cancel() + terminalSubscriptionStartTask = Task { @MainActor [weak self] in + let subscribed = await self?.requestTerminalEventSubscription( + client: client, + reason: "start", + topics: topics + ) ?? false + guard let self else { return } + guard !Task.isCancelled, self.terminalEventListenerID == listenerID else { return } + self.terminalSubscriptionStartTask = nil + guard subscribed else { + MobileDebugLog.anchormux("sync.subscribe_failed reason=start") + self.diagnosticLog?.record(DiagnosticEvent(.error)) + self.stopTerminalRefreshPolling() + self.markMacConnectionUnavailable() + return + } + self.markMacConnectionHealthy() + MobileDebugLog.anchormux("sync.subscribe_ok topics=\(topics.count) transport=\(transport)") + } + } + private func handleTerminalEventStreamEnded(listenerID: UUID, client: MobileCoreRPCClient) { guard !Task.isCancelled, terminalEventListenerID == listenerID, @@ -3838,14 +3903,16 @@ public final class MobileShellComposite: MobileTerminalOutputSinking { /// ticks independently and, on each tick, hops to the main actor to compare /// `lastTerminalEventAt` against `renderGridLivenessSilenceThreshold`. While /// events keep arriving, `lastTerminalEventAt` stays fresh and every tick is a - /// no-op, so an actively-streaming connection never triggers recovery; only a - /// genuinely silent stream crosses the threshold. + /// no-op. A threshold crossing is treated as a SUSPICION, not a verdict: an + /// idle terminal pushes no events, so the tick first probes the host with a + /// bounded `mobile.host.status` round-trip and only recovers when the probe + /// fails (see ``checkRenderGridLiveness(listenerID:)``). private func startRenderGridLivenessWatchdog(listenerID: UUID) { stopRenderGridLivenessWatchdog(listenerID: nil) renderGridLivenessListenerID = listenerID // Reset the window so a freshly-armed subscription gets the full silence // budget before it can be judged dead. - lastTerminalEventAt = runtime?.now() ?? Date() + recordTerminalEventStreamLiveness() // DispatchSourceTimer is the allowed low-level primitive for periodic // event delivery. It fires on the MAIN queue on purpose: the handler is // inferred @MainActor (it touches main-actor store state), and a timer on @@ -3882,6 +3949,21 @@ public final class MobileShellComposite: MobileTerminalOutputSinking { renderGridLivenessTimer?.cancel() renderGridLivenessTimer = nil renderGridLivenessListenerID = nil + renderGridLivenessProbeTask?.cancel() + renderGridLivenessProbeTask = nil + } + + /// Single ownership point for the liveness clock the watchdog reads. + /// + /// Stamped by (1) every envelope the listener loop actually consumes, + /// (2) a successful host probe (positive proof the channel is alive while + /// the terminal is merely idle), and (3) the arming of a new watchdog + /// generation, as the clean generation reset. The watchdog compares this + /// single record against `renderGridLivenessSilenceThreshold`. The only + /// other write is `resetTerminalOutputTracking` clearing it to nil when + /// the connection context is torn down entirely. + private func recordTerminalEventStreamLiveness() { + lastTerminalEventAt = runtime?.now() ?? Date() } #if DEBUG @@ -3897,24 +3979,92 @@ public final class MobileShellComposite: MobileTerminalOutputSinking { /// One watchdog tick on the main actor: if the subscription generation still /// matches, the store is connected, and the stream has been silent past the - /// threshold, tear down + re-subscribe + replay via the existing resync path. + /// threshold, verify the silence with a bounded host probe and only tear + /// down + re-subscribe + replay (via the existing resync path) when the + /// probe fails. + /// + /// The probe step exists because silence is ambiguous: a healthy idle + /// terminal emits nothing (the Mac dedupes unchanged render-grid frames by + /// row signature and stateSeq), which is indistinguishable by wall clock + /// from the half-dead transport this watchdog was built to catch. Treating + /// silence alone as death made the watchdog tear down and full-grid-replay + /// every healthy idle subscription every ~10.5s, forever (the 2026-06-10 + /// Release-sim bisect finding). `mobile.host.status` is the cheapest + /// possible probe: auth-exempt (no token mint), side-effect free on the + /// host, and answered over the same transport the events ride on, so a + /// completed round-trip is positive proof the channel is alive. private func checkRenderGridLiveness(listenerID: UUID) { guard renderGridLivenessListenerID == listenerID else { return } - guard remoteClient != nil, connectionState == .connected else { return } + guard let client = remoteClient, connectionState == .connected else { return } guard terminalEventListenerID == listenerID else { return } let now = runtime?.now() ?? Date() let last = lastTerminalEventAt ?? now let silent = now.timeIntervalSince(last) guard silent >= Self.renderGridLivenessSilenceThreshold else { return } - let silentMs = Int(silent * 1000) - MobileDebugLog.anchormux("sync.liveness re-subscribe silentMs=\(silentMs)") - diagnosticLog?.record(DiagnosticEvent(.livenessResubscribe, ms: UInt32(clamping: silentMs))) - mobileShellLog.info("render-grid stream silent for \(silentMs, privacy: .public)ms, re-subscribing") - // resyncTerminalOutput(restartEventStream: true) stops the wedged listener - // (which cancels this watchdog via stopTerminalRefreshPolling) and starts a - // fresh subscription + watchdog, then replays every surface so the phone - // catches up on the deltas it missed while the stream was silent. - resyncTerminalOutput(reason: "liveness", restartEventStream: true) + guard renderGridLivenessProbeTask == nil else { return } + let probeTimeoutNanoseconds = runtime?.livenessProbeTimeoutNanoseconds + ?? 3_000_000_000 + renderGridLivenessProbeTask = Task { @MainActor [weak self] in + let alive = await Self.probeConnectionLiveness( + client: client, + timeoutNanoseconds: probeTimeoutNanoseconds + ) + guard let self else { return } + self.renderGridLivenessProbeTask = nil + guard !Task.isCancelled, + self.renderGridLivenessListenerID == listenerID, + self.terminalEventListenerID == listenerID, + self.remoteClient === client, + self.connectionState == .connected else { return } + if alive { + // The host answered over the event channel: the stream is + // healthy but quiet. Count the round-trip as the liveness + // evidence so the silence window restarts from this proof. + self.recordTerminalEventStreamLiveness() + MobileDebugLog.anchormux("sync.liveness probe_ok silentMs=\(Int(silent * 1000))") + return + } + // Events may have resumed while the probe was in flight; a fresh + // stamp means the stream already proved itself, so no recovery. + let recheckNow = self.runtime?.now() ?? Date() + let recheckLast = self.lastTerminalEventAt ?? recheckNow + guard recheckNow.timeIntervalSince(recheckLast) >= Self.renderGridLivenessSilenceThreshold else { + return + } + let silentMs = Int(recheckNow.timeIntervalSince(recheckLast) * 1000) + MobileDebugLog.anchormux("sync.liveness re-subscribe silentMs=\(silentMs)") + self.diagnosticLog?.record(DiagnosticEvent(.livenessResubscribe, ms: UInt32(clamping: silentMs))) + mobileShellLog.info("render-grid stream silent for \(silentMs, privacy: .public)ms and host probe failed, re-subscribing") + // resyncTerminalOutput(restartEventStream: true) stops the wedged + // listener (which cancels this watchdog via stopTerminalRefreshPolling) + // and starts a fresh subscription + watchdog, then replays every + // surface so the phone catches up on the deltas it missed while the + // stream was dead. + self.resyncTerminalOutput(reason: "liveness", restartEventStream: true) + } + } + + /// Bounded positive-liveness probe over the event channel. Only a completed + /// `mobile.host.status` round-trip counts as alive; any error (timeout, + /// closed connection, rpc failure) reports dead and lets the watchdog run + /// its recovery. The request is auth-exempt, so the probe never blocks on + /// a token mint and exercises exactly the transport the events use. + private static func probeConnectionLiveness( + client: MobileCoreRPCClient, + timeoutNanoseconds: UInt64 + ) async -> Bool { + guard let request = try? MobileCoreRPCClient.requestData( + method: "mobile.host.status", + params: [:] + ) else { + return false + } + do { + _ = try await client.sendRequest(request, timeoutNanoseconds: timeoutNanoseconds) + return true + } catch { + return false + } } private func resyncTerminalOutput( @@ -4363,6 +4513,8 @@ public final class MobileShellComposite: MobileTerminalOutputSinking { terminalEventListenerTask?.cancel() terminalEventListenerTask = nil terminalEventListenerID = nil + terminalSubscriptionStartTask?.cancel() + terminalSubscriptionStartTask = nil stopRenderGridLivenessWatchdog(listenerID: nil) } From f6116e7c33abe133cc3b604f287d3f8f1f2a5236 Mon Sep 17 00:00:00 2001 From: lawrencecchen <54008264+lawrencecchen@users.noreply.github.com> Date: Wed, 10 Jun 2026 18:44:38 -0700 Subject: [PATCH 03/10] iOS render-grid liveness: probe by idempotent re-subscribe, not host.status Autoreview P1 on the previous probe design: a mobile.host.status answer proves the RPC channel but not the event subscription, so a dropped or wedged server-side registration behind a live channel would reset the liveness clock forever and mask the staleness. The probe is now an idempotent mobile.events.subscribe for the SAME stream id and current topics. A completed round-trip proves the transport the events ride on is alive AND (re)installs the registration, and the host's subscription tracker re-evaluates producer demand on every replace, so the failure mode the finding describes self-heals instead of being masked. The probe still restarts nothing: no listener teardown, no replay, no stream interruption. The deadline bounds the whole attempt including pre-wire token work, so a wedged token provider cannot pin the single-flight probe slot. Tests updated to the final semantics: the teardown signals on a healthy idle stream are now listener restarts (a second mobile.host.status capability resolve) and replay traffic, since the probe itself legitimately re-sends mobile.events.subscribe. Both liveness tests remain red against the pre-fix sources (verified by checking out the commit-1 code under the amended tests: idle stream restarts plus replayCount 2, and the establishment-window event never delivered). Co-Authored-By: Claude Fable 5 --- .../CmuxMobileRPC/MobileSyncRuntime.swift | 14 +-- .../MobileShellComposite.swift | 86 ++++++++++++------- .../MobileShellRenderGridLivenessTests.swift | 75 ++++++++++------ 3 files changed, 112 insertions(+), 63 deletions(-) diff --git a/Packages/CmuxMobileRPC/Sources/CmuxMobileRPC/MobileSyncRuntime.swift b/Packages/CmuxMobileRPC/Sources/CmuxMobileRPC/MobileSyncRuntime.swift index 7c14dbefc097..26ecea5ed388 100644 --- a/Packages/CmuxMobileRPC/Sources/CmuxMobileRPC/MobileSyncRuntime.swift +++ b/Packages/CmuxMobileRPC/Sources/CmuxMobileRPC/MobileSyncRuntime.swift @@ -31,17 +31,17 @@ public protocol MobileSyncRuntime: Sendable { /// consume responses intended for foreground methods. var supportsServerPushEvents: Bool { get } /// Bounded deadline, in nanoseconds, for the render-grid liveness - /// watchdog's host probe (`mobile.host.status`). A healthy idle terminal - /// legitimately pushes no events, so the watchdog verifies prolonged - /// silence with this probe before declaring the stream dead; the deadline - /// bounds how long a dead transport can stall that verdict. + /// watchdog's subscription probe (an idempotent `mobile.events.subscribe` + /// re-assert). A healthy idle terminal legitimately pushes no events, so + /// the watchdog verifies prolonged silence with this probe before + /// declaring the stream dead; the deadline bounds how long a dead + /// transport can stall that verdict. var livenessProbeTimeoutNanoseconds: UInt64 { get } } public extension MobileSyncRuntime { - /// Default probe deadline: generous against a momentarily loaded Mac (the - /// status request is auth-exempt and answered without side effects), while - /// keeping dead-stream recovery within a few seconds of the silence + /// Default probe deadline: generous against a momentarily loaded Mac, + /// while keeping dead-stream recovery within a few seconds of the silence /// threshold instead of the full ``rpcRequestTimeoutNanoseconds``. var livenessProbeTimeoutNanoseconds: UInt64 { 3_000_000_000 } } diff --git a/Packages/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite.swift b/Packages/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite.swift index e166e94e0ba9..791fe29b32e2 100644 --- a/Packages/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite.swift +++ b/Packages/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite.swift @@ -466,9 +466,10 @@ public final class MobileShellComposite: MobileTerminalOutputSinking { // stream and compares "now" against the last received event to detect // prolonged silence. Silence alone is NOT death: a healthy idle terminal // pushes nothing (the Mac dedupes unchanged render-grid frames), so a - // silence-threshold crossing first runs a bounded `mobile.host.status` - // probe and only tears down + re-subscribes + replays when the host fails - // to answer it. + // silence-threshold crossing first runs a bounded idempotent + // `mobile.events.subscribe` probe (same stream id, current topics) and + // only tears down + re-subscribes + replays when the host fails to answer + // it. private var renderGridLivenessTimer: (any DispatchSourceTimer)? private var renderGridLivenessListenerID: UUID? /// The in-flight liveness probe spawned by a silence-threshold crossing. @@ -3904,9 +3905,10 @@ public final class MobileShellComposite: MobileTerminalOutputSinking { /// `lastTerminalEventAt` against `renderGridLivenessSilenceThreshold`. While /// events keep arriving, `lastTerminalEventAt` stays fresh and every tick is a /// no-op. A threshold crossing is treated as a SUSPICION, not a verdict: an - /// idle terminal pushes no events, so the tick first probes the host with a - /// bounded `mobile.host.status` round-trip and only recovers when the probe - /// fails (see ``checkRenderGridLiveness(listenerID:)``). + /// idle terminal pushes no events, so the tick first re-asserts the + /// subscription with a bounded idempotent `mobile.events.subscribe` + /// round-trip and only recovers when that probe fails (see + /// ``checkRenderGridLiveness(listenerID:)``). private func startRenderGridLivenessWatchdog(listenerID: UUID) { stopRenderGridLivenessWatchdog(listenerID: nil) renderGridLivenessListenerID = listenerID @@ -3989,10 +3991,17 @@ public final class MobileShellComposite: MobileTerminalOutputSinking { /// from the half-dead transport this watchdog was built to catch. Treating /// silence alone as death made the watchdog tear down and full-grid-replay /// every healthy idle subscription every ~10.5s, forever (the 2026-06-10 - /// Release-sim bisect finding). `mobile.host.status` is the cheapest - /// possible probe: auth-exempt (no token mint), side-effect free on the - /// host, and answered over the same transport the events ride on, so a - /// completed round-trip is positive proof the channel is alive. + /// Release-sim bisect finding). + /// + /// The probe is an idempotent `mobile.events.subscribe` for the SAME + /// stream id and current topics, not a generic ping: a completed + /// round-trip proves the transport the events ride on is alive AND that + /// the server-side registration is (re)installed, and the host's + /// subscription tracker re-evaluates producer demand on every replace. A + /// generic `mobile.host.status` answer could mask a dropped registration + /// behind a live RPC channel forever. Unlike the resync recovery, the + /// probe restarts nothing: no listener teardown, no replay, no stream + /// interruption. private func checkRenderGridLiveness(listenerID: UUID) { guard renderGridLivenessListenerID == listenerID else { return } guard let client = remoteClient, connectionState == .connected else { return } @@ -4004,11 +4013,13 @@ public final class MobileShellComposite: MobileTerminalOutputSinking { guard renderGridLivenessProbeTask == nil else { return } let probeTimeoutNanoseconds = runtime?.livenessProbeTimeoutNanoseconds ?? 3_000_000_000 + let topics = terminalOutputTransport.eventTopics renderGridLivenessProbeTask = Task { @MainActor [weak self] in - let alive = await Self.probeConnectionLiveness( + let alive = await self?.probeEventSubscriptionLiveness( client: client, + topics: topics, timeoutNanoseconds: probeTimeoutNanoseconds - ) + ) ?? false guard let self else { return } self.renderGridLivenessProbeTask = nil guard !Task.isCancelled, @@ -4017,9 +4028,10 @@ public final class MobileShellComposite: MobileTerminalOutputSinking { self.remoteClient === client, self.connectionState == .connected else { return } if alive { - // The host answered over the event channel: the stream is - // healthy but quiet. Count the round-trip as the liveness - // evidence so the silence window restarts from this proof. + // The host accepted the re-subscribe over the event channel: + // the stream is healthy but quiet. Count the round-trip as the + // liveness evidence so the silence window restarts from this + // proof. self.recordTerminalEventStreamLiveness() MobileDebugLog.anchormux("sync.liveness probe_ok silentMs=\(Int(silent * 1000))") return @@ -4034,7 +4046,7 @@ public final class MobileShellComposite: MobileTerminalOutputSinking { let silentMs = Int(recheckNow.timeIntervalSince(recheckLast) * 1000) MobileDebugLog.anchormux("sync.liveness re-subscribe silentMs=\(silentMs)") self.diagnosticLog?.record(DiagnosticEvent(.livenessResubscribe, ms: UInt32(clamping: silentMs))) - mobileShellLog.info("render-grid stream silent for \(silentMs, privacy: .public)ms and host probe failed, re-subscribing") + mobileShellLog.info("render-grid stream silent for \(silentMs, privacy: .public)ms and subscription probe failed, re-subscribing") // resyncTerminalOutput(restartEventStream: true) stops the wedged // listener (which cancels this watchdog via stopTerminalRefreshPolling) // and starts a fresh subscription + watchdog, then replays every @@ -4044,27 +4056,37 @@ public final class MobileShellComposite: MobileTerminalOutputSinking { } } - /// Bounded positive-liveness probe over the event channel. Only a completed - /// `mobile.host.status` round-trip counts as alive; any error (timeout, - /// closed connection, rpc failure) reports dead and lets the watchdog run - /// its recovery. The request is auth-exempt, so the probe never blocks on - /// a token mint and exercises exactly the transport the events use. - private static func probeConnectionLiveness( + /// Bounded positive-liveness probe: re-assert the event subscription and + /// only count a completed round-trip as alive. Any failure (timeout, + /// closed connection, rpc rejection) reports dead and lets the watchdog + /// run its recovery. + /// + /// The deadline bounds the WHOLE attempt, including any Stack token work + /// that precedes the wire write inside `sendRequest`; an unbounded hang + /// there would otherwise pin the single-flight probe slot and disable the + /// watchdog for the rest of the generation. + private func probeEventSubscriptionLiveness( client: MobileCoreRPCClient, + topics: [String], timeoutNanoseconds: UInt64 ) async -> Bool { - guard let request = try? MobileCoreRPCClient.requestData( - method: "mobile.host.status", - params: [:] - ) else { - return false + let probe = Task { @MainActor [weak self] in + await self?.requestTerminalEventSubscription( + client: client, + reason: "liveness_probe", + topics: topics + ) ?? false } - do { - _ = try await client.sendRequest(request, timeoutNanoseconds: timeoutNanoseconds) - return true - } catch { - return false + // Bounded deadline (allowed: cancellation-wired timeout, not a polling + // sleep). Cancelling the probe task surfaces inside + // requestTerminalEventSubscription as a cancelled request -> false. + let deadline = Task { + try? await Task.sleep(nanoseconds: timeoutNanoseconds) + probe.cancel() } + let alive = await probe.value + deadline.cancel() + return alive } private func resyncTerminalOutput( diff --git a/Packages/CmuxMobileShell/Tests/CmuxMobileShellTests/MobileShellRenderGridLivenessTests.swift b/Packages/CmuxMobileShell/Tests/CmuxMobileShellTests/MobileShellRenderGridLivenessTests.swift index d4038727f6e8..1011e5ca2632 100644 --- a/Packages/CmuxMobileShell/Tests/CmuxMobileShellTests/MobileShellRenderGridLivenessTests.swift +++ b/Packages/CmuxMobileShell/Tests/CmuxMobileShellTests/MobileShellRenderGridLivenessTests.swift @@ -73,6 +73,8 @@ private actor LivenessHostRouter { private var recorded: [RecordedRequest] = [] private var hostStatusRequestCount = 0 private var heldHostStatusRequestNumbers: Set = [] + private var subscribeRequestCount = 0 + private var heldSubscribeRequestNumbers: Set = [] private var holdSubscribe = false private var heldContinuations: [CheckedContinuation] = [] @@ -95,11 +97,18 @@ private actor LivenessHostRouter { heldHostStatusRequestNumbers.insert(number) } + /// Hold the Nth `mobile.events.subscribe` request (1-based) forever, + /// modeling a dead push path whose probe never completes. + func holdSubscribeRequest(number: Int) { + heldSubscribeRequestNumbers.insert(number) + } + /// Resume every held request so parked continuations do not leak past the /// end of the test. func releaseAllHeld() { holdSubscribe = false heldHostStatusRequestNumbers = [] + heldSubscribeRequestNumbers = [] let continuations = heldContinuations heldContinuations = [] for continuation in continuations { @@ -140,7 +149,8 @@ private actor LivenessHostRouter { "capabilities": ["events.v1", "terminal.render_grid.v1", "terminal.replay.v1"], ]) case "mobile.events.subscribe": - if holdSubscribe { + subscribeRequestCount += 1 + if holdSubscribe || heldSubscribeRequestNumbers.contains(subscribeRequestCount) { await park() return nil } @@ -419,11 +429,13 @@ private func makeConnectedStore( /// A healthy idle stream produces zero events (the Mac dedupes unchanged /// frames), so silence alone must not tear the subscription down. The -/// watchdog must verify with a host probe and stay quiet when the host -/// answers. Without this, the phone re-subscribed and full-grid re-replayed -/// every ~10.5s forever on any idle terminal. +/// watchdog may verify the silence with a bounded idempotent re-subscribe +/// probe, but when the host answers it must stay quiet: no listener restart +/// (observable as a second `mobile.host.status` capability resolve) and no +/// full-grid replay. Without this, the phone tore down and full-grid +/// re-replayed every ~10.5s forever on any idle terminal. @MainActor -@Test func watchdogDoesNotResubscribeHealthyIdleStream() async throws { +@Test func watchdogDoesNotTearDownHealthyIdleStream() async throws { let clock = TestClock() let router = LivenessHostRouter() let box = TransportBox() @@ -432,30 +444,43 @@ private func makeConnectedStore( let sawSubscribe = try await pollUntil { await router.count(of: "mobile.events.subscribe") >= 1 } #expect(sawSubscribe, "listener must establish the push subscription") + let collector = OutputCollector() + collector.mount(store: store, surfaceID: "live-terminal") + let sawMountReplay = try await pollUntil { await router.count(of: "mobile.terminal.replay") >= 1 } + #expect(sawMountReplay, "mounting a sink arms exactly one cold-attach replay") + // Idle past the silence threshold: no events at all, host healthy. clock.advance(by: 10) store.debugRunRenderGridLivenessCheckForTesting() - // The watchdog may probe the host (mobile.host.status request number 2; - // number 1 was the listener's transport-capability resolve), and the - // router answers it. What it must NOT do is restart the subscription. - let resubscribed = try await pollUntil(attempts: 60) { - await router.count(of: "mobile.events.subscribe") >= 2 + // A teardown would restart the listener, which re-resolves capabilities + // (mobile.host.status request number 2) and re-replays the mounted sink. + let restarted = try await pollUntil(attempts: 60) { + await router.count(of: "mobile.host.status") >= 2 } #expect( - resubscribed == false, - "the watchdog must not re-subscribe a healthy idle stream; the host answered, so silence only means the terminal had nothing to say" + restarted == false, + "the watchdog must not tear down a healthy idle stream; the host answered the probe, so silence only means the terminal had nothing to say" ) // The probe outcome must reset the silence window: an immediate second // evaluation stays quiet too. store.debugRunRenderGridLivenessCheckForTesting() - let resubscribedAfterRecheck = try await pollUntil(attempts: 30) { - await router.count(of: "mobile.events.subscribe") >= 2 + let restartedAfterRecheck = try await pollUntil(attempts: 30) { + await router.count(of: "mobile.host.status") >= 2 } - #expect(resubscribedAfterRecheck == false) + #expect(restartedAfterRecheck == false) let replayCount = await router.count(of: "mobile.terminal.replay") - #expect(replayCount == 0, "no sink is mounted and the stream is healthy, so no replay traffic may be generated") + #expect(replayCount == 1, "a healthy idle stream must not generate replay traffic beyond the mount's cold-attach replay") + + // The stream was never restarted: the original subscription still + // delivers straight into the mounted sink. + let event = try renderGridEventFrame(surfaceID: "live-terminal", seq: 9, text: "still-alive") + let transport = try #require(box.get()) + await transport.deliver(event) + let delivered = try await pollUntil { collector.lines.isEmpty == false } + #expect(delivered, "the original stream must still be consumed after the probe") + collector.unmount() } /// The watchdog's original purpose (the ~85s silent-death hang) must keep @@ -474,19 +499,21 @@ private func makeConnectedStore( let sawSubscribe = try await pollUntil { await router.count(of: "mobile.events.subscribe") >= 1 } #expect(sawSubscribe, "listener must establish the push subscription") - // The host stops answering the next host.status (the watchdog's probe on - // current code's resync path resolves capabilities first, so holding - // request number 2 models a dead host in both worlds). - await router.holdHostStatusRequest(number: 2) + // The host stops answering the next mobile.events.subscribe (the + // watchdog's re-assert probe), modeling a dead push path while the + // request had already left the phone. + await router.holdSubscribeRequest(number: 2) clock.advance(by: 10) store.debugRunRenderGridLivenessCheckForTesting() - let resubscribed = try await pollUntil(attempts: 600) { - await router.count(of: "mobile.events.subscribe") >= 2 + // Recovery restarts the listener, which re-resolves capabilities: a + // second mobile.host.status request is the teardown-and-restart proof. + let restarted = try await pollUntil(attempts: 600) { + await router.count(of: "mobile.host.status") >= 2 } #expect( - resubscribed, - "a stream that is silent past the threshold AND whose host stops answering must still be torn down and re-subscribed" + restarted, + "a stream that is silent past the threshold AND whose host stops answering the subscription probe must still be torn down and re-subscribed" ) await router.releaseAllHeld() } From 9e8c60231519ca49b88facdcab0a1fb3009e9af6 Mon Sep 17 00:00:00 2001 From: lawrencecchen <54008264+lawrencecchen@users.noreply.github.com> Date: Wed, 10 Jun 2026 18:55:36 -0700 Subject: [PATCH 04/10] iOS render-grid liveness: replay mounted surfaces when the probe repairs a lost registration Autoreview P1 round 2: a successful probe could have just REINSTALLED a registration the host had lost, and render-grid deltas emitted during the gap were never delivered, so stamping liveness and returning would leave the phone stale until unrelated activity happened to repaint the affected rows. The Mac's mobile.events.subscribe acknowledgement now reports already_subscribed (whether the stream id was registered on the connection before the idempotent replace). The phone's probe distinguishes the outcomes: already_subscribed=true (or absent, for older Macs whose registrations cannot outlive the connection anyway) is the healthy-idle case, stamp and stay quiet; already_subscribed=false means the probe repaired a lost registration, so it stamps AND requests a catch-up replay for every mounted surface, without restarting the listener (the phone-side stream is intact; only the host-side registration was missing). New regression test probeRepairingLostSubscriptionReplaysMountedSurfaces drives the lost-registration case end to end through the scripted host: replay re-requested, no listener restart, original stream still consuming afterwards. Co-Authored-By: Claude Fable 5 --- .../MobileEventSubscribeResponse.swift | 13 ++++ .../MobileShellComposite.swift | 71 +++++++++++++------ .../MobileShellRenderGridLivenessTests.swift | 54 ++++++++++++++ Sources/Mobile/MobileHostService.swift | 10 ++- 4 files changed, 126 insertions(+), 22 deletions(-) diff --git a/Packages/CmuxMobileRPC/Sources/CmuxMobileRPC/MobileEventSubscribeResponse.swift b/Packages/CmuxMobileRPC/Sources/CmuxMobileRPC/MobileEventSubscribeResponse.swift index 2872d1e928eb..b498d79a1d50 100644 --- a/Packages/CmuxMobileRPC/Sources/CmuxMobileRPC/MobileEventSubscribeResponse.swift +++ b/Packages/CmuxMobileRPC/Sources/CmuxMobileRPC/MobileEventSubscribeResponse.swift @@ -11,8 +11,20 @@ public struct MobileEventSubscribeResponse: Decodable, Sendable { /// Empty or missing when the subscription was not accepted. public let streamID: String + /// Whether the host already had a registration for this `stream_id` on + /// this connection before processing the request (`mobile.events.subscribe` + /// is idempotent per stream id). + /// + /// `false` means this acknowledgement INSTALLED the registration, so any + /// events emitted while it was absent were never delivered; the liveness + /// probe uses that to trigger a catch-up replay. `nil` when the host + /// predates the field, which callers must treat as "unknown, assume it was + /// already active" so older Macs keep today's behavior. + public let alreadySubscribed: Bool? + private enum CodingKeys: String, CodingKey { case streamID = "stream_id" + case alreadySubscribed = "already_subscribed" } /// Decode a subscribe acknowledgement from raw JSON data. @@ -30,5 +42,6 @@ public struct MobileEventSubscribeResponse: Decodable, Sendable { public init(from decoder: any Decoder) throws { let container = try decoder.container(keyedBy: CodingKeys.self) streamID = (try container.decodeIfPresent(String.self, forKey: .streamID)) ?? "" + alreadySubscribed = try container.decodeIfPresent(Bool.self, forKey: .alreadySubscribed) } } diff --git a/Packages/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite.swift b/Packages/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite.swift index 791fe29b32e2..b4392ef1ef26 100644 --- a/Packages/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite.swift +++ b/Packages/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite.swift @@ -3648,11 +3648,27 @@ public final class MobileShellComposite: MobileTerminalOutputSinking { "ios-terminal-events-\(clientID)" } + /// Outcome of a `mobile.events.subscribe` round-trip. + private enum TerminalEventSubscriptionAck { + case failed + /// The host registered (or re-asserted) the subscription. + /// `alreadySubscribed == false` means this acknowledgement INSTALLED + /// the registration, so events emitted while it was absent were never + /// delivered; `nil` means the host predates the field (treat as + /// already active). + case subscribed(alreadySubscribed: Bool?) + + var isSubscribed: Bool { + if case .subscribed = self { return true } + return false + } + } + private func requestTerminalEventSubscription( client: MobileCoreRPCClient, reason: String, topics: [String] - ) async -> Bool { + ) async -> TerminalEventSubscriptionAck { let requestData: Data do { requestData = try MobileCoreRPCClient.requestData( @@ -3664,7 +3680,7 @@ public final class MobileShellComposite: MobileTerminalOutputSinking { ) } catch { mobileShellLog.error("subscribe payload encode failed: \(String(describing: error), privacy: .private)") - return false + return .failed } let responseData: Data do { @@ -3676,7 +3692,7 @@ public final class MobileShellComposite: MobileTerminalOutputSinking { // `requestTimedOut`. Not a host failure: stay quiet so the log // does not report a self-inflicted cancel as a wire timeout. mobileShellLog.info("subscribe cancelled reason=\(reason, privacy: .public)") - return false + return .failed } mobileShellLog.error("subscribe failed reason=\(reason, privacy: .public): \(String(describing: error), privacy: .private)") // Event-stream (re)subscribe is the view-only/foreground-resume path. @@ -3686,17 +3702,17 @@ public final class MobileShellComposite: MobileTerminalOutputSinking { if remoteClient === client { _ = disconnectForAuthorizationFailureIfNeeded(error) } - return false + return .failed } let response = try? MobileEventSubscribeResponse.decode(responseData) guard let streamID = response?.streamID, !streamID.isEmpty else { mobileShellLog.error("subscribe response missing stream_id reason=\(reason, privacy: .public)") - return false + return .failed } #if DEBUG mobileShellLog.info("subscribe active reason=\(reason, privacy: .public) streamID=\(streamID, privacy: .public)") #endif - return true + return .subscribed(alreadySubscribed: response?.alreadySubscribed) } private func resolveTerminalOutputTransport(client: MobileCoreRPCClient) async -> TerminalOutputTransport { @@ -3857,15 +3873,15 @@ public final class MobileShellComposite: MobileTerminalOutputSinking { guard terminalEventListenerID == listenerID else { return } terminalSubscriptionStartTask?.cancel() terminalSubscriptionStartTask = Task { @MainActor [weak self] in - let subscribed = await self?.requestTerminalEventSubscription( + let ack = await self?.requestTerminalEventSubscription( client: client, reason: "start", topics: topics - ) ?? false + ) ?? .failed guard let self else { return } guard !Task.isCancelled, self.terminalEventListenerID == listenerID else { return } self.terminalSubscriptionStartTask = nil - guard subscribed else { + guard ack.isSubscribed else { MobileDebugLog.anchormux("sync.subscribe_failed reason=start") self.diagnosticLog?.record(DiagnosticEvent(.error)) self.stopTerminalRefreshPolling() @@ -4015,11 +4031,11 @@ public final class MobileShellComposite: MobileTerminalOutputSinking { ?? 3_000_000_000 let topics = terminalOutputTransport.eventTopics renderGridLivenessProbeTask = Task { @MainActor [weak self] in - let alive = await self?.probeEventSubscriptionLiveness( + let ack = await self?.probeEventSubscriptionLiveness( client: client, topics: topics, timeoutNanoseconds: probeTimeoutNanoseconds - ) ?? false + ) ?? .failed guard let self else { return } self.renderGridLivenessProbeTask = nil guard !Task.isCancelled, @@ -4027,13 +4043,26 @@ public final class MobileShellComposite: MobileTerminalOutputSinking { self.terminalEventListenerID == listenerID, self.remoteClient === client, self.connectionState == .connected else { return } - if alive { + if case .subscribed(let alreadySubscribed) = ack { // The host accepted the re-subscribe over the event channel: - // the stream is healthy but quiet. Count the round-trip as the - // liveness evidence so the silence window restarts from this - // proof. + // the stream is healthy. Count the round-trip as the liveness + // evidence so the silence window restarts from this proof. self.recordTerminalEventStreamLiveness() - MobileDebugLog.anchormux("sync.liveness probe_ok silentMs=\(Int(silent * 1000))") + if alreadySubscribed == false { + // The registration had been LOST host-side (the probe just + // reinstalled it), so render-grid deltas emitted during the + // gap were never delivered and delta continuity is broken. + // Replay the mounted surfaces to catch up. The phone-side + // listener stream is intact (registration loss is a + // host-side condition), so no listener restart is needed. + MobileDebugLog.anchormux("sync.liveness probe_repaired silentMs=\(Int(silent * 1000))") + mobileShellLog.info("liveness probe reinstalled a lost event subscription, replaying mounted surfaces") + for surfaceID in self.terminalByteContinuationsBySurfaceID.keys { + self.requestTerminalReplay(surfaceID: surfaceID) + } + } else { + MobileDebugLog.anchormux("sync.liveness probe_ok silentMs=\(Int(silent * 1000))") + } return } // Events may have resumed while the probe was in flight; a fresh @@ -4069,24 +4098,24 @@ public final class MobileShellComposite: MobileTerminalOutputSinking { client: MobileCoreRPCClient, topics: [String], timeoutNanoseconds: UInt64 - ) async -> Bool { + ) async -> TerminalEventSubscriptionAck { let probe = Task { @MainActor [weak self] in await self?.requestTerminalEventSubscription( client: client, reason: "liveness_probe", topics: topics - ) ?? false + ) ?? .failed } // Bounded deadline (allowed: cancellation-wired timeout, not a polling // sleep). Cancelling the probe task surfaces inside - // requestTerminalEventSubscription as a cancelled request -> false. + // requestTerminalEventSubscription as a cancelled request -> .failed. let deadline = Task { try? await Task.sleep(nanoseconds: timeoutNanoseconds) probe.cancel() } - let alive = await probe.value + let ack = await probe.value deadline.cancel() - return alive + return ack } private func resyncTerminalOutput( diff --git a/Packages/CmuxMobileShell/Tests/CmuxMobileShellTests/MobileShellRenderGridLivenessTests.swift b/Packages/CmuxMobileShell/Tests/CmuxMobileShellTests/MobileShellRenderGridLivenessTests.swift index 1011e5ca2632..e51850d6022e 100644 --- a/Packages/CmuxMobileShell/Tests/CmuxMobileShellTests/MobileShellRenderGridLivenessTests.swift +++ b/Packages/CmuxMobileShell/Tests/CmuxMobileShellTests/MobileShellRenderGridLivenessTests.swift @@ -76,6 +76,7 @@ private actor LivenessHostRouter { private var subscribeRequestCount = 0 private var heldSubscribeRequestNumbers: Set = [] private var holdSubscribe = false + private var hasActiveSubscription = false private var heldContinuations: [CheckedContinuation] = [] func record(method: String?, topics: [String]?) { @@ -103,6 +104,13 @@ private actor LivenessHostRouter { heldSubscribeRequestNumbers.insert(number) } + /// Forget the host-side registration, modeling a lost subscription behind + /// a live RPC channel: the next subscribe reports + /// `already_subscribed: false`. + func dropSubscription() { + hasActiveSubscription = false + } + /// Resume every held request so parked continuations do not leak past the /// end of the test. func releaseAllHeld() { @@ -154,9 +162,12 @@ private actor LivenessHostRouter { await park() return nil } + let alreadySubscribed = hasActiveSubscription + hasActiveSubscription = true return try? Self.resultFrame(id: id, result: [ "stream_id": "test-stream", "topics": ["workspace.updated", "terminal.render_grid"], + "already_subscribed": alreadySubscribed, ]) case "mobile.events.unsubscribe", "mobile.terminal.replay", "mobile.terminal.viewport": return try? Self.resultFrame(id: id, result: [:]) @@ -483,6 +494,49 @@ private func makeConnectedStore( collector.unmount() } +/// A successful probe that REPAIRED a lost registration (the host reports +/// `already_subscribed: false`) must replay mounted surfaces: render-grid +/// deltas emitted while the registration was absent were never delivered, so +/// delta continuity is broken even though the channel is healthy again. The +/// phone-side listener stream is intact, so the repair must not restart the +/// listener (no second capability resolve). +@MainActor +@Test func probeRepairingLostSubscriptionReplaysMountedSurfaces() async throws { + let clock = TestClock() + let router = LivenessHostRouter() + let box = TransportBox() + let store = try await makeConnectedStore(router: router, box: box, clock: clock) + + let sawSubscribe = try await pollUntil { await router.count(of: "mobile.events.subscribe") >= 1 } + #expect(sawSubscribe, "listener must establish the push subscription") + + let collector = OutputCollector() + collector.mount(store: store, surfaceID: "live-terminal") + let sawMountReplay = try await pollUntil { await router.count(of: "mobile.terminal.replay") >= 1 } + #expect(sawMountReplay, "mounting a sink arms exactly one cold-attach replay") + + // The host loses the registration while the RPC channel stays healthy. + await router.dropSubscription() + clock.advance(by: 10) + store.debugRunRenderGridLivenessCheckForTesting() + + let replayed = try await pollUntil { await router.count(of: "mobile.terminal.replay") >= 2 } + #expect( + replayed, + "a probe that reinstalls a lost registration must request a catch-up replay for mounted surfaces; deltas emitted during the gap were never delivered" + ) + let hostStatusCount = await router.count(of: "mobile.host.status") + #expect(hostStatusCount == 1, "the repair must not restart the listener; the phone-side stream is intact") + + // The repaired stream delivers straight into the still-mounted sink. + let event = try renderGridEventFrame(surfaceID: "live-terminal", seq: 11, text: "repaired") + let transport = try #require(box.get()) + await transport.deliver(event) + let delivered = try await pollUntil { collector.lines.contains { $0.contains("repaired") } } + #expect(delivered, "the original stream must still be consumed after the repair") + collector.unmount() +} + /// The watchdog's original purpose (the ~85s silent-death hang) must keep /// working: silence past the threshold plus a host that stops answering the /// probe must still tear down and re-subscribe. diff --git a/Sources/Mobile/MobileHostService.swift b/Sources/Mobile/MobileHostService.swift index 1abc437ecc6f..a17c5debeca4 100644 --- a/Sources/Mobile/MobileHostService.swift +++ b/Sources/Mobile/MobileHostService.swift @@ -2166,13 +2166,21 @@ actor MobileHostConnection { guard !topics.isEmpty else { return .failure(MobileHostRPCError(code: "invalid_params", message: "topics is required")) } + // Report whether this stream id was already registered BEFORE the + // idempotent replace. The phone's render-grid liveness probe + // re-asserts its subscription on prolonged silence; `false` tells + // it the registration had been lost (events emitted in the gap + // were never delivered), so it requests a catch-up replay instead + // of trusting delta continuity. + let alreadySubscribed = subscriptions[streamID] != nil subscribe(streamID: streamID, topics: topics) #if DEBUG - cmuxDebugLog("mobile.subscribe streamID=\(streamID) topics=\(topics.sorted()) connID=\(self.id.uuidString)") + cmuxDebugLog("mobile.subscribe streamID=\(streamID) topics=\(topics.sorted()) existing=\(alreadySubscribed) connID=\(self.id.uuidString)") #endif return .ok([ "stream_id": streamID, "topics": Array(topics).sorted(), + "already_subscribed": alreadySubscribed, ]) case "mobile.events.unsubscribe": let streamID = request.params["stream_id"] as? String ?? "" From a5889ab731c22099f12f673edada19318b238867 Mon Sep 17 00:00:00 2001 From: lawrencecchen <54008264+lawrencecchen@users.noreply.github.com> Date: Wed, 10 Jun 2026 19:01:57 -0700 Subject: [PATCH 05/10] iOS render-grid liveness: probe slot ownership token, workspace refresh on repair Two autoreview P2 fixes: 1. A cancelled probe from a previous watchdog generation could complete late and nil out a newer generation's in-flight probe handle, breaking the single-flight guard and orphaning the newer probe from lifecycle cancellation. The slot now carries an ownership token (renderGridLivenessProbeID): only the probe holding it may clear the slot; superseded probes return without touching it. 2. The repaired-registration path only replayed terminal surfaces, but the same registration carries workspace.updated, so workspace create/rename/delete events emitted during the gap were missed too. The repair now also re-fetches the authoritative workspace list via the existing scheduleWorkspaceListRefreshFromEvent path, asserted in the repair regression test. Co-Authored-By: Claude Fable 5 --- .../MobileShellComposite.swift | 19 ++++++++++++++++++- .../MobileShellRenderGridLivenessTests.swift | 10 ++++++++++ 2 files changed, 28 insertions(+), 1 deletion(-) diff --git a/Packages/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite.swift b/Packages/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite.swift index b4392ef1ef26..33881ab1cf2e 100644 --- a/Packages/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite.swift +++ b/Packages/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite.swift @@ -473,8 +473,13 @@ public final class MobileShellComposite: MobileTerminalOutputSinking { private var renderGridLivenessTimer: (any DispatchSourceTimer)? private var renderGridLivenessListenerID: UUID? /// The in-flight liveness probe spawned by a silence-threshold crossing. - /// Single-flight: ticks while a probe is pending are no-ops. + /// Single-flight: ticks while a probe is pending are no-ops. The paired + /// `renderGridLivenessProbeID` is the slot's ownership token: only the + /// probe holding it may clear the slot, so a cancelled probe from an older + /// generation completing late cannot free or clobber a newer generation's + /// in-flight slot. private var renderGridLivenessProbeTask: Task? + private var renderGridLivenessProbeID: UUID? private var lastTerminalEventAt: Date? private var terminalSubscriptionRefreshTask: Task? private var createWorkspaceTask: Task? @@ -3969,6 +3974,7 @@ public final class MobileShellComposite: MobileTerminalOutputSinking { renderGridLivenessListenerID = nil renderGridLivenessProbeTask?.cancel() renderGridLivenessProbeTask = nil + renderGridLivenessProbeID = nil } /// Single ownership point for the liveness clock the watchdog reads. @@ -4030,6 +4036,8 @@ public final class MobileShellComposite: MobileTerminalOutputSinking { let probeTimeoutNanoseconds = runtime?.livenessProbeTimeoutNanoseconds ?? 3_000_000_000 let topics = terminalOutputTransport.eventTopics + let probeID = UUID() + renderGridLivenessProbeID = probeID renderGridLivenessProbeTask = Task { @MainActor [weak self] in let ack = await self?.probeEventSubscriptionLiveness( client: client, @@ -4037,7 +4045,12 @@ public final class MobileShellComposite: MobileTerminalOutputSinking { timeoutNanoseconds: probeTimeoutNanoseconds ) ?? .failed guard let self else { return } + // Only the probe that owns the single-flight slot may clear it; a + // superseded probe completing late returns without touching the + // newer generation's in-flight slot. + guard self.renderGridLivenessProbeID == probeID else { return } self.renderGridLivenessProbeTask = nil + self.renderGridLivenessProbeID = nil guard !Task.isCancelled, self.renderGridLivenessListenerID == listenerID, self.terminalEventListenerID == listenerID, @@ -4060,6 +4073,10 @@ public final class MobileShellComposite: MobileTerminalOutputSinking { for surfaceID in self.terminalByteContinuationsBySurfaceID.keys { self.requestTerminalReplay(surfaceID: surfaceID) } + // The same registration carries `workspace.updated`, so + // workspace create/rename/delete events emitted during the + // gap were missed too; re-fetch the authoritative list. + self.scheduleWorkspaceListRefreshFromEvent() } else { MobileDebugLog.anchormux("sync.liveness probe_ok silentMs=\(Int(silent * 1000))") } diff --git a/Packages/CmuxMobileShell/Tests/CmuxMobileShellTests/MobileShellRenderGridLivenessTests.swift b/Packages/CmuxMobileShell/Tests/CmuxMobileShellTests/MobileShellRenderGridLivenessTests.swift index e51850d6022e..5583f2773f00 100644 --- a/Packages/CmuxMobileShell/Tests/CmuxMobileShellTests/MobileShellRenderGridLivenessTests.swift +++ b/Packages/CmuxMobileShell/Tests/CmuxMobileShellTests/MobileShellRenderGridLivenessTests.swift @@ -517,6 +517,8 @@ private func makeConnectedStore( // The host loses the registration while the RPC channel stays healthy. await router.dropSubscription() + let workspaceListsBeforeRepair = await router.count(of: "mobile.workspace.list") + + router.count(of: "workspace.list") clock.advance(by: 10) store.debugRunRenderGridLivenessCheckForTesting() @@ -527,6 +529,14 @@ private func makeConnectedStore( ) let hostStatusCount = await router.count(of: "mobile.host.status") #expect(hostStatusCount == 1, "the repair must not restart the listener; the phone-side stream is intact") + // workspace.updated events were missed during the gap too: the repair must + // re-fetch the authoritative workspace list. + let workspaceRefetched = try await pollUntil { + let current = await router.count(of: "mobile.workspace.list") + + router.count(of: "workspace.list") + return current > workspaceListsBeforeRepair + } + #expect(workspaceRefetched, "the repaired subscription also carries workspace.updated, so the workspace list must be re-fetched") // The repaired stream delivers straight into the still-mounted sink. let event = try renderGridEventFrame(surfaceID: "live-terminal", seq: 11, text: "repaired") From 971cf677abbc855b2a01a257aee2af76eaf2660c Mon Sep 17 00:00:00 2001 From: lawrencecchen <54008264+lawrencecchen@users.noreply.github.com> Date: Wed, 10 Jun 2026 21:54:33 -0700 Subject: [PATCH 06/10] iOS render-grid liveness: probe success restores connection status; subscribe RPCs are non-interactive host activity Review fixes: a successful liveness probe now calls markMacConnectionHealthy() so a transient RPC failure can't leave the status UI stuck unavailable on an idle terminal, and mobile.events.subscribe/unsubscribe no longer count as interactive mobile activity so the ~9s idle probe can't starve host work gated on mobile quiet (TabManager background git/PR refresh). Co-Authored-By: Claude Fable 5 --- .../Sources/CmuxMobileShell/MobileShellComposite.swift | 5 +++++ Sources/Mobile/MobileHostService.swift | 8 +++++++- 2 files changed, 12 insertions(+), 1 deletion(-) diff --git a/Packages/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite.swift b/Packages/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite.swift index 33881ab1cf2e..f9c719619561 100644 --- a/Packages/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite.swift +++ b/Packages/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite.swift @@ -4061,6 +4061,11 @@ public final class MobileShellComposite: MobileTerminalOutputSinking { // the stream is healthy. Count the round-trip as the liveness // evidence so the silence window restarts from this proof. self.recordTerminalEventStreamLiveness() + // The round-trip is also positive proof of the client/host + // connection itself; recover the visible status if a prior + // transient RPC failure marked it unavailable, since an idle + // terminal may never emit another event to flip it back. + self.markMacConnectionHealthy() if alreadySubscribed == false { // The registration had been LOST host-side (the probe just // reinstalled it), so render-grid deltas emitted during the diff --git a/Sources/Mobile/MobileHostService.swift b/Sources/Mobile/MobileHostService.swift index a17c5debeca4..4f324a19b84e 100644 --- a/Sources/Mobile/MobileHostService.swift +++ b/Sources/Mobile/MobileHostService.swift @@ -2196,7 +2196,13 @@ actor MobileHostConnection { private static func isInteractiveMobileRequest(_ method: String) -> Bool { switch method { - case "mobile.host.status", "mobile.terminal.replay", "terminal.replay": + case "mobile.host.status", "mobile.terminal.replay", "terminal.replay", + // Subscription management is plumbing, not user interaction: the + // phone's render-grid liveness watchdog re-asserts its + // subscription on every silence window (~9s when idle), and + // counting that as interactive activity starves host work gated + // on mobile quiet (e.g. TabManager background git/PR refresh). + "mobile.events.subscribe", "mobile.events.unsubscribe": return false default: return true From ea53f480281f5e1d789f5d2c9d5b0346c60c1ee3 Mon Sep 17 00:00:00 2001 From: lawrencecchen <54008264+lawrencecchen@users.noreply.github.com> Date: Thu, 11 Jun 2026 00:16:05 -0700 Subject: [PATCH 07/10] Probe deadline via one-shot DispatchSourceTimer instead of Task.sleep Greptile P2: Task.sleep is a flagged primitive in production Swift; the file's sanctioned timing primitive is DispatchSourceTimer (the watchdog tick already uses it). Same behavior: deadline cancels the in-flight probe, probe completion cancels the deadline. Co-Authored-By: Claude Fable 5 --- .../CmuxMobileShell/MobileShellComposite.swift | 16 +++++++++------- 1 file changed, 9 insertions(+), 7 deletions(-) diff --git a/Packages/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite.swift b/Packages/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite.swift index f9c719619561..3efb71da6df9 100644 --- a/Packages/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite.swift +++ b/Packages/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite.swift @@ -4128,13 +4128,15 @@ public final class MobileShellComposite: MobileTerminalOutputSinking { topics: topics ) ?? .failed } - // Bounded deadline (allowed: cancellation-wired timeout, not a polling - // sleep). Cancelling the probe task surfaces inside - // requestTerminalEventSubscription as a cancelled request -> .failed. - let deadline = Task { - try? await Task.sleep(nanoseconds: timeoutNanoseconds) - probe.cancel() - } + // Bounded deadline via a one-shot DispatchSourceTimer — the same + // sanctioned primitive the watchdog tick uses — with cancellation + // wired to the probe's lifecycle. Cancelling the probe task surfaces + // inside requestTerminalEventSubscription as a cancelled request -> + // .failed. + let deadline = DispatchSource.makeTimerSource(queue: .main) + deadline.schedule(deadline: .now() + .nanoseconds(Int(clamping: timeoutNanoseconds))) + deadline.setEventHandler { probe.cancel() } + deadline.resume() let ack = await probe.value deadline.cancel() return ack From bb1e595272341394d55009a37540f2055321ad21 Mon Sep 17 00:00:00 2001 From: lawrencecchen <54008264+lawrencecchen@users.noreply.github.com> Date: Thu, 11 Jun 2026 14:09:31 -0700 Subject: [PATCH 08/10] Fix workflow-guard file-length budget: split liveness test fixtures, refresh budget MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Split MobileShellRenderGridLivenessTests.swift (583 > 500 new-file threshold) into the 4 regression tests (212 lines) plus MobileShellRenderGridLivenessTestSupport.swift (381 lines: injected clock, scripted host router, transport mocks, connected-store builder). Budget refresh accepts the remaining growth as known debt: MobileShellComposite.swift +238 (the watchdog probe machinery is single-flight state coupled to the composite's private connection state; extracting it would force ~15 members internal — precedent: the composer merge refreshed this same file's budget) and MobileHostService.swift +14 (alreadySubscribed ack field). Co-Authored-By: Claude Fable 5 --- .github/swift-file-length-budget.tsv | 4 +- ...leShellRenderGridLivenessTestSupport.swift | 381 ++++++++++++++++++ .../MobileShellRenderGridLivenessTests.swift | 371 ----------------- 3 files changed, 383 insertions(+), 373 deletions(-) create mode 100644 Packages/CmuxMobileShell/Tests/CmuxMobileShellTests/MobileShellRenderGridLivenessTestSupport.swift diff --git a/.github/swift-file-length-budget.tsv b/.github/swift-file-length-budget.tsv index 47569845bc78..423116151d0d 100644 --- a/.github/swift-file-length-budget.tsv +++ b/.github/swift-file-length-budget.tsv @@ -21,7 +21,7 @@ 6071 Sources/TextBoxInput.swift 5482 cmuxTests/BrowserConfigTests.swift 5448 Sources/cmuxApp.swift -4574 Packages/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite.swift +4812 Packages/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite.swift 4460 Sources/Panels/FilePreviewPanel.swift 4400 cmuxTests/BrowserPanelTests.swift 4227 Sources/BrowserWindowPortal.swift @@ -41,7 +41,7 @@ 2513 cmuxTests/CommandPaletteSearchEngineTests.swift 2509 Sources/KeyboardShortcutSettings.swift 2327 cmuxTests/CJKIMEInputTests.swift -2326 Sources/Mobile/MobileHostService.swift +2340 Sources/Mobile/MobileHostService.swift 2290 Sources/FileExplorerView.swift 2260 Sources/TerminalWindowPortal.swift 2138 Sources/SessionPersistence.swift diff --git a/Packages/CmuxMobileShell/Tests/CmuxMobileShellTests/MobileShellRenderGridLivenessTestSupport.swift b/Packages/CmuxMobileShell/Tests/CmuxMobileShellTests/MobileShellRenderGridLivenessTestSupport.swift new file mode 100644 index 000000000000..1ae69543c729 --- /dev/null +++ b/Packages/CmuxMobileShell/Tests/CmuxMobileShellTests/MobileShellRenderGridLivenessTestSupport.swift @@ -0,0 +1,381 @@ +import CMUXMobileCore +import CmuxMobileRPC +import Foundation +import Testing +@testable import CmuxMobileShell + +// Shared fixtures for the render-grid liveness watchdog tests +// (MobileShellRenderGridLivenessTests.swift): injected clock, scripted +// host router, transport mocks, and the connected-store builder. + +// MARK: - Injected clock + +final class TestClock: @unchecked Sendable { + private let lock = NSLock() + private var current: Date + + init() { + current = Date() + } + + var now: Date { + lock.withLock { current } + } + + func advance(by interval: TimeInterval) { + lock.withLock { current = current.addingTimeInterval(interval) } + } +} + +// MARK: - Runtime double + +struct LivenessTestRuntime: MobileSyncRuntime { + var transportFactory: any CmxByteTransportFactory + var stackAccessTokenProvider: @Sendable () async throws -> String = { "test-stack-token" } + var stackAccessTokenForceRefresher: @Sendable () async throws -> String = { "test-stack-token" } + var rpcRequestTimeoutNanoseconds: UInt64 = 30 * 1_000_000_000 + var now: @Sendable () -> Date + var supportedRouteKinds: [CmxAttachTransportKind] = [.debugLoopback] + var pairingRequestTimeoutNanoseconds: UInt64 = 30 * 1_000_000_000 + var supportsServerPushEvents: Bool = true + /// Bounded deadline for the watchdog's host liveness probe. Short here so + /// the dead-stream test does not wait the production default. + var livenessProbeTimeoutNanoseconds: UInt64 = 200_000_000 +} + +// MARK: - Scripted host (router + transport) + +/// Scripts the Mac side of the persistent RPC connection: answers the +/// connect-time `workspace.list`, the `mobile.host.status` capability and +/// probe requests, `mobile.events.subscribe`, and replay/viewport calls. +/// Individual requests can be held unresolved to model an ack that has not +/// arrived yet (establishment window) or a host that stopped answering +/// (dead stream). +actor LivenessHostRouter { + struct RecordedRequest: Sendable { + var method: String? + var topics: [String]? + } + + private var recorded: [RecordedRequest] = [] + private var hostStatusRequestCount = 0 + private var heldHostStatusRequestNumbers: Set = [] + private var subscribeRequestCount = 0 + private var heldSubscribeRequestNumbers: Set = [] + private var holdSubscribe = false + private var hasActiveSubscription = false + private var heldContinuations: [CheckedContinuation] = [] + + func record(method: String?, topics: [String]?) { + recorded.append(RecordedRequest(method: method, topics: topics)) + } + + func count(of method: String) -> Int { + recorded.filter { $0.method == method }.count + } + + /// Hold every `mobile.events.subscribe` response until released. + func setHoldSubscribe(_ hold: Bool) { + holdSubscribe = hold + } + + /// Hold the Nth `mobile.host.status` request (1-based) forever, modeling + /// a host that stopped answering on a half-dead transport. + func holdHostStatusRequest(number: Int) { + heldHostStatusRequestNumbers.insert(number) + } + + /// Hold the Nth `mobile.events.subscribe` request (1-based) forever, + /// modeling a dead push path whose probe never completes. + func holdSubscribeRequest(number: Int) { + heldSubscribeRequestNumbers.insert(number) + } + + /// Forget the host-side registration, modeling a lost subscription behind + /// a live RPC channel: the next subscribe reports + /// `already_subscribed: false`. + func dropSubscription() { + hasActiveSubscription = false + } + + /// Resume every held request so parked continuations do not leak past the + /// end of the test. + func releaseAllHeld() { + holdSubscribe = false + heldHostStatusRequestNumbers = [] + heldSubscribeRequestNumbers = [] + let continuations = heldContinuations + heldContinuations = [] + for continuation in continuations { + continuation.resume() + } + } + + func response(method: String?, id: String?) async -> Data? { + switch method { + case "workspace.list", "mobile.workspace.list": + return try? Self.resultFrame(id: id, result: [ + "workspaces": [ + [ + "id": "live-workspace", + "title": "Live Workspace", + "current_directory": "/Users/test/project", + "is_selected": true, + "terminals": [ + [ + "id": "live-terminal", + "title": "Terminal", + "current_directory": "/Users/test/project", + "is_ready": true, + "is_focused": true, + ], + ], + ], + ], + ]) + case "mobile.host.status": + hostStatusRequestCount += 1 + if heldHostStatusRequestNumbers.contains(hostStatusRequestCount) { + await park() + return nil + } + return try? Self.resultFrame(id: id, result: [ + "terminal_fidelity": "render_grid", + "capabilities": ["events.v1", "terminal.render_grid.v1", "terminal.replay.v1"], + ]) + case "mobile.events.subscribe": + subscribeRequestCount += 1 + if holdSubscribe || heldSubscribeRequestNumbers.contains(subscribeRequestCount) { + await park() + return nil + } + let alreadySubscribed = hasActiveSubscription + hasActiveSubscription = true + return try? Self.resultFrame(id: id, result: [ + "stream_id": "test-stream", + "topics": ["workspace.updated", "terminal.render_grid"], + "already_subscribed": alreadySubscribed, + ]) + case "mobile.events.unsubscribe", "mobile.terminal.replay", "mobile.terminal.viewport": + return try? Self.resultFrame(id: id, result: [:]) + default: + return try? Self.errorFrame(id: id, message: "Unexpected method \(method ?? "nil")") + } + } + + private func park() async { + await withCheckedContinuation { continuation in + heldContinuations.append(continuation) + } + } + + private static func resultFrame(id: String?, result: [String: Any]) throws -> Data { + let envelope: [String: Any] = [ + "id": id ?? UUID().uuidString, + "ok": true, + "result": result, + ] + return try MobileSyncFrameCodec.encodeFrame(JSONSerialization.data(withJSONObject: envelope)) + } + + private static func errorFrame(id: String?, message: String) throws -> Data { + let envelope: [String: Any] = [ + "id": id ?? UUID().uuidString, + "ok": false, + "error": ["message": message], + ] + return try MobileSyncFrameCodec.encodeFrame(JSONSerialization.data(withJSONObject: envelope)) + } +} + +/// Holds the live transport instance so the test can push unsolicited +/// server-side event frames through the same receive path production uses. +final class TransportBox: @unchecked Sendable { + private let lock = NSLock() + private var transport: LivenessTransport? + + func set(_ transport: LivenessTransport) { + lock.withLock { self.transport = transport } + } + + func get() -> LivenessTransport? { + lock.withLock { transport } + } +} + +struct LivenessTransportFactory: CmxByteTransportFactory { + let router: LivenessHostRouter + let box: TransportBox + + func makeTransport(for route: CmxAttachRoute) throws -> any CmxByteTransport { + let transport = LivenessTransport(router: router) + box.set(transport) + return transport + } +} + +actor LivenessTransport: CmxByteTransport { + private let router: LivenessHostRouter + private var pendingFrames: [Data] = [] + private var receiveWaiters: [CheckedContinuation] = [] + private var isClosed = false + + init(router: LivenessHostRouter) { + self.router = router + } + + func connect() async throws {} + + func receive() async throws -> Data? { + if !pendingFrames.isEmpty { + return pendingFrames.removeFirst() + } + if isClosed { + return nil + } + return await withCheckedContinuation { continuation in + receiveWaiters.append(continuation) + } + } + + func send(_ data: Data) async throws { + var buffer = data + let payloads = try MobileSyncFrameCodec.decodeFrames(from: &buffer) + for payload in payloads { + let parsed = (try? JSONSerialization.jsonObject(with: payload)) as? [String: Any] + let method = parsed?["method"] as? String + let id = parsed?["id"] as? String + let topics = (parsed?["params"] as? [String: Any])?["topics"] as? [String] + await router.record(method: method, topics: topics) + // Answer each request concurrently so one held response cannot + // head-of-line block later RPCs, matching the Mac host's + // per-frame response tasks. + Task { [router, weak self] in + guard let response = await router.response(method: method, id: id) else { + return + } + await self?.deliver(response) + } + } + } + + func close() async { + isClosed = true + let waiters = receiveWaiters + receiveWaiters = [] + for waiter in waiters { + waiter.resume(returning: nil) + } + } + + /// Deliver a frame to the client's read loop. Also used by tests to push + /// unsolicited server-side event envelopes. + func deliver(_ frame: Data) { + if receiveWaiters.isEmpty { + pendingFrames.append(frame) + return + } + let waiter = receiveWaiters.removeFirst() + waiter.resume(returning: frame) + } +} + +// MARK: - Test helpers + +@MainActor +final class OutputCollector { + private(set) var lines: [String] = [] + private var task: Task? + + func mount(store: MobileShellComposite, surfaceID: String) { + task = Task { @MainActor [weak self] in + for await data in store.terminalOutputStream(surfaceID: surfaceID) { + self?.lines.append(String(decoding: data, as: UTF8.self)) + } + } + } + + func unmount() { + task?.cancel() + task = nil + } +} + +func makeTicket(clock: TestClock) throws -> CmxAttachTicket { + let route = try CmxAttachRoute( + id: "debug_loopback", + kind: .debugLoopback, + endpoint: .hostPort(host: "127.0.0.1", port: 56584) + ) + return try CmxAttachTicket( + workspaceID: "live-workspace", + terminalID: "live-terminal", + macDeviceID: "test-mac", + macDisplayName: "Test Mac", + routes: [route], + expiresAt: clock.now.addingTimeInterval(3600) + ) +} + +func attachURL(for ticket: CmxAttachTicket) throws -> String { + let encoder = JSONEncoder() + encoder.dateEncodingStrategy = .iso8601 + let payload = try encoder.encode(ticket) + .base64EncodedString() + .replacingOccurrences(of: "+", with: "-") + .replacingOccurrences(of: "/", with: "_") + .replacingOccurrences(of: "=", with: "") + return "cmux-ios://attach?v=\(ticket.version)&payload=\(payload)" +} + +func renderGridEventFrame(surfaceID: String, seq: UInt64, text: String) throws -> Data { + let frame = try MobileTerminalRenderGridFrame.fromPlainRows( + surfaceID: surfaceID, + stateSeq: seq, + columns: 16, + rows: 4, + text: text + ) + let envelope: [String: Any] = [ + "kind": "event", + "topic": "terminal.render_grid", + "payload": try frame.jsonObject(), + ] + return try MobileSyncFrameCodec.encodeFrame(JSONSerialization.data(withJSONObject: envelope)) +} + +/// Poll until `condition` is true, bounded at `attempts` x 10ms. Returns the +/// final value so tests can assert both presence and (bounded) absence. +@MainActor +func pollUntil( + attempts: Int = 300, + _ condition: @MainActor () async -> Bool +) async throws -> Bool { + for _ in 0.. MobileShellComposite { + let runtime = LivenessTestRuntime( + transportFactory: LivenessTransportFactory(router: router, box: box), + now: { clock.now }, + livenessProbeTimeoutNanoseconds: probeTimeoutNanoseconds + ) + let store = MobileShellComposite.preview(runtime: runtime) + store.signIn() + let ticket = try makeTicket(clock: clock) + let connected = await store.connectPairingURL(try attachURL(for: ticket)) + #expect(connected, "scripted connect must succeed") + return store +} diff --git a/Packages/CmuxMobileShell/Tests/CmuxMobileShellTests/MobileShellRenderGridLivenessTests.swift b/Packages/CmuxMobileShell/Tests/CmuxMobileShellTests/MobileShellRenderGridLivenessTests.swift index 5583f2773f00..a90d34c38505 100644 --- a/Packages/CmuxMobileShell/Tests/CmuxMobileShellTests/MobileShellRenderGridLivenessTests.swift +++ b/Packages/CmuxMobileShell/Tests/CmuxMobileShellTests/MobileShellRenderGridLivenessTests.swift @@ -21,377 +21,6 @@ import Testing // silence alone can never distinguish "idle" from "dead". The watchdog // needs a bounded host probe before it may declare death. -// MARK: - Injected clock - -private final class TestClock: @unchecked Sendable { - private let lock = NSLock() - private var current: Date - - init() { - current = Date() - } - - var now: Date { - lock.withLock { current } - } - - func advance(by interval: TimeInterval) { - lock.withLock { current = current.addingTimeInterval(interval) } - } -} - -// MARK: - Runtime double - -private struct LivenessTestRuntime: MobileSyncRuntime { - var transportFactory: any CmxByteTransportFactory - var stackAccessTokenProvider: @Sendable () async throws -> String = { "test-stack-token" } - var stackAccessTokenForceRefresher: @Sendable () async throws -> String = { "test-stack-token" } - var rpcRequestTimeoutNanoseconds: UInt64 = 30 * 1_000_000_000 - var now: @Sendable () -> Date - var supportedRouteKinds: [CmxAttachTransportKind] = [.debugLoopback] - var pairingRequestTimeoutNanoseconds: UInt64 = 30 * 1_000_000_000 - var supportsServerPushEvents: Bool = true - /// Bounded deadline for the watchdog's host liveness probe. Short here so - /// the dead-stream test does not wait the production default. - var livenessProbeTimeoutNanoseconds: UInt64 = 200_000_000 -} - -// MARK: - Scripted host (router + transport) - -/// Scripts the Mac side of the persistent RPC connection: answers the -/// connect-time `workspace.list`, the `mobile.host.status` capability and -/// probe requests, `mobile.events.subscribe`, and replay/viewport calls. -/// Individual requests can be held unresolved to model an ack that has not -/// arrived yet (establishment window) or a host that stopped answering -/// (dead stream). -private actor LivenessHostRouter { - struct RecordedRequest: Sendable { - var method: String? - var topics: [String]? - } - - private var recorded: [RecordedRequest] = [] - private var hostStatusRequestCount = 0 - private var heldHostStatusRequestNumbers: Set = [] - private var subscribeRequestCount = 0 - private var heldSubscribeRequestNumbers: Set = [] - private var holdSubscribe = false - private var hasActiveSubscription = false - private var heldContinuations: [CheckedContinuation] = [] - - func record(method: String?, topics: [String]?) { - recorded.append(RecordedRequest(method: method, topics: topics)) - } - - func count(of method: String) -> Int { - recorded.filter { $0.method == method }.count - } - - /// Hold every `mobile.events.subscribe` response until released. - func setHoldSubscribe(_ hold: Bool) { - holdSubscribe = hold - } - - /// Hold the Nth `mobile.host.status` request (1-based) forever, modeling - /// a host that stopped answering on a half-dead transport. - func holdHostStatusRequest(number: Int) { - heldHostStatusRequestNumbers.insert(number) - } - - /// Hold the Nth `mobile.events.subscribe` request (1-based) forever, - /// modeling a dead push path whose probe never completes. - func holdSubscribeRequest(number: Int) { - heldSubscribeRequestNumbers.insert(number) - } - - /// Forget the host-side registration, modeling a lost subscription behind - /// a live RPC channel: the next subscribe reports - /// `already_subscribed: false`. - func dropSubscription() { - hasActiveSubscription = false - } - - /// Resume every held request so parked continuations do not leak past the - /// end of the test. - func releaseAllHeld() { - holdSubscribe = false - heldHostStatusRequestNumbers = [] - heldSubscribeRequestNumbers = [] - let continuations = heldContinuations - heldContinuations = [] - for continuation in continuations { - continuation.resume() - } - } - - func response(method: String?, id: String?) async -> Data? { - switch method { - case "workspace.list", "mobile.workspace.list": - return try? Self.resultFrame(id: id, result: [ - "workspaces": [ - [ - "id": "live-workspace", - "title": "Live Workspace", - "current_directory": "/Users/test/project", - "is_selected": true, - "terminals": [ - [ - "id": "live-terminal", - "title": "Terminal", - "current_directory": "/Users/test/project", - "is_ready": true, - "is_focused": true, - ], - ], - ], - ], - ]) - case "mobile.host.status": - hostStatusRequestCount += 1 - if heldHostStatusRequestNumbers.contains(hostStatusRequestCount) { - await park() - return nil - } - return try? Self.resultFrame(id: id, result: [ - "terminal_fidelity": "render_grid", - "capabilities": ["events.v1", "terminal.render_grid.v1", "terminal.replay.v1"], - ]) - case "mobile.events.subscribe": - subscribeRequestCount += 1 - if holdSubscribe || heldSubscribeRequestNumbers.contains(subscribeRequestCount) { - await park() - return nil - } - let alreadySubscribed = hasActiveSubscription - hasActiveSubscription = true - return try? Self.resultFrame(id: id, result: [ - "stream_id": "test-stream", - "topics": ["workspace.updated", "terminal.render_grid"], - "already_subscribed": alreadySubscribed, - ]) - case "mobile.events.unsubscribe", "mobile.terminal.replay", "mobile.terminal.viewport": - return try? Self.resultFrame(id: id, result: [:]) - default: - return try? Self.errorFrame(id: id, message: "Unexpected method \(method ?? "nil")") - } - } - - private func park() async { - await withCheckedContinuation { continuation in - heldContinuations.append(continuation) - } - } - - private static func resultFrame(id: String?, result: [String: Any]) throws -> Data { - let envelope: [String: Any] = [ - "id": id ?? UUID().uuidString, - "ok": true, - "result": result, - ] - return try MobileSyncFrameCodec.encodeFrame(JSONSerialization.data(withJSONObject: envelope)) - } - - private static func errorFrame(id: String?, message: String) throws -> Data { - let envelope: [String: Any] = [ - "id": id ?? UUID().uuidString, - "ok": false, - "error": ["message": message], - ] - return try MobileSyncFrameCodec.encodeFrame(JSONSerialization.data(withJSONObject: envelope)) - } -} - -/// Holds the live transport instance so the test can push unsolicited -/// server-side event frames through the same receive path production uses. -private final class TransportBox: @unchecked Sendable { - private let lock = NSLock() - private var transport: LivenessTransport? - - func set(_ transport: LivenessTransport) { - lock.withLock { self.transport = transport } - } - - func get() -> LivenessTransport? { - lock.withLock { transport } - } -} - -private struct LivenessTransportFactory: CmxByteTransportFactory { - let router: LivenessHostRouter - let box: TransportBox - - func makeTransport(for route: CmxAttachRoute) throws -> any CmxByteTransport { - let transport = LivenessTransport(router: router) - box.set(transport) - return transport - } -} - -private actor LivenessTransport: CmxByteTransport { - private let router: LivenessHostRouter - private var pendingFrames: [Data] = [] - private var receiveWaiters: [CheckedContinuation] = [] - private var isClosed = false - - init(router: LivenessHostRouter) { - self.router = router - } - - func connect() async throws {} - - func receive() async throws -> Data? { - if !pendingFrames.isEmpty { - return pendingFrames.removeFirst() - } - if isClosed { - return nil - } - return await withCheckedContinuation { continuation in - receiveWaiters.append(continuation) - } - } - - func send(_ data: Data) async throws { - var buffer = data - let payloads = try MobileSyncFrameCodec.decodeFrames(from: &buffer) - for payload in payloads { - let parsed = (try? JSONSerialization.jsonObject(with: payload)) as? [String: Any] - let method = parsed?["method"] as? String - let id = parsed?["id"] as? String - let topics = (parsed?["params"] as? [String: Any])?["topics"] as? [String] - await router.record(method: method, topics: topics) - // Answer each request concurrently so one held response cannot - // head-of-line block later RPCs, matching the Mac host's - // per-frame response tasks. - Task { [router, weak self] in - guard let response = await router.response(method: method, id: id) else { - return - } - await self?.deliver(response) - } - } - } - - func close() async { - isClosed = true - let waiters = receiveWaiters - receiveWaiters = [] - for waiter in waiters { - waiter.resume(returning: nil) - } - } - - /// Deliver a frame to the client's read loop. Also used by tests to push - /// unsolicited server-side event envelopes. - func deliver(_ frame: Data) { - if receiveWaiters.isEmpty { - pendingFrames.append(frame) - return - } - let waiter = receiveWaiters.removeFirst() - waiter.resume(returning: frame) - } -} - -// MARK: - Test helpers - -@MainActor -private final class OutputCollector { - private(set) var lines: [String] = [] - private var task: Task? - - func mount(store: MobileShellComposite, surfaceID: String) { - task = Task { @MainActor [weak self] in - for await data in store.terminalOutputStream(surfaceID: surfaceID) { - self?.lines.append(String(decoding: data, as: UTF8.self)) - } - } - } - - func unmount() { - task?.cancel() - task = nil - } -} - -private func makeTicket(clock: TestClock) throws -> CmxAttachTicket { - let route = try CmxAttachRoute( - id: "debug_loopback", - kind: .debugLoopback, - endpoint: .hostPort(host: "127.0.0.1", port: 56584) - ) - return try CmxAttachTicket( - workspaceID: "live-workspace", - terminalID: "live-terminal", - macDeviceID: "test-mac", - macDisplayName: "Test Mac", - routes: [route], - expiresAt: clock.now.addingTimeInterval(3600) - ) -} - -private func attachURL(for ticket: CmxAttachTicket) throws -> String { - let encoder = JSONEncoder() - encoder.dateEncodingStrategy = .iso8601 - let payload = try encoder.encode(ticket) - .base64EncodedString() - .replacingOccurrences(of: "+", with: "-") - .replacingOccurrences(of: "/", with: "_") - .replacingOccurrences(of: "=", with: "") - return "cmux-ios://attach?v=\(ticket.version)&payload=\(payload)" -} - -private func renderGridEventFrame(surfaceID: String, seq: UInt64, text: String) throws -> Data { - let frame = try MobileTerminalRenderGridFrame.fromPlainRows( - surfaceID: surfaceID, - stateSeq: seq, - columns: 16, - rows: 4, - text: text - ) - let envelope: [String: Any] = [ - "kind": "event", - "topic": "terminal.render_grid", - "payload": try frame.jsonObject(), - ] - return try MobileSyncFrameCodec.encodeFrame(JSONSerialization.data(withJSONObject: envelope)) -} - -/// Poll until `condition` is true, bounded at `attempts` x 10ms. Returns the -/// final value so tests can assert both presence and (bounded) absence. -@MainActor -private func pollUntil( - attempts: Int = 300, - _ condition: @MainActor () async -> Bool -) async throws -> Bool { - for _ in 0.. MobileShellComposite { - let runtime = LivenessTestRuntime( - transportFactory: LivenessTransportFactory(router: router, box: box), - now: { clock.now }, - livenessProbeTimeoutNanoseconds: probeTimeoutNanoseconds - ) - let store = MobileShellComposite.preview(runtime: runtime) - store.signIn() - let ticket = try makeTicket(clock: clock) - let connected = await store.connectPairingURL(try attachURL(for: ticket)) - #expect(connected, "scripted connect must succeed") - return store -} // MARK: - Tests From 98dbf2c53df3c75aae43f431001dc9624e002afc Mon Sep 17 00:00:00 2001 From: lawrencecchen <54008264+lawrencecchen@users.noreply.github.com> Date: Thu, 11 Jun 2026 18:11:27 -0700 Subject: [PATCH 09/10] Red test: stream ending before start ack must converge to unavailable Reproduces the ipad CI failure of macConnectionStatusMarksUnavailableWhenEventStreamCloses: with the start handshake made concurrent, a transport that drops before the subscribe ack lands lets the stream-end restart supersede the listener generation; the parked ack's failure verdict is then silently dropped by its generation guard and recovery loops in reconnecting forever instead of reaching unavailable. Co-Authored-By: Claude Fable 5 --- .../MobileShellRenderGridLivenessTests.swift | 33 +++++++++++++++++++ 1 file changed, 33 insertions(+) diff --git a/Packages/CmuxMobileShell/Tests/CmuxMobileShellTests/MobileShellRenderGridLivenessTests.swift b/Packages/CmuxMobileShell/Tests/CmuxMobileShellTests/MobileShellRenderGridLivenessTests.swift index a90d34c38505..f3a68e2726a0 100644 --- a/Packages/CmuxMobileShell/Tests/CmuxMobileShellTests/MobileShellRenderGridLivenessTests.swift +++ b/Packages/CmuxMobileShell/Tests/CmuxMobileShellTests/MobileShellRenderGridLivenessTests.swift @@ -210,3 +210,36 @@ import Testing ) await router.releaseAllHeld() } + +/// A transport that drops before the start handshake completes must converge +/// to `.unavailable`, not livelock in `.reconnecting`: without the guard, the +/// stream-end restart supersedes the listener generation, so the parked start +/// ack's failure verdict is silently dropped by its generation check and the +/// loop re-arms forever (observed as the ipad-only CI failure of +/// `macConnectionStatusMarksUnavailableWhenEventStreamCloses`, where the race +/// occasionally lands the other way on faster simulators). +@MainActor +@Test func streamEndingBeforeStartAckMarksUnavailable() async throws { + let clock = TestClock() + let router = LivenessHostRouter() + let box = TransportBox() + await router.setHoldSubscribe(true) + let store = try await makeConnectedStore(router: router, box: box, clock: clock) + defer { + Task { await router.releaseAllHeld() } + } + + let sawSubscribe = try await pollUntil { await router.count(of: "mobile.events.subscribe") >= 1 } + #expect(sawSubscribe, "listener must send the start subscribe") + + // The transport dies while the enable handshake is still parked. + let transport = try #require(box.get()) + await transport.close() + + let unavailable = try await pollUntil { store.macConnectionStatus == .unavailable } + #expect( + unavailable, + "a stream that ends before its subscribe ack must converge to unavailable, not loop reconnecting" + ) + #expect(store.connectionRecoveryFailed) +} From ab8f139565c4c055406f209b2c9f25a374f47777 Mon Sep 17 00:00:00 2001 From: lawrencecchen <54008264+lawrencecchen@users.noreply.github.com> Date: Thu, 11 Jun 2026 18:11:28 -0700 Subject: [PATCH 10/10] Treat a stream that ends before its subscribe ack as a failed start When the event stream ends while the generation's enable handshake is still in flight, stop and mark the Mac unavailable instead of restarting: the restart superseded the generation, swallowed the handshake's failure verdict, and livelocked reconnecting against a closed transport. A healthy stream cannot hit this path (it only ends when the transport drops), and a mid-life stream end with a completed handshake still takes the existing restart path. MobileShellComposite budget +15 for the guard. Co-Authored-By: Claude Fable 5 --- .github/swift-file-length-budget.tsv | 2 +- .../CmuxMobileShell/MobileShellComposite.swift | 15 +++++++++++++++ 2 files changed, 16 insertions(+), 1 deletion(-) diff --git a/.github/swift-file-length-budget.tsv b/.github/swift-file-length-budget.tsv index 423116151d0d..e50ead440ef4 100644 --- a/.github/swift-file-length-budget.tsv +++ b/.github/swift-file-length-budget.tsv @@ -21,7 +21,7 @@ 6071 Sources/TextBoxInput.swift 5482 cmuxTests/BrowserConfigTests.swift 5448 Sources/cmuxApp.swift -4812 Packages/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite.swift +4827 Packages/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite.swift 4460 Sources/Panels/FilePreviewPanel.swift 4400 cmuxTests/BrowserPanelTests.swift 4227 Sources/BrowserWindowPortal.swift diff --git a/Packages/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite.swift b/Packages/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite.swift index 3efb71da6df9..a3d574b611cb 100644 --- a/Packages/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite.swift +++ b/Packages/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite.swift @@ -3905,6 +3905,21 @@ public final class MobileShellComposite: MobileTerminalOutputSinking { connectionState == .connected else { return } + if terminalSubscriptionStartTask != nil { + // The stream ended while this generation's enable handshake was + // still in flight: the transport dropped before the subscription + // ever delivered. Restarting here would supersede the generation + // and silently swallow the handshake's failure verdict (its ack + // guard sees a newer listenerID), so a closed transport would + // loop `reconnecting` forever. Converge instead: a stream that + // dies before its handshake completes IS a failed start. + 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)) + stopTerminalRefreshPolling() + markMacConnectionUnavailable() + return + } mobileShellLog.info("terminal event stream ended, restarting") MobileDebugLog.anchormux("sync.stream_ended restarting (render-grid push stopped; falling back to poll)") diagnosticLog?.record(DiagnosticEvent(.streamEnded))