diff --git a/Packages/iOS/CmuxMobileDiagnostics/Sources/CmuxMobileDiagnostics/MobileDebugLogSink.swift b/Packages/iOS/CmuxMobileDiagnostics/Sources/CmuxMobileDiagnostics/MobileDebugLogSink.swift index defcf97b2dd9..d7e8e0ca83e8 100644 --- a/Packages/iOS/CmuxMobileDiagnostics/Sources/CmuxMobileDiagnostics/MobileDebugLogSink.swift +++ b/Packages/iOS/CmuxMobileDiagnostics/Sources/CmuxMobileDiagnostics/MobileDebugLogSink.swift @@ -91,16 +91,32 @@ public actor MobileDebugLogSink { /// file logging is configured, the same line is also appended to disk before /// this method returns. public func append(_ message: String) { + appendMessages(CollectionOfOne(message)) + } + + /// Append a batch while holding actor ownership across the whole write. + /// + /// Latency tracing uses this internal path so its single consumer can + /// forward a drained batch without spawning one task per trace line. + func appendBatch(_ messages: [String]) { + appendMessages(messages) + } + + private func appendMessages(_ messages: Messages) + where Messages.Element == String { + guard !messages.isEmpty else { return } let elapsed = String(format: "%9.3f", now().timeIntervalSince(startedAt)) - let line = "[\(elapsed)] \(message)" - buffer.append(line) + let lines = messages.map { "[\(elapsed)] \($0)" } + buffer.append(contentsOf: lines) if buffer.count > capacity { buffer.removeFirst(buffer.count - capacity) } - for continuation in continuations.values { - continuation.yield(line) + for line in lines { + for continuation in continuations.values { + continuation.yield(line) + } + appendToFile(line) } - appendToFile(line) } /// The full buffer as newline-joined text, newest last. diff --git a/Packages/iOS/CmuxMobileDiagnostics/Sources/CmuxMobileDiagnostics/MobileLatencyTrace.swift b/Packages/iOS/CmuxMobileDiagnostics/Sources/CmuxMobileDiagnostics/MobileLatencyTrace.swift new file mode 100644 index 000000000000..471d55aaab58 --- /dev/null +++ b/Packages/iOS/CmuxMobileDiagnostics/Sources/CmuxMobileDiagnostics/MobileLatencyTrace.swift @@ -0,0 +1,118 @@ +import Dispatch +import Foundation + +/// Low-overhead, opt-in latency stamps for DEBUG mobile builds. +public enum MobileLatencyTrace { + #if DEBUG + private static let writer = MobileLatencyTraceWriter(capacity: 4_096) + #endif + + /// Whether latency tracing is enabled for this process. + public static let isEnabled: Bool = { + #if DEBUG + ProcessInfo.processInfo.environment["CMUX_LATENCY_TRACE"] == "1" + || UserDefaults.standard.bool(forKey: "cmux.debug.latency-trace") + #else + false + #endif + }() + + /// Emits one machine-parseable latency stamp to the mobile file sink. + /// + /// - Parameters: + /// - stage: Stable stage token. + /// - fields: Integer or short-token fields, without terminal content. + @inline(__always) + public static func stamp( + _ stage: StaticString, + _ fields: @autoclosure () -> String = "" + ) { + #if DEBUG + guard isEnabled else { return } + write(stage, uptimeMicroseconds: nowUptimeMicroseconds(), fields: fields()) + #endif + } + + /// Captures the monotonic clock only when tracing is enabled. + @inline(__always) + public static func captureTime() -> UInt64? { + #if DEBUG + guard isEnabled else { return nil } + return nowUptimeMicroseconds() + #else + return nil + #endif + } + + /// Emits a stamp at a previously captured time when tracing is enabled. + /// + /// - Parameters: + /// - stage: Stable stage token. + /// - uptimeMicroseconds: Previously captured monotonic uptime. + /// - fields: Integer or short-token fields, without terminal content. + @inline(__always) + public static func stamp( + _ stage: StaticString, + at uptimeMicroseconds: UInt64, + _ fields: @autoclosure () -> String = "" + ) { + #if DEBUG + guard isEnabled else { return } + write(stage, uptimeMicroseconds: uptimeMicroseconds, fields: fields()) + #endif + } + + /// Returns elapsed microseconds from a captured trace start. + /// + /// - Parameter start: Previously captured monotonic uptime. + /// - Returns: Elapsed monotonic microseconds. + @inline(__always) + public static func elapsedMicroseconds(since start: UInt64) -> UInt64 { + #if DEBUG + nowUptimeMicroseconds() &- start + #else + 0 + #endif + } + + #if DEBUG + /// Emits a completion stamp for an optionally captured trace start. + /// + /// - Parameters: + /// - stage: Stable completion-stage token. + /// - start: Captured start, or `nil` when tracing was disabled. + /// - fields: Builds fields from the elapsed microseconds. + @inline(__always) + public static func stampElapsed( + _ stage: StaticString, + since start: UInt64?, + _ fields: (_ elapsedMicroseconds: UInt64) -> String + ) { + guard isEnabled else { return } + guard let start else { return } + let completionTime = nowUptimeMicroseconds() + write( + stage, + uptimeMicroseconds: completionTime, + fields: fields(completionTime &- start) + ) + } + + @inline(__always) + private static func nowUptimeMicroseconds() -> UInt64 { + // Simulator uptime is in the host Mac clock domain, so simulator and + // Mac stamps are directly comparable. A physical iPhone is not. + DispatchTime.now().uptimeNanoseconds / 1_000 + } + + @inline(__always) + private static func write( + _ stage: StaticString, + uptimeMicroseconds: UInt64, + fields: String + ) { + let suffix = fields.isEmpty ? "" : " \(fields)" + writer.enqueue("LAT \(stage) t=\(uptimeMicroseconds)\(suffix)") + } + #endif +} diff --git a/Packages/iOS/CmuxMobileDiagnostics/Sources/CmuxMobileDiagnostics/MobileLatencyTraceWriter.swift b/Packages/iOS/CmuxMobileDiagnostics/Sources/CmuxMobileDiagnostics/MobileLatencyTraceWriter.swift new file mode 100644 index 000000000000..8bac267f0851 --- /dev/null +++ b/Packages/iOS/CmuxMobileDiagnostics/Sources/CmuxMobileDiagnostics/MobileLatencyTraceWriter.swift @@ -0,0 +1,92 @@ +#if DEBUG +import Foundation +internal import os + +/// Fixed-capacity producer buffer for synchronous latency-trace call sites. +final class MobileLatencyTraceWriter: Sendable { + private struct State: Sendable { + var entries: [String?] + var head = 0 + var count = 0 + var droppedCount = 0 + var consumerStarted = false + + init(capacity: Int) { + entries = Array(repeating: nil, count: capacity) + } + + mutating func enqueue(_ line: String) { + guard count < entries.count else { + droppedCount += 1 + return + } + entries[(head + count) % entries.count] = line + count += 1 + } + + mutating func drain(maximumCount: Int) -> [String]? { + guard count > 0 || droppedCount > 0 else { return nil } + var lines: [String] = [] + lines.reserveCapacity(min(count, maximumCount) + (droppedCount > 0 ? 1 : 0)) + if droppedCount > 0 { + let uptimeMicroseconds = DispatchTime.now().uptimeNanoseconds / 1_000 + lines.append( + "LAT trace.dropped t=\(uptimeMicroseconds) n=\(droppedCount) side=ios" + ) + droppedCount = 0 + } + let drainedCount = min(count, maximumCount) + for _ in 0.. + private let signals: AsyncStream + private let signalContinuation: AsyncStream.Continuation + + init(capacity: Int) { + precondition(capacity > 0) + state = OSAllocatedUnfairLock(initialState: State(capacity: capacity)) + (signals, signalContinuation) = AsyncStream.makeStream( + bufferingPolicy: .bufferingNewest(1) + ) + } + + @inline(__always) + func enqueue(_ line: String) { + let shouldStartConsumer = state.withLock { state in + state.enqueue(line) + guard !state.consumerStarted else { return false } + state.consumerStarted = true + return true + } + if shouldStartConsumer { + Task.detached { [self] in + await consume() + } + } + signalContinuation.yield() + } + + private func consume() async { + for await _ in signals { + while let batch = state.withLock({ + $0.drain(maximumCount: Self.batchSize) + }) { + await MobileDebugLog.shared.sink.appendBatch(batch) + } + } + } +} +#endif diff --git a/Packages/iOS/CmuxMobileShell/Sources/CmuxMobileShell/MobileLatencyProbe.swift b/Packages/iOS/CmuxMobileShell/Sources/CmuxMobileShell/MobileLatencyProbe.swift new file mode 100644 index 000000000000..4d145e49f4dd --- /dev/null +++ b/Packages/iOS/CmuxMobileShell/Sources/CmuxMobileShell/MobileLatencyProbe.swift @@ -0,0 +1,53 @@ +#if DEBUG +import Foundation + +/// Process-scoped configuration and claim gate for the DEBUG typing probe. +@MainActor +enum MobileLatencyProbe { + struct Configuration { + let count: Int + let intervalMilliseconds: Int + } + + private static let configuration: Configuration? = { + guard let raw = ProcessInfo.processInfo.environment["CMUX_LATENCY_PROBE"] else { + return nil + } + if raw == "1" { + return Configuration(count: 40, intervalMilliseconds: 250) + } + let parts = raw.split(separator: ":", omittingEmptySubsequences: false) + guard parts.count == 2, + let count = Int(parts[0]), count > 0, + let intervalMilliseconds = Int(parts[1]), intervalMilliseconds > 0 else { + return nil + } + return Configuration(count: count, intervalMilliseconds: intervalMilliseconds) + }() + + private static var hasClaimedProcessRun = false + private static var hasAutoNavigatedProcessRun = false + + static var hasUnclaimedConfiguration: Bool { + configuration != nil && !hasClaimedProcessRun + } + + static func claimAutoNavigation() -> Bool { + guard hasUnclaimedConfiguration, !hasAutoNavigatedProcessRun else { + return false + } + hasAutoNavigatedProcessRun = true + return true + } + + static func claimConfiguration() -> Configuration? { + guard !hasClaimedProcessRun, let configuration else { return nil } + hasClaimedProcessRun = true + return configuration + } + + static func input(at index: Int) -> Data { + Data([UInt8(ascii: "a") + UInt8(index % 26)]) + } +} +#endif diff --git a/Packages/iOS/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite+LatencyProbe.swift b/Packages/iOS/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite+LatencyProbe.swift new file mode 100644 index 000000000000..feb7e52da0d2 --- /dev/null +++ b/Packages/iOS/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite+LatencyProbe.swift @@ -0,0 +1,78 @@ +#if DEBUG +import CmuxMobileDiagnostics +import Foundation + +extension MobileShellComposite { + func startLatencyProbeAutoNavigationIfNeeded() { + guard connectionState == .connected, + latencyProbeAutoNavigationTask == nil, + MobileLatencyProbe.hasUnclaimedConfiguration, + terminalOutputStreamTokensBySurfaceID.isEmpty, + deeplinkWorkspaceNavigationRequest == nil, + workspaces.contains(where: { !$0.terminals.isEmpty }) else { + return + } + latencyProbeAutoNavigationTask = Task { @MainActor [weak self] in + defer { self?.latencyProbeAutoNavigationTask = nil } + do { + try await Task.sleep(for: .seconds(1)) + guard let self, + self.connectionState == .connected, + self.terminalOutputStreamTokensBySurfaceID.isEmpty, + self.deeplinkWorkspaceNavigationRequest == nil, + let workspaceID = self.workspaces.first(where: { + !$0.terminals.isEmpty + })?.id, + MobileLatencyProbe.claimAutoNavigation() else { + return + } + self.navigateToWorkspaceForDeeplink(workspaceID, origin: .external) + } catch { + return + } + } + } + + func startLatencyProbeIfReady() { + guard connectionState == .connected, + latencyProbeTask == nil, + let surfaceID = terminalOutputStreamTokensBySurfaceID.keys.first, + let configuration = MobileLatencyProbe.claimConfiguration() else { + return + } + latencyProbeTask = Task { @MainActor [weak self] in + defer { self?.latencyProbeTask = nil } + do { + try await Task.sleep(for: .seconds(3)) + for index in 0.. 0 { terminalMirrorHydrationNeededSurfaceIDs.remove(renderGrid.surfaceID) } + #if DEBUG + MobileLatencyTrace.stamp( + "gate", + "s=\(renderGrid.surfaceID.prefix(8).lowercased()) " + + "seq=\(renderGrid.stateSeq) out=delivered" + ) + #endif } /// Whether a surface currently has an attached output stream consumer. @@ -235,13 +312,15 @@ extension MobileShellComposite { func deliverTerminalBytes( _ bytes: Data, surfaceID: String, + endSequence: UInt64? = nil, bypassReplayBarrier: Bool = false ) -> Bool { return deliverTerminalOutput( TerminalOutputDelivery( bytes: bytes, replaceable: false, - viewportPolicy: .natural + viewportPolicy: .natural, + endSequence: endSequence ), surfaceID: surfaceID, bypassReplayBarrier: bypassReplayBarrier @@ -369,6 +448,7 @@ extension MobileShellComposite { streamToken: streamToken, viewportPolicy: immediate.viewportPolicy, sourceRenderGridFrame: immediate.sourceRenderGridFrame, + endSequence: immediate.endSequence, requiresVerifiedReplay: requiresVerifiedReplayApplication(for: immediate), terminalConfigTheme: immediate.terminalConfigTheme ) @@ -477,6 +557,7 @@ extension MobileShellComposite { streamToken: streamToken, viewportPolicy: next.viewportPolicy, sourceRenderGridFrame: next.sourceRenderGridFrame, + endSequence: next.endSequence, requiresVerifiedReplay: requiresVerifiedReplayApplication(for: next), terminalConfigTheme: next.terminalConfigTheme )) @@ -529,6 +610,7 @@ extension MobileShellComposite { terminalAlternateRenderGridBaselineSurfaceIDs.remove(surfaceID) terminalFullReplacementSeqBySurfaceID.removeValue(forKey: surfaceID) terminalFullReplacementGenerationBySurfaceID.removeValue(forKey: surfaceID) + cancelTerminalInputAckResubscribeRetry(surfaceID: surfaceID) pendingTerminalByteEndSeqBySurfaceID.removeValue(forKey: surfaceID) pendingTerminalInputDroppedRenderGridSurfaceIDs.remove(surfaceID) terminalReplayBarrierAckStreamTokensBySurfaceID.removeValue(forKey: surfaceID) diff --git a/Packages/iOS/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite+TerminalReplayLifecycle.swift b/Packages/iOS/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite+TerminalReplayLifecycle.swift index ca375a378608..e36dcb691a9c 100644 --- a/Packages/iOS/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite+TerminalReplayLifecycle.swift +++ b/Packages/iOS/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite+TerminalReplayLifecycle.swift @@ -43,6 +43,7 @@ extension MobileShellComposite { } if let pendingSeq = pendingTerminalByteEndSeqBySurfaceID[surfaceID], endSeq >= pendingSeq { + cancelTerminalInputAckResubscribeRetry(surfaceID: surfaceID) pendingTerminalByteEndSeqBySurfaceID.removeValue(forKey: surfaceID) pendingTerminalInputDroppedRenderGridSurfaceIDs.remove(surfaceID) terminalReplayFailureRetryCountsBySurfaceID.removeValue(forKey: surfaceID) @@ -82,6 +83,7 @@ extension MobileShellComposite { // content under a barrier; only the surface-destroying resets clear it. terminalFullReplacementSeqBySurfaceID.removeValue(forKey: surfaceID) terminalFullReplacementGenerationBySurfaceID.removeValue(forKey: surfaceID) + cancelTerminalInputAckResubscribeRetry(surfaceID: surfaceID) pendingTerminalByteEndSeqBySurfaceID.removeValue(forKey: surfaceID) pendingTerminalInputDroppedRenderGridSurfaceIDs.remove(surfaceID) let token = UUID() @@ -267,6 +269,7 @@ extension MobileShellComposite { terminalRenderGridBaselineReplayBarrierTokensBySurfaceID.removeValue(forKey: surfaceID) terminalReplayBarrierTokensInFlightBySurfaceID.removeValue(forKey: surfaceID) restoreTerminalPreBarrierBaselineIfNeeded(surfaceID: surfaceID) + cancelTerminalInputAckResubscribeRetry(surfaceID: surfaceID) pendingTerminalByteEndSeqBySurfaceID.removeValue(forKey: surfaceID) pendingTerminalInputDroppedRenderGridSurfaceIDs.remove(surfaceID) MobileDebugLog.anchormux("terminal.output.replay_barrier_fail_open surface=\(surfaceID) reason=\(reason)") diff --git a/Packages/iOS/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite+TerminalReplayRetry.swift b/Packages/iOS/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite+TerminalReplayRetry.swift index 32ed8f192d3e..b450fb34b60c 100644 --- a/Packages/iOS/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite+TerminalReplayRetry.swift +++ b/Packages/iOS/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite+TerminalReplayRetry.swift @@ -34,6 +34,7 @@ extension MobileShellComposite { // Same fail-open invariant as failOpenTerminalReplayBarrier: once // retry budget is exhausted, the pending-input gate must not keep // live output suppressed forever. + cancelTerminalInputAckResubscribeRetry(surfaceID: surfaceID) pendingTerminalByteEndSeqBySurfaceID.removeValue(forKey: surfaceID) pendingTerminalInputDroppedRenderGridSurfaceIDs.remove(surfaceID) terminalReplayFailureRetryCountsBySurfaceID.removeValue(forKey: surfaceID) @@ -98,6 +99,7 @@ extension MobileShellComposite { resetRetryBudget: Bool = true ) { MobileDebugLog.anchormux("CMUX_REPLAY pending_input_fail_open surface=\(surfaceID) reason=\(reason)") + cancelTerminalInputAckResubscribeRetry(surfaceID: surfaceID) pendingTerminalByteEndSeqBySurfaceID.removeValue(forKey: surfaceID) pendingTerminalInputDroppedRenderGridSurfaceIDs.remove(surfaceID) if resetRetryBudget { diff --git a/Packages/iOS/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite.swift b/Packages/iOS/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite.swift index 2559a3049307..454939582942 100644 --- a/Packages/iOS/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite.swift +++ b/Packages/iOS/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite.swift @@ -123,6 +123,10 @@ public final class MobileShellComposite: MobileTerminalOutputSinking { /// reschedule per received event (an actively-streaming connection just keeps /// failing the silence check because `lastTerminalEventAt` stays fresh). private static let renderGridLivenessCheckInterval: TimeInterval = 2.5 + /// An input ACK only reasserts the event subscription after this long + /// without an event, preserving lost-registration recovery without adding a + /// control-lane round trip to healthy continuous typing. + static let terminalInputAckResubscribeSilenceThreshold: TimeInterval = 2 /// Short background dwells usually preserve the event stream; beyond this, /// the liveness watchdog and normal foreground resync own catch-up. static let foregroundResyncShortBackgroundThreshold: TimeInterval = 30 @@ -147,9 +151,16 @@ public final class MobileShellComposite: MobileTerminalOutputSinking { if connectionState == .connected { restartTerminalLanesForMountedSurfaces() scheduleWorkspaceChangesSummaryRefresh() + #if DEBUG + startLatencyProbeIfReady() + startLatencyProbeAutoNavigationIfNeeded() + #endif } else { deactivateAllTerminalLanes() resetWorkspaceChangesState() + #if DEBUG + cancelLatencyProbe() + #endif } // Intentional teardown (sign-out, hide, switch) must not look like // a network outage: swallow this edge and reset the throttle so a @@ -791,6 +802,11 @@ public final class MobileShellComposite: MobileTerminalOutputSinking { private var renderGridLivenessProbeID: UUID? private var renderGridLivenessConsecutiveProbeFailures = 0 var lastTerminalEventAt: Date? + @ObservationIgnored var terminalInputAckResubscribeRetryTask: Task? + @ObservationIgnored var terminalInputAckResubscribeRetryTaskID: UUID? + @ObservationIgnored var terminalInputAckResubscribeRetrySurfaceID: String? + /// Injected clock for the bounded ACK freshness-window follow-up. + @ObservationIgnored let terminalInputAckResubscribeClock: any Clock var lastBackgroundedAt: Date? var foregroundResumeEpoch: UInt64 = 0 private var terminalSubscriptionRefreshTask: Task? @@ -998,6 +1014,11 @@ public final class MobileShellComposite: MobileTerminalOutputSinking { private var terminalInputRPCPipeline: MobileTerminalInputRPCPipeline private var rawTerminalInputDrainWaiters: [CheckedContinuation] private var isRawTerminalInputDrainLoopRunning: Bool + #if DEBUG + var latencyProbeAutoNavigationTask: Task? + var latencyProbeTask: Task? + private var rawTerminalInputLatencyBatchNumber: UInt64 + #endif private var pairingAttemptID: UUID /// High-level shell phase derived from sign-in and connection state. @@ -1124,6 +1145,7 @@ public final class MobileShellComposite: MobileTerminalOutputSinking { groupCollapseStore: MobileWorkspaceGroupCollapseStore = MobileWorkspaceGroupCollapseStore(), workspaceChangesHintDismissalStore: MobileWorkspaceChangesHintDismissalStore = MobileWorkspaceChangesHintDismissalStore(), workspaceChangesSchedulingClock: any Clock = ContinuousClock(), + terminalInputAckResubscribeClock: any Clock = ContinuousClock(), taskTemplateStore: (any MobileTaskTemplateStoring)? = nil, storedMacReconnectRestoringDeadlineSeconds: Double = 15 ) { @@ -1132,6 +1154,7 @@ public final class MobileShellComposite: MobileTerminalOutputSinking { self.groupCollapseStore = groupCollapseStore self.workspaceChangesHintDismissalStore = workspaceChangesHintDismissalStore self.workspaceChangesSchedulingClock = workspaceChangesSchedulingClock + self.terminalInputAckResubscribeClock = terminalInputAckResubscribeClock self.taskTemplateStore = taskTemplateStore self.storedMacReconnectRestoringDeadlineSeconds = storedMacReconnectRestoringDeadlineSeconds self.pairedMacStore = pairedMacStore @@ -1192,6 +1215,9 @@ public final class MobileShellComposite: MobileTerminalOutputSinking { self.terminalEventListenerTask = nil self.terminalEventListenerID = nil self.terminalSubscriptionRefreshTask = nil + self.terminalInputAckResubscribeRetryTask = nil + self.terminalInputAckResubscribeRetryTaskID = nil + self.terminalInputAckResubscribeRetrySurfaceID = nil self.createWorkspaceTask = nil self.createWorkspaceTaskGroupID = nil self.createWorkspaceTaskSpec = nil @@ -1258,6 +1284,11 @@ public final class MobileShellComposite: MobileTerminalOutputSinking { self.terminalInputRPCPipeline = MobileTerminalInputRPCPipeline() self.rawTerminalInputDrainWaiters = [] self.isRawTerminalInputDrainLoopRunning = false + #if DEBUG + self.latencyProbeAutoNavigationTask = nil + self.latencyProbeTask = nil + self.rawTerminalInputLatencyBatchNumber = 0 + #endif self.pairingAttemptID = UUID() } @@ -1270,6 +1301,7 @@ public final class MobileShellComposite: MobileTerminalOutputSinking { terminalSubscriptionStartTask?.cancel() renderGridLivenessTimer?.cancel() renderGridLivenessProbeTask?.cancel() + terminalInputAckResubscribeRetryTask?.cancel() terminalSubscriptionRefreshTask?.cancel() createWorkspaceTask?.cancel() createTerminalTask?.cancel() @@ -1292,11 +1324,15 @@ public final class MobileShellComposite: MobileTerminalOutputSinking { } } - public static func preview(runtime: (any MobileSyncRuntime)? = nil) -> CMUXMobileShellStore { + public static func preview( + runtime: (any MobileSyncRuntime)? = nil, + terminalInputAckResubscribeClock: any Clock = ContinuousClock() + ) -> CMUXMobileShellStore { CMUXMobileShellStore( runtime: runtime, workspaces: PreviewMobileHost.workspaces, - deliveredNotificationClearer: NoopDeliveredNotificationClearer() + deliveredNotificationClearer: NoopDeliveredNotificationClearer(), + terminalInputAckResubscribeClock: terminalInputAckResubscribeClock ) } @@ -5402,10 +5438,22 @@ public final class MobileShellComposite: MobileTerminalOutputSinking { // Matches MobileIrohTerminalLane.maximumInputByteCount and the Mac lane // router's maximumInputFrameByteCount. while let chunk = rawTerminalInputBuffer.nextBatch(maximumByteCount: 16 * 1_024) { + #if DEBUG + rawTerminalInputLatencyBatchNumber &+= 1 + let latencyBatchNumber = rawTerminalInputLatencyBatchNumber + MobileLatencyTrace.stamp( + "in.send", + "n=\(latencyBatchNumber) bytes=\(chunk.text.utf8.count)" + ) + let latencyBatchNumberForSend: UInt64? = latencyBatchNumber + #else + let latencyBatchNumberForSend: UInt64? = nil + #endif await sendRemoteTerminalInput( chunk.text, workspaceID: chunk.workspaceID, - terminalID: chunk.terminalID + terminalID: chunk.terminalID, + latencyBatchNumber: latencyBatchNumberForSend ) } } @@ -6114,6 +6162,7 @@ public final class MobileShellComposite: MobileTerminalOutputSinking { clearMacUpdateHint() terminalSubscriptionRefreshTask?.cancel() terminalSubscriptionRefreshTask = nil + cancelTerminalInputAckResubscribeRetry() stopRenderGridLivenessWatchdog(listenerID: nil) lastTerminalEventAt = nil } @@ -6678,18 +6727,25 @@ public final class MobileShellComposite: MobileTerminalOutputSinking { #endif return } - await sendRemoteTerminalInput(text, workspaceID: workspaceID, terminalID: terminalID) + _ = await sendRemoteTerminalInput( + text, + workspaceID: workspaceID, + terminalID: terminalID + ) } + /// Sends terminal input and stamps its eventual settlement outcome. private func sendRemoteTerminalInput( _ text: String, workspaceID: MobileWorkspacePreview.ID, - terminalID: MobileTerminalPreview.ID + terminalID: MobileTerminalPreview.ID, + latencyBatchNumber: UInt64? = nil ) async { guard let client = remoteClient else { #if DEBUG mobileShellLog.info("skip remote terminal input remoteClient=0") #endif + Self.stampTerminalInputSettlement(latencyBatchNumber, succeeded: false) return } let generation = connectionGeneration @@ -6720,6 +6776,10 @@ public final class MobileShellComposite: MobileTerminalOutputSinking { // belong to the previous connection. guard generation == connectionGeneration, client === remoteClient else { + Self.stampTerminalInputSettlement( + latencyBatchNumber, + succeeded: false + ) return } if terminalInputRPCPipeline.hasAmbiguousFailure( @@ -6745,11 +6805,13 @@ public final class MobileShellComposite: MobileTerminalOutputSinking { } switch laneResult { case .sent: + Self.stampTerminalInputSettlement(latencyBatchNumber, succeeded: true) return case .failed: mobileShellLog.error( "independent terminal input failed surface=\(terminalID.rawValue, privacy: .public)" ) + Self.stampTerminalInputSettlement(latencyBatchNumber, succeeded: false) return case .unavailable: break @@ -6776,9 +6838,13 @@ public final class MobileShellComposite: MobileTerminalOutputSinking { ) }, settlementHandler: { [weak self, weak client] result in - guard let self, let client else { return } switch result { case let .success(responseData): + Self.stampTerminalInputSettlement( + latencyBatchNumber, + succeeded: true + ) + guard let self, let client else { return } guard self.isCurrentRemoteOperation( client: client, generation: generation @@ -6788,6 +6854,11 @@ public final class MobileShellComposite: MobileTerminalOutputSinking { surfaceID: terminalID.rawValue ) case let .failure(error): + Self.stampTerminalInputSettlement( + latencyBatchNumber, + succeeded: false + ) + guard let self, let client else { return } self.handleTerminalInputFailure( error, client: client, @@ -6796,11 +6867,13 @@ public final class MobileShellComposite: MobileTerminalOutputSinking { } } ) + return } catch { // A generation change mid-enqueue (pipeline clear) surfaces as // CancellationError; that is a benign teardown, not an // operational failure, regardless of whether the caller also // rotated connectionGeneration. + Self.stampTerminalInputSettlement(latencyBatchNumber, succeeded: false) if error is CancellationError { return } handleTerminalInputFailure( error, @@ -6820,9 +6893,14 @@ public final class MobileShellComposite: MobileTerminalOutputSinking { params: params ) ) - guard isCurrentRemoteOperation(client: client, generation: generation) else { return } + guard isCurrentRemoteOperation(client: client, generation: generation) else { + Self.stampTerminalInputSettlement(latencyBatchNumber, succeeded: false) + return + } handleTerminalInputResponse(responseData, surfaceID: terminalID.rawValue) + Self.stampTerminalInputSettlement(latencyBatchNumber, succeeded: true) } catch { + Self.stampTerminalInputSettlement(latencyBatchNumber, succeeded: false) handleTerminalInputFailure( error, client: client, @@ -6831,6 +6909,20 @@ public final class MobileShellComposite: MobileTerminalOutputSinking { } } + @inline(__always) + private static func stampTerminalInputSettlement( + _ latencyBatchNumber: UInt64?, + succeeded: Bool + ) { + #if DEBUG + guard let latencyBatchNumber else { return } + MobileLatencyTrace.stamp( + "in.settled", + "n=\(latencyBatchNumber) ok=\(succeeded ? 1 : 0)" + ) + #endif + } + private func terminalInputParameters( text: String, workspaceID: MobileWorkspacePreview.ID, @@ -6890,7 +6982,8 @@ public final class MobileShellComposite: MobileTerminalOutputSinking { /// Deliver a composed block to the Mac surface via `terminal.paste`: a /// bracketed paste (so multi-line text is inserted as one literal block) - /// followed by an optional submit key. Mirrors ``sendRemoteTerminalInput(_:workspaceID:terminalID:)`` + /// followed by an optional submit key. Mirrors + /// ``sendRemoteTerminalInput(_:workspaceID:terminalID:latencyBatchNumber:)`` /// but takes the dedicated paste path instead of the raw `terminal.input` /// path, which rewrites newlines to carriage returns. /// @@ -7329,6 +7422,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.cancelTerminalInputAckResubscribeRetry() self.recordTerminalEventStreamLiveness() self.markMacConnectionHealthy() if event.topic == "workspace.updated" { @@ -7606,17 +7700,9 @@ public final class MobileShellComposite: MobileTerminalOutputSinking { // 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.requestAuthoritativeTerminalResync( - surfaceID: surfaceID, - reason: "liveness_probe_repaired" - ) - } - // The same registration carries `workspace.updated` and - // `mobile.sync.delta`, so list changes emitted during the - // gap were missed too; repair through the mode-appropriate - // authoritative path. - self.repairMissedEventWindow() + self.repairLostTerminalEventSubscription( + reason: "liveness_probe_repaired" + ) } else { MobileDebugLog.anchormux("sync.liveness probe_ok silentMs=\(Int(silent * 1000))") } @@ -7722,6 +7808,12 @@ public final class MobileShellComposite: MobileTerminalOutputSinking { let remoteSeq = payload.terminalSeq else { return } + #if DEBUG + MobileLatencyTrace.stamp( + "in.resp", + "s=\(surfaceID.prefix(8).lowercased()) ack_seq=\(remoteSeq)" + ) + #endif let localSeq = deliveredTerminalByteEndSeqBySurfaceID[surfaceID] ?? 0 guard remoteSeq > localSeq else { return } let canRenderGridAdvancePendingSeq = terminalOutputTransport == .renderGrid @@ -7749,7 +7841,20 @@ public final class MobileShellComposite: MobileTerminalOutputSinking { } pendingTerminalByteEndSeqBySurfaceID[surfaceID] = targetSeq MobileDebugLog.anchormux("sync.input_seq_wait surface=\(surfaceID) local=\(localSeq) pending=\(targetSeq) remote=\(remoteSeq)") - refreshTerminalEventSubscription(reason: "input_seq_wait") + let now = runtime?.now() ?? Date() + if lastTerminalEventAt.map({ + now.timeIntervalSince($0) >= Self.terminalInputAckResubscribeSilenceThreshold + }) ?? true { + cancelTerminalInputAckResubscribeRetry() + refreshTerminalEventSubscription(reason: "input_seq_wait") + } else if let lastTerminalEventAt { + scheduleTerminalInputAckResubscribeRetry( + surfaceID: surfaceID, + pendingSeq: targetSeq, + lastTerminalEventAt: lastTerminalEventAt, + now: now + ) + } return } MobileDebugLog.anchormux("sync.input_seq_behind surface=\(surfaceID) local=\(localSeq) remote=\(remoteSeq)") @@ -7767,6 +7872,100 @@ public final class MobileShellComposite: MobileTerminalOutputSinking { ) } + /// Coalesces the freshness-guard follow-up into one cancellable delay. + private func scheduleTerminalInputAckResubscribeRetry( + surfaceID: String, + pendingSeq: UInt64, + lastTerminalEventAt: Date, + now: Date + ) { + cancelTerminalInputAckResubscribeRetry() + guard let listenerID = terminalEventListenerID, + terminalEventListenerTask != nil else { + return + } + let elapsed = max(0, now.timeIntervalSince(lastTerminalEventAt)) + let delay = max( + 0, + Self.terminalInputAckResubscribeSilenceThreshold - elapsed + ) + let taskID = UUID() + terminalInputAckResubscribeRetryTaskID = taskID + terminalInputAckResubscribeRetrySurfaceID = surfaceID + let clock = terminalInputAckResubscribeClock + terminalInputAckResubscribeRetryTask = Task { @MainActor [weak self] in + try? await clock.sleep(for: .seconds(delay)) + guard !Task.isCancelled, let self else { return } + guard self.terminalInputAckResubscribeRetryTaskID == taskID else { return } + self.terminalInputAckResubscribeRetryTask = nil + self.terminalInputAckResubscribeRetryTaskID = nil + self.terminalInputAckResubscribeRetrySurfaceID = nil + guard self.connectionState == .connected, + let client = self.remoteClient, + self.terminalEventListenerID == listenerID, + self.terminalEventListenerTask != nil, + self.lastTerminalEventAt == lastTerminalEventAt, + self.pendingTerminalByteEndSeqBySurfaceID[surfaceID] == pendingSeq, + (self.deliveredTerminalByteEndSeqBySurfaceID[surfaceID] ?? 0) < pendingSeq + else { + return + } + let ack = await self.requestTerminalEventSubscription( + client: client, + reason: "input_seq_wait_retry", + topics: self.terminalOutputTransport.eventTopics + ) + guard !Task.isCancelled, + self.connectionState == .connected, + self.remoteClient === client, + self.terminalEventListenerID == listenerID, + self.terminalEventListenerTask != nil, + self.lastTerminalEventAt == lastTerminalEventAt, + self.pendingTerminalByteEndSeqBySurfaceID[surfaceID] == pendingSeq, + (self.deliveredTerminalByteEndSeqBySurfaceID[surfaceID] ?? 0) < pendingSeq, + case .subscribed(let alreadySubscribed) = ack + else { + return + } + if alreadySubscribed == false { + self.repairLostTerminalEventSubscription( + reason: "input_seq_wait_retry" + ) + } else { + self.requestAuthoritativeTerminalResync( + surfaceID: surfaceID, + reason: "input_seq_wait_retry" + ) + } + } + } + + /// Repairs every collection carried by a host registration that was absent. + private func repairLostTerminalEventSubscription(reason: String) { + for surfaceID in terminalByteContinuationsBySurfaceID.keys { + requestAuthoritativeTerminalResync( + surfaceID: surfaceID, + reason: reason + ) + } + // The same registration carries `workspace.updated` and + // `mobile.sync.delta`, so list changes emitted during the gap were + // missed too; repair through the mode-appropriate authoritative path. + repairMissedEventWindow() + } + + /// Cancels the one-shot ACK retry, optionally only for its owning surface. + func cancelTerminalInputAckResubscribeRetry(surfaceID: String? = nil) { + if let surfaceID, + terminalInputAckResubscribeRetrySurfaceID != surfaceID { + return + } + terminalInputAckResubscribeRetryTask?.cancel() + terminalInputAckResubscribeRetryTask = nil + terminalInputAckResubscribeRetryTaskID = nil + terminalInputAckResubscribeRetrySurfaceID = nil + } + private static func terminalSnapshotReplacementBytes(_ snapshotBytes: Data) -> Data { var bytes = Data("\u{1B}c\u{1B}[H\u{1B}[2J\u{1B}[3J".utf8) bytes.append(snapshotBytes) @@ -7784,15 +7983,18 @@ public final class MobileShellComposite: MobileTerminalOutputSinking { terminalOutputQueuesBySurfaceID[surfaceID] = TerminalOutputDeliveryQueue() deliveredTerminalByteEndSeqBySurfaceID.removeValue(forKey: surfaceID) terminalPreBarrierDeliveredEndSeqBySurfaceID.removeValue(forKey: surfaceID) + terminalRenderGridHistoryContinuityBySurfaceID.removeValue(forKey: surfaceID) terminalRenderGridBaselineReplayRequestCountsBySurfaceID.removeValue(forKey: surfaceID) terminalRenderGridBaselineReplayBarrierTokensBySurfaceID.removeValue(forKey: surfaceID) terminalAlternateRenderGridBaselineSurfaceIDs.remove(surfaceID) terminalFullReplacementSeqBySurfaceID.removeValue(forKey: surfaceID) terminalFullReplacementGenerationBySurfaceID.removeValue(forKey: surfaceID) + cancelTerminalInputAckResubscribeRetry(surfaceID: surfaceID) pendingTerminalByteEndSeqBySurfaceID.removeValue(forKey: surfaceID) pendingTerminalInputDroppedRenderGridSurfaceIDs.remove(surfaceID) #if DEBUG mobileShellLog.info("CMUX_REPLAY register sink surface=\(surfaceID, privacy: .public) connected=\(self.connectionState == .connected, privacy: .public) hasClient=\(self.remoteClient != nil, privacy: .public) workspaceCount=\(self.workspaces.count, privacy: .public)") + startLatencyProbeIfReady() #endif requestColdAttachTerminalReplay(surfaceID: surfaceID) ensureTerminalLane(surfaceID: surfaceID) @@ -7842,6 +8044,7 @@ public final class MobileShellComposite: MobileTerminalOutputSinking { terminalAlternateRenderGridBaselineSurfaceIDs.remove(surfaceID) terminalFullReplacementSeqBySurfaceID.removeValue(forKey: surfaceID) terminalFullReplacementGenerationBySurfaceID.removeValue(forKey: surfaceID) + cancelTerminalInputAckResubscribeRetry(surfaceID: surfaceID) pendingTerminalByteEndSeqBySurfaceID.removeValue(forKey: surfaceID) pendingTerminalInputDroppedRenderGridSurfaceIDs.remove(surfaceID) terminalActiveScreenBySurfaceID.removeValue(forKey: surfaceID) @@ -8186,6 +8389,7 @@ public final class MobileShellComposite: MobileTerminalOutputSinking { } } self.recordTerminalRenderGridDelivery(renderGrid) + self.recordTerminalRenderGridHistoryContinuity(renderGrid) self.rebaseTerminalReplayStaleFloor(surfaceID: surfaceID) // A delivered grid is progress even if the payload omitted // its sequence; fall back to the frame's own sequence so @@ -8230,6 +8434,7 @@ public final class MobileShellComposite: MobileTerminalOutputSinking { let accepted = self.deliverTerminalBytes( deliverBytes, surfaceID: surfaceID, + endSequence: replaySeq, bypassReplayBarrier: replayBarrierTokenForRequest != nil ) if accepted, @@ -8331,6 +8536,9 @@ public final class MobileShellComposite: MobileTerminalOutputSinking { guard let json = event.payloadJSON else { return } + #if DEBUG + let latencyReceiveTime = MobileLatencyTrace.captureTime() + #endif // The frame may arrive nested under `render_grid` or as the bare payload; // try the wrapper first, then fall back to decoding the whole payload. let renderGridDTO = try? MobileTerminalRenderGridEvent.decode(json) @@ -8339,6 +8547,15 @@ public final class MobileShellComposite: MobileTerminalOutputSinking { return } #if DEBUG + if let latencyReceiveTime { + let decodeDuration = MobileLatencyTrace.elapsedMicroseconds(since: latencyReceiveTime) + MobileLatencyTrace.stamp( + "ev.grid", + at: latencyReceiveTime, + "s=\(renderGrid.surfaceID.prefix(8).lowercased()) seq=\(renderGrid.stateSeq) " + + "bytes=\(json.count) dec_us=\(decodeDuration)" + ) + } mobileShellLog.info("CMUX_REPLAY live render_grid surface=\(renderGrid.surfaceID, privacy: .public) full=\(renderGrid.full, privacy: .public) spans=\(renderGrid.rowSpans.count, privacy: .public) cleared=\(renderGrid.clearedRows.count, privacy: .public) seq=\(renderGrid.stateSeq, privacy: .public) hasSink=true") #endif deliverAuthoritativeTerminalRenderGrid(renderGrid, source: "event") @@ -8430,7 +8647,7 @@ public final class MobileShellComposite: MobileTerminalOutputSinking { b: Int(clamping: seq) )) mobileShellLog.info("terminal byte gap surface=\(surfaceID, privacy: .public) deliveredSeq=\(deliveredSeq, privacy: .public) nextSeq=\(seq, privacy: .public)") - guard deliverTerminalBytes(bytes, surfaceID: surfaceID) else { return } + guard deliverTerminalBytes(bytes, surfaceID: surfaceID, endSequence: endSeq) else { return } markTerminalBytesDelivered(surfaceID: surfaceID, endSeq: endSeq) if terminalReplaySurfaceIDsInFlight.contains(surfaceID) { cancelTerminalReplayInFlight(surfaceID: surfaceID) @@ -8447,7 +8664,7 @@ public final class MobileShellComposite: MobileTerminalOutputSinking { } let overlap = deliveredSeq - seq let deliverBytes = Data(bytes.dropFirst(Int(overlap))) - guard deliverTerminalBytes(deliverBytes, surfaceID: surfaceID) else { return } + guard deliverTerminalBytes(deliverBytes, surfaceID: surfaceID, endSequence: endSeq) else { return } markTerminalBytesDelivered(surfaceID: surfaceID, endSeq: endSeq) return } @@ -8461,12 +8678,12 @@ public final class MobileShellComposite: MobileTerminalOutputSinking { if seq < floorSeq { let overlap = floorSeq - seq let deliverBytes = Data(bytes.dropFirst(Int(overlap))) - guard deliverTerminalBytes(deliverBytes, surfaceID: surfaceID) else { return } + guard deliverTerminalBytes(deliverBytes, surfaceID: surfaceID, endSequence: endSeq) else { return } markTerminalBytesDelivered(surfaceID: surfaceID, endSeq: endSeq) return } } - guard deliverTerminalBytes(bytes, surfaceID: surfaceID) else { return } + guard deliverTerminalBytes(bytes, surfaceID: surfaceID, endSequence: endSeq) else { return } markTerminalBytesDelivered(surfaceID: surfaceID, endSeq: endSeq) } @@ -8565,6 +8782,7 @@ public final class MobileShellComposite: MobileTerminalOutputSinking { } func stopTerminalRefreshPolling() { + cancelTerminalInputAckResubscribeRetry() terminalEventListenerTask?.cancel() terminalEventListenerTask = nil terminalEventListenerID = nil @@ -8600,6 +8818,9 @@ public final class MobileShellComposite: MobileTerminalOutputSinking { ) setForegroundWorkspaceState( workspaces: remoteWorkspaces, groups: groups, merge: mergeExistingWorkspaces) + #if DEBUG + startLatencyProbeAutoNavigationIfNeeded() + #endif reconcileWorkspaceChangesSummaryStateWithForeground() let changesSummaryWorkspaceIDs = changesSummaryRefreshScope.workspaceIDs( fullSnapshotWorkspaceIDs: response.workspaces.map(\.id) diff --git a/Packages/iOS/CmuxMobileShell/Sources/CmuxMobileShell/TerminalOutputDelivery.swift b/Packages/iOS/CmuxMobileShell/Sources/CmuxMobileShell/TerminalOutputDelivery.swift index bf8d628d6220..b98fd32a164c 100644 --- a/Packages/iOS/CmuxMobileShell/Sources/CmuxMobileShell/TerminalOutputDelivery.swift +++ b/Packages/iOS/CmuxMobileShell/Sources/CmuxMobileShell/TerminalOutputDelivery.swift @@ -20,6 +20,7 @@ struct TerminalOutputDelivery: Equatable, Sendable { private var payload: Payload var replacementScope: ReplacementScope? var viewportPolicy: MobileTerminalOutputViewportPolicy? + var endSequence: UInt64? var replaceable: Bool { replacementScope != nil @@ -29,17 +30,20 @@ struct TerminalOutputDelivery: Equatable, Sendable { bytes: Data, replaceable: Bool, replacementScope: ReplacementScope? = nil, - viewportPolicy: MobileTerminalOutputViewportPolicy? = nil + viewportPolicy: MobileTerminalOutputViewportPolicy? = nil, + endSequence: UInt64? = nil ) { self.payload = .bytes(bytes) self.replacementScope = replaceable ? (replacementScope ?? .byteViewport) : nil self.viewportPolicy = viewportPolicy + self.endSequence = endSequence } init(theme frame: MobileTerminalRenderGridFrame) { self.payload = .theme(frame) self.replacementScope = .terminalTheme self.viewportPolicy = nil + self.endSequence = nil } init( @@ -51,6 +55,7 @@ struct TerminalOutputDelivery: Equatable, Sendable { self.payload = .renderGrid(frame) self.replacementScope = replaceable ? (replacementScope ?? .renderGridViewport) : nil self.viewportPolicy = viewportPolicy + self.endSequence = frame.stateSeq } var bytes: Data { diff --git a/Packages/iOS/CmuxMobileShell/Tests/CmuxMobileShellTests/InputAckRetryClock.swift b/Packages/iOS/CmuxMobileShell/Tests/CmuxMobileShellTests/InputAckRetryClock.swift new file mode 100644 index 000000000000..90cd653af078 --- /dev/null +++ b/Packages/iOS/CmuxMobileShell/Tests/CmuxMobileShellTests/InputAckRetryClock.swift @@ -0,0 +1,94 @@ +import Foundation + +/// Hand-advanced clock for the input-ACK resubscribe freshness delay. +/// +/// Safety: every mutable field is accessed under `lock`. +final class InputAckRetryClock: Clock, @unchecked Sendable { + struct Instant: InstantProtocol { + var offset: Duration + + func advanced(by duration: Duration) -> Instant { + Instant(offset: offset + duration) + } + + func duration(to other: Instant) -> Duration { + other.offset - offset + } + + static func < (lhs: Instant, rhs: Instant) -> Bool { + lhs.offset < rhs.offset + } + } + + private struct Sleeper { + let id: UUID + let deadline: Instant + let continuation: UnsafeContinuation + } + + private let lock = NSLock() + private var current = Instant(offset: .zero) + private var sleepers: [Sleeper] = [] + private var preCancelledIDs: Set = [] + + var now: Instant { + lock.withLock { current } + } + + var minimumResolution: Duration { .zero } + + var sleeperCount: Int { + lock.withLock { sleepers.count } + } + + func sleep(until deadline: Instant, tolerance _: Duration?) async throws { + let id = UUID() + try await withTaskCancellationHandler { + try await withUnsafeThrowingContinuation { + (continuation: UnsafeContinuation) in + lock.lock() + if preCancelledIDs.remove(id) != nil { + lock.unlock() + continuation.resume(throwing: CancellationError()) + } else if deadline <= current { + lock.unlock() + continuation.resume() + } else { + sleepers.append(Sleeper( + id: id, + deadline: deadline, + continuation: continuation + )) + lock.unlock() + } + } + } onCancel: { + cancelSleeper(id: id) + } + } + + func advance(by duration: Duration) { + lock.lock() + current = current.advanced(by: duration) + let due = sleepers + .filter { $0.deadline <= current } + .sorted { $0.deadline < $1.deadline } + sleepers.removeAll { $0.deadline <= current } + lock.unlock() + for sleeper in due { + sleeper.continuation.resume() + } + } + + private func cancelSleeper(id: UUID) { + lock.lock() + guard let index = sleepers.firstIndex(where: { $0.id == id }) else { + preCancelledIDs.insert(id) + lock.unlock() + return + } + let sleeper = sleepers.remove(at: index) + lock.unlock() + sleeper.continuation.resume(throwing: CancellationError()) + } +} diff --git a/Packages/iOS/CmuxMobileShell/Tests/CmuxMobileShellTests/MobileShellRenderGridEventFrameFixtures.swift b/Packages/iOS/CmuxMobileShell/Tests/CmuxMobileShellTests/MobileShellRenderGridEventFrameFixtures.swift index 33b0780a073c..77c1d90496e0 100644 --- a/Packages/iOS/CmuxMobileShell/Tests/CmuxMobileShellTests/MobileShellRenderGridEventFrameFixtures.swift +++ b/Packages/iOS/CmuxMobileShell/Tests/CmuxMobileShellTests/MobileShellRenderGridEventFrameFixtures.swift @@ -15,7 +15,10 @@ func renderGridEventFrame( columns: Int = 80, rows: Int = 4, activeScreen: MobileTerminalRenderGridFrame.Screen = .primary, - full: Bool = true + full: Bool = true, + anchor: MobileTerminalRenderGridFrame.Anchor = .viewport, + historyRows: UInt64? = nil, + deltaBaseHistoryRows: UInt64? = nil ) throws -> Data { let frame = try MobileTerminalRenderGridFrame( surfaceID: surfaceID, @@ -31,7 +34,10 @@ func renderGridEventFrame( text: text ), ], - activeScreen: activeScreen + activeScreen: activeScreen, + anchor: anchor, + historyRows: historyRows, + deltaBaseHistoryRows: deltaBaseHistoryRows ) let envelope: [String: Any] = [ "kind": "event", diff --git a/Packages/iOS/CmuxMobileShell/Tests/CmuxMobileShellTests/MobileShellRenderGridFrameTestSupport.swift b/Packages/iOS/CmuxMobileShell/Tests/CmuxMobileShellTests/MobileShellRenderGridFrameTestSupport.swift index 78a5aec0a882..660a30355eaa 100644 --- a/Packages/iOS/CmuxMobileShell/Tests/CmuxMobileShellTests/MobileShellRenderGridFrameTestSupport.swift +++ b/Packages/iOS/CmuxMobileShell/Tests/CmuxMobileShellTests/MobileShellRenderGridFrameTestSupport.swift @@ -10,7 +10,10 @@ func renderGridFrame( seq: UInt64, text: String, activeScreen: MobileTerminalRenderGridFrame.Screen = .primary, - full: Bool = true + full: Bool = true, + anchor: MobileTerminalRenderGridFrame.Anchor = .viewport, + historyRows: UInt64? = nil, + deltaBaseHistoryRows: UInt64? = nil ) throws -> MobileTerminalRenderGridFrame { try MobileTerminalRenderGridFrame( surfaceID: surfaceID, @@ -26,6 +29,9 @@ func renderGridFrame( text: text ), ], - activeScreen: activeScreen + activeScreen: activeScreen, + anchor: anchor, + historyRows: historyRows, + deltaBaseHistoryRows: deltaBaseHistoryRows ) } diff --git a/Packages/iOS/CmuxMobileShell/Tests/CmuxMobileShellTests/MobileShellRenderGridInputCatchUpTests.swift b/Packages/iOS/CmuxMobileShell/Tests/CmuxMobileShellTests/MobileShellRenderGridInputCatchUpTests.swift index ec51ea6834ee..c9bcbef0232b 100644 --- a/Packages/iOS/CmuxMobileShell/Tests/CmuxMobileShellTests/MobileShellRenderGridInputCatchUpTests.swift +++ b/Packages/iOS/CmuxMobileShell/Tests/CmuxMobileShellTests/MobileShellRenderGridInputCatchUpTests.swift @@ -5,10 +5,122 @@ import Testing @testable import CmuxMobileShell @MainActor -@Test func renderGridInputAcksDoNotReplayWhileWaitingForCatchUp() async throws { +@Test func renderGridInputAckRetriesSubscriptionWhenFreshnessWindowExpires() async throws { let clock = TestClock() + let retryClock = InputAckRetryClock() let router = LivenessHostRouter() await router.setCapabilities(["events.v1", "terminal.render_grid.v1", "terminal.replay.v1"]) + await router.enqueueTerminalInputSequences([100]) + let box = TransportBox() + let store = try await makeConnectedStore( + router: router, + box: box, + clock: clock, + inputAckRetryClock: retryClock + ) + + let collector = OutputCollector() + collector.mount(store: store, surfaceID: "live-terminal") + let transport = try #require(box.get()) + await transport.deliver(try renderGridEventFrame( + surfaceID: "live-terminal", + seq: 1, + text: "fresh-event" + )) + let freshEventDelivered = try await pollUntil { + collector.lines.contains { $0.contains("fresh-event") } + } + #expect(freshEventDelivered) + let subscribeCount = await router.count(of: "mobile.events.subscribe") + + await store.submitTerminalRawInput(Data("a".utf8), surfaceID: "live-terminal") + let retryArmed = try await pollUntil { retryClock.sleeperCount == 1 } + #expect(retryArmed, "a fresh ACK must arm one deferred subscription retry") + #expect(await router.count(of: "mobile.events.subscribe") == subscribeCount) + let replayCountBeforeRetry = await router.count(of: "mobile.terminal.replay") + await router.dropSubscription() + + clock.advance(by: MobileShellComposite.terminalInputAckResubscribeSilenceThreshold) + retryClock.advance( + by: .seconds(MobileShellComposite.terminalInputAckResubscribeSilenceThreshold) + ) + let retried = await router.waitForCount( + of: "mobile.events.subscribe", + atLeast: subscribeCount + 1 + ) + #expect(retried, "the deferred retry must re-subscribe when the freshness window expires") + let replayed = await router.waitForCount( + of: "mobile.terminal.replay", + atLeast: replayCountBeforeRetry + 1 + ) + #expect( + replayed, + "a retry that reinstalls a lost registration must request catch-up replay for the missed render-grid frame" + ) + collector.unmount() +} + +@MainActor +@Test func renderGridInputAckRetryIsCancelledByArrivingEvent() async throws { + let clock = TestClock() + let retryClock = InputAckRetryClock() + let router = LivenessHostRouter() + await router.setCapabilities(["events.v1", "terminal.render_grid.v1", "terminal.replay.v1"]) + await router.enqueueTerminalInputSequences([100]) + let box = TransportBox() + let store = try await makeConnectedStore( + router: router, + box: box, + clock: clock, + inputAckRetryClock: retryClock + ) + + let collector = OutputCollector() + collector.mount(store: store, surfaceID: "live-terminal") + let transport = try #require(box.get()) + await transport.deliver(try renderGridEventFrame( + surfaceID: "live-terminal", + seq: 1, + text: "fresh-event" + )) + let freshEventDelivered = try await pollUntil { + collector.lines.contains { $0.contains("fresh-event") } + } + #expect(freshEventDelivered) + let subscribeCount = await router.count(of: "mobile.events.subscribe") + + await store.submitTerminalRawInput(Data("a".utf8), surfaceID: "live-terminal") + let retryArmed = try await pollUntil { retryClock.sleeperCount == 1 } + #expect(retryArmed) + + await transport.deliver(try renderGridEventFrame( + surfaceID: "live-terminal", + seq: 2, + text: "later-event" + )) + let retryCancelled = try await pollUntil { retryClock.sleeperCount == 0 } + #expect(retryCancelled, "any arriving terminal event must cancel the deferred retry") + + clock.advance(by: MobileShellComposite.terminalInputAckResubscribeSilenceThreshold) + retryClock.advance( + by: .seconds(MobileShellComposite.terminalInputAckResubscribeSilenceThreshold) + ) + let extraSubscribe = await router.waitForCount( + of: "mobile.events.subscribe", + atLeast: subscribeCount + 1, + timeoutNanoseconds: 500_000_000, + recordIssueOnTimeout: false + ) + #expect(!extraSubscribe, "a cancelled retry must not issue an extra subscription") + collector.unmount() +} + +@MainActor +@Test func renderGridInputAcksOnlyResubscribeAfterEventStreamSilence() async throws { + let clock = TestClock() + let router = LivenessHostRouter() + await router.setCapabilities(["events.v1", "terminal.render_grid.v1", "terminal.replay.v1"]) + await router.enqueueTerminalInputSequences([100, 101]) let box = TransportBox() let store = try await makeConnectedStore(router: router, box: box, clock: clock) @@ -19,30 +131,54 @@ import Testing let replayCountAfterMount = await router.count(of: "mobile.terminal.replay") let subscribed = await router.waitForCount(of: "mobile.events.subscribe", atLeast: 1) #expect(subscribed, "connected render-grid transport must establish the event subscription") - let subscribeCountAfterMount = await router.count(of: "mobile.events.subscribe") + let transport = try #require(box.get()) + await transport.deliver(try renderGridEventFrame( + surfaceID: "live-terminal", + seq: 1, + text: "fresh-event" + )) + let freshEventDelivered = try await pollUntil { + collector.lines.contains { $0.contains("fresh-event") } + } + #expect(freshEventDelivered, "the active listener must record a fresh terminal event") + let subscribeCountAfterFreshEvent = await router.count(of: "mobile.events.subscribe") await store.submitTerminalRawInput(Data("a".utf8), surfaceID: "live-terminal") let firstInputSent = try await pollUntil { await router.count(of: "terminal.input") >= 1 } #expect(firstInputSent) - let firstRefreshSent = await router.waitForCount( + let freshRefreshSent = await router.waitForCount( of: "mobile.events.subscribe", - atLeast: subscribeCountAfterMount + 1 + atLeast: subscribeCountAfterFreshEvent + 1, + timeoutNanoseconds: 500_000_000, + recordIssueOnTimeout: false + ) + #expect( + !freshRefreshSent, + "an ahead-of-render-grid ACK must not resubscribe while terminal events are fresh" ) - #expect(firstRefreshSent, "the first ahead-of-render-grid ACK should refresh the event subscription") - let subscribeCountAfterFirstAck = await router.count(of: "mobile.events.subscribe") + + await transport.deliver(try renderGridEventFrame( + surfaceID: "live-terminal", + seq: 100, + text: "first-input-caught-up" + )) + let firstInputCaughtUp = try await pollUntil { + collector.lines.contains { $0.contains("first-input-caught-up") } + } + #expect(firstInputCaughtUp, "the target frame must settle the first input catch-up episode") + clock.advance(by: 2) + let subscribeCountBeforeStaleAck = await router.count(of: "mobile.events.subscribe") await store.submitTerminalRawInput(Data("b".utf8), surfaceID: "live-terminal") let inputSent = try await pollUntil { await router.count(of: "terminal.input") >= 2 } #expect(inputSent) - let duplicateRefreshSent = await router.waitForCount( + let staleRefreshSent = await router.waitForCount( of: "mobile.events.subscribe", - atLeast: subscribeCountAfterFirstAck + 1, - timeoutNanoseconds: 500_000_000, - recordIssueOnTimeout: false + atLeast: subscribeCountBeforeStaleAck + 1 ) #expect( - !duplicateRefreshSent, - "duplicate ACKs for the same pending sequence must not enqueue another subscription refresh" + staleRefreshSent, + "an ahead-of-render-grid ACK must resubscribe after two seconds of terminal-event silence" ) let replayRequested = try await pollUntil(attempts: 50) { diff --git a/Packages/iOS/CmuxMobileShell/Tests/CmuxMobileShellTests/MobileShellRenderGridLivenessTestSupport.swift b/Packages/iOS/CmuxMobileShell/Tests/CmuxMobileShellTests/MobileShellRenderGridLivenessTestSupport.swift index ff96cb597ca0..ad4bb10dbd8a 100644 --- a/Packages/iOS/CmuxMobileShell/Tests/CmuxMobileShellTests/MobileShellRenderGridLivenessTestSupport.swift +++ b/Packages/iOS/CmuxMobileShell/Tests/CmuxMobileShellTests/MobileShellRenderGridLivenessTestSupport.swift @@ -47,6 +47,7 @@ actor LivenessHostRouter { private var viewportRequestCount = 0 private var heldViewportRequestNumbers: Set = [] private var hasActiveSubscription = false + private var terminalInputSequences: [UInt64] = [] private var heldContinuations: [CheckedContinuation] = [] private var capabilities = ["events.v1", "terminal.bytes.v1", "terminal.render_grid.v1", "terminal.replay.v1"] // This router models the current authenticated Mac host by default. Tests @@ -228,6 +229,10 @@ actor LivenessHostRouter { emptyReplayResponsesRemaining += count } + func enqueueTerminalInputSequences(_ sequences: [UInt64]) { + terminalInputSequences.append(contentsOf: sequences) + } + /// Hold every `mobile.events.subscribe` response until released. func setHoldSubscribe(_ hold: Bool) { holdSubscribe = hold @@ -435,8 +440,11 @@ actor LivenessHostRouter { if let viewportReport = viewportEffectiveGridOverride ?? viewportReport { result["columns"] = viewportReport.columns; result["rows"] = viewportReport.rows } return try? Self.resultFrame(id: id, result: result) case "terminal.input": + let terminalSequence = terminalInputSequences.isEmpty + ? 100 + : terminalInputSequences.removeFirst() return try? Self.resultFrame(id: id, result: [ - "terminal_seq": 100, + "terminal_seq": terminalSequence, ]) default: return try? Self.errorFrame(id: id, message: "Unexpected method \(method ?? "nil")") @@ -664,14 +672,18 @@ func makeConnectedStore( router: LivenessHostRouter, box: TransportBox, clock: TestClock, - probeTimeoutNanoseconds: UInt64 = 200_000_000 + probeTimeoutNanoseconds: UInt64 = 200_000_000, + inputAckRetryClock: any Clock = ContinuousClock() ) async throws -> MobileShellComposite { let runtime = LivenessTestRuntime( transportFactory: LivenessTransportFactory(router: router, box: box), now: { clock.now }, livenessProbeTimeoutNanoseconds: probeTimeoutNanoseconds ) - let store = MobileShellComposite.preview(runtime: runtime) + let store = MobileShellComposite.preview( + runtime: runtime, + terminalInputAckResubscribeClock: inputAckRetryClock + ) store.signIn() let ticket = try makeTicket(clock: clock) let connected = await store.connectPairingURL(try attachURL(for: ticket)) diff --git a/Packages/iOS/CmuxMobileShell/Tests/CmuxMobileShellTests/MobileShellRenderGridReplayStalenessTests.swift b/Packages/iOS/CmuxMobileShell/Tests/CmuxMobileShellTests/MobileShellRenderGridReplayStalenessTests.swift index 7ef1fed0ed73..648397c7dd82 100644 --- a/Packages/iOS/CmuxMobileShell/Tests/CmuxMobileShellTests/MobileShellRenderGridReplayStalenessTests.swift +++ b/Packages/iOS/CmuxMobileShell/Tests/CmuxMobileShellTests/MobileShellRenderGridReplayStalenessTests.swift @@ -11,6 +11,136 @@ import Testing // same-sequence recovery replays must still apply. // Shared fixtures live in MobileShellRenderGridLivenessTestSupport.swift. +@MainActor +@Test(arguments: [UInt64(0), UInt64(3)]) +func screenAnchoredReplayBaselinesNextLiveDelta(historyRows: UInt64) async throws { + let clock = TestClock() + let router = LivenessHostRouter() + await router.setCapabilities(["events.v1", "terminal.render_grid.v1", "terminal.replay.v1"]) + try await router.enqueueReplayRenderGrid( + renderGridFrame( + surfaceID: "live-terminal", + seq: 10, + text: "replay-baseline-\(historyRows)", + anchor: .screen, + historyRows: historyRows + ) + ) + let box = TransportBox() + let store = try await makeConnectedStore(router: router, box: box, clock: clock) + + let collector = OutputCollector() + collector.mount(store: store, surfaceID: "live-terminal") + let replayDelivered = try await pollUntil { + collector.lines.contains { $0.contains("replay-baseline-\(historyRows)") } + } + #expect(replayDelivered, "the replay response must establish the visible baseline") + let replaySettled = try await pollUntil { + store.terminalReplayBarrierTokensBySurfaceID["live-terminal"] == nil + && !store.terminalReplaySurfaceIDsInFlight.contains("live-terminal") + } + #expect(replaySettled, "the replay barrier must clear before the live delta arrives") + let replayCountAfterBaseline = await router.count(of: "mobile.terminal.replay") + let transport = try #require(box.get()) + + await transport.deliver(try renderGridEventFrame( + surfaceID: "live-terminal", + seq: 11, + text: "live-delta-\(historyRows)", + full: false, + anchor: .screen, + historyRows: historyRows, + deltaBaseHistoryRows: historyRows + )) + + let deltaDelivered = try await pollUntil(attempts: 50) { + collector.lines.contains { $0.contains("live-delta-\(historyRows)") } + } + #expect( + deltaDelivered, + "a live screen-anchored delta linked to the replay history must be delivered" + ) + let replayRequested = await router.waitForCount( + of: "mobile.terminal.replay", + atLeast: replayCountAfterBaseline + 1, + timeoutNanoseconds: 500_000_000, + recordIssueOnTimeout: false + ) + #expect( + !replayRequested, + "a linked live delta must not request another replay" + ) + collector.unmount() +} + +@MainActor +@Test func screenAnchoredFrameWithoutHistoryBreaksReplayContinuity() async throws { + let historyRows = UInt64(3) + let clock = TestClock() + let router = LivenessHostRouter() + await router.setCapabilities(["events.v1", "terminal.render_grid.v1", "terminal.replay.v1"]) + try await router.enqueueReplayRenderGrid( + renderGridFrame( + surfaceID: "live-terminal", + seq: 10, + text: "replay-baseline", + anchor: .screen, + historyRows: historyRows + ) + ) + let box = TransportBox() + let store = try await makeConnectedStore(router: router, box: box, clock: clock) + + let collector = OutputCollector() + collector.mount(store: store, surfaceID: "live-terminal") + let replayDelivered = try await pollUntil { + collector.lines.contains { $0.contains("replay-baseline") } + } + #expect(replayDelivered, "the replay response must establish history continuity") + let replaySettled = try await pollUntil { + store.terminalReplayBarrierTokensBySurfaceID["live-terminal"] == nil + && !store.terminalReplaySurfaceIDsInFlight.contains("live-terminal") + } + #expect(replaySettled, "the replay barrier must clear before live frames arrive") + let transport = try #require(box.get()) + + await transport.deliver(try renderGridEventFrame( + surfaceID: "live-terminal", + seq: 11, + text: "full-without-history", + anchor: .screen + )) + let fullFrameDelivered = try await pollUntil { + collector.lines.contains { $0.contains("full-without-history") } + } + #expect(fullFrameDelivered, "the full live frame must replace the replay baseline") + let replayCountBeforeDelta = await router.count(of: "mobile.terminal.replay") + + await transport.deliver(try renderGridEventFrame( + surfaceID: "live-terminal", + seq: 12, + text: "stale-linked-delta", + full: false, + anchor: .screen, + historyRows: historyRows, + deltaBaseHistoryRows: historyRows + )) + + let replayRequested = await router.waitForCount( + of: "mobile.terminal.replay", + atLeast: replayCountBeforeDelta + 1 + ) + #expect( + replayRequested, + "a delta must request replay after a delivered frame loses history continuity" + ) + #expect( + !collector.lines.contains { $0.contains("stale-linked-delta") }, + "a delta linked to stale history must not be delivered" + ) + collector.unmount() +} + @MainActor @Test func renderGridReplayAtSameSeqDoesNotOverwriteNewerLiveGrid() async throws { let clock = TestClock() diff --git a/Packages/iOS/CmuxMobileShell/Tests/CmuxMobileShellTests/TerminalOutputDeliveryQueueTests.swift b/Packages/iOS/CmuxMobileShell/Tests/CmuxMobileShellTests/TerminalOutputDeliveryQueueTests.swift index 3d4a738296a4..23311e6a1c27 100644 --- a/Packages/iOS/CmuxMobileShell/Tests/CmuxMobileShellTests/TerminalOutputDeliveryQueueTests.swift +++ b/Packages/iOS/CmuxMobileShell/Tests/CmuxMobileShellTests/TerminalOutputDeliveryQueueTests.swift @@ -823,6 +823,7 @@ private func waitForReplayRequestCount( let latestTheme = TerminalOutputDelivery(theme: latestFrame) let latestViewport = TerminalOutputDelivery(renderGrid: latestFrame, replaceable: true) + #expect(latestTheme.endSequence == nil) #expect(queue.enqueue(inFlight) == inFlight) #expect(queue.enqueue(TerminalOutputDelivery(theme: oldFrame)) == nil) #expect(queue.enqueue(TerminalOutputDelivery(renderGrid: oldFrame, replaceable: true)) == nil) diff --git a/Packages/iOS/CmuxMobileShellModel/Sources/CmuxMobileShellModel/MobileTerminalOutputSinking.swift b/Packages/iOS/CmuxMobileShellModel/Sources/CmuxMobileShellModel/MobileTerminalOutputSinking.swift index f41543b6e00f..9ef194cc1b97 100644 --- a/Packages/iOS/CmuxMobileShellModel/Sources/CmuxMobileShellModel/MobileTerminalOutputSinking.swift +++ b/Packages/iOS/CmuxMobileShellModel/Sources/CmuxMobileShellModel/MobileTerminalOutputSinking.swift @@ -27,16 +27,29 @@ public struct MobileTerminalOutputChunk: Sendable { public let viewportPolicy: MobileTerminalOutputViewportPolicy? /// Source grid whose VT replay bytes are carried by this chunk. public let sourceRenderGridFrame: MobileTerminalRenderGridFrame? + /// Terminal byte high-water mark represented by this chunk, when known. + public let endSequence: UInt64? /// Whether nonempty output must pass render-grid verification before display. public let requiresVerifiedReplay: Bool /// Raw Ghostty defaults that must be installed before this chunk's VT replay. public let terminalConfigTheme: TerminalTheme? + /// Creates one backpressured terminal-output chunk. + /// + /// - Parameters: + /// - data: VT or PTY bytes to apply. + /// - streamToken: Identity of the mounted output stream. + /// - viewportPolicy: Optional viewport policy to apply with the bytes. + /// - sourceRenderGridFrame: Source grid represented by the bytes. + /// - endSequence: Terminal byte high-water mark represented by the chunk. + /// - requiresVerifiedReplay: Whether the verified replay path is required. + /// - terminalConfigTheme: Raw Ghostty defaults paired with the bytes. public init( data: Data, streamToken: UUID, viewportPolicy: MobileTerminalOutputViewportPolicy? = nil, sourceRenderGridFrame: MobileTerminalRenderGridFrame? = nil, + endSequence: UInt64? = nil, requiresVerifiedReplay: Bool = false, terminalConfigTheme: TerminalTheme? = nil ) { @@ -44,6 +57,7 @@ public struct MobileTerminalOutputChunk: Sendable { self.streamToken = streamToken self.viewportPolicy = viewportPolicy self.sourceRenderGridFrame = sourceRenderGridFrame + self.endSequence = endSequence self.requiresVerifiedReplay = requiresVerifiedReplay self.terminalConfigTheme = terminalConfigTheme } diff --git a/Packages/iOS/CmuxMobileShellUI/Sources/CmuxMobileShellUI/GhosttySurfaceRepresentable.swift b/Packages/iOS/CmuxMobileShellUI/Sources/CmuxMobileShellUI/GhosttySurfaceRepresentable.swift index 7bcde2556b6b..1467f4bda288 100644 --- a/Packages/iOS/CmuxMobileShellUI/Sources/CmuxMobileShellUI/GhosttySurfaceRepresentable.swift +++ b/Packages/iOS/CmuxMobileShellUI/Sources/CmuxMobileShellUI/GhosttySurfaceRepresentable.swift @@ -334,18 +334,52 @@ struct GhosttySurfaceRepresentable: UIViewRepresentable { guard !Task.isCancelled else { return } guard let self else { return } guard let surfaceView else { return } + #if DEBUG + let latencySequence = chunk.sourceRenderGridFrame?.stateSeq + ?? chunk.endSequence + ?? 0 + MobileLatencyTrace.stamp( + "ap.yield", + "s=\(surfaceID.prefix(8).lowercased()) seq=\(latencySequence)" + ) + let latencyApplyStart = MobileLatencyTrace.captureTime() + #endif switch terminalOutputApplicationPath( for: chunk, expectedSurfaceID: surfaceID ) { case .verifiedReplay: guard let frame = chunk.sourceRenderGridFrame else { return } - await self.applyVerifiedRenderGrid( + let applied = await self.applyVerifiedRenderGrid( frame, chunk: chunk, surfaceView: surfaceView, store: store ) + if applied { + #if DEBUG + MobileLatencyTrace.stampElapsed( + "ap.done", + since: latencyApplyStart + ) { + "s=\(surfaceID.prefix(8).lowercased()) seq=\(frame.stateSeq) " + + "path=verified us=\($0)" + } + // Verified replay has already submitted, read back, + // and revealed its tokened presentation before + // `applyVerifiedRenderGrid` returns. Stamp that exact + // frame here, after `ap.done`, instead of associating + // it with a later ordinary redraw. + MobileLatencyTrace.stamp( + "rd.present", + "s=\(surfaceID.prefix(8).lowercased()) seq=\(frame.stateSeq)" + ) + #endif + store.terminalOutputDidProcess( + surfaceID: surfaceID, + streamToken: chunk.streamToken + ) + } continue case .rejectUnverified: let transactionID = self.verifiedReplayState.rejectUnverifiedOutput() @@ -414,6 +448,16 @@ struct GhosttySurfaceRepresentable: UIViewRepresentable { continue } } + #if DEBUG + surfaceView.markLatencyAppliedSequence(latencySequence) + MobileLatencyTrace.stampElapsed( + "ap.done", + since: latencyApplyStart + ) { + "s=\(surfaceID.prefix(8).lowercased()) seq=\(latencySequence) " + + "path=legacy us=\($0)" + } + #endif store.terminalOutputDidProcess( surfaceID: surfaceID, streamToken: chunk.streamToken @@ -478,16 +522,16 @@ struct GhosttySurfaceRepresentable: UIViewRepresentable { chunk: MobileTerminalOutputChunk, surfaceView: GhosttySurfaceView, store: CMUXMobileShellStore - ) async { + ) async -> Bool { if let chunkConfigTheme = chunk.terminalConfigTheme, chunkConfigTheme != store.terminalConfigTheme(for: surfaceID) { store.terminalOutputDidReset( surfaceID: surfaceID, streamToken: chunk.streamToken ) - return + return false } - await applyThemeMatchedVerifiedRenderGrid( + return await applyThemeMatchedVerifiedRenderGrid( frame, chunk: chunk, surfaceView: surfaceView, @@ -500,33 +544,33 @@ struct GhosttySurfaceRepresentable: UIViewRepresentable { chunk: MobileTerminalOutputChunk, surfaceView: GhosttySurfaceView, store: CMUXMobileShellStore - ) async { + ) async -> Bool { guard case .apply(let transaction) = verifiedReplayState.begin(frame: frame) else { _ = await surfaceView.freezeVerifiedReplayPresentation( transactionID: frame.renderRevision ) - guard !Task.isCancelled else { return } + guard !Task.isCancelled else { return false } requestVerifiedReplayReset(transactionID: nil, chunk: chunk, store: store) - return + return false } let frozen = await surfaceView.freezeVerifiedReplayPresentation( transactionID: transaction.id ) - guard !Task.isCancelled else { return } + guard !Task.isCancelled else { return false } guard frozen else { requestVerifiedReplayReset(transactionID: transaction.id, chunk: chunk, store: store) - return + return false } activeViewportPolicy = .remoteGrid(columns: frame.columns, rows: frame.rows) let resized = await surfaceView.applyViewSizeAndWait( cols: frame.columns, rows: frame.rows ) - guard !Task.isCancelled else { return } + guard !Task.isCancelled else { return false } guard resized else { requestVerifiedReplayReset(transactionID: transaction.id, chunk: chunk, store: store) - return + return false } // Capture reads the post-reflow scrollbar, so Ghostty's resize pin @@ -550,10 +594,10 @@ struct GhosttySurfaceRepresentable: UIViewRepresentable { chunk.data, terminalConfigTheme: chunk.terminalConfigTheme ) - guard !Task.isCancelled else { return } + guard !Task.isCancelled else { return false } guard applied else { requestVerifiedReplayReset(transactionID: transaction.id, chunk: chunk, store: store) - return + return false } } @@ -562,8 +606,8 @@ struct GhosttySurfaceRepresentable: UIViewRepresentable { configuredCursorColor: chunk.terminalConfigTheme?.cursor ?? surfaceView.terminalConfigTheme.cursor ) - guard !Task.isCancelled else { return } - await finishVerifiedReplay( + guard !Task.isCancelled else { return false } + return await finishVerifiedReplay( transactionID: transaction.id, observed: observed, viewportAnchor: replayViewportAnchor, @@ -597,7 +641,7 @@ struct GhosttySurfaceRepresentable: UIViewRepresentable { chunk: MobileTerminalOutputChunk, surfaceView: GhosttySurfaceView, store: CMUXMobileShellStore - ) async { + ) async -> Bool { switch verifiedReplayState.complete( transactionID: transactionID, observedFrame: observed @@ -607,13 +651,13 @@ struct GhosttySurfaceRepresentable: UIViewRepresentable { let restored = await surfaceView.restoreVerifiedReplayViewportAnchor( viewportAnchor ) - guard !Task.isCancelled else { return } + guard !Task.isCancelled else { return false } if restored { pendingReplayViewportAnchor = nil // Restore and re-fence happen under render suppression, // so the renderer identity cannot change before reveal. _ = await surfaceView.presentRestoredVerifiedReplayViewport() - guard !Task.isCancelled else { return } + guard !Task.isCancelled else { return false } } } guard surfaceView.revealVerifiedReplayPresentation( @@ -624,17 +668,15 @@ struct GhosttySurfaceRepresentable: UIViewRepresentable { surfaceID: surfaceID, streamToken: chunk.streamToken ) - return + return false } - store.terminalOutputDidProcess( - surfaceID: surfaceID, - streamToken: chunk.streamToken - ) + return true case .keepFrozenAndRequestReplay, .ignoreStaleCompletion: store.terminalOutputDidReset( surfaceID: surfaceID, streamToken: chunk.streamToken ) + return false } } diff --git a/Packages/iOS/CmuxMobileTerminal/Sources/CmuxMobileTerminal/GhosttySurfaceView+LatencyTrace.swift b/Packages/iOS/CmuxMobileTerminal/Sources/CmuxMobileTerminal/GhosttySurfaceView+LatencyTrace.swift new file mode 100644 index 000000000000..260a1ed0a8be --- /dev/null +++ b/Packages/iOS/CmuxMobileTerminal/Sources/CmuxMobileTerminal/GhosttySurfaceView+LatencyTrace.swift @@ -0,0 +1,10 @@ +#if DEBUG && canImport(UIKit) +extension GhosttySurfaceView { + /// Associates subsequent render completions with the latest applied frame. + /// + /// - Parameter sequence: Applied terminal byte high-water mark. + public func markLatencyAppliedSequence(_ sequence: UInt64) { + latencyLastAppliedSequence = sequence + } +} +#endif diff --git a/Packages/iOS/CmuxMobileTerminal/Sources/CmuxMobileTerminal/GhosttySurfaceView.swift b/Packages/iOS/CmuxMobileTerminal/Sources/CmuxMobileTerminal/GhosttySurfaceView.swift index 40187ae685e7..f9e4f2c55678 100644 --- a/Packages/iOS/CmuxMobileTerminal/Sources/CmuxMobileTerminal/GhosttySurfaceView.swift +++ b/Packages/iOS/CmuxMobileTerminal/Sources/CmuxMobileTerminal/GhosttySurfaceView.swift @@ -391,6 +391,9 @@ public final class GhosttySurfaceView: UIView, TerminalSurfaceHosting { var surface: ghostty_surface_t? var surfaceGeneration: UInt64 = 0 + #if DEBUG + var latencyLastAppliedSequence: UInt64? + #endif private var lastReportedSize: TerminalGridSize? /// Latest natural grid awaiting a debounced report to the Mac. The display /// link sends it only after the grid has held steady for @@ -3194,6 +3197,21 @@ public final class GhosttySurfaceView: UIView, TerminalSurfaceHosting { DispatchQueue.main.async { guard let self else { return } guard self.surfaceGeneration == generation else { return } + #if DEBUG + if let surfaceID = self.hostSurfaceID { + if let sequence = self.latencyLastAppliedSequence { + MobileLatencyTrace.stamp( + "rd.present", + "s=\(surfaceID.prefix(8).lowercased()) seq=\(sequence)" + ) + } else { + MobileLatencyTrace.stamp( + "rd.present", + "s=\(surfaceID.prefix(8).lowercased())" + ) + } + } + #endif self.renderInFlight = false self.renderInFlightSince = nil guard !self.isDismantled else { diff --git a/Sources/App/HostLatencyTrace.swift b/Sources/App/HostLatencyTrace.swift new file mode 100644 index 000000000000..1d3be39a30b9 --- /dev/null +++ b/Sources/App/HostLatencyTrace.swift @@ -0,0 +1,171 @@ +#if DEBUG +import CMUXDebugLog +import Dispatch +import Foundation +import os + +/// Low-overhead, opt-in latency stamps for DEBUG host builds. +enum HostLatencyTrace { + private static let writer = HostLatencyTraceWriter(capacity: 4_096) + + static let isEnabled = + ProcessInfo.processInfo.environment["CMUX_LATENCY_TRACE"] == "1" + || UserDefaults.standard.bool(forKey: "cmux.debug.latency-trace") + + @inline(__always) + static func stamp( + _ stage: StaticString, + _ fields: @autoclosure () -> String = "" + ) { + guard isEnabled else { return } + write(stage, uptimeMicroseconds: nowUptimeMicroseconds(), fields: fields()) + } + + @inline(__always) + static func captureTime() -> UInt64? { + guard isEnabled else { return nil } + return nowUptimeMicroseconds() + } + + @inline(__always) + static func stampElapsed( + _ stage: StaticString, + since start: UInt64?, + _ fields: (_ elapsedMicroseconds: UInt64) -> String + ) { + guard let start else { return } + let completionTime = nowUptimeMicroseconds() + write( + stage, + uptimeMicroseconds: completionTime, + fields: fields(completionTime &- start) + ) + } + + @inline(__always) + private static func nowUptimeMicroseconds() -> UInt64 { + // Simulator uptime is in the host Mac clock domain, so simulator and + // Mac stamps are directly comparable. A physical iPhone is not. + DispatchTime.now().uptimeNanoseconds / 1_000 + } + + @inline(__always) + private static func write( + _ stage: StaticString, + uptimeMicroseconds: UInt64, + fields: String + ) { + let suffix = fields.isEmpty ? "" : " \(fields)" + writer.enqueue("LAT \(stage) t=\(uptimeMicroseconds)\(suffix)") + } + + /// Fixed-capacity producer buffer for synchronous latency-trace call sites. + private final class HostLatencyTraceWriter: Sendable { + private struct State: Sendable { + var entries: [String?] + var head = 0 + var count = 0 + var droppedCount = 0 + var consumerStarted = false + + init(capacity: Int) { + entries = Array(repeating: nil, count: capacity) + } + + mutating func enqueue(_ line: String) { + guard count < entries.count else { + droppedCount += 1 + return + } + entries[(head + count) % entries.count] = line + count += 1 + } + + mutating func drain(maximumCount: Int) -> [String]? { + guard count > 0 || droppedCount > 0 else { return nil } + var lines: [String] = [] + lines.reserveCapacity( + min(count, maximumCount) + (droppedCount > 0 ? 1 : 0) + ) + if droppedCount > 0 { + let uptimeMicroseconds = DispatchTime.now().uptimeNanoseconds / 1_000 + lines.append( + "LAT trace.dropped t=\(uptimeMicroseconds) n=\(droppedCount) side=mac" + ) + droppedCount = 0 + } + let drainedCount = min(count, maximumCount) + for _ in 0.. + private let signals: AsyncStream + private let signalContinuation: AsyncStream.Continuation + + init(capacity: Int) { + precondition(capacity > 0) + state = OSAllocatedUnfairLock(initialState: State(capacity: capacity)) + (signals, signalContinuation) = AsyncStream.makeStream( + bufferingPolicy: .bufferingNewest(1) + ) + } + + @inline(__always) + func enqueue(_ line: String) { + let shouldStartConsumer = state.withLock { state in + state.enqueue(line) + guard !state.consumerStarted else { return false } + state.consumerStarted = true + return true + } + if shouldStartConsumer { + Task.detached { [self] in + await consume() + } + } + signalContinuation.yield() + } + + private func consume() async { + for await _ in signals { + while let batch = state.withLock({ + $0.drain(maximumCount: Self.batchSize) + }) { + appendSynchronously(batch) + } + } + } + + private func appendSynchronously(_ batch: [String]) { + guard let data = (batch.joined(separator: "\n") + "\n").data(using: .utf8) else { + return + } + let fileDescriptor = open( + CMUXDebugLog.DebugEventLog.currentLogPath(), + O_WRONLY | O_APPEND | O_CREAT, + 0o644 + ) + guard fileDescriptor >= 0 else { return } + let handle = FileHandle(fileDescriptor: fileDescriptor, closeOnDealloc: true) + do { + try handle.write(contentsOf: data) + try handle.close() + } catch { + try? handle.close() + } + } + } +} +#endif diff --git a/Sources/Mobile/MobileHostConnectionEventQueue.swift b/Sources/Mobile/MobileHostConnectionEventQueue.swift index 4c47ac8d02bc..be1dd3023fed 100644 --- a/Sources/Mobile/MobileHostConnectionEventQueue.swift +++ b/Sources/Mobile/MobileHostConnectionEventQueue.swift @@ -50,12 +50,15 @@ struct MobileHostEventEnqueueResult: Sendable { /// Surfaces whose queued render-grid frames were shed; the caller must ask /// the producer for a full-frame resync of each. let renderGridResyncSurfaceIDs: Set + /// Queue depth immediately after an admitted append. + let depthAfterEnqueue: Int? static let rejected = MobileHostEventEnqueueResult( admitted: false, startDrain: false, shouldClose: false, - renderGridResyncSurfaceIDs: [] + renderGridResyncSurfaceIDs: [], + depthAfterEnqueue: nil ) } @@ -76,6 +79,7 @@ final class MobileHostConnectionEventQueue: @unchecked Sendable { let topic: String let coalesceKey: String? let frame: Data + let stateSeq: UInt64? } static let defaultMaximumEventCount = 256 @@ -140,6 +144,7 @@ final class MobileHostConnectionEventQueue: @unchecked Sendable { topic: String, coalesceKey: String?, isFullRenderGridFrame: Bool, + stateSeq: UInt64? = nil, frame: Data ) -> MobileHostEventEnqueueResult { lock.lock() @@ -172,7 +177,8 @@ final class MobileHostConnectionEventQueue: @unchecked Sendable { admitted: false, startDrain: false, shouldClose: false, - renderGridResyncSurfaceIDs: resyncSurfaceIDs + renderGridResyncSurfaceIDs: resyncSurfaceIDs, + depthAfterEnqueue: nil ) } guard hasRoomLocked(for: frame) else { @@ -182,7 +188,8 @@ final class MobileHostConnectionEventQueue: @unchecked Sendable { admitted: false, startDrain: false, shouldClose: true, - renderGridResyncSurfaceIDs: resyncSurfaceIDs + renderGridResyncSurfaceIDs: resyncSurfaceIDs, + depthAfterEnqueue: nil ) } if isRenderGrid, let coalesceKey { @@ -199,11 +206,20 @@ final class MobileHostConnectionEventQueue: @unchecked Sendable { admitted: false, startDrain: false, shouldClose: false, - renderGridResyncSurfaceIDs: resyncSurfaceIDs + renderGridResyncSurfaceIDs: resyncSurfaceIDs, + depthAfterEnqueue: nil ) } - queuedEvents.append(QueuedEvent(topic: topic, coalesceKey: coalesceKey, frame: frame)) + queuedEvents.append( + QueuedEvent( + topic: topic, + coalesceKey: coalesceKey, + frame: frame, + stateSeq: stateSeq + ) + ) queuedByteCount += frame.count + let depthAfterEnqueue = queuedEvents.count if isRenderGrid, isFullRenderGridFrame, let coalesceKey { poisonedRenderGridSurfaceIDs.remove(coalesceKey) resyncAfterDrainSurfaceIDs.remove(coalesceKey) @@ -217,7 +233,8 @@ final class MobileHostConnectionEventQueue: @unchecked Sendable { admitted: true, startDrain: startDrain, shouldClose: false, - renderGridResyncSurfaceIDs: resyncSurfaceIDs + renderGridResyncSurfaceIDs: resyncSurfaceIDs, + depthAfterEnqueue: depthAfterEnqueue ) } diff --git a/Sources/Mobile/MobileHostService.swift b/Sources/Mobile/MobileHostService.swift index e087a643e17d..7f3eb3f78670 100644 --- a/Sources/Mobile/MobileHostService.swift +++ b/Sources/Mobile/MobileHostService.swift @@ -443,7 +443,8 @@ final class MobileHostService { /// synchronous bounded queues as every other event. nonisolated static func emitRenderGridEvent( framesByAnchor: [MobileTerminalRenderGridFrame.Anchor: (payloadJSON: Data, isFullFrame: Bool)], - surfaceID: String + surfaceID: String, + stateSeq: UInt64 ) { let topic = MobileHostEventTopicPolicy.renderGridTopic guard !framesByAnchor.isEmpty, @@ -462,7 +463,7 @@ final class MobileHostService { encodedByAnchor[anchor] = (frame, item.isFullFrame) } guard !encodedByAnchor.isEmpty else { return } - deliverEventFrames(topic: topic, coalesceKey: surfaceID) { connection in + deliverEventFrames(topic: topic, coalesceKey: surfaceID, stateSeq: stateSeq) { connection in encodedByAnchor[ MobileTerminalRenderGridAnchorRegistry.shared.anchor(connectionID: connection.connectionID) ] @@ -507,7 +508,7 @@ final class MobileHostService { coalesceKey: String?, isFullRenderGridFrame: Bool ) { - deliverEventFrames(topic: topic, coalesceKey: coalesceKey) { _ in + deliverEventFrames(topic: topic, coalesceKey: coalesceKey, stateSeq: nil) { _ in (frame, isFullRenderGridFrame) } } @@ -521,6 +522,7 @@ final class MobileHostService { nonisolated private static func deliverEventFrames( topic: String, coalesceKey: String?, + stateSeq: UInt64?, frameFor: (MobileHostConnection) -> (frame: Data, isFullRenderGridFrame: Bool)? ) { let connections = MobileHostConnectionRegistry.shared.snapshot() @@ -535,8 +537,22 @@ final class MobileHostService { item.frame, topic: topic, coalesceKey: coalesceKey, - isFullRenderGridFrame: item.isFullRenderGridFrame + isFullRenderGridFrame: item.isFullRenderGridFrame, + stateSeq: stateSeq ) + #if DEBUG + if let stateSeq, + let surfaceID = coalesceKey, + result.admitted, + let depth = result.depthAfterEnqueue { + HostLatencyTrace.stamp( + "host.enq", + "s=\(surfaceID.prefix(8).lowercased()) " + + "conn=\(connection.connectionID.uuidString.prefix(8).lowercased()) " + + "seq=\(stateSeq) depth=\(depth)" + ) + } + #endif resyncSurfaceIDs.formUnion(result.renderGridResyncSurfaceIDs) if result.startDrain { Task { await connection.drainQueuedEvents() } @@ -2492,6 +2508,7 @@ actor MobileHostConnection { coalesceKey: MobileHostService.eventCoalesceKey(topic: topic, payload: payload), isFullRenderGridFrame: topic == MobileHostEventTopicPolicy.renderGridTopic && payload["full"] as? Bool == true, + stateSeq: nil, frame: frame ) if !result.renderGridResyncSurfaceIDs.isEmpty { @@ -2526,12 +2543,14 @@ actor MobileHostConnection { _ frame: Data, topic: String, coalesceKey: String?, - isFullRenderGridFrame: Bool + isFullRenderGridFrame: Bool, + stateSeq: UInt64? ) -> MobileHostEventEnqueueResult { eventQueue.enqueue( topic: topic, coalesceKey: coalesceKey, isFullRenderGridFrame: isFullRenderGridFrame, + stateSeq: stateSeq, frame: frame ) } @@ -2589,10 +2608,26 @@ actor MobileHostConnection { return } guard eventQueue.isSubscribed(topic: event.topic) else { continue } + #if DEBUG + let latencyWriteStart = event.stateSeq == nil ? nil : HostLatencyTrace.captureTime() + #endif guard await deliverQueuedEvent(event) else { eventQueue.abandonDrain() return } + #if DEBUG + if let stateSeq = event.stateSeq, + let surfaceID = event.coalesceKey { + HostLatencyTrace.stampElapsed( + "host.write", + since: latencyWriteStart + ) { + "s=\(surfaceID.prefix(8).lowercased()) " + + "conn=\(id.uuidString.prefix(8).lowercased()) " + + "seq=\(stateSeq) us=\($0)" + } + } + #endif let resyncSurfaceIDs = eventQueue.takeResyncAfterDrainRequests() if !resyncSurfaceIDs.isEmpty { MobileTerminalRenderObserver.requestRenderGridFullResync( diff --git a/Sources/Mobile/MobileStateSync.swift b/Sources/Mobile/MobileStateSync.swift index ced29bbaa770..37a0708644a8 100644 --- a/Sources/Mobile/MobileStateSync.swift +++ b/Sources/Mobile/MobileStateSync.swift @@ -97,6 +97,13 @@ final class MobileStateSyncHost { removedIDs: change.removedIDs ) guard let payload = try? MobileSyncFrameCoder().jsonObject(from: event) else { return } + #if DEBUG + HostLatencyTrace.stamp( + "host.sync.emit", + "coll=\(collection.rawValue) rev=\(change.toRev) " + + "rows=\(change.records.count + change.removedIDs.count)" + ) + #endif MobileHostService.shared.emitEvent(topic: Self.deltaTopic, payload: payload) } diff --git a/Sources/Mobile/MobileTerminalByteTee.swift b/Sources/Mobile/MobileTerminalByteTee.swift index 5184982878d5..16246d26c369 100644 --- a/Sources/Mobile/MobileTerminalByteTee.swift +++ b/Sources/Mobile/MobileTerminalByteTee.swift @@ -179,6 +179,12 @@ final class MobileTerminalByteTee { state.replayBuffer.removeFirst(state.replayBuffer.count - replayBudget) } statesBySurfaceID[surfaceID] = state + #if DEBUG + HostLatencyTrace.stamp( + "host.tee", + "s=\(surfaceID.uuidString.prefix(8).lowercased()) seq=\(state.seq) bytes=\(data.count)" + ) + #endif MobileTerminalRenderObserver.shared.noteTerminalBytes(surfaceID: surfaceID) if let continuations = laneContinuationsBySurfaceID[surfaceID] { diff --git a/Sources/Mobile/MobileTerminalRenderObserver.swift b/Sources/Mobile/MobileTerminalRenderObserver.swift index f1ce3f531d99..4e5aa7f50f6c 100644 --- a/Sources/Mobile/MobileTerminalRenderObserver.swift +++ b/Sources/Mobile/MobileTerminalRenderObserver.swift @@ -209,6 +209,9 @@ final class MobileTerminalRenderObserver { } private func flushTerminalUpdates() { + #if DEBUG + HostLatencyTrace.stamp("host.flush", "pending=\(pendingSurfaceIDs.count)") + #endif isEmitFlushScheduled = false guard hasAnyRenderEventSubscribers else { refreshNotificationDemand() @@ -301,6 +304,9 @@ final class MobileTerminalRenderObserver { var surfaceIDString: String? for anchor in anchors { + #if DEBUG + let latencyExportStart = HostLatencyTrace.captureTime() + #endif guard let emitted = emitRenderGridFrame( surface: surface, surfaceID: surfaceID, @@ -312,6 +318,16 @@ final class MobileTerminalRenderObserver { sharedTheme: &sharedTheme ) else { continue } guard let payloadJSON = try? JSONEncoder().encode(emitted) else { continue } + #if DEBUG + HostLatencyTrace.stampElapsed( + "host.grid", + since: latencyExportStart + ) { + "s=\(surfaceID.uuidString.prefix(8).lowercased()) seq=\(emitted.stateSeq) " + + "exp_us=\($0) bytes=\(payloadJSON.count) " + + "kind=\(emitted.full ? "full" : "delta")" + } + #endif framesByAnchor[anchor] = (payloadJSON, emitted.full) emittedByAnchor[anchor] = emitted surfaceIDString = emitted.surfaceID @@ -319,7 +335,8 @@ final class MobileTerminalRenderObserver { guard !framesByAnchor.isEmpty, let surfaceIDString else { return } MobileHostService.emitRenderGridEvent( framesByAnchor: framesByAnchor, - surfaceID: surfaceIDString + surfaceID: surfaceIDString, + stateSeq: stateSeq ) #if DEBUG for (anchor, frame) in emittedByAnchor { diff --git a/Sources/Mobile/MobileWorkspaceListObserver.swift b/Sources/Mobile/MobileWorkspaceListObserver.swift index f6cbbf3e6d53..0b7d17d5b4fc 100644 --- a/Sources/Mobile/MobileWorkspaceListObserver.swift +++ b/Sources/Mobile/MobileWorkspaceListObserver.swift @@ -312,6 +312,9 @@ final class MobileWorkspaceListObserver { } private func emitIfNeeded(force: Bool) { + #if DEBUG + HostLatencyTrace.stamp("host.sync.observe") + #endif let signpost = MobileWorkspaceObserverSignposts.begin("mobile-workspace-emit-if-needed", "force=\(force)"); defer { MobileWorkspaceObserverSignposts.end(signpost) } guard let tabManager else { return } let hash = Self.summaryHash( diff --git a/Sources/TerminalController.swift b/Sources/TerminalController.swift index cc81aa6a724c..ea09ca7083ad 100644 --- a/Sources/TerminalController.swift +++ b/Sources/TerminalController.swift @@ -14884,6 +14884,12 @@ class TerminalController { let terminalPanel = resolved.workspace.terminalInputTarget(forPanelID: surfaceId)?.panel else { return .err(code: "not_found", message: "Terminal surface not found", data: nil) } + #if DEBUG + HostLatencyTrace.stamp( + "host.in.recv", + "s=\(surfaceId.uuidString.prefix(8).lowercased()) bytes=\(text.utf8.count)" + ) + #endif _ = applyMobileViewportReport(params: params, terminalPanel: terminalPanel) @@ -14916,6 +14922,14 @@ class TerminalController { ] if let seq = MobileTerminalByteTee.shared.currentSequence(surfaceID: surfaceId) { payload["terminal_seq"] = seq + #if DEBUG + if sendResult == .sent { + HostLatencyTrace.stamp( + "host.in.applied", + "s=\(surfaceId.uuidString.prefix(8).lowercased()) seq=\(seq)" + ) + } + #endif } return .ok(payload) } diff --git a/cmux.xcodeproj/project.pbxproj b/cmux.xcodeproj/project.pbxproj index dfaf27399545..76973700d808 100644 --- a/cmux.xcodeproj/project.pbxproj +++ b/cmux.xcodeproj/project.pbxproj @@ -1112,6 +1112,7 @@ C0DE71B10000000000000001 /* AppDelegate+AgentChatNotifications.swift in Sources CD0CFE6000000000CD0CFE60 /* HostAccountFlow.swift in Sources */ = {isa = PBXBuildFile; fileRef = CD0CFE6100000000CD0CFE61 /* HostAccountFlow.swift */; }; B1F0C0030000000000000001 /* HostedInspectorDockControlScript.swift in Sources */ = {isa = PBXBuildFile; fileRef = B1F0C0030000000000000002 /* HostedInspectorDockControlScript.swift */; }; B1F0C0040000000000000001 /* HostedInspectorDockControlScriptTests.swift in Sources */ = {isa = PBXBuildFile; fileRef = B1F0C0040000000000000002 /* HostedInspectorDockControlScriptTests.swift */; }; + F17A00000000000000000002 /* HostLatencyTrace.swift in Sources */ = {isa = PBXBuildFile; fileRef = F17A00000000000000000001 /* HostLatencyTrace.swift */; }; CD0CFE6200000000CD0CFE62 /* HostSettingsActions.swift in Sources */ = {isa = PBXBuildFile; fileRef = CD0CFE6300000000CD0CFE63 /* HostSettingsActions.swift */; }; 809100000000000000000001 /* HostSettingsShortcutNotificationTests.swift in Sources */ = {isa = PBXBuildFile; fileRef = 809100000000000000000002 /* HostSettingsShortcutNotificationTests.swift */; }; F1C1AA21B7E84D10A1C10001 /* InactivePaneFirstClickFocusTests.swift in Sources */ = {isa = PBXBuildFile; fileRef = F1C1AA20B7E84D10A1C10001 /* InactivePaneFirstClickFocusTests.swift */; }; @@ -3598,6 +3599,7 @@ C0DE71B10000000000000002 /* AppDelegate+AgentChatNotifications.swift */ = {isa = CD0CFE6100000000CD0CFE61 /* HostAccountFlow.swift */ = {isa = PBXFileReference; includeInIndex = 1; lastKnownFileType = sourcecode.swift; path = HostAccountFlow.swift; sourceTree = ""; }; B1F0C0030000000000000002 /* HostedInspectorDockControlScript.swift */ = {isa = PBXFileReference; lastKnownFileType = sourcecode.swift; path = App/HostedInspectorDockControlScript.swift; sourceTree = ""; }; B1F0C0040000000000000002 /* HostedInspectorDockControlScriptTests.swift */ = {isa = PBXFileReference; lastKnownFileType = sourcecode.swift; path = HostedInspectorDockControlScriptTests.swift; sourceTree = ""; }; + F17A00000000000000000001 /* HostLatencyTrace.swift */ = {isa = PBXFileReference; lastKnownFileType = sourcecode.swift; path = App/HostLatencyTrace.swift; sourceTree = ""; }; CD0CFE6300000000CD0CFE63 /* HostSettingsActions.swift */ = {isa = PBXFileReference; includeInIndex = 1; lastKnownFileType = sourcecode.swift; path = HostSettingsActions.swift; sourceTree = ""; }; 809100000000000000000002 /* HostSettingsShortcutNotificationTests.swift */ = {isa = PBXFileReference; lastKnownFileType = sourcecode.swift; path = HostSettingsShortcutNotificationTests.swift; sourceTree = ""; }; F1C1AA20B7E84D10A1C10001 /* InactivePaneFirstClickFocusTests.swift */ = {isa = PBXFileReference; lastKnownFileType = sourcecode.swift; path = InactivePaneFirstClickFocusTests.swift; sourceTree = ""; }; @@ -5629,6 +5631,7 @@ C0DE71B10000000000000002 /* AppDelegate+AgentChatNotifications.swift */ = {isa = 895500000000000000000004 /* MainWindowController.swift */, A6D1F0000000000000000001 /* ReleasingWindowController.swift */, A500D011A1B2C3D4E5F60718 /* DebugLogging.swift */, + F17A00000000000000000001 /* HostLatencyTrace.swift */, D35B71010000000000000002 /* StartupBreadcrumbLog.swift */, A78C41010000000000000002 /* MacSentryStartupPolicy.swift */, A50016B0A1B2C3D4E5F60718 /* SessionSnapshotDebugBenchmark.swift */, @@ -8674,6 +8677,7 @@ C0DE71B10000000000000002 /* AppDelegate+AgentChatNotifications.swift */ = {isa = F5320000A1B2C3D4E5F60718 /* HermesAgentIndex.swift in Sources */, CD0CFE6000000000CD0CFE60 /* HostAccountFlow.swift in Sources */, B1F0C0030000000000000001 /* HostedInspectorDockControlScript.swift in Sources */, + F17A00000000000000000002 /* HostLatencyTrace.swift in Sources */, CD0CFE6200000000CD0CFE62 /* HostSettingsActions.swift in Sources */, C0DEA7720000000000000001 /* InternalFlagsWindow.swift in Sources */, 4D8B37E0E11A29997632765F /* InternalTabDragConfiguration.swift in Sources */, diff --git a/scripts/mobile-latency-trace/README.md b/scripts/mobile-latency-trace/README.md new file mode 100644 index 000000000000..aef8b312ddea --- /dev/null +++ b/scripts/mobile-latency-trace/README.md @@ -0,0 +1,70 @@ +# Mobile latency tracing + +Tracing is DEBUG-only and off by default. For a tagged Mac build, enable it for +that bundle and relaunch: + +```bash +defaults write com.cmuxterm.app.debug.slat cmux.debug.latency-trace -bool true +./scripts/reload.sh --tag slat --launch +``` + +`CMUX_LATENCY_TRACE=1` is the equivalent process environment gate. The Mac log +is `/tmp/cmux-debug-slat.log`. + +For an iOS Simulator, enable tracing and the optional typing probe at launch: + +```bash +SIMCTL_CHILD_CMUX_LATENCY_TRACE=1 \ +SIMCTL_CHILD_CMUX_LATENCY_PROBE=40:250 \ +xcrun simctl launch +``` + +The probe waits for a connected shell with a mounted terminal, waits another +three seconds, then sends the configured number of single characters through +the production input path. + +Find the simulator log with: + +```bash +data_dir="$(xcrun simctl get_app_container data)" +ios_log="$data_dir/Library/Application Support/cmux-debug.log" +``` + +Simulator and Mac uptime share a clock domain, so analyze with `--same-clock`: + +```bash +python3 scripts/mobile-latency-trace/analyze.py \ + --mac-log /tmp/cmux-debug-slat.log \ + --ios-log "$ios_log" \ + --same-clock +``` + +Omit `--same-clock` for physical-iPhone captures. Add `--json` for raw joined +duration arrays. Run the embedded fixture check with: + +```bash +python3 scripts/mobile-latency-trace/analyze.py --selftest +``` + +## Trace format + +Every per-surface render/input-ack stamp carries `s=`, where +`` is the first eight lowercase hexadecimal characters of the surface +UUID: + +- Mac: `host.tee`, `host.grid`, `host.enq`, `host.write`, `host.in.recv`, and + `host.in.applied`. +- iOS: `ev.grid`, `gate`, `ap.yield`, `ap.done`, `rd.present`, and `in.resp`. + +Mac `host.enq` and `host.write` stamps also carry `conn=`, the first +eight characters of the subscriber connection ID. The analyzer joins render +and input-ack stages by `(s, seq)`. When the same `(s, seq)` appears repeatedly, +it preserves first-at-or-after time ordering; for host fan-out, the earliest +connection write supplies each wire sample. + +Raw-input batches continue to use their process-local `n=` identity: +`in.send` marks departure from the iOS drain loop and `in.settled` marks the +actual response/error settlement (`ok=1` for success, `ok=0` for failure). +Because `in.resp` and `in.settled` are emitted from the same pipelined +settlement, the analyzer associates the next `in.resp` at or after each +successfully settled `in.send` even if its log timestamp follows `in.settled`. diff --git a/scripts/mobile-latency-trace/analyze.py b/scripts/mobile-latency-trace/analyze.py new file mode 100755 index 000000000000..0363661961f4 --- /dev/null +++ b/scripts/mobile-latency-trace/analyze.py @@ -0,0 +1,812 @@ +#!/usr/bin/env python3 +"""Analyze cmux iOS↔Mac latency trace stamps.""" + +from __future__ import annotations + +import argparse +import json +import math +import re +import sys +from dataclasses import dataclass +from pathlib import Path +from typing import Callable, Iterable + + +LINE_RE = re.compile(r"LAT\s+(?P\S+)\s+t=(?P