diff --git a/.github/swift-file-length-budget.tsv b/.github/swift-file-length-budget.tsv index 63aaac795f23..c4c892d7357b 100644 --- a/.github/swift-file-length-budget.tsv +++ b/.github/swift-file-length-budget.tsv @@ -12,7 +12,7 @@ 9497 cmuxTests/CLINotifyProcessIntegrationRegressionTests.swift 8032 Sources/Panels/BrowserPanelView.swift 8016 CLI/cmux_open.swift -7371 Packages/iOS/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite.swift +7783 Packages/iOS/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite.swift 7366 cmuxTests/WorkspaceUnitTests.swift 7218 cmuxTests/WorkspaceRemoteConnectionTests.swift 6359 cmuxTests/SessionPersistenceTests.swift @@ -24,11 +24,11 @@ 4487 Sources/Panels/FilePreviewPanel.swift 4483 Sources/cmuxApp.swift 4367 cmuxTests/BrowserPanelTests.swift +4264 Packages/iOS/CmuxMobileTerminal/Sources/CmuxMobileTerminal/GhosttySurfaceView.swift 4121 Sources/BrowserWindowPortal.swift 3934 Sources/Feed/FeedPanelView.swift 3926 cmuxTests/TabManagerUnitTests.swift 3896 cmuxTests/WindowAndDragTests.swift -3747 Packages/iOS/CmuxMobileTerminal/Sources/CmuxMobileTerminal/GhosttySurfaceView.swift 3668 cmuxTests/CLIGenericHookPersistenceTests.swift 3397 Sources/CmuxConfig.swift 3364 cmuxTests/TabManagerSessionSnapshotTests.swift @@ -112,6 +112,7 @@ 859 Packages/macOS/CmuxControlSocket/Sources/CmuxControlSocket/Coordinator/Workspace/ControlCommandCoordinator+Workspace.swift 847 cmuxTests/AgentSessionAutoResumeSettingsTests.swift 845 cmuxTests/SSHStartupSignalLifecycleTests.swift +841 Packages/iOS/CmuxMobileShell/Tests/CmuxMobileShellTests/TerminalOutputDeliveryQueueTests.swift 834 Sources/MainWindowFocusController.swift 830 Sources/TaskManagerTypes.swift 825 Packages/iOS/CmuxMobileShellUI/Sources/CmuxMobileShellUI/TerminalComposerView.swift @@ -130,7 +131,7 @@ 760 Packages/macOS/CMUXAgentLaunch/Tests/CMUXAgentLaunchTests/AgentLaunchSanitizerTests.swift 756 Sources/Panels/AgentSessionWebRendererCoordinator.swift 754 Sources/TerminalController+ControlWorkspaceContext.swift -753 Packages/iOS/CmuxMobileShellUI/Sources/CmuxMobileShellUI/WorkspaceDetailView.swift +752 Packages/iOS/CmuxMobileShellUI/Sources/CmuxMobileShellUI/WorkspaceDetailView.swift 752 cmuxUITests/CloseWorkspaceCmdDUITests.swift 739 cmuxTests/CLICodexHookTimeoutRegressionTests.swift 738 Packages/macOS/CMUXProjectModel/Sources/CMUXProjectModel/XcodeProjectAdapter.swift @@ -172,6 +173,7 @@ 620 cmuxTests/FinderFileDropRegressionTests.swift 620 cmuxTests/TerminalNotificationQueueTests.swift 615 Sources/SettingsNavigation.swift +613 Packages/iOS/CmuxMobileShell/Tests/CmuxMobileShellTests/MobileShellRenderGridLivenessTestSupport.swift 608 cmuxUITests/FeedSidebarUITests.swift 607 Packages/macOS/CmuxTerminal/Sources/CmuxTerminal/Surface/TerminalSurface.swift 607 Sources/SessionIndexModels.swift diff --git a/Packages/iOS/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite+TerminalOutputDelivery.swift b/Packages/iOS/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite+TerminalOutputDelivery.swift index 95004923b6e7..6aa0fb42f5cd 100644 --- a/Packages/iOS/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite+TerminalOutputDelivery.swift +++ b/Packages/iOS/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite+TerminalOutputDelivery.swift @@ -1,33 +1,59 @@ import CMUXMobileCore +internal import CmuxMobileDiagnostics import CmuxMobileShellModel public import Foundation extension MobileShellComposite { + func claimTerminalReplayBarrierFollowUp(surfaceID: String) -> Bool { + let followUpCount = terminalReplayBarrierFollowUpCountsBySurfaceID[surfaceID] ?? 0 + guard followUpCount < Self.maxTerminalReplayBarrierFollowUps else { + MobileDebugLog.anchormux( + "terminal.output.replay_followup_cap_reached surface=\(surfaceID) attempts=\(followUpCount)" + ) + terminalReplayBarrierFollowUpCountsBySurfaceID.removeValue(forKey: surfaceID) + return false + } + terminalReplayBarrierFollowUpCountsBySurfaceID[surfaceID] = followUpCount + 1 + return true + } + /// Yield a raw PTY byte chunk to the surface stream, if one is attached. - func deliverTerminalBytes(_ bytes: Data, surfaceID: String) { - deliverTerminalOutput( + @discardableResult + func deliverTerminalBytes( + _ bytes: Data, + surfaceID: String, + bypassReplayBarrier: Bool = false + ) -> Bool { + return deliverTerminalOutput( TerminalOutputDelivery( bytes: bytes, replaceable: false, viewportPolicy: .natural ), - surfaceID: surfaceID + surfaceID: surfaceID, + bypassReplayBarrier: bypassReplayBarrier ) } - func deliverTerminalRenderGrid(_ frame: MobileTerminalRenderGridFrame, surfaceID: String) { - deliverTerminalOutput( + @discardableResult + func deliverTerminalRenderGrid( + _ frame: MobileTerminalRenderGridFrame, + surfaceID: String, + bypassReplayBarrier: Bool = false + ) -> Bool { + return deliverTerminalOutput( TerminalOutputDelivery( renderGrid: frame, replaceable: frame.isReplaceableViewportPatchForMobileDelivery, viewportPolicy: frame.mobileViewportPolicy ), - surfaceID: surfaceID + surfaceID: surfaceID, + bypassReplayBarrier: bypassReplayBarrier ) } func deliverTerminalViewportPolicy(_ policy: MobileTerminalOutputViewportPolicy, surfaceID: String) { - deliverTerminalOutput( + _ = deliverTerminalOutput( TerminalOutputDelivery( bytes: Data(), replaceable: true, @@ -38,12 +64,50 @@ extension MobileShellComposite { ) } - private func deliverTerminalOutput(_ delivery: TerminalOutputDelivery, surfaceID: String) { + private func deliverTerminalOutput( + _ delivery: TerminalOutputDelivery, + surfaceID: String, + bypassReplayBarrier: Bool = false + ) -> Bool { guard let continuation = terminalByteContinuationsBySurfaceID[surfaceID], - let streamToken = terminalOutputStreamTokensBySurfaceID[surfaceID] else { return } + let streamToken = terminalOutputStreamTokensBySurfaceID[surfaceID] else { return false } + if let replayBarrierToken = terminalReplayBarrierTokensBySurfaceID[surfaceID], + !bypassReplayBarrier { + terminalReplayBarrierDroppedOutputSurfaceIDs.insert(surfaceID) + let droppedOutputCount = (terminalReplayBarrierDroppedOutputCountsBySurfaceID[surfaceID] ?? 0) &+ 1 + terminalReplayBarrierDroppedOutputCountsBySurfaceID[surfaceID] = droppedOutputCount + if droppedOutputCount == 1 || droppedOutputCount.isMultiple(of: 32) { + MobileDebugLog.anchormux( + "terminal.output.drop_replay_barrier surface=\(surfaceID) count=\(droppedOutputCount)" + ) + } + if remoteClient != nil, + terminalReplayBarrierAckStreamTokensBySurfaceID[surfaceID] == nil, + !terminalReplaySurfaceIDsInFlight.contains(surfaceID), + !terminalReplayFailureRetryExhausted(surfaceID: surfaceID) { + MobileDebugLog.anchormux("terminal.output.replay_retry_after_drop surface=\(surfaceID)") + requestTerminalReplay( + surfaceID: surfaceID, + replayBarrierToken: replayBarrierToken, + coveredReplayBarrierDroppedOutputCount: droppedOutputCount + ) + } + return false + } var queue = terminalOutputQueuesBySurfaceID[surfaceID] ?? TerminalOutputDeliveryQueue() let immediate = queue.enqueue(delivery) + let pendingCount = queue.pendingCount terminalOutputQueuesBySurfaceID[surfaceID] = queue + if bypassReplayBarrier, + immediate != nil, + terminalReplayBarrierTokensBySurfaceID[surfaceID] != nil { + terminalReplayBarrierAckStreamTokensBySurfaceID[surfaceID] = streamToken + } + if pendingCount >= 32, pendingCount.isMultiple(of: 32) { + MobileDebugLog.anchormux( + "terminal.output.pending surface=\(surfaceID) depth=\(pendingCount)" + ) + } if let immediate { continuation.yield( MobileTerminalOutputChunk( @@ -53,6 +117,7 @@ extension MobileShellComposite { ) ) } + return true } /// Mark the current yielded terminal-output chunk as applied by the iOS surface. @@ -61,6 +126,31 @@ extension MobileShellComposite { var queue = terminalOutputQueuesBySurfaceID[surfaceID] else { return } let next = queue.completeInFlight() terminalOutputQueuesBySurfaceID[surfaceID] = queue + if terminalReplayBarrierAckStreamTokensBySurfaceID[surfaceID] == streamToken { + let coveredDroppedOutputCount = + terminalReplayBarrierAckCoveredDroppedOutputCountsBySurfaceID.removeValue(forKey: surfaceID) + let currentDroppedOutputCount = terminalReplayBarrierDroppedOutputCountsBySurfaceID[surfaceID] ?? 0 + let needsFollowUpReplay = coveredDroppedOutputCount.map { + currentDroppedOutputCount > $0 + } ?? true + terminalReplayBarrierAckStreamTokensBySurfaceID.removeValue(forKey: surfaceID) + terminalReplayBarrierTokensBySurfaceID.removeValue(forKey: surfaceID) + MobileDebugLog.anchormux("terminal.output.replay_barrier_cleared surface=\(surfaceID)") + let droppedOutputDuringBarrier = terminalReplayBarrierDroppedOutputSurfaceIDs.remove(surfaceID) != nil + terminalReplayBarrierDroppedOutputCountsBySurfaceID.removeValue(forKey: surfaceID) + if droppedOutputDuringBarrier, + needsFollowUpReplay, + claimTerminalReplayBarrierFollowUp(surfaceID: surfaceID) { + let replayBarrierToken = beginTerminalReplayBarrier( + surfaceID: surfaceID, + preservingFollowUpCount: true + ) + MobileDebugLog.anchormux("terminal.output.replay_followup surface=\(surfaceID)") + requestTerminalReplay(surfaceID: surfaceID, replayBarrierToken: replayBarrierToken) + return + } + terminalReplayBarrierFollowUpCountsBySurfaceID.removeValue(forKey: surfaceID) + } guard let next, let continuation = terminalByteContinuationsBySurfaceID[surfaceID], terminalOutputStreamTokensBySurfaceID[surfaceID] == streamToken else { @@ -72,4 +162,72 @@ extension MobileShellComposite { viewportPolicy: next.viewportPolicy )) } + + /// Abandon the current yielded terminal-output chunk after the local render + /// surface reset. The abandoned bytes may have been applied to the old + /// Ghostty surface or may still be behind a wedged worker queue, so continuing + /// to drain the old pending queue would replay stale deltas into the rebuilt + /// surface. Reset the queue, invalidate stale acks, then request a fresh + /// authoritative replay from the Mac. + public func terminalOutputDidReset(surfaceID: String, streamToken: UUID) { + guard terminalOutputStreamTokensBySurfaceID[surfaceID] == streamToken, + terminalOutputQueuesBySurfaceID[surfaceID] != nil else { return } + if let replayBarrierToken = terminalReplayBarrierTokensBySurfaceID[surfaceID] { + guard terminalReplayBarrierAckStreamTokensBySurfaceID[surfaceID] == streamToken else { + MobileDebugLog.anchormux("terminal.output.reset_barrier_active surface=\(surfaceID)") + return + } + retryTerminalReplayAfterAckReset( + surfaceID: surfaceID, + replayBarrierToken: replayBarrierToken + ) + return + } + let replayBarrierToken = beginTerminalReplayBarrier(surfaceID: surfaceID) + MobileDebugLog.anchormux("terminal.output.reset surface=\(surfaceID)") + requestTerminalReplay(surfaceID: surfaceID, replayBarrierToken: replayBarrierToken) + } + + private func retryTerminalReplayAfterAckReset( + surfaceID: String, + replayBarrierToken: UUID + ) { + guard terminalReplayBarrierTokensBySurfaceID[surfaceID] == replayBarrierToken else { + return + } + terminalOutputQueuesBySurfaceID[surfaceID] = TerminalOutputDeliveryQueue() + terminalOutputStreamTokensBySurfaceID[surfaceID] = UUID() + deliveredTerminalByteEndSeqBySurfaceID.removeValue(forKey: surfaceID) + pendingTerminalByteEndSeqBySurfaceID.removeValue(forKey: surfaceID) + terminalReplayBarrierAckStreamTokensBySurfaceID.removeValue(forKey: surfaceID) + terminalReplayBarrierAckCoveredDroppedOutputCountsBySurfaceID.removeValue(forKey: surfaceID) + terminalReplayBarrierTokensInFlightBySurfaceID.removeValue(forKey: surfaceID) + guard let retryToken = prepareTerminalReplayFailureRetry( + surfaceID: surfaceID, + replayBarrierToken: replayBarrierToken + ) else { + preserveTerminalReplayBarrierIfCurrent( + surfaceID: surfaceID, + token: replayBarrierToken, + reason: "reset_replay_ack" + ) + return + } + MobileDebugLog.anchormux("terminal.output.reset_replay_ack surface=\(surfaceID)") + requestTerminalReplay( + surfaceID: surfaceID, + replayBarrierToken: retryToken, + coveredReplayBarrierDroppedOutputCount: + terminalReplayBarrierDroppedOutputCountsBySurfaceID[surfaceID] + ) + } + + /// Ask the Mac to replay the authoritative terminal state for a surface. + public func terminalOutputNeedsReplay(surfaceID: String) { + guard terminalByteContinuationsBySurfaceID[surfaceID] != nil else { return } + let replayBarrierToken = beginTerminalReplayBarrier(surfaceID: surfaceID) + MobileDebugLog.anchormux("terminal.output.replay_requested surface=\(surfaceID)") + requestTerminalReplay(surfaceID: surfaceID, replayBarrierToken: replayBarrierToken) + } + } diff --git a/Packages/iOS/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite.swift b/Packages/iOS/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite.swift index e3307cb89405..1e6a6ce8aa6a 100644 --- a/Packages/iOS/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite.swift +++ b/Packages/iOS/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite.swift @@ -34,6 +34,9 @@ public typealias CMUXMobileShellStore = MobileShellComposite @MainActor @Observable public final class MobileShellComposite: MobileTerminalOutputSinking { + private static let maxTerminalReplayFailureRetries = 2 + static let maxTerminalReplayBarrierFollowUps = 1 + private enum TerminalOutputTransport: Equatable { case hybrid case renderGrid @@ -715,10 +718,20 @@ public final class MobileShellComposite: MobileTerminalOutputSinking { /// ``secondaryAggregationTask``, can reject old-team results after awaits. var secondaryAggregationScopeGeneration = 0 private var reportedViewportSizesByTerminalKey: [MobileTerminalViewportKey: MobileTerminalViewportSize] - private var deliveredTerminalByteEndSeqBySurfaceID: [String: UInt64] - private var pendingTerminalByteEndSeqBySurfaceID: [String: UInt64] + var deliveredTerminalByteEndSeqBySurfaceID: [String: UInt64] + var pendingTerminalByteEndSeqBySurfaceID: [String: UInt64] private var terminalActiveScreenBySurfaceID: [String: MobileTerminalRenderGridFrame.Screen] - private var terminalReplaySurfaceIDsInFlight: Set + var terminalReplaySurfaceIDsInFlight: Set + private var terminalReplayRequestIDsInFlightBySurfaceID: [String: UUID] + private var terminalReplayTasksBySurfaceID: [String: Task] + var terminalReplayBarrierTokensInFlightBySurfaceID: [String: UUID] + var terminalReplayBarrierTokensBySurfaceID: [String: UUID] + var terminalReplayBarrierAckStreamTokensBySurfaceID: [String: UUID] + var terminalReplayBarrierDroppedOutputSurfaceIDs: Set + var terminalReplayBarrierDroppedOutputCountsBySurfaceID: [String: UInt64] + var terminalReplayBarrierAckCoveredDroppedOutputCountsBySurfaceID: [String: UInt64] + var terminalReplayFailureRetryCountsBySurfaceID: [String: Int] + var terminalReplayBarrierFollowUpCountsBySurfaceID: [String: Int] private var terminalOutputTransport: TerminalOutputTransport var terminalByteContinuationsBySurfaceID: [String: AsyncStream.Continuation] var terminalOutputStreamTokensBySurfaceID: [String: UUID] @@ -910,6 +923,16 @@ public final class MobileShellComposite: MobileTerminalOutputSinking { self.pendingTerminalByteEndSeqBySurfaceID = [:] self.terminalActiveScreenBySurfaceID = [:] self.terminalReplaySurfaceIDsInFlight = [] + self.terminalReplayRequestIDsInFlightBySurfaceID = [:] + self.terminalReplayTasksBySurfaceID = [:] + self.terminalReplayBarrierTokensInFlightBySurfaceID = [:] + self.terminalReplayBarrierTokensBySurfaceID = [:] + self.terminalReplayBarrierAckStreamTokensBySurfaceID = [:] + self.terminalReplayBarrierDroppedOutputSurfaceIDs = [] + self.terminalReplayBarrierDroppedOutputCountsBySurfaceID = [:] + self.terminalReplayBarrierAckCoveredDroppedOutputCountsBySurfaceID = [:] + self.terminalReplayFailureRetryCountsBySurfaceID = [:] + self.terminalReplayBarrierFollowUpCountsBySurfaceID = [:] self.terminalOutputTransport = .rawBytes self.terminalByteContinuationsBySurfaceID = [:] self.terminalOutputStreamTokensBySurfaceID = [:] @@ -935,6 +958,7 @@ public final class MobileShellComposite: MobileTerminalOutputSinking { createTerminalTask?.cancel() workspaceListRefreshTask?.cancel() pullToRefreshTask?.cancel() + cancelAllTerminalReplayTasks() teardownSecondaryMacSubscriptions() if let remoteClient { Task { await remoteClient.disconnect() } @@ -5215,13 +5239,24 @@ public final class MobileShellComposite: MobileTerminalOutputSinking { workspaceListRefreshTask = nil pullToRefreshTask?.cancel() pullToRefreshTask = nil + cancelAllTerminalReplayTasks() } private func resetTerminalOutputTracking() { + cancelAllTerminalReplayTasks() deliveredTerminalByteEndSeqBySurfaceID = [:] pendingTerminalByteEndSeqBySurfaceID = [:] terminalActiveScreenBySurfaceID = [:] terminalReplaySurfaceIDsInFlight = [] + terminalReplayRequestIDsInFlightBySurfaceID = [:] + terminalReplayBarrierTokensInFlightBySurfaceID = [:] + terminalReplayBarrierTokensBySurfaceID = [:] + terminalReplayBarrierAckStreamTokensBySurfaceID = [:] + terminalReplayBarrierDroppedOutputSurfaceIDs = [] + terminalReplayBarrierDroppedOutputCountsBySurfaceID = [:] + terminalReplayBarrierAckCoveredDroppedOutputCountsBySurfaceID = [:] + terminalReplayFailureRetryCountsBySurfaceID = [:] + terminalReplayBarrierFollowUpCountsBySurfaceID = [:] terminalOutputQueuesBySurfaceID = [:] terminalOutputStreamTokensBySurfaceID = terminalOutputStreamTokensBySurfaceID.mapValues { _ in UUID() } terminalScrollQueueTokensBySurfaceID = [:] @@ -6641,6 +6676,174 @@ public final class MobileShellComposite: MobileTerminalOutputSinking { } } + func beginTerminalReplayBarrier( + surfaceID: String, + preservingFollowUpCount: Bool = false + ) -> UUID { + cancelTerminalReplayInFlight(surfaceID: surfaceID) + terminalOutputQueuesBySurfaceID[surfaceID] = TerminalOutputDeliveryQueue() + terminalOutputStreamTokensBySurfaceID[surfaceID] = UUID() + deliveredTerminalByteEndSeqBySurfaceID.removeValue(forKey: surfaceID) + pendingTerminalByteEndSeqBySurfaceID.removeValue(forKey: surfaceID) + let token = UUID() + terminalReplayBarrierTokensBySurfaceID[surfaceID] = token + terminalReplayBarrierAckStreamTokensBySurfaceID.removeValue(forKey: surfaceID) + terminalReplayBarrierDroppedOutputSurfaceIDs.remove(surfaceID) + terminalReplayBarrierDroppedOutputCountsBySurfaceID.removeValue(forKey: surfaceID) + terminalReplayBarrierAckCoveredDroppedOutputCountsBySurfaceID.removeValue(forKey: surfaceID) + terminalReplayFailureRetryCountsBySurfaceID.removeValue(forKey: surfaceID) + if !preservingFollowUpCount { + terminalReplayBarrierFollowUpCountsBySurfaceID.removeValue(forKey: surfaceID) + } + terminalReplayBarrierTokensInFlightBySurfaceID.removeValue(forKey: surfaceID) + return token + } + + @discardableResult + private func clearTerminalReplayBarrierIfCurrent( + surfaceID: String, + token: UUID?, + reason: String, + preserveDroppedOutput: Bool = false + ) -> Bool { + guard let token, + terminalReplayBarrierTokensBySurfaceID[surfaceID] == token else { + return false + } + if preserveDroppedOutput, + terminalReplayBarrierDroppedOutputSurfaceIDs.contains(surfaceID) { + MobileDebugLog.anchormux("terminal.output.replay_barrier_preserved_\(reason) surface=\(surfaceID)") + return false + } + terminalReplayBarrierAckStreamTokensBySurfaceID.removeValue(forKey: surfaceID) + terminalReplayBarrierTokensBySurfaceID.removeValue(forKey: surfaceID) + terminalReplayBarrierDroppedOutputSurfaceIDs.remove(surfaceID) + terminalReplayBarrierDroppedOutputCountsBySurfaceID.removeValue(forKey: surfaceID) + terminalReplayBarrierAckCoveredDroppedOutputCountsBySurfaceID.removeValue(forKey: surfaceID) + terminalReplayFailureRetryCountsBySurfaceID.removeValue(forKey: surfaceID) + terminalReplayBarrierFollowUpCountsBySurfaceID.removeValue(forKey: surfaceID) + terminalReplayBarrierTokensInFlightBySurfaceID.removeValue(forKey: surfaceID) + MobileDebugLog.anchormux("terminal.output.replay_barrier_cleared_\(reason) surface=\(surfaceID)") + return true + } + + @discardableResult + func preserveTerminalReplayBarrierIfCurrent( + surfaceID: String, + token: UUID?, + reason: String + ) -> Bool { + guard let token, + terminalReplayBarrierTokensBySurfaceID[surfaceID] == token else { + return false + } + terminalReplayBarrierAckStreamTokensBySurfaceID.removeValue(forKey: surfaceID) + terminalReplayBarrierTokensInFlightBySurfaceID.removeValue(forKey: surfaceID) + MobileDebugLog.anchormux("terminal.output.replay_barrier_preserved_\(reason) surface=\(surfaceID)") + return true + } + + func prepareTerminalReplayFailureRetry( + surfaceID: String, + replayBarrierToken: UUID? + ) -> UUID? { + guard let replayBarrierToken, + hasTerminalOutputSink(surfaceID: surfaceID), + terminalReplayBarrierTokensBySurfaceID[surfaceID] == replayBarrierToken else { + return nil + } + let retryCount = terminalReplayFailureRetryCountsBySurfaceID[surfaceID] ?? 0 + guard retryCount < Self.maxTerminalReplayFailureRetries else { + MobileDebugLog.anchormux( + "CMUX_REPLAY retry_exhausted surface=\(surfaceID) attempts=\(retryCount)" + ) + return nil + } + terminalReplayFailureRetryCountsBySurfaceID[surfaceID] = retryCount + 1 + MobileDebugLog.anchormux( + "CMUX_REPLAY retry_after_failure surface=\(surfaceID) attempt=\(retryCount + 1)" + ) + return replayBarrierToken + } + + func terminalReplayFailureRetryExhausted(surfaceID: String) -> Bool { + (terminalReplayFailureRetryCountsBySurfaceID[surfaceID] ?? 0) >= Self.maxTerminalReplayFailureRetries + } + + @discardableResult + private func requestTerminalReplayForCurrentBarrier( + surfaceID: String, + replayBarrierToken: UUID?, + coveredReplayBarrierDroppedOutputCount: UInt64?, + reason: String + ) -> Bool { + guard let replayBarrierToken, + hasTerminalOutputSink(surfaceID: surfaceID), + terminalReplayBarrierTokensBySurfaceID[surfaceID] == replayBarrierToken, + remoteClient != nil else { + return false + } + MobileDebugLog.anchormux("CMUX_REPLAY retry_\(reason) surface=\(surfaceID)") + requestTerminalReplay( + surfaceID: surfaceID, + replayBarrierToken: replayBarrierToken, + coveredReplayBarrierDroppedOutputCount: coveredReplayBarrierDroppedOutputCount + ) + return true + } + + private func markTerminalReplayInFlight( + surfaceID: String, + requestID: UUID, + replayBarrierToken: UUID? + ) { + cancelTerminalReplayInFlight(surfaceID: surfaceID) + terminalReplaySurfaceIDsInFlight.insert(surfaceID) + terminalReplayRequestIDsInFlightBySurfaceID[surfaceID] = requestID + if let replayBarrierToken { + terminalReplayBarrierTokensInFlightBySurfaceID[surfaceID] = replayBarrierToken + } else { + terminalReplayBarrierTokensInFlightBySurfaceID.removeValue(forKey: surfaceID) + } + } + + private func storeTerminalReplayTask( + surfaceID: String, + requestID: UUID, + task: Task + ) { + guard terminalReplayRequestIDsInFlightBySurfaceID[surfaceID] == requestID else { + task.cancel() + return + } + terminalReplayTasksBySurfaceID[surfaceID] = task + } + + private func clearTerminalReplayInFlightIfCurrent(surfaceID: String, requestID: UUID) { + guard terminalReplayRequestIDsInFlightBySurfaceID[surfaceID] == requestID else { return } + terminalReplaySurfaceIDsInFlight.remove(surfaceID) + terminalReplayRequestIDsInFlightBySurfaceID.removeValue(forKey: surfaceID) + terminalReplayTasksBySurfaceID.removeValue(forKey: surfaceID) + terminalReplayBarrierTokensInFlightBySurfaceID.removeValue(forKey: surfaceID) + } + + private func cancelTerminalReplayInFlight(surfaceID: String) { + terminalReplayTasksBySurfaceID.removeValue(forKey: surfaceID)?.cancel() + terminalReplaySurfaceIDsInFlight.remove(surfaceID) + terminalReplayRequestIDsInFlightBySurfaceID.removeValue(forKey: surfaceID) + terminalReplayBarrierTokensInFlightBySurfaceID.removeValue(forKey: surfaceID) + } + + private func cancelAllTerminalReplayTasks() { + for task in terminalReplayTasksBySurfaceID.values { + task.cancel() + } + terminalReplayTasksBySurfaceID = [:] + terminalReplaySurfaceIDsInFlight = [] + terminalReplayRequestIDsInFlightBySurfaceID = [:] + terminalReplayBarrierTokensInFlightBySurfaceID = [:] + } + private enum RenderGridEventDeliveryDecision { case deliver case advisory(requestReplay: Bool, updateTrackedScreen: Bool, deliverViewportPolicy: Bool) @@ -6683,7 +6886,6 @@ public final class MobileShellComposite: MobileTerminalOutputSinking { : .deliver switch deliveryDecision { case .deliver: - terminalActiveScreenBySurfaceID[renderGrid.surfaceID] = renderGrid.activeScreen break case .advisory(let requestReplay, let updateTrackedScreen, let deliverViewportPolicy): if updateTrackedScreen { @@ -6700,8 +6902,9 @@ public final class MobileShellComposite: MobileTerminalOutputSinking { } return } + guard deliverTerminalRenderGrid(renderGrid, surfaceID: renderGrid.surfaceID) else { return } + terminalActiveScreenBySurfaceID[renderGrid.surfaceID] = renderGrid.activeScreen markTerminalBytesDelivered(surfaceID: renderGrid.surfaceID, endSeq: renderGrid.stateSeq) - deliverTerminalRenderGrid(renderGrid, surfaceID: renderGrid.surfaceID) } private static func terminalSnapshotReplacementBytes(_ snapshotBytes: Data) -> Data { @@ -6731,9 +6934,17 @@ public final class MobileShellComposite: MobileTerminalOutputSinking { } private func unregisterTerminalOutput(surfaceID: String) { + cancelTerminalReplayInFlight(surfaceID: surfaceID) terminalByteContinuationsBySurfaceID.removeValue(forKey: surfaceID) terminalOutputStreamTokensBySurfaceID.removeValue(forKey: surfaceID) terminalOutputQueuesBySurfaceID.removeValue(forKey: surfaceID) + terminalReplayBarrierTokensBySurfaceID.removeValue(forKey: surfaceID) + terminalReplayBarrierAckStreamTokensBySurfaceID.removeValue(forKey: surfaceID) + terminalReplayBarrierDroppedOutputSurfaceIDs.remove(surfaceID) + terminalReplayBarrierDroppedOutputCountsBySurfaceID.removeValue(forKey: surfaceID) + terminalReplayBarrierAckCoveredDroppedOutputCountsBySurfaceID.removeValue(forKey: surfaceID) + terminalReplayFailureRetryCountsBySurfaceID.removeValue(forKey: surfaceID) + terminalReplayBarrierFollowUpCountsBySurfaceID.removeValue(forKey: surfaceID) terminalScrollQueueTokensBySurfaceID.removeValue(forKey: surfaceID) terminalScrollQueuesBySurfaceID.removeValue(forKey: surfaceID) terminalScrollbackPrefetchStatesBySurfaceID.removeValue(forKey: surfaceID) @@ -6859,30 +7070,64 @@ public final class MobileShellComposite: MobileTerminalOutputSinking { /// resume. The VT snapshot and raw byte ring remain fallbacks, but neither /// is the target architecture: a byte tail is not a complete screen state /// for TUIs, and a VT export is still a replay stream rather than state. - private func requestTerminalReplay(surfaceID: String) { + func requestTerminalReplay( + surfaceID: String, + replayBarrierToken: UUID? = nil, + coveredReplayBarrierDroppedOutputCount: UInt64? = nil + ) { + let replayBarrierTokenForRequest = replayBarrierToken + ?? terminalReplayBarrierTokensBySurfaceID[surfaceID] + let coveredReplayBarrierDroppedOutputCountForRequest = replayBarrierTokenForRequest == nil + ? nil + : (coveredReplayBarrierDroppedOutputCount + ?? terminalReplayBarrierDroppedOutputCountsBySurfaceID[surfaceID] + ?? 0) guard let client = remoteClient else { + clearTerminalReplayBarrierIfCurrent( + surfaceID: surfaceID, + token: replayBarrierTokenForRequest, + reason: "no_remote_client" + ) #if DEBUG mobileShellLog.error("CMUX_REPLAY skip surface=\(surfaceID, privacy: .public) reason=no_remote_client") #endif return } guard let workspaceID = workspaceID(forTerminalID: surfaceID) else { + clearTerminalReplayBarrierIfCurrent( + surfaceID: surfaceID, + token: replayBarrierTokenForRequest, + reason: "workspace_not_found" + ) #if DEBUG mobileShellLog.error("CMUX_REPLAY skip surface=\(surfaceID, privacy: .public) reason=workspace_not_found") #endif return } let remoteWorkspaceID = remoteWorkspaceID(for: workspaceID) - guard !terminalReplaySurfaceIDsInFlight.contains(surfaceID) else { - #if DEBUG - mobileShellLog.info("CMUX_REPLAY skip surface=\(surfaceID, privacy: .public) reason=in_flight") - #endif - return + if let replayBarrierTokenForRequest { + guard terminalReplayBarrierTokensInFlightBySurfaceID[surfaceID] != replayBarrierTokenForRequest else { + #if DEBUG + mobileShellLog.info("CMUX_REPLAY skip surface=\(surfaceID, privacy: .public) reason=barrier_in_flight") + #endif + return + } + } else { + guard !terminalReplaySurfaceIDsInFlight.contains(surfaceID) else { + #if DEBUG + mobileShellLog.info("CMUX_REPLAY skip surface=\(surfaceID, privacy: .public) reason=in_flight") + #endif + return + } } - terminalReplaySurfaceIDsInFlight.insert(surfaceID) - Task { @MainActor [weak self] in - guard let self else { return } - defer { self.terminalReplaySurfaceIDsInFlight.remove(surfaceID) } + let replayRequestID = UUID() + markTerminalReplayInFlight( + surfaceID: surfaceID, + requestID: replayRequestID, + replayBarrierToken: replayBarrierTokenForRequest + ) + let replayTask = Task { @MainActor [weak self] in + let replayResult: Result do { let request = try MobileCoreRPCClient.requestData( method: "mobile.terminal.replay", @@ -6891,14 +7136,59 @@ public final class MobileShellComposite: MobileTerminalOutputSinking { "surface_id": surfaceID, ] ) - let data = try await client.sendRequest(request) - guard self.remoteClient === client else { return } + replayResult = .success(try await client.sendRequest(request)) + } catch { + replayResult = .failure(error) + } + guard let self else { return } + var transferredInFlightToRetry = false + defer { + if !transferredInFlightToRetry { + self.clearTerminalReplayInFlightIfCurrent( + surfaceID: surfaceID, + requestID: replayRequestID + ) + } + } + switch replayResult { + case .success(let data): + guard self.terminalReplayRequestIDsInFlightBySurfaceID[surfaceID] == replayRequestID else { + MobileDebugLog.anchormux("CMUX_REPLAY stale_request surface=\(surfaceID)") + return + } + guard self.remoteClient === client else { + self.clearTerminalReplayInFlightIfCurrent( + surfaceID: surfaceID, + requestID: replayRequestID + ) + transferredInFlightToRetry = true + guard self.requestTerminalReplayForCurrentBarrier( + surfaceID: surfaceID, + replayBarrierToken: replayBarrierTokenForRequest, + coveredReplayBarrierDroppedOutputCount: nil, + reason: "stale_client" + ) else { + self.clearTerminalReplayBarrierIfCurrent( + surfaceID: surfaceID, + token: replayBarrierTokenForRequest, + reason: "stale_client" + ) + return + } + return + } let payload = try? MobileTerminalReplayResponse.decode(data) let bytes = payload?.dataBase64.flatMap { Data(base64Encoded: $0) } let snapshotBytes = payload?.snapshotBase64.flatMap { Data(base64Encoded: $0) } let decodedRenderGrid = payload?.renderGrid let renderGrid = decodedRenderGrid?.surfaceID == surfaceID ? decodedRenderGrid : nil let replaySeq = renderGrid?.stateSeq ?? payload?.sequence + if let replayBarrierTokenForRequest { + guard self.terminalReplayBarrierTokensBySurfaceID[surfaceID] == replayBarrierTokenForRequest else { + MobileDebugLog.anchormux("CMUX_REPLAY barrier_stale surface=\(surfaceID)") + return + } + } #if DEBUG let seq = replaySeq ?? 0 let cols = payload?.columns ?? -1 @@ -6909,6 +7199,11 @@ public final class MobileShellComposite: MobileTerminalOutputSinking { let deliveredSeq = self.deliveredTerminalByteEndSeqBySurfaceID[surfaceID], deliveredSeq > replaySeq { MobileDebugLog.anchormux("CMUX_REPLAY stale surface=\(surfaceID) delivered=\(deliveredSeq) replay=\(replaySeq)") + self.clearTerminalReplayBarrierIfCurrent( + surfaceID: surfaceID, + token: replayBarrierTokenForRequest, + reason: "stale_sequence" + ) return } let deliverBytes: Data? @@ -6922,28 +7217,145 @@ public final class MobileShellComposite: MobileTerminalOutputSinking { deliverBytes = bytes MobileDebugLog.anchormux("CMUX_REPLAY raw_tail surface=\(surfaceID) bytes=\(bytes?.count ?? -1) seq=\(replaySeq ?? 0)") } - if let replaySeq { - self.markTerminalBytesDelivered(surfaceID: surfaceID, endSeq: replaySeq) - } if let renderGrid { + let accepted = self.deliverTerminalRenderGrid( + renderGrid, + surfaceID: surfaceID, + bypassReplayBarrier: replayBarrierTokenForRequest != nil + ) + guard accepted else { + self.clearTerminalReplayBarrierIfCurrent( + surfaceID: surfaceID, + token: replayBarrierTokenForRequest, + reason: "not_delivered", + preserveDroppedOutput: true + ) + return + } + if self.terminalReplayBarrierAckStreamTokensBySurfaceID[surfaceID] != nil { + if let coveredReplayBarrierDroppedOutputCountForRequest { + self.terminalReplayBarrierAckCoveredDroppedOutputCountsBySurfaceID[surfaceID] = + coveredReplayBarrierDroppedOutputCountForRequest + } else { + self.terminalReplayBarrierAckCoveredDroppedOutputCountsBySurfaceID.removeValue(forKey: surfaceID) + } + } self.terminalActiveScreenBySurfaceID[surfaceID] = renderGrid.activeScreen - self.deliverTerminalRenderGrid(renderGrid, surfaceID: surfaceID) + if let replaySeq { + self.markTerminalBytesDelivered(surfaceID: surfaceID, endSeq: replaySeq) + } return } guard let deliverBytes, !deliverBytes.isEmpty else { + if self.terminalReplayBarrierDroppedOutputSurfaceIDs.contains(surfaceID), + let retryToken = self.prepareTerminalReplayFailureRetry( + surfaceID: surfaceID, + replayBarrierToken: replayBarrierTokenForRequest + ) { + self.clearTerminalReplayInFlightIfCurrent( + surfaceID: surfaceID, + requestID: replayRequestID + ) + transferredInFlightToRetry = true + self.requestTerminalReplay( + surfaceID: surfaceID, + replayBarrierToken: retryToken, + coveredReplayBarrierDroppedOutputCount: + self.terminalReplayBarrierDroppedOutputCountsBySurfaceID[surfaceID] + ) + return + } + self.clearTerminalReplayBarrierIfCurrent( + surfaceID: surfaceID, + token: replayBarrierTokenForRequest, + reason: "empty" + ) + return + } + let accepted = self.deliverTerminalBytes( + deliverBytes, + surfaceID: surfaceID, + bypassReplayBarrier: replayBarrierTokenForRequest != nil + ) + if accepted, + self.terminalReplayBarrierAckStreamTokensBySurfaceID[surfaceID] != nil { + if let coveredReplayBarrierDroppedOutputCountForRequest { + self.terminalReplayBarrierAckCoveredDroppedOutputCountsBySurfaceID[surfaceID] = + coveredReplayBarrierDroppedOutputCountForRequest + } else { + self.terminalReplayBarrierAckCoveredDroppedOutputCountsBySurfaceID.removeValue(forKey: surfaceID) + } + } + if accepted, let replaySeq { + self.markTerminalBytesDelivered(surfaceID: surfaceID, endSeq: replaySeq) + } else if !accepted { + self.clearTerminalReplayBarrierIfCurrent( + surfaceID: surfaceID, + token: replayBarrierTokenForRequest, + reason: "not_delivered", + preserveDroppedOutput: true + ) + } + case .failure(let error): + guard self.terminalReplayRequestIDsInFlightBySurfaceID[surfaceID] == replayRequestID else { + MobileDebugLog.anchormux("CMUX_REPLAY stale_request_failed surface=\(surfaceID)") + return + } + mobileShellLog.error("CMUX_REPLAY failed surface=\(surfaceID, privacy: .public) error=\(String(describing: error), privacy: .private)") + guard self.remoteClient === client else { + self.clearTerminalReplayInFlightIfCurrent( + surfaceID: surfaceID, + requestID: replayRequestID + ) + transferredInFlightToRetry = true + guard self.requestTerminalReplayForCurrentBarrier( + surfaceID: surfaceID, + replayBarrierToken: replayBarrierTokenForRequest, + coveredReplayBarrierDroppedOutputCount: nil, + reason: "stale_client" + ) else { + self.clearTerminalReplayBarrierIfCurrent( + surfaceID: surfaceID, + token: replayBarrierTokenForRequest, + reason: "stale_client" + ) + return + } return } - self.deliverTerminalBytes(deliverBytes, surfaceID: surfaceID) - } catch { - mobileShellLog.error("CMUX_REPLAY failed surface=\(surfaceID, privacy: .public) error=\(String(describing: error), privacy: .public)") // The replay request is the view-only/foreground-resume path. A // definitive auth failure here (after the RPC layer's // force-refresh-and-retry already gave up) must drive the re-auth // prompt instead of silently leaving a stale frame. - guard self.remoteClient === client else { return } - _ = self.disconnectForAuthorizationFailureIfNeeded(error) + guard !self.disconnectForAuthorizationFailureIfNeeded(error) else { return } + if let retryToken = self.prepareTerminalReplayFailureRetry( + surfaceID: surfaceID, + replayBarrierToken: replayBarrierTokenForRequest + ) { + self.clearTerminalReplayInFlightIfCurrent( + surfaceID: surfaceID, + requestID: replayRequestID + ) + transferredInFlightToRetry = true + self.requestTerminalReplay( + surfaceID: surfaceID, + replayBarrierToken: retryToken, + coveredReplayBarrierDroppedOutputCount: coveredReplayBarrierDroppedOutputCountForRequest + ) + return + } + self.preserveTerminalReplayBarrierIfCurrent( + surfaceID: surfaceID, + token: replayBarrierTokenForRequest, + reason: "failed" + ) } } + storeTerminalReplayTask( + surfaceID: surfaceID, + requestID: replayRequestID, + task: replayTask + ) } private func handleTerminalRenderGridEvent(_ event: MobileEventEnvelope) { @@ -7060,11 +7472,11 @@ public final class MobileShellComposite: MobileTerminalOutputSinking { } let overlap = deliveredSeq - seq let deliverBytes = Data(bytes.dropFirst(Int(overlap))) - deliverTerminalBytes(deliverBytes, surfaceID: surfaceID) + guard deliverTerminalBytes(deliverBytes, surfaceID: surfaceID) else { return } markTerminalBytesDelivered(surfaceID: surfaceID, endSeq: endSeq) return } - deliverTerminalBytes(bytes, surfaceID: surfaceID) + guard deliverTerminalBytes(bytes, surfaceID: surfaceID) else { return } markTerminalBytesDelivered(surfaceID: surfaceID, endSeq: endSeq) } diff --git a/Packages/iOS/CmuxMobileShell/Tests/CmuxMobileShellTests/MobileShellRenderGridLivenessTestSupport.swift b/Packages/iOS/CmuxMobileShell/Tests/CmuxMobileShellTests/MobileShellRenderGridLivenessTestSupport.swift index e09f84a16ee8..a0436e19a5ce 100644 --- a/Packages/iOS/CmuxMobileShell/Tests/CmuxMobileShellTests/MobileShellRenderGridLivenessTestSupport.swift +++ b/Packages/iOS/CmuxMobileShell/Tests/CmuxMobileShellTests/MobileShellRenderGridLivenessTestSupport.swift @@ -59,23 +59,110 @@ actor LivenessHostRouter { } private var recorded: [RecordedRequest] = [] + private var countWaiters: [( + id: UUID, + method: String, + expectedCount: Int, + continuation: CheckedContinuation + )] = [] private var hostStatusRequestCount = 0 private var heldHostStatusRequestNumbers: Set = [] private var subscribeRequestCount = 0 private var heldSubscribeRequestNumbers: Set = [] private var holdSubscribe = false + private var replayRequestCount = 0 + private var heldReplayRequestNumbers: Set = [] + private var heldReplayResponsesRemaining = 0 private var hasActiveSubscription = false private var heldContinuations: [CheckedContinuation] = [] private var capabilities = ["events.v1", "terminal.bytes.v1", "terminal.render_grid.v1", "terminal.replay.v1"] + private var replayTexts: [String] = [] + private var replayFailuresRemaining = 0 + private var emptyReplayResponsesRemaining = 0 func record(method: String?, topics: [String]?) { recorded.append(RecordedRequest(method: method, topics: topics)) + resumeSatisfiedCountWaiters() } func count(of method: String) -> Int { recorded.filter { $0.method == method }.count } + @discardableResult + func waitForCount( + of method: String, + atLeast expectedCount: Int, + timeoutNanoseconds: UInt64 = 3_000_000_000, + recordIssueOnTimeout: Bool = true + ) async -> Bool { + let reached = await withTaskGroup(of: Bool.self) { group in + group.addTask { + await self.waitUntilCountReached(of: method, atLeast: expectedCount) + return true + } + group.addTask { + // Test assertion deadline only; request arrival is signaled by record(). + try? await Task.sleep(nanoseconds: timeoutNanoseconds) + return false + } + let reached = await group.next() ?? false + group.cancelAll() + return reached + } + if !reached, recordIssueOnTimeout { + Issue.record("timed out waiting for \(method) count >= \(expectedCount)") + } + return reached + } + + private func waitUntilCountReached(of method: String, atLeast expectedCount: Int) async { + guard count(of: method) < expectedCount else { return } + let waiterID = UUID() + await withTaskCancellationHandler { + await withCheckedContinuation { continuation in + countWaiters.append(( + id: waiterID, + method: method, + expectedCount: expectedCount, + continuation: continuation + )) + resumeSatisfiedCountWaiters() + } + } onCancel: { + Task { await self.cancelCountWaiter(id: waiterID) } + } + } + + private func resumeSatisfiedCountWaiters() { + var remaining: [( + id: UUID, + method: String, + expectedCount: Int, + continuation: CheckedContinuation + )] = [] + var satisfied: [CheckedContinuation] = [] + for waiter in countWaiters { + if count(of: waiter.method) >= waiter.expectedCount { + satisfied.append(waiter.continuation) + } else { + remaining.append(waiter) + } + } + countWaiters = remaining + for continuation in satisfied { + continuation.resume() + } + } + + private func cancelCountWaiter(id: UUID) { + guard let index = countWaiters.firstIndex(where: { $0.id == id }) else { + return + } + let waiter = countWaiters.remove(at: index) + waiter.continuation.resume() + } + func topics(for method: String) -> [[String]] { recorded.compactMap { request in guard request.method == method else { return nil } @@ -87,6 +174,18 @@ actor LivenessHostRouter { self.capabilities = capabilities } + func enqueueReplayTexts(_ texts: [String]) { + replayTexts.append(contentsOf: texts) + } + + func failNextReplay(count: Int = 1) { + replayFailuresRemaining += count + } + + func enqueueEmptyReplayResponses(count: Int = 1) { + emptyReplayResponsesRemaining += count + } + /// Hold every `mobile.events.subscribe` response until released. func setHoldSubscribe(_ hold: Bool) { holdSubscribe = hold @@ -104,6 +203,18 @@ actor LivenessHostRouter { heldSubscribeRequestNumbers.insert(number) } + /// Hold the Nth `mobile.terminal.replay` response (1-based), letting a test + /// swap clients while the old request is still in flight. + func holdReplayRequest(number: Int) { + heldReplayRequestNumbers.insert(number) + } + + /// Hold the next N `mobile.terminal.replay` responses, independent of any + /// replay requests the connect/mount path already used. + func holdNextReplayResponses(count: Int = 1) { + heldReplayResponsesRemaining += count + } + /// Forget the host-side registration, modeling a lost subscription behind /// a live RPC channel: the next subscribe reports /// `already_subscribed: false`. @@ -117,6 +228,8 @@ actor LivenessHostRouter { holdSubscribe = false heldHostStatusRequestNumbers = [] heldSubscribeRequestNumbers = [] + heldReplayRequestNumbers = [] + heldReplayResponsesRemaining = 0 let continuations = heldContinuations heldContinuations = [] for continuation in continuations { @@ -169,7 +282,30 @@ actor LivenessHostRouter { "topics": ["workspace.updated", "terminal.render_grid"], "already_subscribed": alreadySubscribed, ]) - case "mobile.events.unsubscribe", "mobile.terminal.replay", "mobile.terminal.viewport": + case "mobile.terminal.replay": + replayRequestCount += 1 + if heldReplayResponsesRemaining > 0 { + heldReplayResponsesRemaining -= 1 + await park() + } else if heldReplayRequestNumbers.contains(replayRequestCount) { + await park() + } + if replayFailuresRemaining > 0 { + replayFailuresRemaining -= 1 + return try? Self.errorFrame(id: id, message: "replay failed") + } + if emptyReplayResponsesRemaining > 0 { + emptyReplayResponsesRemaining -= 1 + return try? Self.resultFrame(id: id, result: [:]) + } + guard !replayTexts.isEmpty else { + return try? Self.resultFrame(id: id, result: [:]) + } + let text = replayTexts.removeFirst() + return try? Self.resultFrame(id: id, result: [ + "data_b64": Data(text.utf8).base64EncodedString(), + ]) + case "mobile.events.unsubscribe", "mobile.terminal.viewport": return try? Self.resultFrame(id: id, result: [:]) case "terminal.input": return try? Self.resultFrame(id: id, result: [ @@ -454,3 +590,24 @@ func makeConnectedStore( #expect(connected, "scripted connect must succeed") return store } + +@MainActor +func installFreshLivenessRemoteClient( + on store: MobileShellComposite, + router: LivenessHostRouter, + box: TransportBox, + clock: TestClock +) throws { + let runtime = LivenessTestRuntime( + transportFactory: LivenessTransportFactory(router: router, box: box), + now: { clock.now } + ) + let ticket = try makeTicket(clock: clock) + let route = try #require(ticket.routes.first) + store.remoteClient = MobileCoreRPCClient( + runtime: runtime, + route: route, + ticket: ticket, + allowsStackAuthFallback: true + ) +} diff --git a/Packages/iOS/CmuxMobileShell/Tests/CmuxMobileShellTests/MobileShellReplayFallbackScreenTests.swift b/Packages/iOS/CmuxMobileShell/Tests/CmuxMobileShellTests/MobileShellReplayFallbackScreenTests.swift new file mode 100644 index 000000000000..b3bc2265d588 --- /dev/null +++ b/Packages/iOS/CmuxMobileShell/Tests/CmuxMobileShellTests/MobileShellReplayFallbackScreenTests.swift @@ -0,0 +1,81 @@ +import Foundation +import Testing +@testable import CmuxMobileShell + +@MainActor +@Test func replayByteFallbackPreservesAlternateScreenSuppression() async throws { + let clock = TestClock() + let router = LivenessHostRouter() + let box = TransportBox() + let store = try await makeConnectedStore(router: router, box: box, clock: clock) + + await router.enqueueReplayTexts(["cold-replay", "fallback-replay"]) + let collector = OutputCollector() + collector.mount(store: store, surfaceID: "live-terminal") + await router.waitForCount(of: "mobile.terminal.replay", atLeast: 1) + let coldReplayDelivered = try await pollUntil { + collector.lines.contains { $0.contains("cold-replay") } + } + #expect(coldReplayDelivered) + let transport = try #require(box.get()) + + await transport.deliver(try renderGridEventFrame( + surfaceID: "live-terminal", + seq: 3, + text: "alt", + activeScreen: .alternate + )) + let altDelivered = try await pollUntil { + collector.viewportPolicies.last == .remoteGrid(columns: 16, rows: 4) + } + #expect(altDelivered) + + let replayCountAfterAlternate = await router.count(of: "mobile.terminal.replay") + let replayBarrierToken = store.beginTerminalReplayBarrier(surfaceID: "live-terminal") + store.requestTerminalReplay( + surfaceID: "live-terminal", + replayBarrierToken: replayBarrierToken + ) + await router.waitForCount( + of: "mobile.terminal.replay", + atLeast: replayCountAfterAlternate + 1 + ) + let fallbackReplayDelivered = try await pollUntil { + collector.lines.contains { $0.contains("fallback-replay") } + } + #expect(fallbackReplayDelivered) + #expect(store.terminalReplayBarrierTokensBySurfaceID["live-terminal"] == nil) + + await transport.deliver(try terminalBytesEventFrame( + surfaceID: "live-terminal", + seq: 4, + text: "raw-after-fallback" + )) + let rawAfterFallbackDelivered = try await pollUntil(attempts: 50) { + collector.lines.contains { $0.contains("raw-after-fallback") } + } + #expect( + !rawAfterFallbackDelivered, + "byte/snapshot replay fallbacks must not clear alternate-screen raw-byte suppression" + ) + + await transport.deliver(try renderGridEventFrame( + surfaceID: "live-terminal", + seq: 10, + text: "primary-full" + )) + let primaryDelivered = try await pollUntil { + collector.lines.contains { $0.contains("primary-full") } + } + #expect(primaryDelivered) + await transport.deliver(try terminalBytesEventFrame( + surfaceID: "live-terminal", + seq: 10, + text: "raw-after-primary" + )) + let rawAfterPrimaryDelivered = try await pollUntil { + collector.lines.contains { $0.contains("raw-after-primary") } + } + #expect(rawAfterPrimaryDelivered) + collector.unmount() +} diff --git a/Packages/iOS/CmuxMobileShell/Tests/CmuxMobileShellTests/TerminalOutputDeliveryQueueTests.swift b/Packages/iOS/CmuxMobileShell/Tests/CmuxMobileShellTests/TerminalOutputDeliveryQueueTests.swift index 4b8db691f153..3a338ecffbbf 100644 --- a/Packages/iOS/CmuxMobileShell/Tests/CmuxMobileShellTests/TerminalOutputDeliveryQueueTests.swift +++ b/Packages/iOS/CmuxMobileShell/Tests/CmuxMobileShellTests/TerminalOutputDeliveryQueueTests.swift @@ -38,7 +38,655 @@ import Testing store.terminalOutputDidProcess(surfaceID: surfaceID, streamToken: currentChunk.streamToken) let secondChunk = try #require(await currentIterator.next()) - #expect(String(decoding: secondChunk.data, as: UTF8.self) == "new-second") + #expect(String(data: secondChunk.data, encoding: .utf8) == "new-second") +} + +@MainActor +@Test func terminalReplayBarrierDropsStalledBacklogAndInvalidatesOldAcks() async throws { + let store = MobileShellComposite.preview() + let surfaceID = "terminal" + + var iterator = store.terminalOutputStream(surfaceID: surfaceID).makeAsyncIterator() + store.deliverTerminalBytes(Data("stalled-first".utf8), surfaceID: surfaceID) + let stalledChunk = try #require(await iterator.next()) + store.deliverTerminalBytes(Data("stale-second".utf8), surfaceID: surfaceID) + + #expect(store.terminalOutputQueuesBySurfaceID[surfaceID]?.pendingCount == 1) + + let replayBarrierToken = store.beginTerminalReplayBarrier(surfaceID: surfaceID) + + #expect(store.terminalOutputQueuesBySurfaceID[surfaceID]?.isIdle == true) + #expect(store.terminalReplayBarrierTokensBySurfaceID[surfaceID] == replayBarrierToken) + + let liveBeforeReplayAccepted = store.deliverTerminalBytes( + Data("live-before-replay".utf8), + surfaceID: surfaceID + ) + #expect(liveBeforeReplayAccepted == false) + #expect(store.terminalOutputQueuesBySurfaceID[surfaceID]?.isIdle == true) + + store.terminalOutputDidProcess(surfaceID: surfaceID, streamToken: stalledChunk.streamToken) + + store.deliverTerminalBytes( + Data("authoritative-replay".utf8), + surfaceID: surfaceID, + bypassReplayBarrier: true + ) + let replayChunk = try #require(await iterator.next()) + #expect(String(data: replayChunk.data, encoding: .utf8) == "authoritative-replay") + #expect(replayChunk.streamToken != stalledChunk.streamToken) + + let liveBeforeReplayAckAccepted = store.deliverTerminalBytes( + Data("live-before-replay-ack".utf8), + surfaceID: surfaceID + ) + #expect(liveBeforeReplayAckAccepted == false) + #expect(store.terminalOutputQueuesBySurfaceID[surfaceID]?.pendingCount == 0) + + let afterStaleAckAccepted = store.deliverTerminalBytes( + Data("after-stale-ack".utf8), + surfaceID: surfaceID + ) + #expect(afterStaleAckAccepted == false) + #expect(store.terminalOutputQueuesBySurfaceID[surfaceID]?.pendingCount == 0) + + store.terminalOutputDidProcess(surfaceID: surfaceID, streamToken: replayChunk.streamToken) + #expect(store.terminalReplayBarrierTokensBySurfaceID[surfaceID] == nil) + + store.deliverTerminalBytes(Data("after-replay-ack".utf8), surfaceID: surfaceID) + + let afterReplayAck = try #require(await iterator.next()) + #expect(String(data: afterReplayAck.data, encoding: .utf8) == "after-replay-ack") +} + +@MainActor +@Test func terminalOutputResetClearsBarrierWhenReplayCannotStart() async throws { + let store = MobileShellComposite.preview() + let surfaceID = "terminal" + + var iterator = store.terminalOutputStream(surfaceID: surfaceID).makeAsyncIterator() + store.deliverTerminalBytes(Data("stalled-first".utf8), surfaceID: surfaceID) + let stalledChunk = try #require(await iterator.next()) + store.deliverTerminalBytes(Data("stale-second".utf8), surfaceID: surfaceID) + + #expect(store.terminalOutputQueuesBySurfaceID[surfaceID]?.pendingCount == 1) + + store.terminalOutputDidReset(surfaceID: surfaceID, streamToken: stalledChunk.streamToken) + + #expect(store.terminalOutputQueuesBySurfaceID[surfaceID]?.isIdle == true) + #expect(store.terminalReplayBarrierTokensBySurfaceID[surfaceID] == nil) + + store.terminalOutputDidProcess(surfaceID: surfaceID, streamToken: stalledChunk.streamToken) + + let accepted = store.deliverTerminalBytes(Data("after-aborted-replay".utf8), surfaceID: surfaceID) + #expect(accepted == true) + + let afterAbort = try #require(await iterator.next()) + #expect(String(data: afterAbort.data, encoding: .utf8) == "after-aborted-replay") +} + +@MainActor +@Test func terminalReplayBarrierRequestsFollowUpWhenLiveOutputDropsBeforeAck() async throws { + let router = LivenessHostRouter() + let box = TransportBox() + let clock = TestClock() + let store = try await makeConnectedStore(router: router, box: box, clock: clock) + let surfaceID = "live-terminal" + + await router.enqueueReplayTexts(["cold-replay", "first-replay", "follow-up-replay"]) + var iterator = store.terminalOutputStream(surfaceID: surfaceID).makeAsyncIterator() + await router.waitForCount(of: "mobile.terminal.replay", atLeast: 1) + let coldReplayChunk = try #require(await iterator.next()) + #expect(String(data: coldReplayChunk.data, encoding: .utf8) == "cold-replay") + store.terminalOutputDidProcess(surfaceID: surfaceID, streamToken: coldReplayChunk.streamToken) + let coldReplaySettled = await waitForReplayBarrierFailureToSettle { + !store.terminalReplaySurfaceIDsInFlight.contains(surfaceID) + } + #expect(coldReplaySettled, "cold mount replay must settle before starting the held replay") + let replayCountAfterMount = await router.count(of: "mobile.terminal.replay") + + store.deliverTerminalBytes(Data("stalled-first".utf8), surfaceID: surfaceID) + let stalledChunk = try #require(await iterator.next()) + store.deliverTerminalBytes(Data("stale-second".utf8), surfaceID: surfaceID) + + store.terminalOutputDidReset(surfaceID: surfaceID, streamToken: stalledChunk.streamToken) + await router.waitForCount(of: "mobile.terminal.replay", atLeast: replayCountAfterMount + 1) + + let replayChunk = try #require(await iterator.next()) + #expect(String(data: replayChunk.data, encoding: .utf8) == "first-replay") + + let acceptedDuringBarrier = store.deliverTerminalBytes( + Data("live-during-barrier".utf8), + surfaceID: surfaceID + ) + #expect(acceptedDuringBarrier == false) + #expect(store.terminalReplayBarrierDroppedOutputSurfaceIDs.contains(surfaceID)) + + store.terminalOutputDidProcess(surfaceID: surfaceID, streamToken: replayChunk.streamToken) + await router.waitForCount(of: "mobile.terminal.replay", atLeast: replayCountAfterMount + 2) + + let followUpChunk = try #require(await iterator.next()) + #expect(String(data: followUpChunk.data, encoding: .utf8) == "follow-up-replay") + #expect(!store.terminalReplayBarrierDroppedOutputSurfaceIDs.contains(surfaceID)) + + store.terminalOutputDidProcess(surfaceID: surfaceID, streamToken: followUpChunk.streamToken) + #expect(store.terminalReplayBarrierTokensBySurfaceID[surfaceID] == nil) + + store.deliverTerminalBytes(Data("after-follow-up".utf8), surfaceID: surfaceID) + let afterFollowUp = try #require(await iterator.next()) + #expect(String(data: afterFollowUp.data, encoding: .utf8) == "after-follow-up") +} + +@MainActor +@Test func terminalReplayBarrierRetriesAfterReplayFailureWithDroppedOutput() async throws { + let router = LivenessHostRouter() + let box = TransportBox() + let clock = TestClock() + let store = try await makeConnectedStore(router: router, box: box, clock: clock) + let surfaceID = "live-terminal" + + await router.enqueueReplayTexts(["cold-replay", "retry-replay"]) + var iterator = store.terminalOutputStream(surfaceID: surfaceID).makeAsyncIterator() + await router.waitForCount(of: "mobile.terminal.replay", atLeast: 1) + let coldReplayChunk = try #require(await iterator.next()) + #expect(String(data: coldReplayChunk.data, encoding: .utf8) == "cold-replay") + store.terminalOutputDidProcess(surfaceID: surfaceID, streamToken: coldReplayChunk.streamToken) + let replayCountAfterMount = await router.count(of: "mobile.terminal.replay") + + await router.failNextReplay() + _ = store.beginTerminalReplayBarrier(surfaceID: surfaceID) + let firstDropAccepted = store.deliverTerminalBytes( + Data("live-during-failed-replay".utf8), + surfaceID: surfaceID + ) + #expect(firstDropAccepted == false) + + await router.waitForCount(of: "mobile.terminal.replay", atLeast: replayCountAfterMount + 1) + + let retryRequested = await waitForReplayRequestCount( + router, + atLeast: replayCountAfterMount + 2 + ) + #expect(retryRequested, "failed replay with dropped output must request a replacement replay") + #expect(store.terminalReplayBarrierDroppedOutputSurfaceIDs.contains(surfaceID)) + + let retryReplayChunk = try #require(await iterator.next()) + #expect(String(data: retryReplayChunk.data, encoding: .utf8) == "retry-replay") + store.terminalOutputDidProcess(surfaceID: surfaceID, streamToken: retryReplayChunk.streamToken) + #expect(store.terminalReplayBarrierTokensBySurfaceID[surfaceID] == nil) + #expect(!store.terminalReplayBarrierDroppedOutputSurfaceIDs.contains(surfaceID)) +} + +@MainActor +@Test func terminalReplayBarrierRetriesAfterReplayFailureWithoutDroppedOutput() async throws { + let router = LivenessHostRouter() + let box = TransportBox() + let clock = TestClock() + let store = try await makeConnectedStore(router: router, box: box, clock: clock) + let surfaceID = "live-terminal" + + await router.enqueueReplayTexts(["cold-replay", "retry-replay"]) + var iterator = store.terminalOutputStream(surfaceID: surfaceID).makeAsyncIterator() + await router.waitForCount(of: "mobile.terminal.replay", atLeast: 1) + let coldReplayChunk = try #require(await iterator.next()) + #expect(String(data: coldReplayChunk.data, encoding: .utf8) == "cold-replay") + store.terminalOutputDidProcess(surfaceID: surfaceID, streamToken: coldReplayChunk.streamToken) + let replayCountAfterMount = await router.count(of: "mobile.terminal.replay") + + await router.failNextReplay() + store.deliverTerminalBytes(Data("stalled-first".utf8), surfaceID: surfaceID) + let stalledChunk = try #require(await iterator.next()) + store.terminalOutputDidReset(surfaceID: surfaceID, streamToken: stalledChunk.streamToken) + + let retryRequested = await waitForReplayRequestCount( + router, + atLeast: replayCountAfterMount + 2 + ) + #expect(retryRequested, "failed reset replay must retry even without a later live-output drop") + #expect(store.terminalReplayBarrierTokensBySurfaceID[surfaceID] != nil) + + if retryRequested { + let retryReplayChunk = try #require(await iterator.next()) + #expect(String(data: retryReplayChunk.data, encoding: .utf8) == "retry-replay") + store.terminalOutputDidProcess(surfaceID: surfaceID, streamToken: retryReplayChunk.streamToken) + #expect(store.terminalReplayBarrierTokensBySurfaceID[surfaceID] == nil) + } +} + +@MainActor +@Test func terminalReplayBarrierReplaysOnReplacementClientAfterStaleResponse() async throws { + let router = LivenessHostRouter() + let box = TransportBox() + let clock = TestClock() + let store = try await makeConnectedStore(router: router, box: box, clock: clock) + let surfaceID = "live-terminal" + + await router.enqueueReplayTexts(["cold-replay"]) + var iterator = store.terminalOutputStream(surfaceID: surfaceID).makeAsyncIterator() + await router.waitForCount(of: "mobile.terminal.replay", atLeast: 1) + let coldReplayChunk = try #require(await iterator.next()) + #expect(String(data: coldReplayChunk.data, encoding: .utf8) == "cold-replay") + store.terminalOutputDidProcess(surfaceID: surfaceID, streamToken: coldReplayChunk.streamToken) + let replayCountAfterMount = await router.count(of: "mobile.terminal.replay") + await router.holdNextReplayResponses() + + store.deliverTerminalBytes(Data("stalled-first".utf8), surfaceID: surfaceID) + let stalledChunk = try #require(await iterator.next()) + store.terminalOutputDidReset(surfaceID: surfaceID, streamToken: stalledChunk.streamToken) + await router.waitForCount(of: "mobile.terminal.replay", atLeast: replayCountAfterMount + 1) + + let replacementRouter = LivenessHostRouter() + let replacementBox = TransportBox() + await replacementRouter.enqueueReplayTexts(["replacement-replay"]) + try installFreshLivenessRemoteClient( + on: store, + router: replacementRouter, + box: replacementBox, + clock: clock + ) + await router.enqueueReplayTexts(["stale-replay"]) + await router.releaseAllHeld() + + let replacementReplayRequested = await waitForReplayRequestCount( + replacementRouter, + atLeast: 1 + ) + #expect(replacementReplayRequested, "stale replay responses must resync on the replacement client") + + if replacementReplayRequested { + let replacementReplayChunk = try #require(await iterator.next()) + #expect(String(data: replacementReplayChunk.data, encoding: .utf8) == "replacement-replay") + store.terminalOutputDidProcess(surfaceID: surfaceID, streamToken: replacementReplayChunk.streamToken) + #expect(store.terminalReplayBarrierTokensBySurfaceID[surfaceID] == nil) + } +} + +@MainActor +@Test func terminalReplayBarrierReplaysOnReplacementClientAfterStaleFailure() async throws { + let router = LivenessHostRouter() + let box = TransportBox() + let clock = TestClock() + let store = try await makeConnectedStore(router: router, box: box, clock: clock) + let surfaceID = "live-terminal" + + await router.enqueueReplayTexts(["cold-replay"]) + var iterator = store.terminalOutputStream(surfaceID: surfaceID).makeAsyncIterator() + await router.waitForCount(of: "mobile.terminal.replay", atLeast: 1) + let coldReplayChunk = try #require(await iterator.next()) + #expect(String(data: coldReplayChunk.data, encoding: .utf8) == "cold-replay") + store.terminalOutputDidProcess(surfaceID: surfaceID, streamToken: coldReplayChunk.streamToken) + let replayCountAfterMount = await router.count(of: "mobile.terminal.replay") + await router.holdNextReplayResponses() + + store.deliverTerminalBytes(Data("stalled-first".utf8), surfaceID: surfaceID) + let stalledChunk = try #require(await iterator.next()) + store.terminalOutputDidReset(surfaceID: surfaceID, streamToken: stalledChunk.streamToken) + await router.waitForCount(of: "mobile.terminal.replay", atLeast: replayCountAfterMount + 1) + + let replacementRouter = LivenessHostRouter() + let replacementBox = TransportBox() + await replacementRouter.enqueueReplayTexts(["replacement-replay"]) + try installFreshLivenessRemoteClient( + on: store, + router: replacementRouter, + box: replacementBox, + clock: clock + ) + await router.failNextReplay() + await router.releaseAllHeld() + + let replacementReplayRequested = await waitForReplayRequestCount( + replacementRouter, + atLeast: 1 + ) + #expect(replacementReplayRequested, "stale replay failures must resync on the replacement client") + + if replacementReplayRequested { + let replacementReplayChunk = try #require(await iterator.next()) + #expect(String(data: replacementReplayChunk.data, encoding: .utf8) == "replacement-replay") + store.terminalOutputDidProcess(surfaceID: surfaceID, streamToken: replacementReplayChunk.streamToken) + #expect(store.terminalReplayBarrierTokensBySurfaceID[surfaceID] == nil) + } +} + +@MainActor +@Test func terminalReplayBarrierStaysActiveAfterRetryExhaustionWithoutDroppedOutput() async throws { + let router = LivenessHostRouter() + let box = TransportBox() + let clock = TestClock() + let store = try await makeConnectedStore(router: router, box: box, clock: clock) + let surfaceID = "live-terminal" + + await router.enqueueReplayTexts(["cold-replay"]) + var iterator = store.terminalOutputStream(surfaceID: surfaceID).makeAsyncIterator() + await router.waitForCount(of: "mobile.terminal.replay", atLeast: 1) + let coldReplayChunk = try #require(await iterator.next()) + #expect(String(data: coldReplayChunk.data, encoding: .utf8) == "cold-replay") + store.terminalOutputDidProcess(surfaceID: surfaceID, streamToken: coldReplayChunk.streamToken) + let replayCountAfterMount = await router.count(of: "mobile.terminal.replay") + + await router.failNextReplay(count: 3) + store.deliverTerminalBytes(Data("stalled-first".utf8), surfaceID: surfaceID) + let stalledChunk = try #require(await iterator.next()) + store.terminalOutputDidReset(surfaceID: surfaceID, streamToken: stalledChunk.streamToken) + + let exhaustedRetries = await waitForReplayRequestCount( + router, + atLeast: replayCountAfterMount + 3 + ) + #expect(exhaustedRetries, "reset replay should exhaust the initial request plus two retries") + + let failureSettled = await waitForReplayBarrierFailureToSettle { + !store.terminalReplaySurfaceIDsInFlight.contains(surfaceID) + } + #expect(failureSettled) + #expect(store.terminalReplayBarrierTokensBySurfaceID[surfaceID] != nil) + + let replayCountAfterExhaustion = await router.count(of: "mobile.terminal.replay") + let acceptedAfterExhaustion = store.deliverTerminalBytes( + Data("live-after-exhausted-replay".utf8), + surfaceID: surfaceID + ) + #expect(acceptedAfterExhaustion == false) + #expect(store.terminalOutputQueuesBySurfaceID[surfaceID]?.isIdle == true) + let replayRestartedAfterExhaustion = await waitForReplayRequestCount( + router, + atLeast: replayCountAfterExhaustion + 1 + ) + #expect(!replayRestartedAfterExhaustion, "dropped live output must not bypass the replay retry cap") +} + +@MainActor +@Test func terminalReplayInFlightClearsWhenOutputStreamUnmountsBeforeResponse() async throws { + let router = LivenessHostRouter() + let box = TransportBox() + let clock = TestClock() + let store = try await makeConnectedStore(router: router, box: box, clock: clock) + let surfaceID = "live-terminal" + + await router.holdNextReplayResponses() + let collector = OutputCollector() + collector.mount(store: store, surfaceID: surfaceID) + await router.waitForCount(of: "mobile.terminal.replay", atLeast: 1) + #expect(store.terminalReplaySurfaceIDsInFlight.contains(surfaceID)) + + collector.unmount() + let unregistered = await waitForReplayBarrierFailureToSettle { + store.terminalByteContinuationsBySurfaceID[surfaceID] == nil + } + #expect(unregistered) + #expect(!store.terminalReplaySurfaceIDsInFlight.contains(surfaceID)) + + let replayCountAfterUnmount = await router.count(of: "mobile.terminal.replay") + await router.enqueueReplayTexts(["remount-replay"]) + let remountCollector = OutputCollector() + remountCollector.mount(store: store, surfaceID: surfaceID) + + await router.waitForCount( + of: "mobile.terminal.replay", + atLeast: replayCountAfterUnmount + 1 + ) + let replayDelivered = try await pollUntil { + remountCollector.lines.contains("remount-replay") + } + #expect(replayDelivered) + + remountCollector.unmount() + await router.releaseAllHeld() +} + +@MainActor +@Test func staleReplayResponseAfterRemountDoesNotDeliverIntoCurrentStream() async throws { + let router = LivenessHostRouter() + let box = TransportBox() + let clock = TestClock() + let store = try await makeConnectedStore(router: router, box: box, clock: clock) + let surfaceID = "live-terminal" + + await router.holdNextReplayResponses() + let oldCollector = OutputCollector() + oldCollector.mount(store: store, surfaceID: surfaceID) + await router.waitForCount(of: "mobile.terminal.replay", atLeast: 1) + oldCollector.unmount() + + let unregistered = await waitForReplayBarrierFailureToSettle { + store.terminalByteContinuationsBySurfaceID[surfaceID] == nil + } + #expect(unregistered) + let replayCountAfterUnmount = await router.count(of: "mobile.terminal.replay") + + await router.enqueueReplayTexts(["fresh-replay", "stale-replay"]) + let currentCollector = OutputCollector() + currentCollector.mount(store: store, surfaceID: surfaceID) + await router.waitForCount( + of: "mobile.terminal.replay", + atLeast: replayCountAfterUnmount + 1 + ) + + let freshDelivered = try await pollUntil { + currentCollector.lines.contains("fresh-replay") + } + #expect(freshDelivered) + + await router.releaseAllHeld() + let staleDelivered = try await pollUntil(attempts: 50) { + currentCollector.lines.contains("stale-replay") + } + #expect(!staleDelivered, "superseded replay responses must not deliver stale bytes") + + currentCollector.unmount() +} + +@MainActor +@Test func genericReplayRequestReusesPreservedBarrierAfterRetryExhaustion() async throws { + let router = LivenessHostRouter() + let box = TransportBox() + let clock = TestClock() + let store = try await makeConnectedStore(router: router, box: box, clock: clock) + let surfaceID = "live-terminal" + + await router.enqueueReplayTexts(["cold-replay"]) + var iterator = store.terminalOutputStream(surfaceID: surfaceID).makeAsyncIterator() + await router.waitForCount(of: "mobile.terminal.replay", atLeast: 1) + let coldReplayChunk = try #require(await iterator.next()) + #expect(String(data: coldReplayChunk.data, encoding: .utf8) == "cold-replay") + store.terminalOutputDidProcess(surfaceID: surfaceID, streamToken: coldReplayChunk.streamToken) + let replayCountAfterMount = await router.count(of: "mobile.terminal.replay") + + await router.failNextReplay(count: 3) + store.deliverTerminalBytes(Data("stalled-first".utf8), surfaceID: surfaceID) + let stalledChunk = try #require(await iterator.next()) + store.terminalOutputDidReset(surfaceID: surfaceID, streamToken: stalledChunk.streamToken) + + let exhaustedRetries = await waitForReplayRequestCount( + router, + atLeast: replayCountAfterMount + 3 + ) + #expect(exhaustedRetries, "reset replay should exhaust the initial request plus two retries") + + let failureSettled = await waitForReplayBarrierFailureToSettle { + !store.terminalReplaySurfaceIDsInFlight.contains(surfaceID) + } + #expect(failureSettled) + let preservedBarrierToken = try #require(store.terminalReplayBarrierTokensBySurfaceID[surfaceID]) + + let replayCountAfterExhaustion = await router.count(of: "mobile.terminal.replay") + await router.enqueueReplayTexts(["resync-replay"]) + store.requestTerminalReplay(surfaceID: surfaceID) + + let genericReplayRequested = await waitForReplayRequestCount( + router, + atLeast: replayCountAfterExhaustion + 1 + ) + #expect(genericReplayRequested, "generic resync must not be blocked by a preserved barrier") + + if genericReplayRequested { + let replayChunk = try #require(await iterator.next()) + #expect(String(data: replayChunk.data, encoding: .utf8) == "resync-replay") + #expect(store.terminalReplayBarrierTokensBySurfaceID[surfaceID] == preservedBarrierToken) + store.terminalOutputDidProcess(surfaceID: surfaceID, streamToken: replayChunk.streamToken) + #expect(store.terminalReplayBarrierTokensBySurfaceID[surfaceID] == nil) + } +} + +@MainActor +@Test func terminalReplayBarrierRetriesWhenReplayAckChunkResets() async throws { + let router = LivenessHostRouter() + let box = TransportBox() + let clock = TestClock() + let store = try await makeConnectedStore(router: router, box: box, clock: clock) + let surfaceID = "live-terminal" + + await router.enqueueReplayTexts(["cold-replay", "first-replay", "retry-replay"]) + var iterator = store.terminalOutputStream(surfaceID: surfaceID).makeAsyncIterator() + await router.waitForCount(of: "mobile.terminal.replay", atLeast: 1) + let coldReplayChunk = try #require(await iterator.next()) + #expect(String(data: coldReplayChunk.data, encoding: .utf8) == "cold-replay") + store.terminalOutputDidProcess(surfaceID: surfaceID, streamToken: coldReplayChunk.streamToken) + let replayCountAfterMount = await router.count(of: "mobile.terminal.replay") + + store.deliverTerminalBytes(Data("stalled-first".utf8), surfaceID: surfaceID) + let stalledChunk = try #require(await iterator.next()) + store.deliverTerminalBytes(Data("stale-second".utf8), surfaceID: surfaceID) + + store.terminalOutputDidReset(surfaceID: surfaceID, streamToken: stalledChunk.streamToken) + await router.waitForCount(of: "mobile.terminal.replay", atLeast: replayCountAfterMount + 1) + + let replayChunk = try #require(await iterator.next()) + #expect(String(data: replayChunk.data, encoding: .utf8) == "first-replay") + let firstBarrierToken = try #require(store.terminalReplayBarrierTokensBySurfaceID[surfaceID]) + + store.terminalOutputDidReset(surfaceID: surfaceID, streamToken: replayChunk.streamToken) + await router.waitForCount(of: "mobile.terminal.replay", atLeast: replayCountAfterMount + 2) + + let retryReplayChunk = try #require(await iterator.next()) + #expect(String(data: retryReplayChunk.data, encoding: .utf8) == "retry-replay") + let retryBarrierToken = try #require(store.terminalReplayBarrierTokensBySurfaceID[surfaceID]) + #expect(retryBarrierToken == firstBarrierToken) + + store.terminalOutputDidProcess(surfaceID: surfaceID, streamToken: retryReplayChunk.streamToken) + #expect(store.terminalReplayBarrierTokensBySurfaceID[surfaceID] == nil) +} + +@MainActor +@Test func terminalReplayBarrierRequestsFollowUpForDropAfterCoveredRetryResponse() async throws { + let router = LivenessHostRouter() + let box = TransportBox() + let clock = TestClock() + let store = try await makeConnectedStore(router: router, box: box, clock: clock) + let surfaceID = "live-terminal" + + await router.enqueueReplayTexts(["cold-replay", "retry-replay", "follow-up-replay"]) + var iterator = store.terminalOutputStream(surfaceID: surfaceID).makeAsyncIterator() + await router.waitForCount(of: "mobile.terminal.replay", atLeast: 1) + let coldReplayChunk = try #require(await iterator.next()) + #expect(String(data: coldReplayChunk.data, encoding: .utf8) == "cold-replay") + store.terminalOutputDidProcess(surfaceID: surfaceID, streamToken: coldReplayChunk.streamToken) + let replayCountAfterMount = await router.count(of: "mobile.terminal.replay") + + await router.failNextReplay() + let replayBarrierToken = store.beginTerminalReplayBarrier(surfaceID: surfaceID) + let firstDropAccepted = store.deliverTerminalBytes( + Data("live-during-failed-replay".utf8), + surfaceID: surfaceID + ) + #expect(firstDropAccepted == false) + await router.waitForCount(of: "mobile.terminal.replay", atLeast: replayCountAfterMount + 1) + + let failureSettled = await waitForReplayBarrierFailureToSettle { + !store.terminalReplaySurfaceIDsInFlight.contains(surfaceID) + && store.terminalReplayBarrierTokensBySurfaceID[surfaceID] == replayBarrierToken + } + #expect(failureSettled) + + let retryDropAccepted = store.deliverTerminalBytes( + Data("live-before-retry-request".utf8), + surfaceID: surfaceID + ) + #expect(retryDropAccepted == false) + await router.waitForCount(of: "mobile.terminal.replay", atLeast: replayCountAfterMount + 2) + + let retryReplayChunk = try #require(await iterator.next()) + #expect(String(data: retryReplayChunk.data, encoding: .utf8) == "retry-replay") + + let postResponseDropAccepted = store.deliverTerminalBytes( + Data("live-after-retry-response-before-ack".utf8), + surfaceID: surfaceID + ) + #expect(postResponseDropAccepted == false) + + store.terminalOutputDidProcess(surfaceID: surfaceID, streamToken: retryReplayChunk.streamToken) + await router.waitForCount(of: "mobile.terminal.replay", atLeast: replayCountAfterMount + 3) + + let followUpChunk = try #require(await iterator.next()) + #expect(String(data: followUpChunk.data, encoding: .utf8) == "follow-up-replay") + store.terminalOutputDidProcess(surfaceID: surfaceID, streamToken: followUpChunk.streamToken) + #expect(store.terminalReplayBarrierTokensBySurfaceID[surfaceID] == nil) +} + +@MainActor +@Test func terminalReplayBarrierRetriesEmptyReplayAfterDroppedOutput() async throws { + let router = LivenessHostRouter() + let box = TransportBox() + let clock = TestClock() + let store = try await makeConnectedStore(router: router, box: box, clock: clock) + let surfaceID = "live-terminal" + + await router.enqueueReplayTexts(["cold-replay"]) + var iterator = store.terminalOutputStream(surfaceID: surfaceID).makeAsyncIterator() + await router.waitForCount(of: "mobile.terminal.replay", atLeast: 1) + let coldReplayChunk = try #require(await iterator.next()) + #expect(String(data: coldReplayChunk.data, encoding: .utf8) == "cold-replay") + store.terminalOutputDidProcess(surfaceID: surfaceID, streamToken: coldReplayChunk.streamToken) + let replayCountAfterMount = await router.count(of: "mobile.terminal.replay") + + await router.enqueueEmptyReplayResponses() + await router.enqueueReplayTexts(["follow-up-replay"]) + let replayBarrierToken = store.beginTerminalReplayBarrier(surfaceID: surfaceID) + let droppedOutputAccepted = store.deliverTerminalBytes( + Data("live-during-empty-replay".utf8), + surfaceID: surfaceID + ) + #expect(droppedOutputAccepted == false) + #expect(store.terminalReplayBarrierTokensBySurfaceID[surfaceID] == replayBarrierToken) + + await router.waitForCount(of: "mobile.terminal.replay", atLeast: replayCountAfterMount + 1) + let followUpReplayRequested = await waitForReplayRequestCount( + router, + atLeast: replayCountAfterMount + 2 + ) + #expect(followUpReplayRequested, "empty replay with dropped output must immediately request a replacement replay") + + if followUpReplayRequested { + let followUpChunk = try #require(await iterator.next()) + #expect(String(data: followUpChunk.data, encoding: .utf8) == "follow-up-replay") + store.terminalOutputDidProcess(surfaceID: surfaceID, streamToken: followUpChunk.streamToken) + #expect(store.terminalReplayBarrierTokensBySurfaceID[surfaceID] == nil) + } +} + +@MainActor +private func waitForReplayBarrierFailureToSettle( + _ condition: @MainActor () -> Bool +) async -> Bool { + for _ in 0..<300 { + if condition() { + return true + } + try? await Task.sleep(nanoseconds: 10_000_000) + } + return condition() +} + +private func waitForReplayRequestCount( + _ router: LivenessHostRouter, + atLeast expectedCount: Int +) async -> Bool { + await router.waitForCount( + of: "mobile.terminal.replay", + atLeast: expectedCount, + recordIssueOnTimeout: false + ) } @Test func terminalOutputQueueCoalescesReplaceableViewportFramesBehindBackpressure() { diff --git a/Packages/iOS/CmuxMobileShell/Tests/CmuxMobileShellTests/TerminalReplayBarrierAckResetTests.swift b/Packages/iOS/CmuxMobileShell/Tests/CmuxMobileShellTests/TerminalReplayBarrierAckResetTests.swift new file mode 100644 index 000000000000..389cc7d939a1 --- /dev/null +++ b/Packages/iOS/CmuxMobileShell/Tests/CmuxMobileShellTests/TerminalReplayBarrierAckResetTests.swift @@ -0,0 +1,56 @@ +import Foundation +import Testing +@testable import CmuxMobileShell + +@MainActor +@Test func terminalReplayBarrierCapsRepeatedReplayAckResets() async throws { + let router = LivenessHostRouter() + let box = TransportBox() + let clock = TestClock() + let store = try await makeConnectedStore(router: router, box: box, clock: clock) + let surfaceID = "live-terminal" + + await router.enqueueReplayTexts([ + "cold-replay", + "first-replay", + "second-replay", + "third-replay", + "unexpected-replay", + ]) + var iterator = store.terminalOutputStream(surfaceID: surfaceID).makeAsyncIterator() + await router.waitForCount(of: "mobile.terminal.replay", atLeast: 1) + let coldReplayChunk = try #require(await iterator.next()) + #expect(String(data: coldReplayChunk.data, encoding: .utf8) == "cold-replay") + store.terminalOutputDidProcess(surfaceID: surfaceID, streamToken: coldReplayChunk.streamToken) + let replayCountAfterMount = await router.count(of: "mobile.terminal.replay") + + store.deliverTerminalBytes(Data("stalled-first".utf8), surfaceID: surfaceID) + let stalledChunk = try #require(await iterator.next()) + store.terminalOutputDidReset(surfaceID: surfaceID, streamToken: stalledChunk.streamToken) + await router.waitForCount(of: "mobile.terminal.replay", atLeast: replayCountAfterMount + 1) + + let firstReplayChunk = try #require(await iterator.next()) + #expect(String(data: firstReplayChunk.data, encoding: .utf8) == "first-replay") + let barrierToken = try #require(store.terminalReplayBarrierTokensBySurfaceID[surfaceID]) + + store.terminalOutputDidReset(surfaceID: surfaceID, streamToken: firstReplayChunk.streamToken) + await router.waitForCount(of: "mobile.terminal.replay", atLeast: replayCountAfterMount + 2) + let secondReplayChunk = try #require(await iterator.next()) + #expect(String(data: secondReplayChunk.data, encoding: .utf8) == "second-replay") + #expect(store.terminalReplayBarrierTokensBySurfaceID[surfaceID] == barrierToken) + + store.terminalOutputDidReset(surfaceID: surfaceID, streamToken: secondReplayChunk.streamToken) + await router.waitForCount(of: "mobile.terminal.replay", atLeast: replayCountAfterMount + 3) + let thirdReplayChunk = try #require(await iterator.next()) + #expect(String(data: thirdReplayChunk.data, encoding: .utf8) == "third-replay") + #expect(store.terminalReplayBarrierTokensBySurfaceID[surfaceID] == barrierToken) + + store.terminalOutputDidReset(surfaceID: surfaceID, streamToken: thirdReplayChunk.streamToken) + let replayRestartedAfterExhaustion = await router.waitForCount( + of: "mobile.terminal.replay", + atLeast: replayCountAfterMount + 4, + recordIssueOnTimeout: false + ) + #expect(!replayRestartedAfterExhaustion, "replay ack resets must not bypass the replay retry cap") + #expect(store.terminalReplayBarrierTokensBySurfaceID[surfaceID] == barrierToken) +} diff --git a/Packages/iOS/CmuxMobileShell/Tests/CmuxMobileShellTests/TerminalReplayBarrierFollowUpTests.swift b/Packages/iOS/CmuxMobileShell/Tests/CmuxMobileShellTests/TerminalReplayBarrierFollowUpTests.swift new file mode 100644 index 000000000000..f6925369f929 --- /dev/null +++ b/Packages/iOS/CmuxMobileShell/Tests/CmuxMobileShellTests/TerminalReplayBarrierFollowUpTests.swift @@ -0,0 +1,63 @@ +import Foundation +import Testing +@testable import CmuxMobileShell + +@MainActor +@Test func terminalReplayBarrierCapsFollowUpReplaysUnderContinuousOutput() async throws { + let router = LivenessHostRouter() + let box = TransportBox() + let clock = TestClock() + let store = try await makeConnectedStore(router: router, box: box, clock: clock) + let surfaceID = "live-terminal" + + await router.enqueueReplayTexts([ + "cold-replay", + "first-replay", + "follow-up-replay", + "unexpected-second-follow-up", + ]) + var iterator = store.terminalOutputStream(surfaceID: surfaceID).makeAsyncIterator() + await router.waitForCount(of: "mobile.terminal.replay", atLeast: 1) + let coldReplayChunk = try #require(await iterator.next()) + #expect(String(data: coldReplayChunk.data, encoding: .utf8) == "cold-replay") + store.terminalOutputDidProcess(surfaceID: surfaceID, streamToken: coldReplayChunk.streamToken) + let replayCountAfterMount = await router.count(of: "mobile.terminal.replay") + + store.deliverTerminalBytes(Data("stalled-first".utf8), surfaceID: surfaceID) + let stalledChunk = try #require(await iterator.next()) + store.terminalOutputDidReset(surfaceID: surfaceID, streamToken: stalledChunk.streamToken) + await router.waitForCount(of: "mobile.terminal.replay", atLeast: replayCountAfterMount + 1) + + let replayChunk = try #require(await iterator.next()) + #expect(String(data: replayChunk.data, encoding: .utf8) == "first-replay") + let firstDropAccepted = store.deliverTerminalBytes( + Data("live-during-first-barrier".utf8), + surfaceID: surfaceID + ) + #expect(firstDropAccepted == false) + + store.terminalOutputDidProcess(surfaceID: surfaceID, streamToken: replayChunk.streamToken) + await router.waitForCount(of: "mobile.terminal.replay", atLeast: replayCountAfterMount + 2) + + let followUpChunk = try #require(await iterator.next()) + #expect(String(data: followUpChunk.data, encoding: .utf8) == "follow-up-replay") + let followUpDropAccepted = store.deliverTerminalBytes( + Data("live-during-follow-up-barrier".utf8), + surfaceID: surfaceID + ) + #expect(followUpDropAccepted == false) + + store.terminalOutputDidProcess(surfaceID: surfaceID, streamToken: followUpChunk.streamToken) + let secondFollowUpRequested = await router.waitForCount( + of: "mobile.terminal.replay", + atLeast: replayCountAfterMount + 3, + timeoutNanoseconds: 200_000_000, + recordIssueOnTimeout: false + ) + #expect(!secondFollowUpRequested, "continuous output must not keep replaying indefinitely") + #expect(store.terminalReplayBarrierTokensBySurfaceID[surfaceID] == nil) + + store.deliverTerminalBytes(Data("after-bounded-replay".utf8), surfaceID: surfaceID) + let afterBoundedReplay = try #require(await iterator.next()) + #expect(String(data: afterBoundedReplay.data, encoding: .utf8) == "after-bounded-replay") +} diff --git a/Packages/iOS/CmuxMobileShellModel/Sources/CmuxMobileShellModel/MobileTerminalOutputSinking.swift b/Packages/iOS/CmuxMobileShellModel/Sources/CmuxMobileShellModel/MobileTerminalOutputSinking.swift index baeec82e0e1c..7b2bf990f9f4 100644 --- a/Packages/iOS/CmuxMobileShellModel/Sources/CmuxMobileShellModel/MobileTerminalOutputSinking.swift +++ b/Packages/iOS/CmuxMobileShellModel/Sources/CmuxMobileShellModel/MobileTerminalOutputSinking.swift @@ -49,4 +49,16 @@ public protocol MobileTerminalOutputSinking: Sendable { /// - Parameter surfaceID: The terminal surface identifier. /// - Parameter streamToken: The token carried by the yielded chunk. @MainActor func terminalOutputDidProcess(surfaceID: String, streamToken: UUID) + + /// Abandon the current yielded chunk after the local renderer was reset. + /// + /// The sink must drop stale pending output, invalidate the old stream token, + /// and request an authoritative replay for the same surface. + /// - Parameter surfaceID: The terminal surface identifier. + /// - Parameter streamToken: The token carried by the abandoned chunk. + @MainActor func terminalOutputDidReset(surfaceID: String, streamToken: UUID) + + /// Request an authoritative replay without an abandoned in-flight chunk. + /// - Parameter surfaceID: The terminal surface identifier. + @MainActor func terminalOutputNeedsReplay(surfaceID: String) } diff --git a/Packages/iOS/CmuxMobileShellUI/Sources/CmuxMobileShellUI/GhosttySurfaceRepresentable.swift b/Packages/iOS/CmuxMobileShellUI/Sources/CmuxMobileShellUI/GhosttySurfaceRepresentable.swift index 8163f4540a0d..16c04614b5bf 100644 --- a/Packages/iOS/CmuxMobileShellUI/Sources/CmuxMobileShellUI/GhosttySurfaceRepresentable.swift +++ b/Packages/iOS/CmuxMobileShellUI/Sources/CmuxMobileShellUI/GhosttySurfaceRepresentable.swift @@ -3,6 +3,7 @@ import CMUXMobileCore import CmuxMobileDiagnostics import CmuxMobileShell import CmuxMobileShellModel +import CmuxMobileSupport import CmuxMobileTerminal import SwiftUI import UIKit @@ -49,7 +50,10 @@ struct GhosttySurfaceRepresentable: UIViewRepresentable { fallback.numberOfLines = 0 fallback.textColor = .white fallback.backgroundColor = UIColor(red: 0x27/255.0, green: 0x28/255.0, blue: 0x22/255.0, alpha: 1) - fallback.text = "Ghostty runtime failed to initialise:\n\(error.localizedDescription)" + fallback.text = L10n.string( + "mobile.terminal.rendererFailed", + defaultValue: "Terminal renderer failed to start." + ) return fallback } let view = GhosttySurfaceView( @@ -145,20 +149,41 @@ struct GhosttySurfaceRepresentable: UIViewRepresentable { if chunk.data.isEmpty { surfaceView.useNaturalViewSize() } else { - await surfaceView.useNaturalViewSizeAndWait() + let applied = await surfaceView.useNaturalViewSizeAndWait() + guard applied else { + store.terminalOutputDidReset( + surfaceID: surfaceID, + streamToken: chunk.streamToken + ) + continue + } } case .remoteGrid(let columns, let rows): self.activeViewportPolicy = .remoteGrid(columns: columns, rows: rows) if chunk.data.isEmpty { surfaceView.applyViewSize(cols: columns, rows: rows) } else { - await surfaceView.applyViewSizeAndWait(cols: columns, rows: rows) + let applied = await surfaceView.applyViewSizeAndWait(cols: columns, rows: rows) + guard applied else { + store.terminalOutputDidReset( + surfaceID: surfaceID, + streamToken: chunk.streamToken + ) + continue + } } case nil: break } if !chunk.data.isEmpty { - await surfaceView.processOutputAndWait(chunk.data) + let applied = await surfaceView.processOutputAndWait(chunk.data) + guard applied else { + store.terminalOutputDidReset( + surfaceID: surfaceID, + streamToken: chunk.streamToken + ) + continue + } } store.terminalOutputDidProcess( surfaceID: surfaceID, @@ -402,6 +427,13 @@ struct GhosttySurfaceRepresentable: UIViewRepresentable { } } + func ghosttySurfaceViewDidResetRenderPipeline(_ surfaceView: GhosttySurfaceView) { + Task { @MainActor [weak self, weak store, surfaceID] in + guard let self, self.surfaceView === surfaceView else { return } + store?.terminalOutputNeedsReplay(surfaceID: surfaceID) + } + } + /// Walk up from `view` to the nearest owning `UIViewController`, then to /// its top-most presented controller, so a sheet presents above whatever /// is already on screen. diff --git a/Packages/iOS/CmuxMobileShellUI/Sources/CmuxMobileShellUI/WorkspaceDetailView.swift b/Packages/iOS/CmuxMobileShellUI/Sources/CmuxMobileShellUI/WorkspaceDetailView.swift index 2b11af258a36..159048f6ad1b 100644 --- a/Packages/iOS/CmuxMobileShellUI/Sources/CmuxMobileShellUI/WorkspaceDetailView.swift +++ b/Packages/iOS/CmuxMobileShellUI/Sources/CmuxMobileShellUI/WorkspaceDetailView.swift @@ -484,8 +484,7 @@ struct WorkspaceDetailView: View { #if DEBUG Button(action: copyDebugLogsFromMenu) { - // DEV-only debug tooling; not shipped, so not localized. - Label("Copy Debug Logs", systemImage: "doc.on.clipboard") + Label(L10n.string("mobile.debug.copyLogs", defaultValue: "Copy Debug Logs"), systemImage: "doc.on.clipboard") } .accessibilityIdentifier("MobileCopyDebugLogsMenuItem") #endif @@ -512,8 +511,8 @@ struct WorkspaceDetailView: View { private func copyDebugLogsFromMenu() { // Include "what the user sees" (the visible terminal text) above the // debug log so a pasted bug report shows the on-screen content too. - let terminalText = GhosttySurfaceView.visibleTerminalSnapshot() Task { @MainActor in + let terminalText = await GhosttySurfaceView.visibleTerminalSnapshot() let count = await MobileDebugLog.shared.copyToPasteboard(prepending: terminalText) UINotificationFeedbackGenerator().notificationOccurred(.success) NSLog("cmux.terminal copied %d debug log lines + visible terminal to pasteboard", count) @@ -639,10 +638,10 @@ struct WorkspaceDetailView: View { // Only the agent path reads the terminal/debug snapshots; reading them is // cheap and harmless on the email path, but skip the work when unused. // `visibleTerminalSnapshot()` reads off the output queue with a bounded - // wait (never a main-thread `ghostty_surface_read_text`, which blanks the + // async deadline (never a main-thread `ghostty_surface_read_text`, which blanks the // terminal). The debug-log snapshot is awaited from its actor. - let terminalText = routesToAgent ? GhosttySurfaceView.visibleTerminalSnapshot() : "" Task { @MainActor in + let terminalText = routesToAgent ? await GhosttySurfaceView.visibleTerminalSnapshot() : "" let debugLogText = routesToAgent ? await MobileDebugLog.shared.sink.snapshotWithCount().1 : "" let outcome = await store.submitFeedback( message: note, diff --git a/Packages/iOS/CmuxMobileTerminal/Sources/CmuxMobileTerminal/GhosttySurfaceRegistry.swift b/Packages/iOS/CmuxMobileTerminal/Sources/CmuxMobileTerminal/GhosttySurfaceRegistry.swift index 2adb524e703a..aee3702c839f 100644 --- a/Packages/iOS/CmuxMobileTerminal/Sources/CmuxMobileTerminal/GhosttySurfaceRegistry.swift +++ b/Packages/iOS/CmuxMobileTerminal/Sources/CmuxMobileTerminal/GhosttySurfaceRegistry.swift @@ -53,16 +53,16 @@ extension GhosttySurfaceView { /// on the serial `outputQueue` because `ghostty_surface_read_text` takes /// the surface lock that `process_output` holds during a render storm, so /// a main-thread read would stall the present and blank the terminal. - /// Unlike that synchronous DEV path there is no bounded semaphore wait - /// here — the caller awaits, so a busy queue just resumes the continuation - /// late while the sheet shows its loading state. + /// The await is bounded by the surface's display-link deadline, like the + /// diagnostic visible snapshot path, so a wedged output queue resumes nil + /// instead of leaving the copy sheet loading forever. /// /// The continuation body enqueues on `outputQueue` synchronously while /// still on the main actor, so the read is FIFO-ordered before any /// later-enqueued `disposeSurface` free of the same pointer — the same /// lifetime argument `visibleTerminalSnapshot()` relies on. /// - /// The read is bounded at the source: iOS surfaces are created with + /// The bytes read are bounded at the source: iOS surfaces are created with /// `scrollback-limit = 2000000` (see `GhosttyRuntime.applyiOSDefaults`), /// so the SCREEN range can never materialize more than ~2MB of text no /// matter how long the session ran. The sheet's 5000-line budget is then @@ -94,29 +94,8 @@ extension GhosttySurfaceView { && candidate.window != nil && !candidate.isHidden && candidate.alpha > 0.01 } - guard let surface = matchingView?.surface else { return nil } - let handle = CopyableTextSurfaceHandle(surface: surface) - return await withCheckedContinuation { continuation in - outputQueue.async { - // SCREEN = scrollback + all written rows. Fall back to the - // viewport-only read if the screen read fails outright. - let text = surfaceText(handle.surface, pointTag: GHOSTTY_POINT_SCREEN) - ?? surfaceText(handle.surface, pointTag: GHOSTTY_POINT_VIEWPORT) - continuation.resume(returning: text) - } - } + guard let matchingView, + let surface = matchingView.surface else { return nil } + return await matchingView.copyableTextForCurrentSurface(surface: surface) } } - -/// Carrier for the "View as Text" sheet's surface pointer across the hop to -/// `GhosttySurfaceView.outputQueue`. Same safety argument as -/// `VisibleSnapshotRequest` in `GhosttySurfaceView.swift`: the pointer is only -/// dereferenced on the queue that owns `process_output` and is FIFO-ordered -/// before any queued free — hence `@unchecked Sendable`. -/// -/// Deliberately `private` to this file: it holds the class's raw -/// `ghostty_surface_t`, which must not escape `GhosttySurfaceView`'s -/// queue/lifetime discipline into the wider module. -private struct CopyableTextSurfaceHandle: @unchecked Sendable { - let surface: ghostty_surface_t -} diff --git a/Packages/iOS/CmuxMobileTerminal/Sources/CmuxMobileTerminal/GhosttySurfaceView.swift b/Packages/iOS/CmuxMobileTerminal/Sources/CmuxMobileTerminal/GhosttySurfaceView.swift index 2d285bf5b12f..67567dc949f2 100644 --- a/Packages/iOS/CmuxMobileTerminal/Sources/CmuxMobileTerminal/GhosttySurfaceView.swift +++ b/Packages/iOS/CmuxMobileTerminal/Sources/CmuxMobileTerminal/GhosttySurfaceView.swift @@ -5,6 +5,7 @@ import CmuxMobileSupport import CmuxMobileTerminalKit import GhosttyKit import OSLog +import Synchronization import UIKit private let log = Logger(subsystem: "ai.manaflow.cmux.ios", category: "ghostty.surface") @@ -75,6 +76,9 @@ public protocol GhosttySurfaceViewDelegate: AnyObject { /// reveal-after-hide and the present-while-suppressed paths so the draft and its /// focus return together. Optional. func ghosttySurfaceViewDidRequestComposerFocus(_ surfaceView: GhosttySurfaceView) + /// The local Ghostty render pipeline was rebuilt after a stuck render/output + /// operation. The host should replay authoritative terminal state. + func ghosttySurfaceViewDidResetRenderPipeline(_ surfaceView: GhosttySurfaceView) } public extension GhosttySurfaceViewDelegate { @@ -86,6 +90,8 @@ public extension GhosttySurfaceViewDelegate { func ghosttySurfaceViewDidRequestComposerToggle(_ surfaceView: GhosttySurfaceView) {} /// Default no-op so hosts without a composer can ignore the focus request. func ghosttySurfaceViewDidRequestComposerFocus(_ surfaceView: GhosttySurfaceView) {} + /// Default no-op so hosts without terminal-output replay can ignore renderer resets. + func ghosttySurfaceViewDidResetRenderPipeline(_ surfaceView: GhosttySurfaceView) {} } @MainActor @@ -570,8 +576,8 @@ public final class GhosttySurfaceView: UIView, TerminalSurfaceHosting { private var zoomOverlayShown = false /// Media time of the last zoom interaction (pinch step, zoom button, or HUD /// tap). The display link fades the HUD once this is older than - /// `zoomOverlayVisibleDuration`. Time-based off the per-frame callback, not a - /// timer/`Task.sleep`, so it honors the no-sleep rule and tracks real + /// `zoomOverlayVisibleDuration`. Time-based off the per-frame callback, not + /// a sleeping timer task, so it honors the no-sleep rule and tracks real /// elapsed time regardless of frame rate. private var zoomOverlayLastInteraction: CFTimeInterval = 0 private static let zoomOverlayVisibleDuration: CFTimeInterval = 2.5 @@ -579,7 +585,7 @@ public final class GhosttySurfaceView: UIView, TerminalSurfaceHosting { /// reset/save/restore actions. Owned by the surface (constructed at init) /// rather than reached through a singleton, so it is injectable in tests. private let zoomPreference = MobileTerminalZoomPreference() - private let bridge = GhosttySurfaceBridge() + private var bridge = GhosttySurfaceBridge() private let prefersSnapshotFallbackRendering = false var onFocusInputRequestedForTesting: (() -> Void)? private var surfaceTitle: String? @@ -607,6 +613,7 @@ public final class GhosttySurfaceView: UIView, TerminalSurfaceHosting { /// symptom. Coalescing caps the backlog: while a render is in flight, mark /// `needsAnotherRender` and re-enqueue exactly one when it completes. private var renderInFlight: Bool = false + private var renderInFlightSince: CFTimeInterval? private var needsAnotherRender: Bool = false /// True while the app is inactive/backgrounded. On iOS `render_now` /// produces a frame synchronously on `outputQueue` and acquires a @@ -647,10 +654,25 @@ public final class GhosttySurfaceView: UIView, TerminalSurfaceHosting { /// Internal (not private) so the copyable-text extension in /// `GhosttySurfaceCopyableText.swift` can enqueue its surface read with /// the same FIFO-before-dispose ordering discipline. - static let outputQueue = DispatchQueue( - label: "dev.cmux.GhosttySurfaceView.output", - qos: .userInitiated - ) + var outputQueue = GhosttySurfaceWorkQueue(generation: 0) + private var outputQueueGeneration: UInt64 = 0 + private var pendingSurfaceFreeCount = 0 + private var renderPipelineRecoveryPaused = false + private var lastRecoveryPausedDropLogTime: CFTimeInterval = 0 + private static let renderPipelineStallDeadline: CFTimeInterval = 2.0 + private static let outputApplyTimeout: CFTimeInterval = 2.0 + private static let visibleSnapshotTimeout: CFTimeInterval = 0.6 + private static let copyableTextTimeout: CFTimeInterval = 2.0 + private static let maxPendingSurfaceFrees = 1 + private var nextSurfaceOperationID: UInt64 = 0 + private var pendingOutputApply: PendingSurfaceOperation? + private var pendingGeometryApply: PendingSurfaceOperation? + private var pendingVisibleSnapshot: PendingVisibleSnapshot? + private var pendingCopyableTextRead: PendingCopyableTextRead? + private var hasPendingSurfaceOperationDeadline: Bool { + pendingOutputApply != nil || pendingGeometryApply != nil || pendingVisibleSnapshot != nil + || pendingCopyableTextRead != nil + } private static let scrollMechanicsContentHeight: CGFloat = 1_000_000 private var scrollMechanicsIsRecentering = false private var lastScrollMechanicsOffsetY: CGFloat? @@ -788,6 +810,7 @@ public final class GhosttySurfaceView: UIView, TerminalSurfaceHosting { }() private(set) var surface: ghostty_surface_t? + private var surfaceGeneration: UInt64 = 0 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 @@ -1096,6 +1119,8 @@ public final class GhosttySurfaceView: UIView, TerminalSurfaceHosting { /// `willResignActive` and `didEnterBackground`. private func suspendRendering() { renderingSuspended = true + skipPendingVisibleSnapshot() + skipPendingCopyableTextRead() stopDisplayLink() guard let surface else { return } ghostty_surface_set_occlusion(surface, false) // false = occluded; drawFrame skips @@ -1112,6 +1137,7 @@ public final class GhosttySurfaceView: UIView, TerminalSurfaceHosting { private func resumeRendering() { renderingSuspended = false renderInFlight = false + renderInFlightSince = nil needsAnotherRender = false guard let surface, window != nil else { return } ghostty_surface_set_occlusion(surface, true) // true = visible @@ -1558,7 +1584,7 @@ public final class GhosttySurfaceView: UIView, TerminalSurfaceHosting { /// toolbar stays visible regardless, so closing the composer just collapses its /// band. Gating the re-focus on `keyboardVisible` makes both directions /// correct. - /// No sleep / `asyncAfter`: the `become` is issued synchronously here. + /// No deferred timer task: the `become` is issued synchronously here. public func setComposerActive(_ active: Bool) { guard composerActive != active else { return } composerActive = active @@ -2058,7 +2084,7 @@ public final class GhosttySurfaceView: UIView, TerminalSurfaceHosting { // An absolute `set_font_size:` keeps libghostty in lockstep // with `liveFontSize`, which we keep inside [minimumSize, maximumSize]. let action = "set_font_size:\(target)" - Self.outputQueue.async { + outputQueue.async { action.withCString { pointer in _ = ghostty_surface_binding_action(surface, pointer, UInt(action.utf8.count)) } @@ -2258,20 +2284,194 @@ public final class GhosttySurfaceView: UIView, TerminalSurfaceHosting { /// per-surface backpressure without blocking the main actor while Ghostty /// consumes the chunk. /// - Parameter data: VT or PTY bytes to feed into the surface. - public func processOutputAndWait(_ data: Data) async { - await withCheckedContinuation { continuation in - processOutput(data) { - continuation.resume() + /// - Returns: `true` when the bytes reached the current surface generation, + /// or `false` when the caller should reset its delivery queue and replay. + @discardableResult + public func processOutputAndWait(_ data: Data) async -> Bool { + return await withCheckedContinuation { continuation in + let operationID = registerPendingOutputApply( + byteCount: data.count, + continuation: continuation + ) + processOutput(data) { [weak self] applied in + self?.completePendingOutputApply(id: operationID, returning: applied) } } } + private func makeSurfaceOperationID() -> UInt64 { + nextSurfaceOperationID &+= 1 + return nextSurfaceOperationID + } + + private func ensureSurfaceOperationDeadlinePump() { + guard window != nil, displayLink == nil, !renderingSuspended, !renderPipelineRecoveryPaused else { return } + startDisplayLink() + } + + private func registerPendingOutputApply( + byteCount: Int, + continuation: CheckedContinuation + ) -> UInt64 { + let operationID = makeSurfaceOperationID() + if let existing = pendingOutputApply { + pendingOutputApply = nil + let elapsedMs = Int((CACurrentMediaTime() - existing.startedAt) * 1000) + MobileDebugLog.anchormux("output.apply.OVERLAP elapsedMs=\(elapsedMs)") + existing.continuation.resume(returning: false) + } + pendingOutputApply = PendingSurfaceOperation( + id: operationID, + startedAt: CACurrentMediaTime(), + byteCount: byteCount, + continuation: continuation + ) + ensureSurfaceOperationDeadlinePump() + return operationID + } + + @discardableResult + private func completePendingOutputApply(id: UInt64, returning result: Bool) -> Bool { + guard let pending = pendingOutputApply, pending.id == id else { return false } + pendingOutputApply = nil + pending.continuation.resume(returning: result) + return true + } + + private func registerPendingGeometryApply( + continuation: CheckedContinuation + ) -> UInt64 { + let operationID = makeSurfaceOperationID() + if let existing = pendingGeometryApply { + pendingGeometryApply = nil + let elapsedMs = Int((CACurrentMediaTime() - existing.startedAt) * 1000) + MobileDebugLog.anchormux("geometry.apply.OVERLAP elapsedMs=\(elapsedMs)") + existing.continuation.resume(returning: false) + } + pendingGeometryApply = PendingSurfaceOperation( + id: operationID, + startedAt: CACurrentMediaTime(), + byteCount: nil, + continuation: continuation + ) + ensureSurfaceOperationDeadlinePump() + return operationID + } + + @discardableResult + private func completePendingGeometryApply(id: UInt64, returning result: Bool) -> Bool { + guard let pending = pendingGeometryApply, pending.id == id else { return false } + pendingGeometryApply = nil + pending.continuation.resume(returning: result) + return true + } + + @discardableResult + private func checkSurfaceOperationDeadlines(now: CFTimeInterval) -> Bool { + if let pending = pendingOutputApply, + now - pending.startedAt >= Self.outputApplyTimeout { + pendingOutputApply = nil + let elapsedMs = Int((now - pending.startedAt) * 1000) + MobileDebugLog.anchormux( + "output.apply.TIMEOUT bytes=\(pending.byteCount ?? 0) elapsedMs=\(elapsedMs)" + ) + let recovered = recoverRenderPipeline( + reason: "output_timeout", + stalledMs: elapsedMs, + replay: .callerWillRequestReplay + ) + pending.continuation.resume(returning: false) + return recovered + } + + if let pending = pendingGeometryApply, + now - pending.startedAt >= Self.outputApplyTimeout { + pendingGeometryApply = nil + let elapsedMs = Int((now - pending.startedAt) * 1000) + MobileDebugLog.anchormux("geometry.apply.TIMEOUT elapsedMs=\(elapsedMs)") + let recovered = recoverRenderPipeline( + reason: "geometry_timeout", + stalledMs: elapsedMs, + replay: .callerWillRequestReplay + ) + pending.continuation.resume(returning: false) + return recovered + } + + if let pending = pendingVisibleSnapshot, + now - pending.startedAt >= Self.visibleSnapshotTimeout { + pendingVisibleSnapshot = nil + pending.continuation.resume(returning: nil) + } + + if let pending = pendingCopyableTextRead, + now - pending.startedAt >= Self.copyableTextTimeout { + pendingCopyableTextRead = nil + pending.cancel() + pending.continuation.resume(returning: nil) + } + return false + } + + @discardableResult + private func completePendingSurfaceOperations(returning result: Bool) -> Bool { + var completed = false + if let pending = pendingOutputApply { + pendingOutputApply = nil + pending.continuation.resume(returning: result) + completed = true + } + if let pending = pendingGeometryApply { + pendingGeometryApply = nil + pending.continuation.resume(returning: result) + completed = true + } + skipPendingVisibleSnapshot() + skipPendingCopyableTextRead() + return completed + } + + private func skipPendingVisibleSnapshot() { + guard let pending = pendingVisibleSnapshot else { return } + pendingVisibleSnapshot = nil + pending.continuation.resume(returning: nil) + } + + private func skipPendingCopyableTextRead() { + guard let pending = pendingCopyableTextRead else { return } + pendingCopyableTextRead = nil + pending.cancel() + pending.continuation.resume(returning: nil) + } + + @discardableResult + private func completePendingVisibleSnapshot(id: UInt64, returning section: String?) -> Bool { + guard let pending = pendingVisibleSnapshot, pending.id == id else { return false } + pendingVisibleSnapshot = nil + pending.continuation.resume(returning: section) + return true + } + + @discardableResult + private func completePendingCopyableTextRead(id: UInt64, returning text: String?) -> Bool { + guard let pending = pendingCopyableTextRead, pending.id == id else { return false } + pendingCopyableTextRead = nil + pending.cancel() + pending.continuation.resume(returning: text) + return true + } + private func processOutput( _ data: Data, - completion: (@MainActor @Sendable () -> Void)? + completion: (@MainActor @Sendable (Bool) -> Void)? ) { + guard !renderPipelineRecoveryPaused else { + logRecoveryPausedDrop(kind: "output", byteCount: data.count) + completion?(false) + return + } guard let surface, !isDismantled else { - completion?() + completion?(true) return } #if DEBUG @@ -2289,6 +2489,7 @@ public final class GhosttySurfaceView: UIView, TerminalSurfaceHosting { } #endif let forwarded = Self.forwardDaemonOutputBytes(data) + let generation = surfaceGeneration // Track the host's cursor-visible mode (DECTCEM) straight from the VT // bytes the surface is about to apply, so the cursor overlay can match a // TUI that hides the cursor. nil = this delta carried no DECTCEM, so the @@ -2301,7 +2502,8 @@ public final class GhosttySurfaceView: UIView, TerminalSurfaceHosting { // scene-update watchdog (0x8BADF00D) kills the app. It must run off // the main thread. Feed it on a serial background queue (order // preserved) and hop back to main only for the Swift-side UI state. - Self.outputQueue.async { [weak self] in + let workQueue = outputQueue + workQueue.async { [weak self] in forwarded.withUnsafeBytes { buffer in guard let baseAddress = buffer.baseAddress else { return } let pointer = baseAddress.assumingMemoryBound(to: CChar.self) @@ -2319,14 +2521,18 @@ public final class GhosttySurfaceView: UIView, TerminalSurfaceHosting { // main. Off-main reads can never trip the main-thread watchdog. var accessibilityText: String? let a11yNow = CACurrentMediaTime() - if a11yNow - Self.lastAccessibilityTextTime > 0.5 { - Self.lastAccessibilityTextTime = a11yNow + if a11yNow - workQueue.lastAccessibilityTextTime > 0.5 { + workQueue.lastAccessibilityTextTime = a11yNow accessibilityText = Self.accessibilitySurfaceText(surface) } #endif DispatchQueue.main.async { guard let self, !self.isDismantled else { - completion?() + completion?(true) + return + } + guard self.surfaceGeneration == generation else { + completion?(false) return } self.needsDraw = true @@ -2355,7 +2561,7 @@ public final class GhosttySurfaceView: UIView, TerminalSurfaceHosting { } self.onOutputProcessedForTesting?() #endif - completion?() + completion?(true) } } } @@ -2373,7 +2579,7 @@ public final class GhosttySurfaceView: UIView, TerminalSurfaceHosting { // `process_output` also preserves ordering. The return was already // discarded. let action = "scroll_to_bottom" - Self.outputQueue.async { + outputQueue.async { action.withCString { pointer in _ = ghostty_surface_binding_action(surface, pointer, UInt(action.utf8.count)) } @@ -2438,6 +2644,10 @@ public final class GhosttySurfaceView: UIView, TerminalSurfaceHosting { /// tree. Does not set ``isDismantled`` so a transient detach can re-attach /// and resume; only ``prepareForDismantle()`` marks the surface dead. private func prepareForReuseAfterDetach() { + completePendingSurfaceOperations(returning: false) + renderInFlight = false + renderInFlightSince = nil + needsAnotherRender = false resignInput() stopKeyboardHeightAnimation() stopDisplayLink() @@ -2517,11 +2727,6 @@ public final class GhosttySurfaceView: UIView, TerminalSurfaceHosting { } } - /// Throttle stamp for the off-main accessibility-label read in - /// `processOutput`. Accessed only on the serial `outputQueue`, so the - /// unchecked mutation is safe. - nonisolated(unsafe) fileprivate static var lastAccessibilityTextTime: CFTimeInterval = 0 - /// Off-main equivalent of ``accessibilityRenderedTextForTesting()`` that /// reads via the raw surface handle so it can run on the serial output queue /// (alongside `process_output`) instead of the main thread. See the call @@ -2555,6 +2760,55 @@ public final class GhosttySurfaceView: UIView, TerminalSurfaceHosting { return String(decoding: Data(bytes: ptr, count: Int(text.text_len)), as: UTF8.self) } + func copyableTextForCurrentSurface(surface expectedSurface: ghostty_surface_t) async -> String? { + let generation = surfaceGeneration + guard surface == expectedSurface, + !isDismantled, + !renderPipelineRecoveryPaused, + !renderingSuspended else { + return nil + } + return await withCheckedContinuation { continuation in + let operationID = makeSurfaceOperationID() + if let existing = pendingCopyableTextRead { + pendingCopyableTextRead = nil + existing.cancel() + existing.continuation.resume(returning: nil) + } + let cancellation = SurfaceOperationCancellationToken() + pendingCopyableTextRead = PendingCopyableTextRead( + id: operationID, + startedAt: CACurrentMediaTime(), + cancellation: cancellation, + continuation: continuation + ) + ensureSurfaceOperationDeadlinePump() + let read = CopyableTextRead( + surface: expectedSurface, + generation: generation, + cancellation: cancellation + ) + outputQueue.async { [weak self] in + guard !read.cancellation.isCancelled else { return } + // SCREEN = scrollback + all written rows. Fall back to the + // viewport-only read if the screen read fails outright. + let screenText = Self.surfaceText(read.surface, pointTag: GHOSTTY_POINT_SCREEN) + guard !read.cancellation.isCancelled else { return } + let text = screenText ?? Self.surfaceText(read.surface, pointTag: GHOSTTY_POINT_VIEWPORT) + guard !read.cancellation.isCancelled else { return } + Task { @MainActor [weak self] in + guard let self else { return } + guard self.surface == read.surface, + self.surfaceGeneration == read.generation else { + self.completePendingCopyableTextRead(id: operationID, returning: nil) + return + } + self.completePendingCopyableTextRead(id: operationID, returning: text) + } + } + } + } + func renderedHTMLForTesting(pointTag: ghostty_point_tag_e = GHOSTTY_POINT_VIEWPORT) -> String? { _ = pointTag // ghostty_surface_read_text_html not available in this build @@ -2571,7 +2825,9 @@ public final class GhosttySurfaceView: UIView, TerminalSurfaceHosting { guard let surface else { return } GhosttySurfaceView.unregister(surface: surface) self.surface = nil - bridge.detach() + let currentBridge = bridge + let currentQueue = outputQueue + currentBridge.detach() // Free on the SAME serial `outputQueue` that runs `process_output`, // `render_now`, and `binding_action` (all of which capture this C // surface pointer), not a separate queue. FIFO ordering guarantees the @@ -2582,11 +2838,127 @@ public final class GhosttySurfaceView: UIView, TerminalSurfaceHosting { // work from being enqueued once `surface` is nil, so only the bounded // backlog drains before the free. (Retain the bridge across the hop; it // owns the userdata libghostty still references until the free.) + enqueueSurfaceFree(surface, bridge: currentBridge, on: currentQueue) + } + + private func enqueueSurfaceFree( + _ surface: ghostty_surface_t, + bridge: GhosttySurfaceBridge, + on queue: GhosttySurfaceWorkQueue, + completion: (@MainActor @Sendable () -> Void)? = nil + ) { let retainedBridge = Unmanaged.passRetained(bridge) - Self.outputQueue.async { + queue.async { ghostty_surface_free(surface) retainedBridge.release() + if let completion { + Task { @MainActor in completion() } + } + } + } + + private func logRecoveryPausedDrop(kind: String, byteCount: Int? = nil) { + let now = CACurrentMediaTime() + guard now - lastRecoveryPausedDropLogTime >= 1 else { return } + lastRecoveryPausedDropLogTime = now + MobileDebugLog.anchormux( + "render.recover.paused_drop kind=\(kind) bytes=\(byteCount ?? 0) pendingFrees=\(pendingSurfaceFreeCount)" + ) + } + + @discardableResult + private func pauseRenderPipelineRecovery( + reason: String, + stalledMs: Int + ) -> Bool { + MobileDebugLog.anchormux( + "render.recover.paused reason=\(reason) stalledMs=\(stalledMs) pendingFrees=\(pendingSurfaceFreeCount)" + ) + renderPipelineRecoveryPaused = true + stopDisplayLink() + _ = completePendingSurfaceOperations(returning: false) + renderInFlight = false + renderInFlightSince = nil + needsAnotherRender = false + needsDraw = false + return true + } + + private func resumePausedRenderPipelineRecoveryIfPossible() { + guard renderPipelineRecoveryPaused, + !isDismantled, + surface != nil, + pendingSurfaceFreeCount < Self.maxPendingSurfaceFrees else { return } + MobileDebugLog.anchormux( + "render.recover.resuming pendingFrees=\(pendingSurfaceFreeCount)" + ) + renderPipelineRecoveryPaused = false + _ = recoverRenderPipeline( + reason: "free_drained", + stalledMs: 0, + replay: .delegateWhenNoCaller + ) + } + + @discardableResult + private func recoverRenderPipeline( + reason: String, + stalledMs: Int, + replay: RenderPipelineRecoveryReplay + ) -> Bool { + guard !isDismantled, + surface != nil else { + return false + } + guard !renderPipelineRecoveryPaused else { + return pauseRenderPipelineRecovery(reason: reason, stalledMs: stalledMs) + } + guard pendingSurfaceFreeCount < Self.maxPendingSurfaceFrees else { + return pauseRenderPipelineRecovery(reason: reason, stalledMs: stalledMs) + } + MobileDebugLog.anchormux( + "render.recover reason=\(reason) stalledMs=\(stalledMs) generation=\(surfaceGeneration) pendingFrees=\(pendingSurfaceFreeCount)" + ) + let completedFailedOperation = completePendingSurfaceOperations(returning: false) + + stopDisplayLink() + let oldSurface = surface + let oldBridge = bridge + let oldQueue = outputQueue + oldBridge.detach() + if let oldSurface { + GhosttySurfaceView.unregister(surface: oldSurface) + pendingSurfaceFreeCount += 1 + enqueueSurfaceFree(oldSurface, bridge: oldBridge, on: oldQueue) { [weak self] in + guard let self else { return } + self.pendingSurfaceFreeCount = max(0, self.pendingSurfaceFreeCount - 1) + MobileDebugLog.anchormux( + "render.recover.free_drained pendingFrees=\(self.pendingSurfaceFreeCount)" + ) + self.resumePausedRenderPipelineRecoveryIfPossible() + } + } + + surface = nil + renderInFlight = false + renderInFlightSince = nil + needsAnotherRender = false + needsDraw = true + cellPixelSize = .zero + lastRenderRect = .zero + lastAppliedContentScale = 0 + + surfaceGeneration &+= 1 + outputQueueGeneration &+= 1 + outputQueue = GhosttySurfaceWorkQueue(generation: outputQueueGeneration) + bridge = GhosttySurfaceBridge() + bridge.attach(to: self) + + initializeSurface() + if replay == .delegateWhenNoCaller && !completedFailedOperation { + delegate?.ghosttySurfaceViewDidResetRenderPipeline(self) } + return true } private var preferredScreenScale: CGFloat { @@ -2660,7 +3032,16 @@ public final class GhosttySurfaceView: UIView, TerminalSurfaceHosting { } @objc func handleDisplayLinkFire() { - guard let surface else { return } + let now = CACurrentMediaTime() + if checkSurfaceOperationDeadlines(now: now) { + return + } + guard surface != nil else { + if !hasPendingSurfaceOperationDeadline { + stopDisplayLink() + } + return + } #if DEBUG // Main-thread liveness heartbeat + presented-surface state. Time-gated, // no behavior change. The `contents`/size fields let an IDLE blank be @@ -2670,7 +3051,7 @@ public final class GhosttySurfaceView: UIView, TerminalSurfaceHosting { // content itself is empty (sync/producer). `sinceOutput` ties a blank // to a render-grid stream gap or rules it out. CALayer reads only — no // libghostty call, so no futex/main-thread-wedge risk. - let nowHeartbeat = CACurrentMediaTime() + let nowHeartbeat = now if nowHeartbeat - lastHeartbeatTime >= 2.0 { lastHeartbeatTime = nowHeartbeat let renderLayer = (layer.sublayers ?? []).first(where: { isGhosttyRendererLayer($0) }) @@ -2687,6 +3068,17 @@ public final class GhosttySurfaceView: UIView, TerminalSurfaceHosting { ) } #endif + if let renderInFlightSince { + let stalledMs = Int((now - renderInFlightSince) * 1000) + if stalledMs >= Int(Self.renderPipelineStallDeadline * 1000), + recoverRenderPipeline( + reason: "render_in_flight", + stalledMs: stalledMs, + replay: .delegateWhenNoCaller + ) { + return + } + } // Apply at most one coalesced zoom per frame. This only changes the // font; the geometry resync is deferred until zoom settles. let appliedZoom = applyPendingFontSizeIfNeeded() @@ -2715,7 +3107,6 @@ public final class GhosttySurfaceView: UIView, TerminalSurfaceHosting { pendingGeometryReassert = false syncSurfaceGeometry(shouldReassertNaturalSize: reassert) } - let now = CACurrentMediaTime() let blinkChanged = cursorBlinkState.advance(now: now) // Draw on content/cursor changes, and for a short bounded burst after // any geometry change. iOS has no renderer-side vsync, so a frame is @@ -2767,7 +3158,7 @@ public final class GhosttySurfaceView: UIView, TerminalSurfaceHosting { // Fade the zoom HUD once interaction has been quiet. Uses real elapsed // time off the continuous display link (no timer / sleep). if zoomOverlayShown, - CACurrentMediaTime() - zoomOverlayLastInteraction > Self.zoomOverlayVisibleDuration { + now - zoomOverlayLastInteraction > Self.zoomOverlayVisibleDuration { fadeOutZoomOverlay() } } @@ -2795,7 +3186,10 @@ public final class GhosttySurfaceView: UIView, TerminalSurfaceHosting { // now bounded in libghostty (so a foreground stall self-heals as a // skipped frame the display link re-drives), but we still gate on // suspension; `resumeRendering` clears it on the next active transition. - guard !renderingSuspended, let surface, !isDismantled else { return } + guard !renderPipelineRecoveryPaused, + !renderingSuspended, + let surface, + !isDismantled else { return } // Coalesce: never let more than one render_now sit on the serial queue. // (Called on main from the display link.) if renderInFlight { @@ -2803,8 +3197,10 @@ public final class GhosttySurfaceView: UIView, TerminalSurfaceHosting { return } renderInFlight = true + renderInFlightSince = CACurrentMediaTime() + let generation = surfaceGeneration let enqueuedAt = CACurrentMediaTime() - Self.outputQueue.async { [weak self] in + outputQueue.async { [weak self] in // Queue LAG = how long this render waited behind other ops. If this // climbs into hundreds of ms the queue is backlogged (the freeze). let lagMs = (CACurrentMediaTime() - enqueuedAt) * 1000 @@ -2812,7 +3208,9 @@ public final class GhosttySurfaceView: UIView, TerminalSurfaceHosting { ghostty_surface_render_now(surface) DispatchQueue.main.async { guard let self else { return } + guard self.surfaceGeneration == generation else { return } self.renderInFlight = false + self.renderInFlightSince = nil guard !self.isDismantled else { self.needsAnotherRender = false return @@ -2978,11 +3376,18 @@ public final class GhosttySurfaceView: UIView, TerminalSurfaceHosting { applyViewSize(cols: cols, rows: rows, confirmedViewportEcho: false) } - public func applyViewSizeAndWait(cols: Int, rows: Int) async { + /// Apply the daemon's authoritative rendering grid and wait until libghostty + /// accepts the geometry for the current surface generation. + /// - Parameter cols: The authoritative terminal column count. + /// - Parameter rows: The authoritative terminal row count. + /// - Returns: `false` when the surface reset before the geometry applied. + @discardableResult + public func applyViewSizeAndWait(cols: Int, rows: Int) async -> Bool { let changed = updateEffectiveGrid(cols: cols, rows: rows, confirmedViewportEcho: false) if changed || needsGeometrySync { - await syncSurfaceGeometryAndWait(shouldReassertNaturalSize: false) + return await syncSurfaceGeometryAndWait(shouldReassertNaturalSize: false) } + return true } public func applyConfirmedViewSize(cols: Int, rows: Int) { @@ -3020,11 +3425,16 @@ public final class GhosttySurfaceView: UIView, TerminalSurfaceHosting { setNeedsGeometrySync(reassertNaturalSize: false) } - public func useNaturalViewSizeAndWait() async { + /// Return to the phone's natural viewport capacity and wait until libghostty + /// accepts the geometry for the current surface generation. + /// - Returns: `false` when the surface reset before the geometry applied. + @discardableResult + public func useNaturalViewSizeAndWait() async -> Bool { let changed = clearEffectiveGrid() if changed || needsGeometrySync { - await syncSurfaceGeometryAndWait(shouldReassertNaturalSize: false) + return await syncSurfaceGeometryAndWait(shouldReassertNaturalSize: false) } + return true } private func clearEffectiveGrid() -> Bool { @@ -3081,22 +3491,28 @@ public final class GhosttySurfaceView: UIView, TerminalSurfaceHosting { let pinnedSize: CGSize? } - private func syncSurfaceGeometryAndWait(shouldReassertNaturalSize: Bool = true) async { + private func syncSurfaceGeometryAndWait(shouldReassertNaturalSize: Bool = true) async -> Bool { needsGeometrySync = false pendingGeometryReassert = false - await withCheckedContinuation { continuation in - syncSurfaceGeometry(shouldReassertNaturalSize: shouldReassertNaturalSize) { - continuation.resume() + return await withCheckedContinuation { continuation in + let operationID = registerPendingGeometryApply(continuation: continuation) + syncSurfaceGeometry(shouldReassertNaturalSize: shouldReassertNaturalSize) { [weak self] applied in + self?.completePendingGeometryApply(id: operationID, returning: applied) } } } private func syncSurfaceGeometry( shouldReassertNaturalSize: Bool = true, - completion: (@MainActor @Sendable () -> Void)? = nil + completion: (@MainActor @Sendable (Bool) -> Void)? = nil ) { + guard !renderPipelineRecoveryPaused else { + logRecoveryPausedDrop(kind: "geometry") + completion?(false) + return + } guard let surface else { - completion?() + completion?(true) return } @@ -3144,8 +3560,9 @@ public final class GhosttySurfaceView: UIView, TerminalSurfaceHosting { let eff = effectiveGrid let pushContentScale = abs(lastAppliedContentScale - scale) > 0.001 if pushContentScale { lastAppliedContentScale = scale } + let generation = surfaceGeneration - Self.outputQueue.async { [weak self] in + outputQueue.async { [weak self] in if pushContentScale { ghostty_surface_set_content_scale(surface, scale, scale) } @@ -3185,14 +3602,22 @@ public final class GhosttySurfaceView: UIView, TerminalSurfaceHosting { ) let result = GeometryResult(cellPixelSize: cell, naturalSize: natural, pinnedSize: pinnedSize) DispatchQueue.main.async { - self?.applyGeometryResult( + guard let self else { + completion?(true) + return + } + guard self.surfaceGeneration == generation else { + completion?(false) + return + } + self.applyGeometryResult( result, scale: scale, containerW: containerW, containerH: containerH, shouldReassertNaturalSize: shouldReassertNaturalSize ) - completion?() + completion?(true) } } } @@ -3378,7 +3803,7 @@ public final class GhosttySurfaceView: UIView, TerminalSurfaceHosting { ios: ghostty_platform_ios_s(uiview: Unmanaged.passUnretained(self).toOpaque()) ) surfaceConfig.scale_factor = preferredScreenScale - surfaceConfig.font_size = fontSize + surfaceConfig.font_size = liveFontSize surfaceConfig.context = GHOSTTY_SURFACE_CONTEXT_WINDOW surfaceConfig.io_mode = GHOSTTY_SURFACE_IO_MANUAL surfaceConfig.io_write_cb = { userdata, buf, len in @@ -3600,7 +4025,7 @@ public final class GhosttySurfaceView: UIView, TerminalSurfaceHosting { /// terminal surface, for the DEV "Copy Debug Logs" action so a bug report /// pairs the on-screen content with the debug log. Reads the VIEWPORT /// (visible grid only, not scrollback) via libghostty. - public static func visibleTerminalSnapshot() -> String { + public static let visibleTerminalSnapshot: @MainActor () async -> String = { registeredSurfaceViews = registeredSurfaceViews.filter { $0.value.value != nil } // Collect the main-actor state + surface pointers first, then read the // viewport text on the serial output queue. `ghostty_surface_read_text` @@ -3608,46 +4033,80 @@ public final class GhosttySurfaceView: UIView, TerminalSurfaceHosting { // reading it on the MAIN thread here contends that lock during a render // storm and stalls the present — tapping Copy Debug Logs would itself // blank the terminal. The output queue is never concurrent with - // `process_output`, so the read can't wedge. No `main.sync` runs on that - // queue, so this `.sync` cannot deadlock. + // `process_output`, so the read can't wedge. The await is bounded by + // the surface's display-link deadline so this diagnostic path does not + // add a sleeping timer task or block the main actor. var pending: [VisibleSnapshotRequest] = [] for view in registeredSurfaceViews.values.compactMap(\.value) { guard view.window != nil, !view.isHidden, view.alpha > 0.01, let surface = view.surface else { continue } let grid = view.effectiveGrid.map { "\($0.cols)x\($0.rows)" } ?? "?" - pending.append(VisibleSnapshotRequest(grid: grid, font: Int(view.liveFontSize), surface: surface)) + pending.append(VisibleSnapshotRequest( + view: view, + grid: grid, + font: Int(view.liveFontSize), + surface: surface, + generation: view.surfaceGeneration + )) } if pending.isEmpty { return "===== visible terminal: (no on-screen surface) =====" } - // Read on the output queue, but bound the wait. If a render wedge has the - // queue stuck mid-`process_output`, a plain `.sync` here would freeze the - // whole app exactly when the user taps Copy Debug Logs to capture that - // bug. Time out and ship the logs without the snapshot instead. - let holder = VisibleSnapshotHolder() - // This synchronous DEV-only "Copy Debug Logs" path reads the viewport off - // the serial output queue and must give up after a deadline if a render - // wedge holds it; an actor/await cannot express the bounded synchronous - // wait the synchronous caller needs. - // carve-out justification: one-shot cross-queue completion signal with a - // bounded wait, not a lock guarding shared state. - let done = DispatchSemaphore(value: 0) - outputQueue.async { - var built: [String] = [] - for item in pending { - let text = surfaceText(item.surface, pointTag: GHOSTTY_POINT_VIEWPORT) ?? "(unavailable)" - built.append( - "===== visible terminal · grid=\(item.grid) · font=\(item.font) =====\n" - + text - ) + var sections: [String] = [] + for item in pending { + guard let section = await item.view.visibleSnapshotSection( + surface: item.surface, + generation: item.generation, + grid: item.grid, + font: item.font + ) else { + return "===== visible terminal: (snapshot skipped — render busy) =====" } - holder.sections = built - done.signal() + sections.append(section) } - if done.wait(timeout: .now() + 0.6) == .timedOut { - return "===== visible terminal: (snapshot skipped — render busy) =====" + return sections.joined(separator: "\n\n") + } + + private func visibleSnapshotSection( + surface: ghostty_surface_t, + generation: UInt64, + grid: String, + font: Int + ) async -> String? { + guard self.surface == surface, + surfaceGeneration == generation, + !renderPipelineRecoveryPaused else { + return nil + } + return await withCheckedContinuation { continuation in + let operationID = makeSurfaceOperationID() + if let existing = pendingVisibleSnapshot { + pendingVisibleSnapshot = nil + existing.continuation.resume(returning: nil) + } + pendingVisibleSnapshot = PendingVisibleSnapshot( + id: operationID, + startedAt: CACurrentMediaTime(), + continuation: continuation + ) + ensureSurfaceOperationDeadlinePump() + let queue = outputQueue + let read = VisibleSnapshotRead(surface: surface, generation: generation, grid: grid, font: font) + queue.async { + let text = Self.surfaceText(read.surface, pointTag: GHOSTTY_POINT_VIEWPORT) ?? "(unavailable)" + let section = "===== visible terminal · grid=\(read.grid) · font=\(read.font) =====\n" + + text + Task { @MainActor [weak self] in + guard let view = self else { return } + guard view.surface == read.surface, + view.surfaceGeneration == read.generation else { + view.completePendingVisibleSnapshot(id: operationID, returning: nil) + return + } + view.completePendingVisibleSnapshot(id: operationID, returning: section) + } + } } - return holder.sections.joined(separator: "\n\n") } private func handleBell() { @@ -3710,26 +4169,84 @@ extension GhosttySurfaceView: UIScrollViewDelegate { } } +nonisolated private enum RenderPipelineRecoveryReplay { + case callerWillRequestReplay + case delegateWhenNoCaller +} + +/// One output/geometry operation awaiting either its output-queue completion or +/// the display-link deadline that rebuilds the stalled render pipeline. +nonisolated private struct PendingSurfaceOperation { + let id: UInt64 + let startedAt: CFTimeInterval + let byteCount: Int? + let continuation: CheckedContinuation +} + +/// One visible-terminal snapshot read awaiting output-queue completion or its +/// display-link deadline. A timeout skips only the diagnostic snapshot. +nonisolated private struct PendingVisibleSnapshot { + let id: UInt64 + let startedAt: CFTimeInterval + let continuation: CheckedContinuation +} + +/// One "View as Text" read awaiting output-queue completion or deadline. +nonisolated private struct PendingCopyableTextRead { + let id: UInt64 + let startedAt: CFTimeInterval + let cancellation: SurfaceOperationCancellationToken + let continuation: CheckedContinuation + + func cancel() { + cancellation.cancel() + } +} + /// One surface's request for the bounded visible-terminal snapshot. -/// -/// The `ghostty_surface_t` is a C pointer that the snapshot only dereferences on -/// `GhosttySurfaceView.outputQueue` (the queue that owns `process_output`) and -/// never mutates, so carrying it across the queue hop is safe — hence -/// `@unchecked Sendable`. -private struct VisibleSnapshotRequest: @unchecked Sendable { +nonisolated private struct VisibleSnapshotRequest { + let view: GhosttySurfaceView let grid: String let font: Int let surface: ghostty_surface_t + let generation: UInt64 } -/// Carrier for the snapshot text produced off `GhosttySurfaceView.outputQueue`. +/// Raw surface read payload captured by the off-main output queue. /// -/// `sections` is written exactly once on that queue before its semaphore is -/// signaled and read by the caller only after the matching wait, so the two -/// accesses never overlap — hence `@unchecked Sendable`. On the timeout path the -/// caller never reads it, leaving the queue task the sole accessor. -private final class VisibleSnapshotHolder: @unchecked Sendable { - var sections: [String] = [] +/// The C surface pointer is dereferenced only on `GhosttySurfaceWorkQueue`, +/// which is the same FIFO queue that owns `process_output` and surface free. +nonisolated private struct VisibleSnapshotRead: @unchecked Sendable { + let surface: ghostty_surface_t + let generation: UInt64 + let grid: String + let font: Int +} + +/// Raw full-text read payload captured by the off-main output queue. +/// +/// The C surface pointer is dereferenced only on `GhosttySurfaceWorkQueue`, +/// which is the same FIFO queue that owns `process_output` and surface free. +nonisolated private struct CopyableTextRead: @unchecked Sendable { + let surface: ghostty_surface_t + let generation: UInt64 + let cancellation: SurfaceOperationCancellationToken +} + +nonisolated private final class SurfaceOperationCancellationToken: Sendable { + // lint:allow lock - tiny cross-queue cancellation flag for already-enqueued + // libghostty work; actor hops would put the serial surface queue back behind + // the main actor and defeat the stale-read fast path. + private let cancelled: Mutex + = .init(false) + + var isCancelled: Bool { + cancelled.withLock { $0 } + } + + func cancel() { + cancelled.withLock { $0 = true } + } } private class DisplayLinkProxy { diff --git a/Packages/iOS/CmuxMobileTerminal/Sources/CmuxMobileTerminal/GhosttySurfaceWorkQueue.swift b/Packages/iOS/CmuxMobileTerminal/Sources/CmuxMobileTerminal/GhosttySurfaceWorkQueue.swift new file mode 100644 index 000000000000..999dec0ba9c0 --- /dev/null +++ b/Packages/iOS/CmuxMobileTerminal/Sources/CmuxMobileTerminal/GhosttySurfaceWorkQueue.swift @@ -0,0 +1,23 @@ +import Foundation + +/// Owns the serial libghostty work queue for one surface generation. +/// All mutable state is accessed only from `queue`; main-actor code replaces whole instances on recovery. +final class GhosttySurfaceWorkQueue: @unchecked Sendable { + let queue: DispatchQueue + #if DEBUG + /// Accessed only from ``queue`` while producing DEBUG accessibility snapshots. + var lastAccessibilityTextTime: CFTimeInterval = 0 + #endif + + init(generation: UInt64) { + // carve-out justification: serial event-delivery queue for low-level libghostty C calls; not used as a lock. + queue = DispatchQueue( + label: "dev.cmux.GhosttySurfaceView.output.\(generation)", + qos: .userInitiated + ) + } + + func async(_ work: @escaping @Sendable () -> Void) { + queue.async(execute: work) + } +} diff --git a/ios/cmux/Resources/Localizable.xcstrings b/ios/cmux/Resources/Localizable.xcstrings index a7e8b26cc5e7..938f3f49ee74 100644 --- a/ios/cmux/Resources/Localizable.xcstrings +++ b/ios/cmux/Resources/Localizable.xcstrings @@ -2296,6 +2296,23 @@ } } }, + "mobile.debug.copyLogs": { + "extractionState": "manual", + "localizations": { + "en": { + "stringUnit": { + "state": "translated", + "value": "Copy Debug Logs" + } + }, + "ja": { + "stringUnit": { + "state": "translated", + "value": "デバッグログをコピー" + } + } + } + }, "mobile.feedback.cancel": { "extractionState": "manual", "localizations": { @@ -5968,6 +5985,23 @@ } } }, + "mobile.terminal.rendererFailed": { + "extractionState": "manual", + "localizations": { + "en": { + "stringUnit": { + "state": "translated", + "value": "Terminal renderer failed to start." + } + }, + "ja": { + "stringUnit": { + "state": "translated", + "value": "ターミナルレンダラーを起動できませんでした。" + } + } + } + }, "mobile.terminal.viewAsText": { "extractionState": "manual", "localizations": {