diff --git a/.github/workflows/test-depot.yml b/.github/workflows/test-depot.yml index d394291686fd..dd7ecbb7c163 100644 --- a/.github/workflows/test-depot.yml +++ b/.github/workflows/test-depot.yml @@ -165,7 +165,11 @@ jobs: suite_output=$(run_unit_suite "-only-testing:cmuxTests/$suite") || suite_status=$? printf '%s\n' "$suite_output" if [ "$suite_status" -ne 0 ]; then return "$suite_status"; fi - if ! printf '%s\n' "$suite_output" | grep -Eq 'Test run with [1-9][0-9]* tests|Executed [1-9][0-9]* tests'; then + # A pipeline here makes `set -o pipefail` treat grep's normal + # early exit as a printf SIGPIPE when the xcodebuild log is + # large. Feed grep with a here-string so a passing suite cannot + # be reported as empty because of that logging artifact. + if ! grep -Eq 'Test run with [1-9][0-9]* tests|Executed [1-9][0-9]* tests' <<<"$suite_output"; then echo "No tests executed for $suite" >&2 return 1 fi diff --git a/Packages/Shared/CmuxIrohTransport/Sources/CmuxIrohTransport/CmxIrohClientSession.swift b/Packages/Shared/CmuxIrohTransport/Sources/CmuxIrohTransport/CmxIrohClientSession.swift index 6ebd5548b09b..fc8f991d5de5 100644 --- a/Packages/Shared/CmuxIrohTransport/Sources/CmuxIrohTransport/CmxIrohClientSession.swift +++ b/Packages/Shared/CmuxIrohTransport/Sources/CmuxIrohTransport/CmxIrohClientSession.swift @@ -157,7 +157,7 @@ public actor CmxIrohClientSession { priority: Int32 ) async throws -> CmxIrohBidirectionalStream { switch lane { - case .terminal, .artifact, .simulatorStream: + case .terminal, .terminalInput, .artifact, .simulatorStream: break case .control, .serverEvents: throw CmxIrohClientSessionError.invalidOutgoingLane diff --git a/Packages/Shared/CmuxIrohTransport/Sources/CmuxIrohTransport/CmxIrohLane.swift b/Packages/Shared/CmuxIrohTransport/Sources/CmuxIrohTransport/CmxIrohLane.swift index 9d7544a97ea2..995bc83a39e4 100644 --- a/Packages/Shared/CmuxIrohTransport/Sources/CmuxIrohTransport/CmxIrohLane.swift +++ b/Packages/Shared/CmuxIrohTransport/Sources/CmuxIrohTransport/CmxIrohLane.swift @@ -9,6 +9,12 @@ public enum CmxIrohLane: Equatable, Sendable { /// One terminal's ordered stream resumed after the optional byte cursor. case terminal(resourceID: CmxIrohResourceID, cursor: UInt64?) + /// A terminal input-only lane. The host sends one empty replay envelope + /// as a readiness baseline, then the stream carries only input frames. + /// Render-grid sessions use this lane so keystrokes never wait on RPC + /// settlement or compete with the authoritative output event stream. + case terminalInput(resourceID: CmxIrohResourceID) + /// A low-priority artifact stream resumed at an exact byte offset. case artifact(resourceID: CmxIrohResourceID, offset: UInt64) diff --git a/Packages/Shared/CmuxIrohTransport/Sources/CmuxIrohTransport/CmxIrohServerSession.swift b/Packages/Shared/CmuxIrohTransport/Sources/CmuxIrohTransport/CmxIrohServerSession.swift index c83a76d054f2..e1299e0ad077 100644 --- a/Packages/Shared/CmuxIrohTransport/Sources/CmuxIrohTransport/CmxIrohServerSession.swift +++ b/Packages/Shared/CmuxIrohTransport/Sources/CmuxIrohTransport/CmxIrohServerSession.swift @@ -197,7 +197,7 @@ public actor CmxIrohServerSession { from: stream.receiveStream ) switch decoded.header.lane { - case .terminal, .artifact, .simulatorStream: + case .terminal, .terminalInput, .artifact, .simulatorStream: break case .control, .serverEvents: throw CmxIrohServerSessionError.invalidPeerLane @@ -235,7 +235,7 @@ public actor CmxIrohServerSession { break case .artifact: throw CmxIrohServerSessionError.applicationLanesUnavailable - case .control, .terminal, .simulatorStream: + case .control, .terminal, .terminalInput, .simulatorStream: throw CmxIrohServerSessionError.invalidServerLane } let stream = try await connection.openSendStream() diff --git a/Packages/Shared/CmuxIrohTransport/Sources/CmuxIrohTransport/CmxIrohStreamHeaderCodec.swift b/Packages/Shared/CmuxIrohTransport/Sources/CmuxIrohTransport/CmxIrohStreamHeaderCodec.swift index 22f9fb216385..928a252bacd2 100644 --- a/Packages/Shared/CmuxIrohTransport/Sources/CmuxIrohTransport/CmxIrohStreamHeaderCodec.swift +++ b/Packages/Shared/CmuxIrohTransport/Sources/CmuxIrohTransport/CmxIrohStreamHeaderCodec.swift @@ -78,6 +78,12 @@ public struct CmxIrohStreamHeaderCodec: Sendable { append(cursor, to: &payload) } + case let .terminalInput(resourceID): + laneCode = 6 + credentialCode = 0 + flags = 0 + try appendLengthPrefixedString(resourceID.value, lengthByteCount: 1, to: &payload) + case let .artifact(resourceID, offset): laneCode = 4 credentialCode = 0 @@ -201,6 +207,15 @@ public struct CmxIrohStreamHeaderCodec: Sendable { } let resourceID = try readResourceID(payload: &payload) return try CmxIrohStreamHeader(lane: .simulatorStream(resourceID: resourceID)) + case 6: + guard flags == 0 else { + throw CmxIrohStreamHeaderCodecError.invalidFlags(flags) + } + guard credentialCode == 0 else { + throw CmxIrohStreamHeaderCodecError.invalidCredentialKind(credentialCode) + } + let resourceID = try readResourceID(payload: &payload) + return try CmxIrohStreamHeader(lane: .terminalInput(resourceID: resourceID)) default: throw CmxIrohStreamHeaderCodecError.unknownLane(laneCode) } diff --git a/Packages/Shared/CmuxIrxTransport/Sources/CmuxIrxTransport/IrxProtocol.swift b/Packages/Shared/CmuxIrxTransport/Sources/CmuxIrxTransport/IrxProtocol.swift index 3f6708985571..a5980248881d 100644 --- a/Packages/Shared/CmuxIrxTransport/Sources/CmuxIrxTransport/IrxProtocol.swift +++ b/Packages/Shared/CmuxIrxTransport/Sources/CmuxIrxTransport/IrxProtocol.swift @@ -87,6 +87,7 @@ public enum IrxLaneKind: String, Codable, Sendable { case keepalive case events case terminal + case terminalInput = "terminal_input" case artifact case simulatorStream = "simulator_stream" } diff --git a/Packages/iOS/CmuxMobileRPC/Sources/CmuxMobileRPC/MobileSyncRuntime.swift b/Packages/iOS/CmuxMobileRPC/Sources/CmuxMobileRPC/MobileSyncRuntime.swift index cac68bfb5cca..7448614c5c4a 100644 --- a/Packages/iOS/CmuxMobileRPC/Sources/CmuxMobileRPC/MobileSyncRuntime.swift +++ b/Packages/iOS/CmuxMobileRPC/Sources/CmuxMobileRPC/MobileSyncRuntime.swift @@ -42,6 +42,9 @@ public protocol MobileSyncRuntime: Sendable { /// Optional source for one independent, sequence-aware terminal lane per /// mounted surface. A nil provider preserves control/event delivery. var terminalLaneProvider: MobileTerminalLaneProvider? { get } + /// Optional source for a terminal input-only lane. It carries one empty + /// replay baseline, then fire-and-forget input frames without output. + var terminalInputLaneProvider: MobileTerminalLaneProvider? { get } /// Optional source for low-priority raw artifact bytes on an admitted Iroh peer. var artifactLaneProvider: MobileArtifactLaneProvider? { get } /// Optional source for one dedicated simulator-stream v2 video lane per @@ -69,6 +72,7 @@ public protocol MobileSyncRuntime: Sendable { public extension MobileSyncRuntime { var independentEventByteStreamProvider: CmxIndependentEventByteStreamProvider? { nil } var terminalLaneProvider: MobileTerminalLaneProvider? { nil } + var terminalInputLaneProvider: MobileTerminalLaneProvider? { nil } var artifactLaneProvider: MobileArtifactLaneProvider? { nil } var simulatorStreamLaneProvider: MobileSimulatorStreamLaneProvider? { nil } diff --git a/Packages/iOS/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite+Capabilities.swift b/Packages/iOS/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite+Capabilities.swift index 235351525777..d0331397c583 100644 --- a/Packages/iOS/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite+Capabilities.swift +++ b/Packages/iOS/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite+Capabilities.swift @@ -42,6 +42,14 @@ extension MobileShellComposite { && supportedHostCapabilities.contains(Self.terminalVerifiedReplayCapability) } + /// Hybrid sessions subscribe to render-grid events only for screen-state + /// tracking and alternate-screen recovery. Primary-screen painting stays + /// on the sequence-aware byte lane, so an advisory grid must never impose + /// its shared viewport dimensions on the local natural surface. + public var usesHybridTerminalOutput: Bool { + terminalOutputTransport == .hybrid + } + /// Screen-anchored render-grid sessions receive active-area-anchored /// frames whose deltas carry exact scrolled-row counts, so this device /// keeps a deep local scrollback and scrolls the primary screen locally diff --git a/Packages/iOS/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite+TerminalLane.swift b/Packages/iOS/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite+TerminalLane.swift index c498f8197a0b..58dddb01d7f1 100644 --- a/Packages/iOS/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite+TerminalLane.swift +++ b/Packages/iOS/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite+TerminalLane.swift @@ -10,13 +10,21 @@ extension MobileShellComposite { guard !demonstrationOwnsSurface(surfaceID) else { return } guard let terminalLaneCoordinator, connectionState == .connected, - terminalOutputTransport != .renderGrid, terminalByteContinuationsBySurfaceID[surfaceID] != nil, let activeRoute, activeRoute.kind == .iroh, let activeTicket else { return } + let laneMode: MobileTerminalLaneCoordinator.LaneMode + if terminalOutputTransport == .renderGrid { + // Render-grid owns terminal output. Use a separate input-only + // stream so every key bypasses ordered RPC settlement. + guard runtime?.terminalInputLaneProvider != nil else { return } + laneMode = .inputOnly + } else { + laneMode = .output + } // A Direct lane request can redial the peer session, so it must carry // the same method-pinned allowlist as the control dial or it could // ride relay or discovered paths the method forbids. Tailscale @@ -36,6 +44,7 @@ extension MobileShellComposite { let configuration = MobileTerminalLaneCoordinator.Configuration( request: request, surfaceID: surfaceID, + mode: laneMode, cursor: { @MainActor [weak self] in guard let self, self.connectionGeneration == connectionGeneration, @@ -70,7 +79,7 @@ extension MobileShellComposite { func resumeTerminalLaneIfSuspended(surfaceID: String) { guard let terminalLaneCoordinator, connectionState == .connected, - terminalOutputTransport != .renderGrid else { return } + terminalReplayBarrierTokensBySurfaceID[surfaceID] == nil else { return } Task { await terminalLaneCoordinator.resume(surfaceID: surfaceID) } } @@ -105,10 +114,22 @@ extension MobileShellComposite { } func reconcileTerminalLanesForOutputTransport() { - if terminalOutputTransport == .renderGrid { - deactivateAllTerminalLanes() - } else { - restartTerminalLanesForMountedSurfaces() + // Render-grid keeps its authoritative event stream for output, but + // retains an input-only lane for fire-and-forget keystrokes. + guard let terminalLaneCoordinator else { return } + terminalLaneLifecycleID = UUID() + let lifecycleID = terminalLaneLifecycleID + terminalLaneOutputReadySurfaceIDs.removeAll() + let mountedSurfaceIDs = Array(terminalByteContinuationsBySurfaceID.keys) + Task { @MainActor [weak self] in + await terminalLaneCoordinator.deactivateAll() + guard let self, + self.terminalLaneLifecycleID == lifecycleID, + self.connectionState == .connected else { return } + for surfaceID in mountedSurfaceIDs + where self.terminalByteContinuationsBySurfaceID[surfaceID] != nil { + self.ensureTerminalLane(surfaceID: surfaceID) + } } } @@ -119,6 +140,11 @@ extension MobileShellComposite { guard terminalByteContinuationsBySurfaceID[surfaceID] != nil else { return .stop } + if terminalOutputTransport == .renderGrid { + // The input-only lane's baseline only gates readiness. Its output + // half is deliberately ignored because render-grid is authoritative. + return .accepted(outputReady: true) + } if terminalOutputTransport == .hybrid, terminalActiveScreenBySurfaceID[surfaceID] == .alternate { return .accepted(outputReady: true) diff --git a/Packages/iOS/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite+TerminalOutputDelivery.swift b/Packages/iOS/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite+TerminalOutputDelivery.swift index 1c24da81db43..64f2fb87da0b 100644 --- a/Packages/iOS/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite+TerminalOutputDelivery.swift +++ b/Packages/iOS/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite+TerminalOutputDelivery.swift @@ -604,6 +604,10 @@ extension MobileShellComposite { // so the floor restore is the truthful baseline hand-back. restoreTerminalPreBarrierBaselineIfNeeded(surfaceID: surfaceID) terminalReplayBarrierFollowUpCountsBySurfaceID.removeValue(forKey: surfaceID) + // Admission updates the cursor before the renderer acknowledges + // the replay. Reopen a backpressured lane only after that ACK + // releases the barrier, so its next frame is not dropped again. + resumeTerminalLaneIfSuspended(surfaceID: surfaceID) } } guard let next, diff --git a/Packages/iOS/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite+TerminalReplayLifecycle.swift b/Packages/iOS/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite+TerminalReplayLifecycle.swift index e0bbb5d55f53..504dcf5114c3 100644 --- a/Packages/iOS/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite+TerminalReplayLifecycle.swift +++ b/Packages/iOS/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite+TerminalReplayLifecycle.swift @@ -226,6 +226,10 @@ extension MobileShellComposite { } terminalRenderGridBaselineReplayBarrierTokensBySurfaceID.removeValue(forKey: surfaceID) terminalReplayBarrierTokensInFlightBySurfaceID.removeValue(forKey: surfaceID) + // A terminal lane may have paused on a replay-barrier backpressure + // response. The barrier is now resolved, so let it reopen from the + // delivered sequence instead of leaving input on the RPC fallback. + resumeTerminalLaneIfSuspended(surfaceID: surfaceID) MobileDebugLog.anchormux("terminal.output.replay_barrier_cleared_\(reason) surface=\(surfaceID)") return true } @@ -342,6 +346,9 @@ extension MobileShellComposite { cancelTerminalInputAckResubscribeRetry(surfaceID: surfaceID) pendingTerminalByteEndSeqBySurfaceID.removeValue(forKey: surfaceID) pendingTerminalInputDroppedRenderGridSurfaceIDs.remove(surfaceID) + // Fail-open also releases a lane paused behind a replay that could not + // settle. Its next attach will request a fresh bounded cursor replay. + resumeTerminalLaneIfSuspended(surfaceID: surfaceID) MobileDebugLog.anchormux("terminal.output.replay_barrier_fail_open surface=\(surfaceID) reason=\(reason)") return true } diff --git a/Packages/iOS/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite+TerminalViewport.swift b/Packages/iOS/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite+TerminalViewport.swift index 5820c9979b83..3e2edea287b5 100644 --- a/Packages/iOS/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite+TerminalViewport.swift +++ b/Packages/iOS/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite+TerminalViewport.swift @@ -95,6 +95,22 @@ extension MobileShellComposite { let requestGeneration = (viewportReportGenerationsBySequenceKey[sequenceKey] ?? 0) + 1 viewportReportGenerationsBySequenceKey[sequenceKey] = requestGeneration + // The output sink is registered immediately after the first geometry + // callback. Keep this generation marked until its acknowledgement + // settles so registration can let that acknowledgement own the cold + // replay instead of racing it with a second request. + // A geometry callback can supersede the preparation that was already + // in flight after the output sink mounted. Carry that sink's deferred + // cold replay onto the newer generation so the replacement + // acknowledgement still owns the initial replay. + if terminalViewportDeferredColdReplayGenerationsBySequenceKey + .removeValue(forKey: sequenceKey) != nil + { + terminalViewportDeferredColdReplayGenerationsBySequenceKey[ + sequenceKey + ] = requestGeneration + } + terminalViewportPreparationGenerationsBySequenceKey[sequenceKey] = requestGeneration return MobileTerminalViewportPreparation( workspaceID: workspaceID, ownerKey: ownerKey, @@ -128,6 +144,30 @@ extension MobileShellComposite { workspaceID: preparedWorkspaceID, terminalID: MobileTerminalPreview.ID(rawValue: surfaceID) ) + func finishPreparation( + requestColdReplay: Bool = false, + replayAlreadyRequested: Bool = false + ) { + guard terminalViewportPreparationGenerationsBySequenceKey[sequenceKey] + == requestGeneration else { return } + terminalViewportPreparationGenerationsBySequenceKey.removeValue( + forKey: sequenceKey + ) + let deferredColdReplay = !replayAlreadyRequested + && terminalViewportDeferredColdReplayGenerationsBySequenceKey.removeValue( + forKey: sequenceKey + ) == requestGeneration + if replayAlreadyRequested { + terminalViewportDeferredColdReplayGenerationsBySequenceKey.removeValue( + forKey: sequenceKey + ) + } + if (requestColdReplay || deferredColdReplay), + hasTerminalOutputSink(surfaceID: surfaceID) { + requestColdAttachTerminalReplay(surfaceID: surfaceID) + } + } + guard foregroundMacKey == preparation.ownerKey, workspaceID(forTerminalID: surfaceID) == preparedWorkspaceID, reportedViewportSizesByTerminalKey[key] == reportedGrid, @@ -137,6 +177,7 @@ extension MobileShellComposite { correlationID: surfaceID, failure: .superseded ) + finishPreparation() return nil } // Demonstration surfaces answer the viewport report locally with the @@ -153,6 +194,7 @@ extension MobileShellComposite { correlationID: surfaceID, count: columns * rows ) + finishPreparation() return ( columns: columns, rows: rows, @@ -166,6 +208,7 @@ extension MobileShellComposite { correlationID: surfaceID, failure: .offline ) + finishPreparation(requestColdReplay: true) return nil } let previousReportedGrid = reportedTerminalViewportSizesBySurfaceID[surfaceID] @@ -199,6 +242,7 @@ extension MobileShellComposite { correlationID: surfaceID, failure: .superseded ) + finishPreparation(requestColdReplay: true) return nil } guard viewportReportGenerationsBySequenceKey[sequenceKey] == requestGeneration else { @@ -208,11 +252,14 @@ extension MobileShellComposite { correlationID: surfaceID, failure: .superseded ) + // A newer preparation owns the replay decision and its + // generation marker. Do not issue a fallback for this stale + // response. return nil } guard let payload = try? MobileTerminalViewportResponse.decode(data), let grid = payload.effectiveGrid else { - finishPrearmedTerminalViewportBarrierWithoutResize( + let replayRequested = finishPrearmedTerminalViewportBarrierWithoutResize( surfaceID: surfaceID, token: prearmedReplayBarrierToken, reason: "viewport_missing_grid" @@ -222,6 +269,10 @@ extension MobileShellComposite { correlationID: surfaceID, failure: .protocolViolation ) + finishPreparation( + requestColdReplay: !replayRequested, + replayAlreadyRequested: replayRequested + ) return nil } reportedTerminalViewportSizesBySurfaceID[surfaceID] = reportedGrid @@ -229,6 +280,7 @@ extension MobileShellComposite { let previousGrid = effectiveViewportSizesBySurfaceID[surfaceID] effectiveViewportSizesBySurfaceID[surfaceID] = effectiveGrid let shouldRequestReplay = previousGrid.map { $0 != effectiveGrid } ?? true + var replayRequested = false if shouldRequestReplay, hasTerminalOutputSink(surfaceID: surfaceID) { let replayBarrierToken: UUID @@ -253,6 +305,7 @@ extension MobileShellComposite { "terminal.output.viewport_resync surface=\(surfaceID) grid=\(effectiveGrid.columns)x\(effectiveGrid.rows)" ) requestTerminalReplay(surfaceID: surfaceID, replayBarrierToken: replayBarrierToken) + replayRequested = true } else if prearmedReplayBarrierToken == nil, terminalReplayBarrierTokensBySurfaceID[surfaceID] != nil, terminalReplayFailureRetryExhausted(surfaceID: surfaceID), @@ -266,8 +319,9 @@ extension MobileShellComposite { let replayBarrierToken = beginTerminalReplayBarrierCarryingReplacedWork(surfaceID: surfaceID) MobileDebugLog.anchormux("terminal.output.viewport_rearm_exhausted surface=\(surfaceID)") requestTerminalReplay(surfaceID: surfaceID, replayBarrierToken: replayBarrierToken) + replayRequested = true } else { - finishPrearmedTerminalViewportBarrierWithoutResize( + replayRequested = finishPrearmedTerminalViewportBarrierWithoutResize( surfaceID: surfaceID, token: prearmedReplayBarrierToken, reason: "viewport_unchanged" @@ -278,6 +332,12 @@ extension MobileShellComposite { correlationID: surfaceID, count: grid.columns * grid.rows ) + // Sink registration defers the cold attach replay while this + // acknowledgement is in flight. If the effective grid is + // unchanged, the normal resync branches above do not request a + // replay, so fulfill that deferred registration before clearing + // its preparation marker. + finishPreparation(replayAlreadyRequested: replayRequested) return ( columns: grid.columns, rows: grid.rows, @@ -303,11 +363,15 @@ extension MobileShellComposite { ) return nil } - finishPrearmedTerminalViewportBarrierWithoutResize( + let replayRequested = finishPrearmedTerminalViewportBarrierWithoutResize( surfaceID: surfaceID, token: prearmedReplayBarrierToken, reason: "viewport_failed" ) + finishPreparation( + requestColdReplay: !replayRequested, + replayAlreadyRequested: replayRequested + ) terminalViewportLog.error("viewport report failed surface=\(surfaceID, privacy: .public) error=\(String(describing: error), privacy: .public)") recordAppEvent( .terminalViewportReportFailed, @@ -322,6 +386,16 @@ extension MobileShellComposite { /// detach). Fire-and-forget; the Mac also clears on connection close. public func clearTerminalViewport(surfaceID: String) { recordAppEvent(.terminalViewportClearStarted, correlationID: surfaceID) + let sequenceKey = MobileTerminalViewportSequenceKey( + ownerKey: foregroundMacKey, + surfaceID: surfaceID + ) + terminalViewportPreparationGenerationsBySequenceKey.removeValue( + forKey: sequenceKey + ) + terminalViewportDeferredColdReplayGenerationsBySequenceKey.removeValue( + forKey: sequenceKey + ) let workspaceID = workspaceID(forTerminalID: surfaceID) // A clear releases the presentation's full local viewport lease. Any // replay, input, or paste that races after this point must not carry @@ -346,10 +420,6 @@ extension MobileShellComposite { // applying after detach and blocks generation reuse across re-attach. // Warm focus swaps keep the peer connection alive, so its fence must // outlive the focused role. The account boundary clears all sequences. - let sequenceKey = MobileTerminalViewportSequenceKey( - ownerKey: foregroundMacKey, - surfaceID: surfaceID - ) let clearGeneration = (viewportReportGenerationsBySequenceKey[sequenceKey] ?? 0) + 1 viewportReportGenerationsBySequenceKey[sequenceKey] = clearGeneration @@ -420,25 +490,27 @@ extension MobileShellComposite { return replayBarrierToken } + @discardableResult private func finishPrearmedTerminalViewportBarrierWithoutResize( surfaceID: String, token: UUID?, reason: String - ) { - guard let token else { return } + ) -> Bool { + guard let token else { return false } terminalViewportReplayBarrierPendingAckTokensBySurfaceID.removeValue(forKey: surfaceID) - guard terminalReplayBarrierTokensBySurfaceID[surfaceID] == token else { return } + guard terminalReplayBarrierTokensBySurfaceID[surfaceID] == token else { return false } if terminalReplayBarrierDroppedOutputSurfaceIDs.contains(surfaceID), hasTerminalOutputSink(surfaceID: surfaceID), remoteClient != nil { MobileDebugLog.anchormux("terminal.output.viewport_replay_after_\(reason) surface=\(surfaceID)") requestTerminalReplay(surfaceID: surfaceID, replayBarrierToken: token) - return + return true } clearTerminalReplayBarrierIfCurrent( surfaceID: surfaceID, token: token, reason: reason ) + return false } } diff --git a/Packages/iOS/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite.swift b/Packages/iOS/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite.swift index b3338886cad8..1cf2ddc85e93 100644 --- a/Packages/iOS/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite.swift +++ b/Packages/iOS/CmuxMobileShell/Sources/CmuxMobileShell/MobileShellComposite.swift @@ -1562,6 +1562,13 @@ public final class MobileShellComposite: MobileTerminalOutputSinking { var terminalReplayBarrierDroppedOutputCountsBySurfaceID: [String: UInt64] var terminalReplayBarrierAckCoveredDroppedOutputCountsBySurfaceID: [String: UInt64] var terminalViewportReplayBarrierPendingAckTokensBySurfaceID: [String: UUID] + /// Viewport reports committed by the surface before its output sink is + /// registered. The first authoritative replay is scheduled by the + /// viewport acknowledgement, so cold attach must not race it. + var terminalViewportPreparationGenerationsBySequenceKey: + [MobileTerminalViewportSequenceKey: UInt64] + var terminalViewportDeferredColdReplayGenerationsBySequenceKey: + [MobileTerminalViewportSequenceKey: UInt64] var terminalReplayFailureRetryCountsBySurfaceID: [String: Int] var terminalReplayBarrierFollowUpCountsBySurfaceID: [String: Int] var terminalColdAttachReplayBarrierTokensBySurfaceID: [String: UUID] @@ -1963,6 +1970,8 @@ public final class MobileShellComposite: MobileTerminalOutputSinking { self.terminalReplayBarrierDroppedOutputCountsBySurfaceID = [:] self.terminalReplayBarrierAckCoveredDroppedOutputCountsBySurfaceID = [:] self.terminalViewportReplayBarrierPendingAckTokensBySurfaceID = [:] + self.terminalViewportPreparationGenerationsBySequenceKey = [:] + self.terminalViewportDeferredColdReplayGenerationsBySequenceKey = [:] self.terminalReplayFailureRetryCountsBySurfaceID = [:] self.terminalReplayBarrierFollowUpCountsBySurfaceID = [:] self.terminalColdAttachReplayBarrierTokensBySurfaceID = [:] @@ -1972,9 +1981,11 @@ public final class MobileShellComposite: MobileTerminalOutputSinking { self.terminalOutputStreamTokensBySurfaceID = [:] self.terminalOutputConsumerOwnerIDsBySurfaceID = [:] self.terminalOutputQueuesBySurfaceID = [:] - if let terminalLaneProvider = runtime?.terminalLaneProvider { + if runtime?.terminalLaneProvider != nil + || runtime?.terminalInputLaneProvider != nil { self.terminalLaneCoordinator = MobileTerminalLaneCoordinator( - provider: terminalLaneProvider + provider: runtime?.terminalLaneProvider, + inputOnlyProvider: runtime?.terminalInputLaneProvider ) } else { self.terminalLaneCoordinator = nil @@ -11363,6 +11374,8 @@ public final class MobileShellComposite: MobileTerminalOutputSinking { terminalReplayBarrierDroppedOutputCountsBySurfaceID = [:] terminalReplayBarrierAckCoveredDroppedOutputCountsBySurfaceID = [:] terminalViewportReplayBarrierPendingAckTokensBySurfaceID = [:] + terminalViewportPreparationGenerationsBySequenceKey = [:] + terminalViewportDeferredColdReplayGenerationsBySequenceKey = [:] terminalReplayFailureRetryCountsBySurfaceID = [:] terminalReplayBarrierFollowUpCountsBySurfaceID = [:] terminalColdAttachReplayBarrierTokensBySurfaceID = [:] @@ -14484,7 +14497,28 @@ public final class MobileShellComposite: MobileTerminalOutputSinking { mobileShellLog.info("CMUX_REPLAY register sink surface=\(surfaceID, privacy: .public) connected=\(self.connectionState == .connected, privacy: .public) hasClient=\(self.remoteClient != nil, privacy: .public) workspaceCount=\(self.workspaces.count, privacy: .public)") startLatencyProbeIfReady() #endif - requestColdAttachTerminalReplay(surfaceID: surfaceID) + // The first viewport callback commits a generation before this sink is + // registered. Let its viewport acknowledgement schedule the one + // authoritative replay, avoiding a cold attach request that races the + // geometry RPC and gets immediately superseded. + let preparationSequenceKey = MobileTerminalViewportSequenceKey( + ownerKey: foregroundMacKey, + surfaceID: surfaceID + ) + if terminalViewportPreparationGenerationsBySequenceKey[ + preparationSequenceKey + ] == nil { + requestColdAttachTerminalReplay(surfaceID: surfaceID) + } else { + terminalViewportDeferredColdReplayGenerationsBySequenceKey[ + preparationSequenceKey + ] = terminalViewportPreparationGenerationsBySequenceKey[ + preparationSequenceKey + ] + MobileDebugLog.anchormux( + "terminal.output.defer_cold_replay surface=\(surfaceID)" + ) + } ensureTerminalLane(surfaceID: surfaceID) return streamToken } @@ -14514,6 +14548,18 @@ public final class MobileShellComposite: MobileTerminalOutputSinking { terminalScrollbackPrefetchStatesBySurfaceID.removeValue(forKey: surfaceID) effectiveViewportSizesBySurfaceID.removeValue(forKey: surfaceID); reportedTerminalViewportSizesBySurfaceID.removeValue(forKey: surfaceID) terminalViewportReplayBarrierPendingAckTokensBySurfaceID.removeValue(forKey: surfaceID) + terminalViewportPreparationGenerationsBySequenceKey.removeValue( + forKey: MobileTerminalViewportSequenceKey( + ownerKey: foregroundMacKey, + surfaceID: surfaceID + ) + ) + terminalViewportDeferredColdReplayGenerationsBySequenceKey.removeValue( + forKey: MobileTerminalViewportSequenceKey( + ownerKey: foregroundMacKey, + surfaceID: surfaceID + ) + ) deliveredTerminalByteEndSeqBySurfaceID.removeValue(forKey: surfaceID) terminalPreBarrierDeliveredEndSeqBySurfaceID.removeValue(forKey: surfaceID) terminalRenderGridBaselineReplayRequestCountsBySurfaceID.removeValue(forKey: surfaceID) diff --git a/Packages/iOS/CmuxMobileShell/Sources/CmuxMobileShell/MobileTerminalLaneCoordinator.swift b/Packages/iOS/CmuxMobileShell/Sources/CmuxMobileShell/MobileTerminalLaneCoordinator.swift index c6cd13686417..a9603ebf3e4e 100644 --- a/Packages/iOS/CmuxMobileShell/Sources/CmuxMobileShell/MobileTerminalLaneCoordinator.swift +++ b/Packages/iOS/CmuxMobileShell/Sources/CmuxMobileShell/MobileTerminalLaneCoordinator.swift @@ -4,6 +4,11 @@ import Foundation /// Owns independent terminal lanes keyed by peer and mounted surface. actor MobileTerminalLaneCoordinator { + enum LaneMode: Equatable, Sendable { + case output + case inputOnly + } + enum FrameDisposition: Sendable { case accepted(outputReady: Bool) case suspendUntilAuthoritativeOutput @@ -26,9 +31,26 @@ actor MobileTerminalLaneCoordinator { struct Configuration: Sendable { let request: CmxByteTransportRequest let surfaceID: String + let mode: LaneMode let cursor: @Sendable () async -> UInt64? let consume: @Sendable (MobileTerminalLaneOutputFrame) async -> FrameDisposition let readinessChanged: @Sendable (Bool) async -> Void + + init( + request: CmxByteTransportRequest, + surfaceID: String, + mode: LaneMode = .output, + cursor: @escaping @Sendable () async -> UInt64?, + consume: @escaping @Sendable (MobileTerminalLaneOutputFrame) async -> FrameDisposition, + readinessChanged: @escaping @Sendable (Bool) async -> Void + ) { + self.request = request + self.surfaceID = surfaceID + self.mode = mode + self.cursor = cursor + self.consume = consume + self.readinessChanged = readinessChanged + } } private struct LaneKey: Hashable, Sendable { @@ -75,12 +97,17 @@ actor MobileTerminalLaneCoordinator { private static let maximumOpenAttempts = 3 - private let provider: MobileTerminalLaneProvider + private let provider: MobileTerminalLaneProvider? + private let inputOnlyProvider: MobileTerminalLaneProvider? private var entriesByKey: [LaneKey: Entry] = [:] private var focusedKeyBySurfaceID: [String: LaneKey] = [:] - init(provider: @escaping MobileTerminalLaneProvider) { + init( + provider: MobileTerminalLaneProvider?, + inputOnlyProvider: MobileTerminalLaneProvider? = nil + ) { self.provider = provider + self.inputOnlyProvider = inputOnlyProvider } func ensure(_ configuration: Configuration) async { @@ -196,9 +223,27 @@ actor MobileTerminalLaneCoordinator { while openAttempt < Self.maximumOpenAttempts, !Task.isCancelled { guard let entry = entriesByKey[key], entry.id == id else { return } let configuration = entry.configuration - let requestedCursor = await configuration.cursor() + // Input-only lanes carry an empty replay baseline solely to gate + // readiness. Their host baseline is sampled independently from + // the output event stream, so validating it against the output + // cursor would reject a healthy lane whenever output advanced + // between those two operations. + let requestedCursor = configuration.mode == .inputOnly + ? nil + : await configuration.cursor() do { - let lane = try await provider( + let laneProvider: MobileTerminalLaneProvider? + switch configuration.mode { + case .output: + laneProvider = provider + case .inputOnly: + laneProvider = inputOnlyProvider ?? provider + } + guard let laneProvider else { + await markFailed(key: key, id: id) + return + } + let lane = try await laneProvider( configuration.request, configuration.surfaceID, requestedCursor @@ -227,7 +272,17 @@ actor MobileTerminalLaneCoordinator { } switch disposition { case let .accepted(outputReady): - await setOutputReady(outputReady, key: key, id: id) + if outputReady { + await setOutputReady(true, key: key, id: id) + } else { + // A consumer can reject a frame temporarily while + // an authoritative replay barrier owns the output + // sink. Stop reading immediately. Draining the + // next chunk would discard the replay baseline and + // turn backpressure into a false sequence gap. + await suspend(key: key, id: id, lane: lane) + return + } case .suspendUntilAuthoritativeOutput: await suspend(key: key, id: id, lane: lane) return diff --git a/Packages/iOS/CmuxMobileShell/Tests/CmuxMobileShellTests/ComposerSubmitRoutingTestSupport.swift b/Packages/iOS/CmuxMobileShell/Tests/CmuxMobileShellTests/ComposerSubmitRoutingTestSupport.swift index 48d746f72ec7..5bed1b77f97d 100644 --- a/Packages/iOS/CmuxMobileShell/Tests/CmuxMobileShellTests/ComposerSubmitRoutingTestSupport.swift +++ b/Packages/iOS/CmuxMobileShell/Tests/CmuxMobileShellTests/ComposerSubmitRoutingTestSupport.swift @@ -16,6 +16,7 @@ import Testing struct RoutingTestRuntime: MobileSyncRuntime { var transportFactory: any CmxByteTransportFactory var terminalLaneProvider: MobileTerminalLaneProvider? = nil + var terminalInputLaneProvider: MobileTerminalLaneProvider? = nil var stackAccessTokenProvider: @Sendable () async throws -> String = { "test-stack-token" } var stackAccessTokenForceRefresher: @Sendable () async throws -> String = { "test-stack-token" } var rpcRequestTimeoutNanoseconds: UInt64 = 30 * 1_000_000_000 diff --git a/Packages/iOS/CmuxMobileShell/Tests/CmuxMobileShellTests/MobileShellTerminalLaneConfigurationTests.swift b/Packages/iOS/CmuxMobileShell/Tests/CmuxMobileShellTests/MobileShellTerminalLaneConfigurationTests.swift new file mode 100644 index 000000000000..ce1d9b048632 --- /dev/null +++ b/Packages/iOS/CmuxMobileShell/Tests/CmuxMobileShellTests/MobileShellTerminalLaneConfigurationTests.swift @@ -0,0 +1,29 @@ +import CMUXMobileCore +import CmuxMobileRPC +import Foundation +import Testing +@testable import CmuxMobileShell + +@MainActor +@Suite +struct MobileShellTerminalLaneConfigurationTests { + @Test + func inputOnlyProviderAloneCreatesLaneCoordinator() { + let runtime = RoutingTestRuntime( + transportFactory: LaneConfigurationTestTransportFactory(), + terminalInputLaneProvider: { _, _, _ in + fatalError("the provider is only used after mounting a terminal") + } + ) + + let store = MobileShellComposite(runtime: runtime) + + #expect(store.terminalLaneCoordinator != nil) + } +} + +private struct LaneConfigurationTestTransportFactory: CmxByteTransportFactory { + func makeTransport(for _: CmxAttachRoute) throws -> any CmxByteTransport { + fatalError("the transport is not used by this configuration test") + } +} diff --git a/Packages/iOS/CmuxMobileShell/Tests/CmuxMobileShellTests/MobileTerminalLaneCoordinatorTests.swift b/Packages/iOS/CmuxMobileShell/Tests/CmuxMobileShellTests/MobileTerminalLaneCoordinatorTests.swift index 2d0810e6eca6..0e70a9549ab4 100644 --- a/Packages/iOS/CmuxMobileShell/Tests/CmuxMobileShellTests/MobileTerminalLaneCoordinatorTests.swift +++ b/Packages/iOS/CmuxMobileShell/Tests/CmuxMobileShellTests/MobileTerminalLaneCoordinatorTests.swift @@ -98,6 +98,61 @@ struct MobileTerminalLaneCoordinatorTests { #expect(await coordinator.sendInput("fallback", surfaceID: Self.surfaceID) == .unavailable) } + @Test + func inputOnlyLaneDoesNotValidateAgainstOutputCursor() async throws { + let lane = TerminalLaneTestConnection( + frames: [Self.frame(kind: .replay, sequence: 12, bytes: "")], + waitsAfterFrames: true + ) + let provider = TerminalLaneTestProvider(lanes: [lane]) + let coordinator = MobileTerminalLaneCoordinator { request, surfaceID, cursor in + try await provider.callAsFunction(request, surfaceID, cursor: cursor) + } + let readiness = TerminalLaneReadinessRecorder() + var readinessIterator = await readiness.stream().makeAsyncIterator() + + await coordinator.ensure(Self.configuration( + mode: .inputOnly, + providerRequest: try Self.request(), + cursor: { 5 }, + consume: { _ in .accepted(outputReady: true) }, + readinessChanged: { await readiness.append($0) } + )) + + #expect(await readinessIterator.next() == true) + #expect(await provider.requestedCursors() == [nil]) + #expect(await coordinator.sendInput("echo fast\\n", surfaceID: Self.surfaceID) == .sent) + await coordinator.deactivateAll() + } + + @Test + func outputLaneDoesNotFallBackToInputOnlyProvider() async throws { + let inputProvider = TerminalLaneTestProvider(lanes: [ + TerminalLaneTestConnection( + frames: [Self.frame(kind: .replay, sequence: 0, bytes: "")], + waitsAfterFrames: true + ), + ]) + let coordinator = MobileTerminalLaneCoordinator( + provider: nil, + inputOnlyProvider: { request, surfaceID, cursor in + try await inputProvider.callAsFunction(request, surfaceID, cursor: cursor) + } + ) + + await coordinator.ensure(Self.configuration( + providerRequest: try Self.request(), + cursor: { nil }, + consume: { _ in .accepted(outputReady: true) }, + readinessChanged: { _ in } + )) + try await Task.sleep(for: .milliseconds(10)) + + #expect(await inputProvider.requestCount() == 0) + #expect(await coordinator.isOutputReady(surfaceID: Self.surfaceID) == false) + await coordinator.deactivateAll() + } + @Test func sequenceGapSuspendsUntilAuthoritativeCursorThenReopens() async throws { let firstLane = TerminalLaneTestConnection( @@ -142,6 +197,56 @@ struct MobileTerminalLaneCoordinatorTests { await coordinator.deactivateAll() } + @Test + func consumerBackpressureSuspendsBeforeDrainingNextChunk() async throws { + let firstLane = TerminalLaneTestConnection( + frames: [ + Self.frame(kind: .replay, sequence: 0, bytes: "baseline"), + Self.frame(kind: .chunk, sequence: 8, bytes: "must-not-drain"), + ], + waitsAfterFrames: true + ) + let secondLane = TerminalLaneTestConnection( + frames: [Self.frame(kind: .replay, sequence: 0, bytes: "baseline")], + waitsAfterFrames: true + ) + let provider = TerminalLaneTestProvider(lanes: [firstLane, secondLane]) + let coordinator = MobileTerminalLaneCoordinator { request, surfaceID, cursor in + try await provider.callAsFunction(request, surfaceID, cursor: cursor) + } + let readiness = TerminalLaneReadinessRecorder() + var readinessIterator = await readiness.stream().makeAsyncIterator() + let consumed = TerminalLaneFrameRecorder() + let shouldRejectFirstFrame = TerminalLaneFlag(value: true) + + await coordinator.ensure(Self.configuration( + providerRequest: try Self.request(), + cursor: { nil }, + consume: { frame in + await consumed.append(frame) + if await shouldRejectFirstFrame.value() { + await shouldRejectFirstFrame.setValue(false) + return .accepted(outputReady: false) + } + return .accepted(outputReady: true) + }, + readinessChanged: { await readiness.append($0) } + )) + + for _ in 0..<100 { + if await firstLane.closeCount() == 1 { break } + try await Task.sleep(for: .milliseconds(1)) + } + #expect(await firstLane.closeCount() == 1) + #expect(await consumed.frames().map(\.kind) == [.replay]) + + await coordinator.resume(surfaceID: Self.surfaceID) + #expect(await readinessIterator.next() == true) + #expect(await consumed.frames().map(\.kind) == [.replay, .replay]) + #expect(await provider.requestCount() == 2) + await coordinator.deactivateAll() + } + @Test func replayCursorMismatchNeverBecomesReadyOrAcceptsInput() async throws { let mismatchedLane = TerminalLaneTestConnection( @@ -209,6 +314,7 @@ struct MobileTerminalLaneCoordinatorTests { } private static func configuration( + mode: MobileTerminalLaneCoordinator.LaneMode = .output, providerRequest: CmxByteTransportRequest, cursor: @escaping @Sendable () async -> UInt64?, consume: @escaping @Sendable (MobileTerminalLaneOutputFrame) async -> MobileTerminalLaneCoordinator.FrameDisposition, @@ -217,6 +323,7 @@ struct MobileTerminalLaneCoordinatorTests { MobileTerminalLaneCoordinator.Configuration( request: providerRequest, surfaceID: surfaceID, + mode: mode, cursor: cursor, consume: consume, readinessChanged: readinessChanged @@ -335,3 +442,14 @@ private actor TerminalLaneCursor { func value() -> UInt64? { storedValue } func setValue(_ value: UInt64?) { storedValue = value } } + +private actor TerminalLaneFlag { + private var storedValue: Bool + + init(value: Bool) { + self.storedValue = value + } + + func value() -> Bool { storedValue } + func setValue(_ value: Bool) { storedValue = value } +} diff --git a/Packages/iOS/CmuxMobileShell/Tests/CmuxMobileShellTests/MobileTerminalViewportSequenceKeyTests.swift b/Packages/iOS/CmuxMobileShell/Tests/CmuxMobileShellTests/MobileTerminalViewportSequenceKeyTests.swift new file mode 100644 index 000000000000..17d787f271c7 --- /dev/null +++ b/Packages/iOS/CmuxMobileShell/Tests/CmuxMobileShellTests/MobileTerminalViewportSequenceKeyTests.swift @@ -0,0 +1,24 @@ +import Testing +@testable import CmuxMobileShell + +@Suite("Terminal viewport preparation ownership") +struct MobileTerminalViewportSequenceKeyTests { + @Test("preparations for two Mac app instances do not collide") + func appInstanceOwnsPreparationMarker() { + let surfaceID = "terminal" + let first = MobileTerminalViewportSequenceKey( + ownerKey: MacPairingKey(macDeviceID: "mac", instanceTag: "first"), + surfaceID: surfaceID + ) + let second = MobileTerminalViewportSequenceKey( + ownerKey: MacPairingKey(macDeviceID: "mac", instanceTag: "second"), + surfaceID: surfaceID + ) + var markers: [MobileTerminalViewportSequenceKey: UInt64] = [first: 1] + + #expect(first != second) + #expect(markers[second] == nil) + markers[second] = 1 + #expect(markers.count == 2) + } +} diff --git a/Packages/iOS/CmuxMobileShell/Tests/CmuxMobileShellTests/TerminalLaneReplayBarrierTests.swift b/Packages/iOS/CmuxMobileShell/Tests/CmuxMobileShellTests/TerminalLaneReplayBarrierTests.swift new file mode 100644 index 000000000000..f3f915451f67 --- /dev/null +++ b/Packages/iOS/CmuxMobileShell/Tests/CmuxMobileShellTests/TerminalLaneReplayBarrierTests.swift @@ -0,0 +1,117 @@ +import CmuxMobileRPC +import Foundation +import Testing + +@testable import CmuxMobileShell + +@MainActor +@Test func terminalLaneResumesOnlyAfterReplayAcknowledgement() async throws { + let surfaceID = RoutingHostRouter.terminalA + let firstLane = ReplayBarrierTestLane(sequence: 0, bytes: "old") + let secondLane = ReplayBarrierTestLane(sequence: 5, bytes: "") + let provider = ReplayBarrierTestLaneProvider(lanes: [firstLane, secondLane]) + let router = RoutingHostRouter() + let store = try await makeRoutingConnectedStore( + router: router, + routeKind: .iroh, + terminalLaneProvider: { _, _, cursor in + try await provider.open(cursor: cursor) + } + ) + let output = store.terminalOutputStream(surfaceID: surfaceID) + defer { withExtendedLifetime(output) {} } + var iterator = output.makeAsyncIterator() + let barrier = store.beginTerminalReplayBarrier(surfaceID: surfaceID) + // Hold the authoritative request while the independent lane arrives. + store.terminalReplaySurfaceIDsInFlight.insert(surfaceID) + store.terminalReplayBarrierTokensInFlightBySurfaceID[surfaceID] = barrier + + for _ in 0..<1_000 { + if await firstLane.closeCount() == 1 { break } + await Task.yield() + } + #expect(await firstLane.closeCount() == 1) + #expect(await provider.cursors() == [nil]) + + #expect(store.deliverTerminalBytes( + Data("fresh".utf8), + surfaceID: surfaceID, + endSequence: 5, + bypassReplayBarrier: true + )) + store.terminalReplayBarrierAckCoveredDroppedOutputCountsBySurfaceID[surfaceID] = + store.terminalReplayBarrierDroppedOutputCountsBySurfaceID[surfaceID] + store.markTerminalBytesDelivered(surfaceID: surfaceID, endSeq: 5) + let replay = try #require(await iterator.next()) + for _ in 0..<100 { await Task.yield() } + #expect(await provider.cursors() == [nil]) + #expect(!store.terminalLaneOutputReadySurfaceIDs.contains(surfaceID)) + + store.terminalOutputDidProcess(surfaceID: surfaceID, streamToken: replay.streamToken) + for _ in 0..<1_000 { + if store.terminalLaneOutputReadySurfaceIDs.contains(surfaceID) { break } + await Task.yield() + } + #expect(store.terminalReplayBarrierTokensBySurfaceID[surfaceID] == nil) + #expect(await provider.cursors() == [nil, 5]) + #expect(store.terminalLaneOutputReadySurfaceIDs.contains(surfaceID)) + + await store.submitTerminalRawInput(Data("a".utf8), surfaceID: surfaceID) + #expect(await secondLane.inputs() == ["a"]) + #expect(await router.recordedTerminalInputs().isEmpty) + await store.terminalLaneCoordinator?.deactivateAll() +} + +private actor ReplayBarrierTestLane: MobileTerminalLaneConnection { + private var pendingFrame: MobileTerminalLaneOutputFrame? + private var waiter: CheckedContinuation? + private var closedCount = 0 + private var sentInputs: [String] = [] + + init(sequence: UInt64, bytes: String) { + let data = Data(bytes.utf8) + pendingFrame = MobileTerminalLaneOutputFrame( + kind: .replay, + retainedBaseSequence: sequence, + sequence: sequence, + currentSequence: sequence + UInt64(data.count), + bytes: data + ) + } + + func receiveOutput() async -> MobileTerminalLaneOutputFrame? { + if let frame = pendingFrame { + pendingFrame = nil + return frame + } + if closedCount > 0 { return nil } + return await withCheckedContinuation { waiter = $0 } + } + + func sendInput(_ input: String) { sentInputs.append(input) } + + func close() { + closedCount += 1 + waiter?.resume(returning: nil) + waiter = nil + } + + func closeCount() -> Int { closedCount } + func inputs() -> [String] { sentInputs } +} + +private actor ReplayBarrierTestLaneProvider { + private let lanes: [ReplayBarrierTestLane] + private var requestedCursors: [UInt64?] = [] + + init(lanes: [ReplayBarrierTestLane]) { self.lanes = lanes } + + func open(cursor: UInt64?) throws -> any MobileTerminalLaneConnection { + let index = requestedCursors.count + requestedCursors.append(cursor) + guard index < lanes.count else { throw CancellationError() } + return lanes[index] + } + + func cursors() -> [UInt64?] { requestedCursors } +} diff --git a/Packages/iOS/CmuxMobileShell/Tests/CmuxMobileShellTests/TerminalViewportResyncTests.swift b/Packages/iOS/CmuxMobileShell/Tests/CmuxMobileShellTests/TerminalViewportResyncTests.swift index 66e322a04d98..68854dcf32bd 100644 --- a/Packages/iOS/CmuxMobileShell/Tests/CmuxMobileShellTests/TerminalViewportResyncTests.swift +++ b/Packages/iOS/CmuxMobileShell/Tests/CmuxMobileShellTests/TerminalViewportResyncTests.swift @@ -755,6 +755,192 @@ import Testing #expect(String(data: liveChunk.data, encoding: .utf8) == "live-after-deferred-replay") } +@MainActor +@Test func terminalViewportAcknowledgementCompletesDeferredColdAttachReplay() 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" + + // Establish the shared grid before a sink is mounted. The raced + // acknowledgement below will therefore report the same effective grid, + // which is the path where the deferred cold replay must be fulfilled. + let baselineGrid = await store.updateTerminalViewport( + surfaceID: surfaceID, + columns: 80, + rows: 48 + ) + #expect(baselineGrid?.columns == 80) + #expect(baselineGrid?.rows == 48) + + await router.holdViewportRequest(number: 2) + await router.enqueueReplayTexts(["deferred-cold-replay"]) + let preparation = try #require(store.prepareTerminalViewport( + surfaceID: surfaceID, + columns: 80, + rows: 48 + )) + let viewportTask = Task { + await store.updatePreparedTerminalViewport(preparation) + } + let viewportRequested = await router.waitForCount( + of: "mobile.terminal.viewport", + atLeast: 2 + ) + #expect(viewportRequested) + guard viewportRequested else { + await router.releaseAllHeld() + _ = await viewportTask.value + return + } + + // The output sink can mount while its first viewport acknowledgement is + // in flight. Registration defers the cold replay to that acknowledgement, + // but the successful response still has to fulfill the deferred request. + var iterator = store.terminalOutputStream(surfaceID: surfaceID).makeAsyncIterator() + let replayBeforeAcknowledgement = await router.waitForCount( + of: "mobile.terminal.replay", + atLeast: 1, + timeoutNanoseconds: 200_000_000, + recordIssueOnTimeout: false + ) + #expect( + !replayBeforeAcknowledgement, + "mounting during viewport preparation must defer the cold replay until the acknowledgement" + ) + + await router.releaseAllHeld() + let acknowledgedGrid = await viewportTask.value + #expect(acknowledgedGrid?.columns == 80) + #expect(acknowledgedGrid?.rows == 48) + + let replayRequested = await router.waitForCount( + of: "mobile.terminal.replay", + atLeast: 1 + ) + #expect( + replayRequested, + "a successful viewport acknowledgement must fulfill a deferred cold replay" + ) + guard replayRequested else { return } + let replayChunk = try #require(await iterator.next()) + #expect(String(data: replayChunk.data, encoding: .utf8) == "deferred-cold-replay") + store.terminalOutputDidProcess(surfaceID: surfaceID, streamToken: replayChunk.streamToken) +} + +@MainActor +@Test func terminalViewportSupersededPreparationPreservesDeferredColdAttachReplay() 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" + + // Establish the shared grid before the sink mounts. Both preparations + // below report that same grid, so only the deferred cold replay can cause + // the replay request. + let baselineGrid = await store.updateTerminalViewport( + surfaceID: surfaceID, + columns: 80, + rows: 48 + ) + #expect(baselineGrid?.columns == 80) + #expect(baselineGrid?.rows == 48) + + await router.holdViewportRequest(number: 2) + await router.enqueueReplayTexts(["superseded-deferred-cold-replay"]) + let firstPreparation = try #require(store.prepareTerminalViewport( + surfaceID: surfaceID, + columns: 80, + rows: 48 + )) + let firstViewportTask = Task { + await store.updatePreparedTerminalViewport(firstPreparation) + } + #expect(await router.waitForCount(of: "mobile.terminal.viewport", atLeast: 2)) + + // Mount while the first acknowledgement is parked, then supersede that + // preparation with another same-size report. The second acknowledgement + // must inherit the first registration's deferred replay. + var iterator = store.terminalOutputStream(surfaceID: surfaceID).makeAsyncIterator() + let replayBeforeAcknowledgement = await router.waitForCount( + of: "mobile.terminal.replay", + atLeast: 1, + timeoutNanoseconds: 200_000_000, + recordIssueOnTimeout: false + ) + #expect(!replayBeforeAcknowledgement) + + let secondPreparation = try #require(store.prepareTerminalViewport( + surfaceID: surfaceID, + columns: 80, + rows: 48 + )) + let secondViewportTask = Task { + await store.updatePreparedTerminalViewport(secondPreparation) + } + #expect(await router.waitForCount(of: "mobile.terminal.viewport", atLeast: 3)) + + let acknowledgedGrid = await secondViewportTask.value + #expect(acknowledgedGrid?.columns == 80) + #expect(acknowledgedGrid?.rows == 48) + let replayRequested = await router.waitForCount( + of: "mobile.terminal.replay", + atLeast: 1 + ) + #expect( + replayRequested, + "a superseding same-size acknowledgement must fulfill the deferred cold replay" + ) + guard replayRequested else { + await router.releaseAllHeld() + _ = await firstViewportTask.value + return + } + let replayChunk = try #require(await iterator.next()) + #expect(String(data: replayChunk.data, encoding: .utf8) == "superseded-deferred-cold-replay") + store.terminalOutputDidProcess(surfaceID: surfaceID, streamToken: replayChunk.streamToken) + + await router.releaseAllHeld() + _ = await firstViewportTask.value +} + +@MainActor +@Test func terminalViewportClearDropsDeferredColdAttachReplay() 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 store.updateTerminalViewport(surfaceID: surfaceID, columns: 80, rows: 48) + await router.holdViewportRequest(number: 2) + let preparation = try #require(store.prepareTerminalViewport( + surfaceID: surfaceID, + columns: 80, + rows: 48 + )) + let viewportTask = Task { + await store.updatePreparedTerminalViewport(preparation) + } + #expect(await router.waitForCount(of: "mobile.terminal.viewport", atLeast: 2)) + + _ = store.terminalOutputStream(surfaceID: surfaceID) + #expect( + store.terminalViewportDeferredColdReplayGenerationsBySequenceKey.isEmpty == false, + "a sink mounted during preparation must leave a deferred replay marker" + ) + store.clearTerminalViewport(surfaceID: surfaceID) + #expect( + store.terminalViewportDeferredColdReplayGenerationsBySequenceKey.isEmpty, + "clearing a terminal must discard its deferred cold replay marker" + ) + + await router.releaseAllHeld() + _ = await viewportTask.value +} + @MainActor @Test func terminalPipelineResetDuringViewportAckDefersToAckReplay() async throws { let router = LivenessHostRouter() diff --git a/Packages/iOS/CmuxMobileShellUI/Sources/CmuxMobileShellUI/GhosttySurfaceRepresentable.swift b/Packages/iOS/CmuxMobileShellUI/Sources/CmuxMobileShellUI/GhosttySurfaceRepresentable.swift index f25b81227e90..a608bf53ff92 100644 --- a/Packages/iOS/CmuxMobileShellUI/Sources/CmuxMobileShellUI/GhosttySurfaceRepresentable.swift +++ b/Packages/iOS/CmuxMobileShellUI/Sources/CmuxMobileShellUI/GhosttySurfaceRepresentable.swift @@ -546,6 +546,20 @@ struct GhosttySurfaceRepresentable: UIViewRepresentable { // would strand the sink until a later UIKit callback. guard self.surfaceView === surfaceView else { return } self.armOutputConsumerStabilityReset(generation: taskGeneration) + if let frame = chunk.sourceRenderGridFrame, + store.usesHybridTerminalOutput, + !frame.full, + frame.activeScreen == .primary { + // Hybrid uses partial render-grid primary frames as + // advisory state only. Full frames still apply so a + // transition back from the alternate screen cannot + // leave the byte lane showing stale TUI content. + store.terminalOutputDidProcess( + surfaceID: surfaceID, + streamToken: chunk.streamToken + ) + continue + } #if DEBUG let latencySequence = chunk.sourceRenderGridFrame?.stateSeq ?? chunk.endSequence @@ -621,37 +635,49 @@ struct GhosttySurfaceRepresentable: UIViewRepresentable { case .legacy: break } - switch chunk.viewportPolicy { - case .natural: - self.activeViewportPolicy = .natural - if chunk.data.isEmpty { - surfaceView.useNaturalViewSize() - } else { - let applied = await surfaceView.useNaturalViewSizeAndWait() - guard applied else { - store.terminalOutputDidReset( - surfaceID: surfaceID, - streamToken: chunk.streamToken - ) - continue + let directPrimaryDelta = chunk.sourceRenderGridFrame.map { + !$0.full && $0.anchor == .screen && $0.activeScreen == .primary + } ?? false + // Screen-anchored primary deltas are relative to the + // established replay grid. Their advisory `.natural` + // policy describes the phone's preferred layout, but + // applying it before every delta can reflow one extra row + // while the host is still emitting the established grid. + // Keep the baseline until a full frame or an explicit + // viewport report reconciles it. + if !directPrimaryDelta { + switch chunk.viewportPolicy { + case .natural: + self.activeViewportPolicy = .natural + if chunk.data.isEmpty { + surfaceView.useNaturalViewSize() + } else { + 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 { - let applied = await surfaceView.applyViewSizeAndWait(cols: columns, rows: rows) - 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 { + let applied = await surfaceView.applyViewSizeAndWait(cols: columns, rows: rows) + guard applied else { + store.terminalOutputDidReset( + surfaceID: surfaceID, + streamToken: chunk.streamToken + ) + continue + } } + case nil: + break } - case nil: - break } if self.shouldResetForConfigThemeMismatch(chunk, store: store) { store.terminalOutputDidReset( @@ -677,8 +703,7 @@ struct GhosttySurfaceRepresentable: UIViewRepresentable { // alternate-screen frames retain the exact // dimension check used by replay safety. requiresSurfaceDimensionCheck: !( - store.usesScreenAnchoredRenderGrid - && !$0.full + !$0.full && $0.anchor == .screen && $0.activeScreen == .primary ) diff --git a/Packages/iOS/CmuxMobileShellUI/Sources/CmuxMobileShellUI/TerminalArtifactChipCountState.swift b/Packages/iOS/CmuxMobileShellUI/Sources/CmuxMobileShellUI/TerminalArtifactChipCountState.swift index d7b0abff1279..7207014c02be 100644 --- a/Packages/iOS/CmuxMobileShellUI/Sources/CmuxMobileShellUI/TerminalArtifactChipCountState.swift +++ b/Packages/iOS/CmuxMobileShellUI/Sources/CmuxMobileShellUI/TerminalArtifactChipCountState.swift @@ -55,6 +55,15 @@ struct TerminalArtifactChipCountState: Sendable { private var stateGeneration: UInt64 = 0 private var inFlight: Request? private var trailing: Pending? + /// A count-only scan is keyed by the visible count that caused it. The + /// viewport snapshot can change on every render-grid frame while the + /// artifact count stays the same; issuing another session RPC for those + /// snapshots needlessly competes with terminal input and output traffic. + private var lastRequestedLocalCount: Int? + /// Refresh the authoritative count after a bounded run of unchanged + /// render-grid generations so an artifact created outside the viewport is + /// eventually reflected without reopening a scan for every frame. + private var lastRequestedSurfaceGeneration: UInt64? private var consecutiveRearmCount = 0 /// Last successful gallery total (or positive legacy session total), held /// across transient scan failures so @@ -68,11 +77,14 @@ struct TerminalArtifactChipCountState: Sendable { private var lastAuthoritativeSessionID: String? static let maxConsecutiveRearms = 3 + static let maxDedupeSurfaceGenerationGap: UInt64 = 120 mutating func reset() { stateGeneration &+= 1 inFlight = nil trailing = nil + lastRequestedLocalCount = nil + lastRequestedSurfaceGeneration = nil consecutiveRearmCount = 0 lastAuthoritativeTotal = nil lastAuthoritativeSessionID = nil @@ -92,6 +104,31 @@ struct TerminalArtifactChipCountState: Sendable { surfaceGeneration: surfaceGeneration ) let pending = Pending(surfaceGeneration: surfaceGeneration, localCount: localCount) + let isWithinDedupeWindow: Bool + if let requestedGeneration = lastRequestedSurfaceGeneration, + surfaceGeneration >= requestedGeneration { + isWithinDedupeWindow = + surfaceGeneration - requestedGeneration < Self.maxDedupeSurfaceGenerationGap + } else { + isWithinDedupeWindow = false + } + if lastRequestedLocalCount == localCount, isWithinDedupeWindow { + if inFlight != nil, trailing?.localCount == localCount { + // Keep a queued count tied to the freshest render-grid + // generation. The original request may settle after several + // frames, and an old generation would otherwise discard the + // follow-up as stale. + trailing = pending + } + // Keep the chip's local observation current, but do not re-open + // the count RPC until the visible count changes. A matching + // trailing request still represents a count that has not been + // scanned yet, so repeated render-grid observations must retain + // it until the in-flight request completes. + return .provisionalReport(provisional) + } + lastRequestedLocalCount = localCount + lastRequestedSurfaceGeneration = surfaceGeneration guard inFlight == nil else { trailing = pending return .provisionalReport(provisional) @@ -174,13 +211,43 @@ struct TerminalArtifactChipCountState: Sendable { if let trailing { self.trailing = nil - if trailing.surfaceGeneration == currentSurfaceGeneration { - let nextRequest = makeRequest(trailing) + // A queued observation was captured after the completed request, + // but the caller's current render-grid snapshot can lag that + // observation. Preserve the queued refresh whenever it is still + // newer than the completed request, and pin it to the freshest + // generation known by either side instead of dropping it on an + // exact-generation mismatch. + let trailingGenerationIsNewer = + trailing.surfaceGeneration > request.surfaceGeneration + let trailingCountChanged = trailing.localCount != request.localCount + if trailing.surfaceGeneration >= request.surfaceGeneration, + trailingGenerationIsNewer || trailingCountChanged { + let nextRequest = makeRequest(Pending( + surfaceGeneration: max( + trailing.surfaceGeneration, + currentSurfaceGeneration + ), + localCount: trailing.localCount + )) inFlight = nextRequest + // The promoted request is now the request that dedupe must + // compare against. A trailing observation may have a newer + // generation than the last trigger marker, so leaving the + // markers behind can reopen or suppress the wrong scan. + lastRequestedLocalCount = nextRequest.localCount + lastRequestedSurfaceGeneration = nextRequest.surfaceGeneration return Completion(outcome: outcome, nextRequest: nextRequest) } } + if !scanSucceeded { + // A failed request must not permanently claim this visible count. + // Transient relay or RPC errors commonly leave the viewport + // unchanged, including after a queued count briefly changed and + // returned to the original value. + lastRequestedLocalCount = nil + } + guard outcome == .droppedForSurfaceGenerationMismatch, consecutiveRearmCount < Self.maxConsecutiveRearms else { return Completion(outcome: outcome, nextRequest: nil) @@ -191,6 +258,8 @@ struct TerminalArtifactChipCountState: Sendable { localCount: freshestLocalCount )) inFlight = nextRequest + lastRequestedLocalCount = nextRequest.localCount + lastRequestedSurfaceGeneration = nextRequest.surfaceGeneration return Completion(outcome: outcome, nextRequest: nextRequest) } diff --git a/Packages/iOS/CmuxMobileShellUI/Tests/CmuxMobileShellUITests/TerminalArtifactChipCountStateTests.swift b/Packages/iOS/CmuxMobileShellUI/Tests/CmuxMobileShellUITests/TerminalArtifactChipCountStateTests.swift index 9a1e5c6ff5bb..43cfea733883 100644 --- a/Packages/iOS/CmuxMobileShellUI/Tests/CmuxMobileShellUITests/TerminalArtifactChipCountStateTests.swift +++ b/Packages/iOS/CmuxMobileShellUI/Tests/CmuxMobileShellUITests/TerminalArtifactChipCountStateTests.swift @@ -190,6 +190,211 @@ struct TerminalArtifactChipCountStateTests { #expect(upgraded == .init(count: 12, surfaceGeneration: 5)) } + @Test("unchanged visible counts do not reissue the session scan") + func unchangedVisibleCountDoesNotRearmScan() throws { + var state = TerminalArtifactChipCountState() + let first = try request(from: state.trigger( + localCount: 0, + surfaceGeneration: 1, + supportsSessionCount: true + )) + #expect(state.trigger( + localCount: 0, + surfaceGeneration: 2, + supportsSessionCount: true + ) == .provisionalReport(.init(count: 0, surfaceGeneration: 2))) + let completion = state.complete( + first, + galleryRowTotal: 0, + sessionTotal: 0, + currentSurfaceGeneration: 1, + freshestLocalCount: 0 + ) + #expect(completion.outcome == .reported(.init(count: 0, surfaceGeneration: 1))) + #expect(completion.nextRequest == nil) + + // Render-grid snapshots can advance their surface generation without + // changing the visible artifact count. That observation only updates + // the chip and must not open another RPC. + #expect(state.trigger( + localCount: 0, + surfaceGeneration: 2, + supportsSessionCount: true + ) == .provisionalReport(.init(count: 0, surfaceGeneration: 2))) + } + + @Test("an unchanged visible count refreshes after the dedupe window") + func unchangedVisibleCountRefreshesAfterDedupeWindow() throws { + var state = TerminalArtifactChipCountState() + let first = try request(from: state.trigger( + localCount: 0, + surfaceGeneration: 1, + supportsSessionCount: true + )) + _ = state.complete( + first, + sessionTotal: 0, + currentSurfaceGeneration: 1, + freshestLocalCount: 0 + ) + + let refreshGeneration = TerminalArtifactChipCountState.maxDedupeSurfaceGenerationGap + 1 + guard case .reportAndRequest(_, let refresh) = state.trigger( + localCount: 0, + surfaceGeneration: refreshGeneration, + supportsSessionCount: true + ) else { + Issue.record("Expected the unchanged count to refresh after the dedupe window") + throw UnexpectedAction() + } + #expect(refresh.localCount == 0) + #expect(refresh.surfaceGeneration == refreshGeneration) + } + + @Test("a repeated queued count keeps its follow-up scan") + func repeatedQueuedCountKeepsTrailingScan() throws { + var state = TerminalArtifactChipCountState() + let first = try request(from: state.trigger( + localCount: 1, + surfaceGeneration: 1, + supportsSessionCount: true + )) + #expect(state.trigger( + localCount: 2, + surfaceGeneration: 1, + supportsSessionCount: true + ) == .provisionalReport(.init(count: 2, surfaceGeneration: 1))) + + // Render-grid can report the same changed count repeatedly while the + // first session scan is in flight. The queued follow-up must survive + // those observations so the post-change count is eventually scanned. + #expect(state.trigger( + localCount: 2, + surfaceGeneration: 2, + supportsSessionCount: true + ) == .provisionalReport(.init(count: 2, surfaceGeneration: 2))) + + let completion = state.complete( + first, + sessionTotal: 1, + currentSurfaceGeneration: 1, + freshestLocalCount: 2 + ) + guard let nextRequest = completion.nextRequest else { + Issue.record("Expected the changed count to remain queued") + throw UnexpectedAction() + } + #expect(nextRequest.localCount == 2) + #expect(nextRequest.surfaceGeneration == 2) + } + + @Test("a promoted trailing scan owns the dedupe generation") + func promotedTrailingScanUpdatesDedupeMarkers() throws { + var state = TerminalArtifactChipCountState() + let first = try request(from: state.trigger( + localCount: 1, + surfaceGeneration: 1, + supportsSessionCount: true + )) + #expect(state.trigger( + localCount: 2, + surfaceGeneration: 2, + supportsSessionCount: true + ) == .provisionalReport(.init(count: 2, surfaceGeneration: 2))) + #expect(state.trigger( + localCount: 2, + surfaceGeneration: 3, + supportsSessionCount: true + ) == .provisionalReport(.init(count: 2, surfaceGeneration: 3))) + + let completion = state.complete( + first, + sessionTotal: 1, + currentSurfaceGeneration: 1, + freshestLocalCount: 2 + ) + let promoted = try #require(completion.nextRequest) + #expect(promoted.localCount == 2) + #expect(promoted.surfaceGeneration == 3) + + // Generation 122 is just outside the old trigger marker (2), but it + // remains inside the dedupe window of the promoted request (3). + guard case .provisionalReport(let dedupedReport) = state.trigger( + localCount: 2, + surfaceGeneration: 122, + supportsSessionCount: true + ) else { + Issue.record("Expected the promoted request to own the dedupe window") + throw UnexpectedAction() + } + #expect(dedupedReport.surfaceGeneration == 122) + } + + @Test("a failed scan can retry at the same visible count") + func failedScanAllowsSameCountRetry() throws { + var state = TerminalArtifactChipCountState() + let first = try request(from: state.trigger( + localCount: 2, + surfaceGeneration: 1, + supportsSessionCount: true + )) + #expect(state.complete( + first, + sessionTotal: nil, + scanSucceeded: false, + currentSurfaceGeneration: 1, + freshestLocalCount: 2 + ).outcome == .reported(.init(count: 2, surfaceGeneration: 1))) + + guard case .reportAndRequest(_, let retry) = state.trigger( + localCount: 2, + surfaceGeneration: 1, + supportsSessionCount: true + ) else { + Issue.record("Expected the failed scan to be retried") + throw UnexpectedAction() + } + #expect(retry.localCount == 2) + } + + @Test("a failed scan retries after a queued count returns to the original") + func failedScanRetriesAfterQueuedCountRoundTrip() throws { + var state = TerminalArtifactChipCountState() + let first = try request(from: state.trigger( + localCount: 1, + surfaceGeneration: 1, + supportsSessionCount: true + )) + #expect(state.trigger( + localCount: 2, + surfaceGeneration: 1, + supportsSessionCount: true + ) == .provisionalReport(.init(count: 2, surfaceGeneration: 1))) + #expect(state.trigger( + localCount: 1, + surfaceGeneration: 1, + supportsSessionCount: true + ) == .provisionalReport(.init(count: 1, surfaceGeneration: 1))) + + _ = state.complete( + first, + sessionTotal: nil, + scanSucceeded: false, + currentSurfaceGeneration: 1, + freshestLocalCount: 1 + ) + + guard case .reportAndRequest(_, let retry) = state.trigger( + localCount: 1, + surfaceGeneration: 1, + supportsSessionCount: true + ) else { + Issue.record("Expected the failed round-trip scan to be retried") + throw UnexpectedAction() + } + #expect(retry.localCount == 1) + } + @Test("a failed scan with fresh local evidence drops a held authoritative zero") func failedScanDropsHeldZero() throws { var state = TerminalArtifactChipCountState() diff --git a/Packages/iOS/CmuxMobileTerminal/Sources/CmuxMobileTerminal/GhosttySurfaceView.swift b/Packages/iOS/CmuxMobileTerminal/Sources/CmuxMobileTerminal/GhosttySurfaceView.swift index acb3aa7fcde7..010c9ff74666 100644 --- a/Packages/iOS/CmuxMobileTerminal/Sources/CmuxMobileTerminal/GhosttySurfaceView.swift +++ b/Packages/iOS/CmuxMobileTerminal/Sources/CmuxMobileTerminal/GhosttySurfaceView.swift @@ -5092,7 +5092,6 @@ public final class GhosttySurfaceView: UIView, TerminalSurfaceHosting { let bottomInsetPx = UInt32(max(0, Int((bottomInsetPts * scale).rounded(.down)))) let appliedBottomInsetPts = CGFloat(bottomInsetPx) / scale let eff = effectiveGrid - let requiresExactEffectiveGrid = verifiedReplayRenderSuppressed let pushContentScale = abs(lastAppliedContentScale - scale) > 0.001 if pushContentScale { lastAppliedContentScale = scale } let generation = surfaceGeneration @@ -5117,13 +5116,19 @@ public final class GhosttySurfaceView: UIView, TerminalSurfaceHosting { var pinnedSize: CGSize? if let eff, eff.cols > 0, eff.rows > 0, cell.width > 0, cell.height > 0 { let fillsNaturalGrid = eff.cols >= Int(measured.columns) && eff.rows >= Int(measured.rows) - let withinOneCell = (Int(measured.columns) - eff.cols) <= 1 && (Int(measured.rows) - eff.rows) <= 1 let exactGridFitsInsideNatural = eff.cols <= Int(measured.columns) && eff.rows <= Int(measured.rows) let pinnedW = CGFloat(eff.cols) * cell.width / scale let pinnedH = CGFloat(eff.rows) * cell.height / scale + // The producer's effective grid is the contract for every + // authoritative render-grid replay. Even a one-row/column + // difference must be fitted locally, otherwise the apply + // fence rejects every replay and the lane keeps reopening + // behind a fresh recovery cycle. Keep the fit bounded to + // grids that actually fit inside the measured surface; a + // larger effective grid still needs a normal geometry pass. let shouldFitEffectiveGrid = !fillsNaturalGrid - && (!withinOneCell || requiresExactEffectiveGrid && exactGridFitsInsideNatural) + && exactGridFitsInsideNatural if shouldFitEffectiveGrid, pinnedW + 0.5 < containerW || pinnedH + 0.5 < containerH { let fitted = Self.fitSurfaceToGrid(surface, cols: eff.cols, rows: eff.rows, cellPixelSize: cell) diff --git a/Sources/Mobile/MobileHostIrohApplicationLaneRouter.swift b/Sources/Mobile/MobileHostIrohApplicationLaneRouter.swift index a8f8258ae4d6..48450a222979 100644 --- a/Sources/Mobile/MobileHostIrohApplicationLaneRouter.swift +++ b/Sources/Mobile/MobileHostIrohApplicationLaneRouter.swift @@ -605,7 +605,7 @@ actor MobileHostIrohApplicationLaneRouter { ) async { let laneClass: MobileHostIrohApplicationLaneQuota.LaneClass switch lane { - case .terminal: + case .terminal, .terminalInput: laneClass = .terminal case .artifact: laneClass = .artifact @@ -631,6 +631,11 @@ actor MobileHostIrohApplicationLaneRouter { cursor: cursor, stream: stream ) + case let .terminalInput(resourceID): + await Self.handleTerminalInputLane( + resourceID: resourceID, + stream: stream + ) case let .artifact(resourceID, offset): let didTakeOwnership = await artifactHandler.handleArtifactLane( resourceID: resourceID, @@ -702,6 +707,46 @@ actor MobileHostIrohApplicationLaneRouter { await stream.receiveStream.stop(errorCode: 0) } + /// Serves render-grid input without opening a second byte-output stream. + /// The empty replay envelope establishes readiness and the input half then + /// stays open for fire-and-forget length-prefixed frames. + private nonisolated static func handleTerminalInputLane( + resourceID: CmxIrohResourceID, + stream: CmxIrohBidirectionalStream + ) async { + guard let surfaceID = terminalSurfaceID(resourceID), + await MainActor.run(body: { + GhosttyApp.terminalSurfaceRegistry.terminalSurface(id: surfaceID) != nil + }) else { + await reject(stream, errorCode: ErrorCode.unsupportedResource) + return + } + do { + let currentSequence = await MainActor.run { + MobileTerminalByteTee.shared.replayState(surfaceID: surfaceID)?.seq ?? 0 + } + let baseline = try CmxIrohTerminalOutputEnvelope( + kind: .replay, + retainedBaseSequence: currentSequence, + sequence: currentSequence, + currentSequence: currentSequence, + payload: Data() + ) + try await stream.sendStream.send( + CmxIrohTerminalOutputEnvelopeCodec().encode(baseline) + ) + _ = await receiveTerminalInput( + surfaceID: surfaceID, + stream: stream + ) + } catch is CancellationError { + await stream.sendStream.reset(errorCode: 0) + } catch { + await reject(stream, errorCode: ErrorCode.invalidInput) + } + await stream.receiveStream.stop(errorCode: 0) + } + /// Returns `true` when the complete lane should close. A clean input-side /// finish returns false because the client may intentionally retain an /// output-only terminal stream. @@ -722,7 +767,10 @@ actor MobileHostIrohApplicationLaneRouter { return true } for input in try decodeTerminalInputFrames(from: &buffer) { - guard await sendTerminalInput(input, surfaceID: surfaceID) else { + guard await sendTerminalInput( + input, + surfaceID: surfaceID + ) else { await reject(stream, errorCode: ErrorCode.invalidInput) return true } @@ -845,7 +893,10 @@ actor MobileHostIrohApplicationLaneRouter { } switch surface.sendInputResult(input) { case .sent: - surface.forceRefresh(reason: "mobileHost.irohTerminalLaneInput") + // PTY output is observed by MobileTerminalByteTee, which + // schedules the normal render tick. A refresh here would + // emit a duplicate full frame before the echo and make every + // key compete with the output lane's replay fence. return true case .queued: return true diff --git a/Sources/Mobile/MobileHostIrxRuntime.swift b/Sources/Mobile/MobileHostIrxRuntime.swift index fe9e8c406891..5be8519208bc 100644 --- a/Sources/Mobile/MobileHostIrxRuntime.swift +++ b/Sources/Mobile/MobileHostIrxRuntime.swift @@ -813,7 +813,7 @@ final class MobileHostIrxRuntime { artifactRegistry: MobileHostIrohArtifactTransferRegistry, journal: IrxJournal ) async { - var terminalLaneCount = 0 + let terminalLaneQuota = MobileHostIrxTerminalLaneQuota() while !Task.isCancelled { guard let lane = await irx.acceptLane() else { return } journal.record( @@ -827,12 +827,11 @@ final class MobileHostIrxRuntime { case .keepalive: _ = irx.respondKeepalive(on: lane) case .terminal: - guard terminalLaneCount < 4 else { + guard await terminalLaneQuota.reserve() else { await lane.writer.reset(errorCode: 3) await lane.reader.stop(errorCode: 3) continue } - terminalLaneCount += 1 let resource = lane.descriptor.resource ?? "" let cursor = lane.descriptor.cursor Task { @@ -842,6 +841,22 @@ final class MobileHostIrxRuntime { stream: lane.bidirectional(), journal: journal ) + await terminalLaneQuota.release() + } + case .terminalInput: + guard await terminalLaneQuota.reserve() else { + await lane.writer.reset(errorCode: 3) + await lane.reader.stop(errorCode: 3) + continue + } + let resource = lane.descriptor.resource ?? "" + Task { + await MobileHostIrxTerminalLaneServer.serveInputOnly( + resourceID: resource, + stream: lane.bidirectional(), + journal: journal + ) + await terminalLaneQuota.release() } case .artifact: guard let resource = try? CmxIrohResourceID(lane.descriptor.resource ?? "") @@ -892,6 +907,25 @@ final class MobileHostIrxRuntime { } } +/// Tracks active IRX terminal lanes rather than cumulative opens. Replay +/// barriers intentionally close and reopen lanes, so a connection must return +/// its credit when a serving task finishes or the fast input lane eventually +/// becomes permanently unavailable after four reopen cycles. +private actor MobileHostIrxTerminalLaneQuota { + private static let maximum = 4 + private var activeCount = 0 + + func reserve() -> Bool { + guard activeCount < Self.maximum else { return false } + activeCount += 1 + return true + } + + func release() { + activeCount = max(0, activeCount - 1) + } +} + /// Synchronous trust-snapshot reader for the admission path (no actor hop, /// no network): reads the JSON the broker service persists. The caller passes /// the per-bundle, per-broker state directory computed at activation so diff --git a/Sources/Mobile/MobileHostIrxTerminalLaneServer.swift b/Sources/Mobile/MobileHostIrxTerminalLaneServer.swift index 4aef59d0de64..0ae369226aef 100644 --- a/Sources/Mobile/MobileHostIrxTerminalLaneServer.swift +++ b/Sources/Mobile/MobileHostIrxTerminalLaneServer.swift @@ -58,6 +58,49 @@ enum MobileHostIrxTerminalLaneServer { journal.record("host-terminal", "lane-closed", ["surface": surfaceID.uuidString]) } + /// Serves render-grid input without opening a second byte-output stream. + /// The empty replay envelope establishes readiness and the input half then + /// stays open for fire-and-forget length-prefixed frames. + static func serveInputOnly( + resourceID: String, + stream: CmxIrohBidirectionalStream, + journal: IrxJournal + ) async { + guard let surfaceID = terminalSurfaceID(resourceID), + await MainActor.run(body: { + GhosttyApp.terminalSurfaceRegistry.terminalSurface(id: surfaceID) != nil + }) + else { + await reject(stream, errorCode: ErrorCode.unsupportedResource) + return + } + do { + let currentSequence = await MainActor.run { + MobileTerminalByteTee.shared.replayState(surfaceID: surfaceID)?.seq ?? 0 + } + let baseline = try CmxIrohTerminalOutputEnvelope( + kind: .replay, + retainedBaseSequence: currentSequence, + sequence: currentSequence, + currentSequence: currentSequence, + payload: Data() + ) + try await stream.sendStream.send( + CmxIrohTerminalOutputEnvelopeCodec().encode(baseline) + ) + _ = await receiveInput( + surfaceID: surfaceID, + stream: stream + ) + } catch is CancellationError { + await stream.sendStream.reset(errorCode: 0) + } catch { + await reject(stream, errorCode: ErrorCode.invalidInput) + } + await stream.receiveStream.stop(errorCode: 0) + journal.record("host-terminal", "input-lane-closed", ["surface": surfaceID.uuidString]) + } + private static func sendOutput( surfaceID: UUID, cursor: UInt64?, @@ -181,7 +224,10 @@ enum MobileHostIrxTerminalLaneServer { for input in try MobileHostIrohApplicationLaneRouter .decodeTerminalInputFrames(from: &buffer) { - guard await deliverInput(input, surfaceID: surfaceID) else { + guard await deliverInput( + input, + surfaceID: surfaceID + ) else { await reject(stream, errorCode: ErrorCode.invalidInput) return true } @@ -200,7 +246,10 @@ enum MobileHostIrxTerminalLaneServer { } } - private static func deliverInput(_ input: String, surfaceID: UUID) async -> Bool { + private static func deliverInput( + _ input: String, + surfaceID: UUID + ) async -> Bool { await MainActor.run { guard let surface = GhosttyApp.terminalSurfaceRegistry.terminalSurface( @@ -208,7 +257,10 @@ enum MobileHostIrxTerminalLaneServer { else { return false } switch surface.sendInputResult(input) { case .sent: - surface.forceRefresh(reason: "mobileHost.irxTerminalLaneInput") + // PTY output is observed by MobileTerminalByteTee, which + // schedules the normal render tick. A refresh here would + // emit a duplicate full frame before the echo and make every + // key compete with the output lane's replay fence. return true case .queued: return true diff --git a/Sources/Mobile/MobileTerminalRenderObserver.swift b/Sources/Mobile/MobileTerminalRenderObserver.swift index 4e5aa7f50f6c..2b2be6006f7b 100644 --- a/Sources/Mobile/MobileTerminalRenderObserver.swift +++ b/Sources/Mobile/MobileTerminalRenderObserver.swift @@ -20,8 +20,8 @@ final class MobileTerminalRenderObserver { private var isEmitFlushScheduled = false private var renderGridStatesBySurfaceID: [UUID: [MobileTerminalRenderGridFrame.Anchor: MobileTerminalRenderGridEmissionState]] = [:] - private var terminalThemesBySurfaceID: [UUID: TerminalTheme] = [:] - private var terminalConfigThemesBySurfaceID: [UUID: TerminalTheme] = [:] + var terminalThemesBySurfaceID: [UUID: TerminalTheme] = [:] + var terminalConfigThemesBySurfaceID: [UUID: TerminalTheme] = [:] private var runtimeSurfaceGenerationsBySurfaceID: [UUID: UInt64] = [:] private var reconciledSurfaceTopologyGeneration: UInt64? private var cachedTerminalTheme: TerminalTheme = .monokai @@ -424,8 +424,9 @@ final class MobileTerminalRenderObserver { themedFrame.terminalTheme = resolvedTheme.theme themedFrame.terminalThemeRevision = resolvedTheme.revision + let previousEmissionState = renderGridStatesBySurfaceID[surfaceID]?[anchor] guard let emission = try? themedFrame.renderGridEmission( - comparedTo: renderGridStatesBySurfaceID[surfaceID]?[anchor], + comparedTo: previousEmissionState, fullScrollbackTarget: fullScrollbackTarget, allowScrollbackRequest: allowScrollbackRequest ) else { return nil } @@ -464,6 +465,12 @@ final class MobileTerminalRenderObserver { func adoptReplayBaseline(_ frame: MobileTerminalRenderGridFrame, surfaceID: UUID) { guard frame.anchor == .screen else { return } renderGridStatesBySurfaceID[surfaceID, default: [:]][.screen] = frame.emissionState + // Replay uses the same decorated theme/config state as the phone. Keep + // the resolver caches in that state as well, otherwise the next live + // capture resolves missing theme fields from a different source and + // promotes the delta to a full frame. + terminalThemesBySurfaceID[surfaceID] = frame.terminalTheme + terminalConfigThemesBySurfaceID[surfaceID] = frame.terminalConfigTheme } func decorateReplayFrame(_ frame: MobileTerminalRenderGridFrame) -> MobileTerminalRenderGridFrame { diff --git a/Sources/TerminalController+MobileScrollPrefetch.swift b/Sources/TerminalController+MobileScrollPrefetch.swift index e53aacfb8f50..c5ee84c4f434 100644 --- a/Sources/TerminalController+MobileScrollPrefetch.swift +++ b/Sources/TerminalController+MobileScrollPrefetch.swift @@ -45,8 +45,13 @@ extension TerminalController { scrollbackLines: scrollbackLines, anchor: anchor )?.frame else { return nil } - MobileTerminalRenderObserver.shared.adoptReplayBaseline(frame, surfaceID: surfaceID) - return MobileTerminalRenderObserver.shared.decorateReplayFrame(frame) + // The phone applies the decorated frame, so rebase the producer cache + // with that same theme/config state. Rebasing the raw snapshot first + // makes the next event look like a theme change and promotes every + // keystroke to another verified full replay. + let decorated = MobileTerminalRenderObserver.shared.decorateReplayFrame(frame) + MobileTerminalRenderObserver.shared.adoptReplayBaseline(decorated, surfaceID: surfaceID) + return decorated } /// Captures a render grid from the canonical socket-bound runtime surface. diff --git a/Sources/TerminalController.swift b/Sources/TerminalController.swift index 6dd1505bfd88..a697f603f9fa 100644 --- a/Sources/TerminalController.swift +++ b/Sources/TerminalController.swift @@ -16005,7 +16005,12 @@ class TerminalController { let sendResult = terminalTarget.sendInputResult(text) switch sendResult { case .sent: - terminalTarget.forceRefresh(reason: "mobileHost.terminalInput") + // PTY output is already observed by MobileTerminalByteTee, which + // schedules the post-parser render-grid tick. Forcing a refresh + // here emits a full frame before the echoed bytes arrive, sending + // screen-anchored iOS clients through the verified replay path on + // every keystroke. + break case .queued: break case .inputQueueFull: diff --git a/cmuxTests/MobileHostTerminalThemeTests.swift b/cmuxTests/MobileHostTerminalThemeTests.swift index 73b734ae8419..2857ecf0fb6e 100644 --- a/cmuxTests/MobileHostTerminalThemeTests.swift +++ b/cmuxTests/MobileHostTerminalThemeTests.swift @@ -187,6 +187,40 @@ import Testing #expect(resolved?.boldColor == "#4e2a84") } + + @MainActor + @Test func replayBaselineClearsRemovedThemeCaches() throws { + let observer = MobileTerminalRenderObserver.shared + observer.stop() + defer { observer.stop() } + let surfaceID = UUID() + let seeded = try MobileTerminalRenderGridFrame( + surfaceID: surfaceID.uuidString, + stateSeq: 1, + columns: 2, + rows: 1, + rowSpans: [], + terminalTheme: .monokai, + terminalConfigTheme: .monokai, + anchor: .screen + ) + observer.adoptReplayBaseline(seeded, surfaceID: surfaceID) + #expect(observer.terminalThemesBySurfaceID[surfaceID] != nil) + #expect(observer.terminalConfigThemesBySurfaceID[surfaceID] != nil) + + let cleared = try MobileTerminalRenderGridFrame( + surfaceID: surfaceID.uuidString, + stateSeq: 2, + columns: 2, + rows: 1, + rowSpans: [], + anchor: .screen + ) + observer.adoptReplayBaseline(cleared, surfaceID: surfaceID) + + #expect(observer.terminalThemesBySurfaceID[surfaceID] == nil) + #expect(observer.terminalConfigThemesBySurfaceID[surfaceID] == nil) + } } private final class ThemeInvalidationTestClock: Clock, @unchecked Sendable { diff --git a/ios/cmux/cmuxApp.swift b/ios/cmux/cmuxApp.swift index 21022ca3d3f2..fcc6deb8bef7 100644 --- a/ios/cmux/cmuxApp.swift +++ b/ios/cmux/cmuxApp.swift @@ -139,6 +139,20 @@ struct cmuxApp: App { cursor: cursor ) }, + terminalInputLaneProvider: { request, surfaceID, _ in + guard let surfaceUUID = UUID(uuidString: surfaceID) else { + throw MobileIrohTerminalLaneError.invalidSurfaceID + } + return irxEnabled + ? try await irx.openTerminalInputLane( + for: request, + surfaceID: surfaceUUID + ) + : try await iroh.openTerminalInputLane( + for: request, + surfaceID: surfaceUUID + ) + }, artifactLaneProvider: { request, resourceID, offset in irxEnabled ? try await irx.openArtifactLane( diff --git a/ios/cmuxPackage/Sources/cmuxFeature/CMUXMobileRuntime.swift b/ios/cmuxPackage/Sources/cmuxFeature/CMUXMobileRuntime.swift index a5da120f7bf4..9f943af06e5a 100644 --- a/ios/cmuxPackage/Sources/cmuxFeature/CMUXMobileRuntime.swift +++ b/ios/cmuxPackage/Sources/cmuxFeature/CMUXMobileRuntime.swift @@ -35,6 +35,7 @@ public struct CMUXMobileRuntime: Sendable, MobileSyncRuntime { public var supportsServerPushEvents: Bool public var independentEventByteStreamProvider: CmxIndependentEventByteStreamProvider? public var terminalLaneProvider: MobileTerminalLaneProvider? + public var terminalInputLaneProvider: MobileTerminalLaneProvider? public var simulatorStreamLaneProvider: MobileSimulatorStreamLaneProvider? public var artifactLaneProvider: MobileArtifactLaneProvider? @@ -146,6 +147,7 @@ public struct CMUXMobileRuntime: Sendable, MobileSyncRuntime { supportsServerPushEvents: Bool = true, independentEventByteStreamProvider: CmxIndependentEventByteStreamProvider? = nil, terminalLaneProvider: MobileTerminalLaneProvider? = nil, + terminalInputLaneProvider: MobileTerminalLaneProvider? = nil, artifactLaneProvider: MobileArtifactLaneProvider? = nil, simulatorStreamLaneProvider: MobileSimulatorStreamLaneProvider? = nil ) { @@ -161,6 +163,7 @@ public struct CMUXMobileRuntime: Sendable, MobileSyncRuntime { self.supportsServerPushEvents = supportsServerPushEvents self.independentEventByteStreamProvider = independentEventByteStreamProvider self.terminalLaneProvider = terminalLaneProvider + self.terminalInputLaneProvider = terminalInputLaneProvider self.artifactLaneProvider = artifactLaneProvider self.simulatorStreamLaneProvider = simulatorStreamLaneProvider } @@ -177,6 +180,7 @@ public struct CMUXMobileRuntime: Sendable, MobileSyncRuntime { supportsServerPushEvents: Bool = true, independentEventByteStreamProvider: CmxIndependentEventByteStreamProvider? = nil, terminalLaneProvider: MobileTerminalLaneProvider? = nil, + terminalInputLaneProvider: MobileTerminalLaneProvider? = nil, artifactLaneProvider: MobileArtifactLaneProvider? = nil, simulatorStreamLaneProvider: MobileSimulatorStreamLaneProvider? = nil ) { @@ -191,6 +195,7 @@ public struct CMUXMobileRuntime: Sendable, MobileSyncRuntime { self.supportsServerPushEvents = supportsServerPushEvents self.independentEventByteStreamProvider = independentEventByteStreamProvider self.terminalLaneProvider = terminalLaneProvider + self.terminalInputLaneProvider = terminalInputLaneProvider self.artifactLaneProvider = artifactLaneProvider self.simulatorStreamLaneProvider = simulatorStreamLaneProvider self.now = now diff --git a/ios/cmuxPackage/Sources/cmuxFeature/MobileIrohRuntimeComposition.swift b/ios/cmuxPackage/Sources/cmuxFeature/MobileIrohRuntimeComposition.swift index e978c73c8797..f9bf469fc13e 100644 --- a/ios/cmuxPackage/Sources/cmuxFeature/MobileIrohRuntimeComposition.swift +++ b/ios/cmuxPackage/Sources/cmuxFeature/MobileIrohRuntimeComposition.swift @@ -757,6 +757,23 @@ public final class MobileIrohRuntimeComposition: return MobileIrohTerminalLane(stream: stream) } + /// Opens a terminal input-only lane. Render-grid output remains on the + /// ordered event stream, while keystrokes use an independent QUIC stream + /// whose writes do not wait for an RPC response. + public func openTerminalInputLane( + for request: CmxByteTransportRequest, + surfaceID: UUID, + priority: Int32 = 0 + ) async throws -> MobileIrohTerminalLane { + let resourceID = try CmxIrohResourceID("terminal:\(surfaceID.uuidString.lowercased())") + let stream = try await openBidirectionalLane( + for: request, + lane: .terminalInput(resourceID: resourceID), + priority: priority + ) + return MobileIrohTerminalLane(stream: stream) + } + /// Opens a simulator-stream v2 lane for one Mac simulator panel. The /// phone-to-Mac half carries start/ack/input messages, so it rides above /// terminal typing (tiny messages, interaction-critical); the Mac sets diff --git a/ios/cmuxPackage/Sources/cmuxFeature/MobileIrxRuntimeComposition.swift b/ios/cmuxPackage/Sources/cmuxFeature/MobileIrxRuntimeComposition.swift index 0dd585e19465..b94f5b0114d1 100644 --- a/ios/cmuxPackage/Sources/cmuxFeature/MobileIrxRuntimeComposition.swift +++ b/ios/cmuxPackage/Sources/cmuxFeature/MobileIrxRuntimeComposition.swift @@ -1396,6 +1396,29 @@ public actor MobileIrxRuntimeComposition { return MobileIrohTerminalLane(stream: lane.bidirectional()) } + /// Opens a terminal input-only lane. Render-grid output remains on the + /// ordered event stream, while keystrokes use an independent QUIC stream + /// whose writes do not wait for an RPC response. + public func openTerminalInputLane( + for request: CmxByteTransportRequest, + surfaceID: UUID + ) async throws -> MobileIrohTerminalLane { + let peerHex = try peerTarget(for: request) + let session = try await engine(forPeer: peerHex) + .ensureSession(trigger: "terminal-input-lane") + let lane = try await session.connection.openLane( + IrxLaneDescriptor( + lane: .terminalInput, + resource: "terminal:\(surfaceID.uuidString.lowercased())" + ) + ) + Self.journal.record( + "client-terminal-input", "lane-opened", + ["surface": surfaceID.uuidString.lowercased()] + ) + return MobileIrohTerminalLane(stream: lane.bidirectional()) + } + public func openArtifactLane( for request: CmxByteTransportRequest, resourceID: String,