diff --git a/.github/swift-file-length-budget.tsv b/.github/swift-file-length-budget.tsv index 07944f504270..a97f24753258 100644 --- a/.github/swift-file-length-budget.tsv +++ b/.github/swift-file-length-budget.tsv @@ -102,8 +102,8 @@ 905 Sources/CmuxSSHURLRequest.swift 899 Sources/Panels/MarkdownWebRenderer.swift 885 Sources/Panels/TerminalPanel.swift +948 Packages/Shared/CmuxAgentChat/Tests/CmuxAgentChatTests/ChatConversationStoreTests.swift 885 cmuxTests/SidebarWorkspaceDropPlannerTests.swift -877 Packages/Shared/CmuxAgentChat/Tests/CmuxAgentChatTests/ChatConversationStoreTests.swift 871 cmuxTests/ClaudeHookSurfaceResolutionSwiftTests.swift 868 Sources/Panels/BrowserScreenshotSnapshotter.swift 859 Packages/macOS/CmuxControlSocket/Sources/CmuxControlSocket/Coordinator/Workspace/ControlCommandCoordinator+Workspace.swift @@ -129,8 +129,8 @@ 752 cmuxUITests/CloseWorkspaceCmdDUITests.swift 739 cmuxTests/CLICodexHookTimeoutRegressionTests.swift 738 Packages/macOS/CMUXProjectModel/Sources/CMUXProjectModel/XcodeProjectAdapter.swift -724 Sources/Panels/BrowserPopupWindowController.swift -722 Packages/Shared/CmuxAgentChat/Sources/CmuxAgentChat/Store/ChatConversationStore.swift +725 Sources/Panels/BrowserPopupWindowController.swift +764 Packages/Shared/CmuxAgentChat/Sources/CmuxAgentChat/Store/ChatConversationStore.swift 718 Packages/Shared/CmuxAuthRuntime/Sources/CmuxAuthRuntime/Coordinator/AuthCoordinator.swift 716 Sources/TaskManagerSnapshot.swift 715 Sources/AppleScriptSupport.swift @@ -214,6 +214,7 @@ 526 Packages/macOS/CMUXAgentLaunch/Sources/CMUXAgentLaunch/Workstream/WorkstreamStore.swift 525 Packages/macOS/CmuxSettings/Sources/CmuxSettings/SocketControl/SocketControlSettings.swift 524 CLI/CMUXCLI+AutoNaming.swift +524 Sources/Mobile/AgentChat/AgentChatTranscriptService.swift 520 CLI/CMUXCLI+AmpExtension.swift 520 cmuxTests/MainWindowVisibilityControllerTests.swift 519 Packages/macOS/CmuxSwiftRender/Tests/CmuxSwiftRenderTests/Corpus/stress-two-column-cockpit-sidebar.swift @@ -225,11 +226,11 @@ 509 Packages/macOS/CMUXAgentLaunch/Sources/CMUXAgentLaunch/AgentLaunchSanitizerAdditionalPolicies.swift 507 Sources/TerminalControllerTopSupport.swift 506 Sources/App/MainWindowVisibilityController.swift +518 Packages/iOS/CmuxMobileShellUI/Sources/CmuxMobileShellUI/CMUXMobileRootView.swift 505 cmuxUITests/DisplayResolutionRegressionUITests.swift 504 cmuxTests/TerminalNotificationSocketActionTests.swift 503 Sources/Settings/ConfigSource.swift 502 Sources/CmuxEventPublishing.swift 502 Sources/RemoteTmuxSessionMirror.swift -500 Packages/iOS/CmuxMobileShellUI/Sources/CmuxMobileShellUI/CMUXMobileRootView.swift 500 Packages/macOS/CmuxControlSocket/Tests/CmuxControlSocketTests/ControlCommandContextTestStubs.swift 500 Sources/KeyboardShortcutRecorder.swift diff --git a/Packages/Shared/CmuxAgentChat/Sources/CmuxAgentChat/Parsing/AgentChatProseScreenExtractor.swift b/Packages/Shared/CmuxAgentChat/Sources/CmuxAgentChat/Parsing/AgentChatProseScreenExtractor.swift new file mode 100644 index 000000000000..324d48fb9c8f --- /dev/null +++ b/Packages/Shared/CmuxAgentChat/Sources/CmuxAgentChat/Parsing/AgentChatProseScreenExtractor.swift @@ -0,0 +1,352 @@ +import Foundation + +/// Extracts the agent's in-progress prose from a snapshot of the terminal's +/// rendered screen, for the live streaming preview. +/// +/// The agent CLIs paint their turn with a cursor-addressed TUI and never write +/// token-level deltas to their JSONL transcript, so the only token-grained +/// source of a streaming answer is the emulated screen grid. This extractor is +/// deliberately conservative and **best-effort**: the preview it returns is +/// always superseded by the authoritative JSONL line when the turn settles, so +/// a transient mis-extraction self-corrects within one turn. It returns `nil` +/// whenever it cannot confidently locate an actively-streaming answer, which the +/// caller treats as "show nothing" rather than guess. +/// +/// Strategy (the "spinner anchor"): while a turn is in flight the agent renders +/// a working/status line carrying an elapsed timer (`(4s · ↓ 21 tokens)`, +/// `Thinking… (esc to interrupt)`). That line sits directly below the streaming +/// answer and above the input box, so it is a stable local landmark that needs +/// no knowledge of the prompt text or the input-box format. Everything at or +/// below it is chrome; the contiguous text block immediately above it, up to the +/// previous committed block, is the in-progress answer. +public struct AgentChatProseScreenExtractor: Sendable { + /// Hard cap on how many lines above the anchor are considered, so a screen + /// with no committed-block boundary can't fold the whole scrollback into one + /// preview. + private static let maxAnswerLines = 200 + + public init() {} + + /// Extracts the current streaming answer from rendered screen rows. + /// + /// - Parameters: + /// - lines: Rendered screen rows, top to bottom (e.g. a render-grid + /// snapshot's plain rows). Trailing whitespace per row is ignored. + /// - agentKind: Selects per-agent boundary markers. + /// - Returns: The cleaned in-progress prose, or `nil` when no actively + /// streaming answer is present. + public func extract(lines: [String], agentKind: ChatAgentKind) -> String? { + let rows = lines.map { Self.trimTrailing($0) } + guard let anchor = Self.statusLineIndex(in: rows) else { return nil } + guard anchor > 0 else { return nil } + + let lowerBound = max(0, anchor - Self.maxAnswerLines) + let answerTops = Self.answerTopBullets(for: agentKind) + // Agents that bullet their live answer (Claude's "⏺ ") require that bullet + // to be reached, so the early "thinking" screen — where the spinner sits + // directly under the wrapped user prompt and no answer exists yet — yields + // nil instead of leaking the prompt's tail as a fake answer. + let requireAnswerTop = !answerTops.isEmpty + var collected: [String] = [] + var index = anchor - 1 + var foundAnswerTop = false + while index >= lowerBound { + let row = rows[index] + if let first = row.trimmingCharacters(in: .whitespaces).first, + answerTops.contains(first) { + // Inclusive top: the answer's own leading bullet. Include it + // stripped and stop — anything above belongs to an earlier block. + collected.append(Self.strippingLeadingBullet(row, agentKind: agentKind)) + foundAnswerTop = true + break + } + if Self.isBoundary(row, agentKind: agentKind) { break } + collected.append(row) + index -= 1 + } + if requireAnswerTop && !foundAnswerTop { return nil } + collected.reverse() + + // Strip a leading committed-block bullet if the answer just committed + // on screen (e.g. Claude prefixes a finalized block with "⏺ "). + if let first = collected.first { + collected[0] = Self.strippingLeadingBullet(first, agentKind: agentKind) + } + // Claude wraps the answer under a 2-space hanging indent aligned past the + // "⏺ " bullet. Drop it from continuation rows so wrapped lines read as one + // flowing paragraph rather than an indented block. + if requireAnswerTop { + for i in collected.indices where i > 0 { + collected[i] = Self.strippingHangingIndent(collected[i]) + } + } + + let cleaned = Self.collapsingBlankRuns(collected) + .joined(separator: "\n") + .trimmingCharacters(in: .whitespacesAndNewlines) + return cleaned.isEmpty ? nil : cleaned + } + + /// Removes up to two leading spaces (Claude's hanging-indent width) from a + /// wrapped answer continuation row; blank rows are returned unchanged. + static func strippingHangingIndent(_ row: String) -> String { + var working = Substring(row) + var removed = 0 + while removed < 2, working.first == " " { + working = working.dropFirst() + removed += 1 + } + return String(working) + } + + // MARK: - Anchoring + + /// The row to anchor on: the agent's working/status line that sits directly + /// above the streaming answer. + /// + /// Two tiers, because Claude 2.1 renders *two* working signals at once: the + /// spinner line with the elapsed timer (e.g. `✻ Forming… (4s · ↓ 21 tokens)`) + /// directly above the answer, and a persistent bottom mode bar that carries + /// `esc to interrupt` *below* the input box. Anchoring on the lower of the two + /// (the mode bar) would fold the input box and dividers into the preview, so + /// the timer line is strongly preferred; the interrupt hint is only a fallback + /// for layouts/agents that render no timer line (e.g. Codex). + static func statusLineIndex(in rows: [String]) -> Int? { + // Tier 1: the spinner line carrying an elapsed timer and throughput, + // directly above the answer once tokens start (`✻ … (3s · ↓ 1 tokens)`). + for index in stride(from: rows.count - 1, through: 0, by: -1) { + if isTimerStatusLine(rows[index]) { return index } + } + // Tier 2: the gerund spinner line before the timer appears + // (`✻ Nebulizing… ` during the first seconds), matched by its leading + // animated glyph + the trailing ellipsis so the post-turn `Brewed for 3s` + // summary (no ellipsis) is excluded. + for index in stride(from: rows.count - 1, through: 0, by: -1) { + if isGerundWorkingLine(rows[index]) { return index } + } + // Tier 3: an explicit interrupt hint *on the working line itself* + // (Codex's `Working (3s • Esc to interrupt)`). The persistent Claude mode + // bar also carries that phrase but sits below the input box, so the footer + // form is excluded — anchoring there would fold in the input box chrome. + for index in stride(from: rows.count - 1, through: 0, by: -1) { + if isInterruptHintLine(rows[index]), !isModeFooterLine(rows[index]) { return index } + } + return nil + } + + /// Whether a row is a status line by any signal. Retained for callers/tests + /// that ask the question without caring which tier matched. + static func isStatusLine(_ row: String) -> Bool { + isTimerStatusLine(row) || isGerundWorkingLine(row) + || (isInterruptHintLine(row) && !isModeFooterLine(row)) + } + + /// Whether a row carries an explicit interrupt hint (`esc to interrupt` / + /// `esc to cancel`). + static func isInterruptHintLine(_ row: String) -> Bool { + let lower = row.lowercased() + return lower.contains("esc to interrupt") || lower.contains("esc to cancel") + } + + /// Whether a row is the persistent bottom mode/footer bar rather than a + /// working line. Identified by its stable footer phrases. + static func isModeFooterLine(_ row: String) -> Bool { + let lower = row.lowercased() + return lower.contains("shift+tab") || lower.contains("for agents") + || lower.contains("auto mode") || lower.contains("⏵⏵") + } + + /// Whether a row is Claude's gerund spinner line before the elapsed timer + /// renders: a leading animated spinner glyph and a trailing `…`. Excludes the + /// post-turn `Brewed for Ns` summary, which carries no ellipsis. + static func isGerundWorkingLine(_ row: String) -> Bool { + let trimmed = row.trimmingCharacters(in: .whitespaces) + guard let first = trimmed.first, Self.spinnerLeadGlyphs.contains(first) else { return false } + return trimmed.contains("…") + } + + /// Whether a row is the spinner/elapsed-timer working line that sits directly + /// above the streaming answer. Matched on stable signals rather than the + /// (randomized, localized) gerund: an elapsed timer paired with a spinner + /// glyph or the token/throughput markers Claude shows alongside it, so a + /// parenthesized `(3s ...)` inside prose is not mistaken for the anchor. + /// + /// An *active* turn always pairs the timer with either a parenthesis (`(4s`) + /// or a throughput marker (`↓ 21 tokens`). The post-turn summary Claude leaves + /// on screen, `✻ Brewed for 3s`, has a bare timer and neither, so it reads as + /// settled (no anchor) rather than a still-streaming line. + static func isTimerStatusLine(_ row: String) -> Bool { + let trimmed = row.trimmingCharacters(in: .whitespaces) + guard !trimmed.isEmpty else { return false } + let lower = trimmed.lowercased() + let hasSpinner = trimmed.contains(where: { Self.spinnerGlyphs.contains($0) }) + let hasThroughput = lower.contains("token") || trimmed.contains("↓") || trimmed.contains("↑") + guard hasSpinner || hasThroughput else { return false } + if hasThroughput { + // The "running stop hooks… 0/3 · 3s · ↓ 56 tokens" form drops the + // paren around the timer, so accept the bare form when throughput + // markers confirm the turn is live. + return Self.containsElapsedTimer(trimmed) + } + return Self.containsParenthesizedTimer(trimmed) + } + + /// Whether the row contains an elapsed-time token like `(4s`, `(12s`, + /// `(1m05s`, or the bare `· 3s ·` form Claude switches to once it starts + /// running stop hooks (`(running stop hooks… 0/3 · 3s · ↓ 56 tokens)`). The + /// digit run must start at a word boundary (preceded by a non-alphanumeric) + /// and the `s` must not be followed by a letter, so neither `0/3` nor a + /// version like `2.1.191s` is mistaken for a timer. Hand-scanned to avoid a + /// regex literal, whose `/.../ ` parse is ambiguous next to division. + static func containsElapsedTimer(_ text: String) -> Bool { + let chars = Array(text) + var index = 0 + while index < chars.count { + guard chars[index].isNumber else { index += 1; continue } + // The digit run must begin at a word boundary so "0/3" or a mid-token + // digit can't anchor a false match. + if index > 0 { + let prev = chars[index - 1] + if prev.isLetter || prev.isNumber { index += 1; continue } + } + var cursor = index + while cursor < chars.count, chars[cursor].isNumber { cursor += 1 } + // optional minutes group: m + if cursor < chars.count, chars[cursor] == "m" { + let afterM = cursor + 1 + if afterM < chars.count, chars[afterM].isNumber { + cursor = afterM + while cursor < chars.count, chars[cursor].isNumber { cursor += 1 } + } + } + if cursor < chars.count, chars[cursor] == "s" { + let afterS = cursor + 1 + if afterS >= chars.count || !chars[afterS].isLetter { + return true + } + } + index = cursor + } + return false + } + + /// Whether the row contains a *parenthesized* elapsed-time token of the form + /// `(s` or `(ms`, e.g. `(4s`, `(12s`, `(1m05s`. This + /// is the stricter form used to tell a live working line (`Forming… (9s)`) + /// from the post-turn `Brewed for 3s` summary, which has a bare timer. + static func containsParenthesizedTimer(_ text: String) -> Bool { + let chars = Array(text) + var index = 0 + while index < chars.count { + guard chars[index] == "(" else { index += 1; continue } + var cursor = index + 1 + var sawDigits = false + while cursor < chars.count, chars[cursor].isNumber { cursor += 1; sawDigits = true } + if sawDigits, cursor < chars.count, chars[cursor] == "m" { + cursor += 1 + while cursor < chars.count, chars[cursor].isNumber { cursor += 1 } + } + if sawDigits, cursor < chars.count, chars[cursor] == "s" { + return true + } + index += 1 + } + return false + } + + /// Glyphs Claude/Codex cycle through for the working spinner. + private static let spinnerGlyphs: Set = [ + "✢", "✶", "✻", "✽", "✳", "·", "∗", "⟢", "✦", "✧", "◐", "◓", "◑", "◒", + ] + + /// Animated spinner glyphs that *lead* the gerund working line. Excludes "·" + /// (a mid-line separator in the mode bar and prose) so only a genuine spinner + /// at the start of a row qualifies as the gerund anchor. + private static let spinnerLeadGlyphs: Set = [ + "✢", "✶", "✻", "✽", "✳", "∗", "⟢", "✦", "✧", "◐", "◓", "◑", "◒", + ] + + // MARK: - Boundaries + + /// Whether a row marks the top boundary of the current answer: a previous + /// committed block (tool call / earlier answer) or a user-prompt line. The + /// streaming answer is the uncommitted text between the boundary and the + /// status line. + static func isBoundary(_ row: String, agentKind: ChatAgentKind) -> Bool { + guard let first = row.trimmingCharacters(in: .whitespaces).first else { + return false + } + return boundaryLeadingGlyphs(for: agentKind).contains(first) + } + + /// Leading glyphs that begin a committed block or prompt line for an agent. + /// These are *exclusive* boundaries: collection stops before the row. + static func boundaryLeadingGlyphs(for agentKind: ChatAgentKind) -> Set { + switch agentKind { + case .claude, .other: + // ● tool bullet, ⎿ tool-result continuation, ❯/> user prompt echo, + // │ prompt-box border. (⏺ is handled as an *inclusive* answer top in + // answerTopBullets, so it is not listed here.) + return ["●", "⎿", "❯", ">", "│"] + case .codex: + // Codex marks user turns with "user" headers and tool calls with + // bullets; ">" / box borders are the reliable cross-version anchors. + return ["•", "›", "❯", ">", "│", "⎿"] + } + } + + /// Leading glyphs that mark the *inclusive* top of the in-progress answer: + /// the bullet the agent prefixes onto the streaming block itself. The row is + /// kept (with the bullet stripped) and collection stops there, so an earlier + /// committed block above it is excluded. Claude prefixes the live answer with + /// `⏺ `; Codex prose carries no per-block bullet in v1. + static func answerTopBullets(for agentKind: ChatAgentKind) -> Set { + switch agentKind { + case .claude, .other: + return ["⏺"] + case .codex: + return [] + } + } + + /// Removes a leading committed-block bullet ("⏺ ", "● ", "• ") from a row. + static func strippingLeadingBullet(_ row: String, agentKind: ChatAgentKind) -> String { + var working = row + let leading = Set(["⏺", "●", "•", "›"]) + if let first = working.first, leading.contains(first) { + working.removeFirst() + if working.first == " " { working.removeFirst() } + } + return working + } + + // MARK: - Cleanup + + private static func trimTrailing(_ row: String) -> String { + var scalars = Array(row.unicodeScalars) + while let last = scalars.last, last == " " || last == "\t" { + scalars.removeLast() + } + return String(String.UnicodeScalarView(scalars)) + } + + /// Trims leading/trailing blank rows and collapses runs of 2+ blank rows to + /// a single blank, so paragraph spacing survives but TUI padding does not. + static func collapsingBlankRuns(_ rows: [String]) -> [String] { + var out: [String] = [] + var previousBlank = false + for row in rows { + let isBlank = row.trimmingCharacters(in: .whitespaces).isEmpty + if isBlank { + if previousBlank { continue } + previousBlank = true + } else { + previousBlank = false + } + out.append(row) + } + while out.first?.trimmingCharacters(in: .whitespaces).isEmpty == true { out.removeFirst() } + while out.last?.trimmingCharacters(in: .whitespaces).isEmpty == true { out.removeLast() } + return out + } +} diff --git a/Packages/Shared/CmuxAgentChat/Sources/CmuxAgentChat/Source/FixtureChatEventSource.swift b/Packages/Shared/CmuxAgentChat/Sources/CmuxAgentChat/Source/FixtureChatEventSource.swift index 403472f89e68..478f76c98cd8 100644 --- a/Packages/Shared/CmuxAgentChat/Sources/CmuxAgentChat/Source/FixtureChatEventSource.swift +++ b/Packages/Shared/CmuxAgentChat/Sources/CmuxAgentChat/Source/FixtureChatEventSource.swift @@ -147,6 +147,10 @@ public actor FixtureChatEventSource: ChatEventSource { case .reset: backlog = [] terminalBacklog = [] + case .streamingProse: + // A live preview is transient and not part of history; forward it to + // subscribers but never fold it into the backlog. + break case .unknown: break case .stateChanged, .descriptorChanged: diff --git a/Packages/Shared/CmuxAgentChat/Sources/CmuxAgentChat/Store/ChatConversationStore.swift b/Packages/Shared/CmuxAgentChat/Sources/CmuxAgentChat/Store/ChatConversationStore.swift index 3fc839fed135..0031e20cff29 100644 --- a/Packages/Shared/CmuxAgentChat/Sources/CmuxAgentChat/Store/ChatConversationStore.swift +++ b/Packages/Shared/CmuxAgentChat/Sources/CmuxAgentChat/Store/ChatConversationStore.swift @@ -60,6 +60,14 @@ public final class ChatConversationStore { @ObservationIgnored private var messages: [ChatMessage] = [] @ObservationIgnored private var pending: [ChatPendingOutbound] = [] + /// Live, not-yet-committed preview of the agent's in-progress prose for the + /// current turn, scraped from the rendered terminal screen. Held outside + /// ``messages`` (so it never collides with window dedup/paging/seq) and + /// rendered as a trailing agent bubble. Cleared the instant the authoritative + /// agent prose lands via ``ChatSessionEvent/appended`` or an explicit + /// ``ChatSessionEvent/streamingProse`` `nil`, so it never duplicates a real + /// message. + @ObservationIgnored private var streamingMessage: ChatMessage? @ObservationIgnored private var firstUnreadSeq: Int? /// Terminal command-blocks for a `.terminal`-kind session, upserted by /// id; `terminalBlockOrder` preserves arrival order. Unused (and the @@ -458,6 +466,13 @@ public final class ChatConversationStore { switch event { case .appended(let newMessages): reconcilePending(against: newMessages) + // The authoritative agent prose for the in-flight turn just landed: + // drop the live preview so the committed message takes over with no + // duplicate, even if the host's explicit clear is delayed or lost. + if streamingMessage != nil, + newMessages.contains(where: { $0.role == .agent && Self.isProse($0) }) { + streamingMessage = nil + } // A live append whose seq regresses below the window tail means // the transcript was truncated/replaced and the tailer reset; // appending would corrupt window ordering. Re-anchor instead. @@ -499,6 +514,14 @@ public final class ChatConversationStore { // optimistic pending row it came from so it doesn't linger or leak. reconcileTerminalPending(against: blocks) reproject() + case .streamingProse(let message): + // The preview is a whole-value replace; an agent session only. A + // terminal session has no agent prose, so ignore it there. + guard descriptor.kind != .terminal else { break } + let next = message.flatMap { Self.isProse($0) ? $0 : nil } + guard next != streamingMessage else { break } + streamingMessage = next + reproject() case .reset: // The transcript was truncated/replaced on the Mac (tailer // re-read from scratch). The window's seq space is void; clear @@ -507,6 +530,8 @@ public final class ChatConversationStore { // their retry and in-flight sends may still land in the new // transcript and reconcile normally. messages = [] + // The preview belongs to the old seq space; drop it on re-anchor. + streamingMessage = nil // Terminal blocks must clear here too: the terminal reproject() // does not consult `messages`, so without this the synchronous // reproject below would re-render stale blocks (and they'd persist @@ -713,10 +738,27 @@ public final class ChatConversationStore { + pending.map(ChatTranscriptRow.pendingOutbound) return } + // The live preview renders as a trailing agent bubble after the + // committed window. Appending it to the projector input lets it group + // with adjacent agent prose exactly like a real message; it carries no + // window identity (never paged, deduped, or reconciled by id). + let projected: [ChatMessage] + if let streamingMessage, !messages.contains(where: { $0.id == streamingMessage.id }) { + projected = messages + [streamingMessage] + } else { + projected = messages + } rows = projector.rows( - messages: messages, + messages: projected, pending: pending, firstUnreadSeq: firstUnreadSeq ) } + + /// Whether a message is renderable agent/user prose (used to settle the + /// live preview against the authoritative transcript line). + private static func isProse(_ message: ChatMessage) -> Bool { + if case .prose = message.kind { return true } + return false + } } diff --git a/Packages/Shared/CmuxAgentChat/Sources/CmuxAgentChat/Store/ChatSessionListReducer.swift b/Packages/Shared/CmuxAgentChat/Sources/CmuxAgentChat/Store/ChatSessionListReducer.swift index d1d6cfa98305..c23db35bc47f 100644 --- a/Packages/Shared/CmuxAgentChat/Sources/CmuxAgentChat/Store/ChatSessionListReducer.swift +++ b/Packages/Shared/CmuxAgentChat/Sources/CmuxAgentChat/Store/ChatSessionListReducer.swift @@ -67,7 +67,7 @@ public struct ChatSessionListReducer: Sendable { // conversation's `ChatConversationStore` still consumes `stateChanged` // directly for its own live state (it is not version-reconciled). return sessions - case .appended, .updated, .terminalBlocks, .reset, .unknown: + case .appended, .updated, .terminalBlocks, .streamingProse, .reset, .unknown: // Transcript-content frames don't affect the session list. return sessions } diff --git a/Packages/Shared/CmuxAgentChat/Sources/CmuxAgentChat/Wire/ChatSessionEvent.swift b/Packages/Shared/CmuxAgentChat/Sources/CmuxAgentChat/Wire/ChatSessionEvent.swift index 9ad1e35914f1..e9fa75ef80ca 100644 --- a/Packages/Shared/CmuxAgentChat/Sources/CmuxAgentChat/Wire/ChatSessionEvent.swift +++ b/Packages/Shared/CmuxAgentChat/Sources/CmuxAgentChat/Wire/ChatSessionEvent.swift @@ -17,6 +17,14 @@ public enum ChatSessionEvent: Sendable, Equatable { /// reconnect is idempotent. case terminalBlocks([TerminalCommandBlock]) + /// A live, not-yet-committed preview of the agent's in-progress prose for + /// the current turn, scraped from the terminal's rendered screen while the + /// authoritative JSONL line has not been written yet. The payload replaces + /// any prior preview wholesale; `nil` clears it. It lives outside the + /// message window and is superseded the instant the authoritative agent + /// prose lands via ``appended``, so it never duplicates a real message. + case streamingProse(ChatMessage?) + /// The producing transcript was truncated or replaced; the session's /// seq space restarted and clients must re-anchor from history. case reset @@ -30,6 +38,7 @@ extension ChatSessionEvent: Codable { private enum CodingKeys: String, CodingKey { case event case messages + case message case state case descriptor case blocks @@ -41,6 +50,7 @@ extension ChatSessionEvent: Codable { case stateChanged = "state_changed" case descriptorChanged = "descriptor_changed" case terminalBlocks = "terminal_blocks" + case streamingProse = "streaming_prose" case reset } @@ -58,6 +68,8 @@ extension ChatSessionEvent: Codable { self = .descriptorChanged(try container.decode(ChatSessionDescriptor.self, forKey: .descriptor)) case .terminalBlocks: self = .terminalBlocks(try container.decode([TerminalCommandBlock].self, forKey: .blocks)) + case .streamingProse: + self = .streamingProse(try container.decodeIfPresent(ChatMessage.self, forKey: .message)) case .reset: self = .reset case .none: @@ -86,6 +98,9 @@ extension ChatSessionEvent: Codable { case .terminalBlocks(let blocks): try container.encode(EventName.terminalBlocks.rawValue, forKey: .event) try container.encode(blocks, forKey: .blocks) + case .streamingProse(let message): + try container.encode(EventName.streamingProse.rawValue, forKey: .event) + try container.encodeIfPresent(message, forKey: .message) case .reset: try container.encode(EventName.reset.rawValue, forKey: .event) case .unknown(let raw): diff --git a/Packages/Shared/CmuxAgentChat/Tests/CmuxAgentChatTests/AgentChatProseScreenExtractorTests.swift b/Packages/Shared/CmuxAgentChat/Tests/CmuxAgentChatTests/AgentChatProseScreenExtractorTests.swift new file mode 100644 index 000000000000..131cd9f301de --- /dev/null +++ b/Packages/Shared/CmuxAgentChat/Tests/CmuxAgentChatTests/AgentChatProseScreenExtractorTests.swift @@ -0,0 +1,260 @@ +import Foundation +import Testing + +@testable import CmuxAgentChat + +/// Fixtures mirror the rendered viewport of Claude Code 2.1 / Codex while a turn +/// streams: an answer block above a working/status line, with the input box and +/// footer below it. The extractor must isolate the answer and return `nil` when +/// no turn is actively streaming. +@Suite("AgentChatProseScreenExtractor") +struct AgentChatProseScreenExtractorTests { + private let extractor = AgentChatProseScreenExtractor() + + private static let rule = String(repeating: "─", count: 48) + + /// A Claude streaming viewport: prior tool block, the in-progress answer + /// (introduced by the "⏺ " bullet and wrapped under a 2-space hanging indent, + /// as the real TUI renders it), the spinner/status line, then the input box + /// and the bottom mode bar (which carries "esc to interrupt" while working). + private func claudeStreamingScreen(answer: [String]) -> [String] { + var rows = [ + "❯ Reply with three short sentences about the color blue.", + "", + "⏺ Read(notes.md)", + " ⎿ Read 12 lines", + "", + ] + for (offset, line) in answer.enumerated() { + rows.append(offset == 0 ? "⏺ \(line)" : " \(line)") + } + rows.append(contentsOf: [ + "", + "✢ Forming… (4s · ↓ 21 tokens)", + Self.rule, + "❯ ", + Self.rule, + " ⏵⏵ auto mode on (shift+tab to cycle) · esc to interrupt · ← for agents", + ]) + return rows + } + + @Test("isolates the in-progress answer above the status line") + func isolatesAnswer() { + let answer = [ + "The sky owes its blue to how air scatters sunlight.", + "Blue is often linked with calm, depth, and quiet focus.", + "From sapphires to deep ocean water, it is everywhere.", + ] + let result = extractor.extract(lines: claudeStreamingScreen(answer: answer), agentKind: .claude) + #expect(result == answer.joined(separator: "\n")) + } + + @Test("keeps paragraph breaks but drops padding blank runs") + func keepsParagraphBreaks() { + let answer = [ + "First paragraph.", + "", + "", + "Second paragraph.", + ] + let result = extractor.extract(lines: claudeStreamingScreen(answer: answer), agentKind: .claude) + #expect(result == "First paragraph.\n\nSecond paragraph.") + } + + @Test("returns nil when no turn is actively streaming") + func nilWhenSettled() { + // No status line: the turn has ended and the answer is committed. + let rows = [ + "⏺ The sky is blue because of Rayleigh scattering.", + "", + Self.rule, + "❯ ", + Self.rule, + "⏵⏵ auto mode", + ] + #expect(extractor.extract(lines: rows, agentKind: .claude) == nil) + } + + @Test("returns nil when the status line has no answer above it") + func nilWhenNoAnswer() { + let rows = [ + "⏺ Read(notes.md)", + " ⎿ Read 12 lines", + "✶ Thinking… (2s · esc to interrupt)", + Self.rule, + "❯ ", + ] + #expect(extractor.extract(lines: rows, agentKind: .claude) == nil) + } + + @Test("anchors on an esc-to-interrupt working line without a timer glyph") + func anchorsOnInterruptHint() { + // Codex renders the interrupt hint on the working line itself and does not + // bullet its answer, so the bullet-less body above the hint is the answer. + let rows = [ + "Streaming answer body line one.", + "Streaming answer body line two.", + " Thinking… esc to interrupt", + String(repeating: "─", count: 20), + "❯ ", + ] + let result = extractor.extract(lines: rows, agentKind: .codex) + #expect(result == "Streaming answer body line one.\nStreaming answer body line two.") + } + + @Test("a Codex working screen isolates its answer") + func codexScreen() { + let rows = [ + "› summarize the file", + "", + "Here is the summary you asked for.", + "It spans two lines of streaming prose.", + "Working (3s • Esc to interrupt)", + "▌", + ] + let result = extractor.extract(lines: rows, agentKind: .codex) + #expect(result == "Here is the summary you asked for.\nIt spans two lines of streaming prose.") + } + + @Test("elapsed-timer scanner matches seconds and minutes forms") + func elapsedTimer() { + #expect(AgentChatProseScreenExtractor.containsElapsedTimer("(4s")) + #expect(AgentChatProseScreenExtractor.containsElapsedTimer("foo (12s · bar)")) + #expect(AgentChatProseScreenExtractor.containsElapsedTimer("(1m05s)")) + // Bare form (no paren), as in the "running stop hooks… 0/3 · 3s" status. + #expect(AgentChatProseScreenExtractor.containsElapsedTimer("running stop hooks… 0/3 · 3s · ↓ 56 tokens")) + #expect(!AgentChatProseScreenExtractor.containsElapsedTimer("(no timer here)")) + #expect(!AgentChatProseScreenExtractor.containsElapsedTimer("plain text")) + // "0/3" alone is not a timer. + #expect(!AgentChatProseScreenExtractor.containsElapsedTimer("progress 0/3 done")) + } + + @Test("parenthesized-timer scanner rejects the bare Brewed-for summary") + func parenthesizedTimer() { + #expect(AgentChatProseScreenExtractor.containsParenthesizedTimer("✢ Forming… (9s)")) + #expect(AgentChatProseScreenExtractor.containsParenthesizedTimer("(1m05s)")) + // The post-turn summary has a bare timer, so it is not a live anchor. + #expect(!AgentChatProseScreenExtractor.containsParenthesizedTimer("✻ Brewed for 3s")) + } + + // MARK: - Real Claude Code 2.1.191 frames + + // The synthetic fixtures above missed two things the live TUI does: the + // in-progress answer is itself prefixed with "⏺ ", and the bottom mode bar + // carries "esc to interrupt" *while working* (below the input box). These + // frames are captured verbatim from a live `claude` turn via the debug + // socket's read-screen, then replayed so the extractor is pinned to the real + // rendering, not an idealized one. + + private static let realModeBarWorking = + " ⏵⏵ auto mode on (shift+tab to cycle) · esc to interrupt · ← for agents" + private static let realModeBarSettled = + " ⏵⏵ auto mode on (shift+tab to cycle) · ← for agents" + + /// A faithful Claude Code 2.1.191 viewport: welcome box, the echoed (wrapped) + /// prompt, the answer body, the spinner/timer line, then the input box and the + /// bottom mode bar (which carries "esc to interrupt" only while `working`). + private func realClaudeScreen(answerBody: [String], status: String, working: Bool) -> [String] { + var rows = [ + "Last login: Thu Jun 25 20:40:07 on ttys099", + "claude", + "╭─── Claude Code v2.1.191 ──────────────────────────╮", + "│ Welcome back Aziz! │", + "╰───────────────────────────────────────────────────╯", + "", + "", + "❯ Reply with exactly three short sentences about the color blue. No preamble, no lists, just", + " three sentences.", + " ", + ] + rows.append(contentsOf: answerBody) + rows.append("") + rows.append(status) + rows.append("") + rows.append(Self.rule) + rows.append("❯ ") + rows.append(Self.rule) + rows.append(working ? Self.realModeBarWorking : Self.realModeBarSettled) + return rows + } + + @Test("real frame: mid-stream partial sentence is isolated, not the mode bar") + func realPartialFrame() { + // Frame 19: the answer is cut mid-sentence and the bottom mode bar shows + // "esc to interrupt". Anchoring on that bar would yield chrome; the + // extractor must anchor on the spinner/timer line above the answer. + let rows = realClaudeScreen( + answerBody: [ + "⏺ The sky owes its blue to sunlight scattering across the atmosphere. Blue is often linked to", + ], + status: "✻ Nebulizing… (3s · ↓ 1 tokens)", + working: true + ) + let result = extractor.extract(lines: rows, agentKind: .claude) + #expect(result == "The sky owes its blue to sunlight scattering across the atmosphere. Blue is often linked to") + } + + @Test("real frame: full answer captured while still running stop hooks") + func realFullFrame() { + // Frame 21: full three-sentence answer, status switched to the bare-timer + // "running stop hooks… 0/3 · 3s · ↓ 56 tokens" form (no paren around 3s). + let rows = realClaudeScreen( + answerBody: [ + "⏺ The sky owes its blue to sunlight scattering across the atmosphere. Blue is often linked to", + " calm, depth, and quiet trust. From sapphires to deep oceans, it spans some of nature's most", + " striking sights.", + ], + status: "✻ Nebulizing… (running stop hooks… 0/3 · 3s · ↓ 56 tokens)", + working: true + ) + let result = extractor.extract(lines: rows, agentKind: .claude) + // The 2-space hanging indent under "⏺ " is stripped so the wrapped lines + // read as one flowing answer. + #expect(result == """ + The sky owes its blue to sunlight scattering across the atmosphere. Blue is often linked to + calm, depth, and quiet trust. From sapphires to deep oceans, it spans some of nature's most + striking sights. + """) + } + + @Test("real frame: empty answer (only the ⏺ bullet) yields nil") + func realEmptyAnswerFrame() { + // Frame 17: the block bullet has rendered but no words yet. + let rows = realClaudeScreen( + answerBody: ["⏺ "], + status: "✢ Nebulizing… (2s · ↓ 1 tokens)", + working: true + ) + #expect(extractor.extract(lines: rows, agentKind: .claude) == nil) + } + + @Test("real frame: settled turn (Brewed for 3s summary) yields nil") + func realSettledFrame() { + // Frame 22: turn done. The spinner line is replaced by the "Brewed for 3s" + // summary (bare timer, no throughput) and the mode bar drops "esc to + // interrupt", so the extractor must report no active stream. + let rows = realClaudeScreen( + answerBody: [ + "⏺ The sky owes its blue to sunlight scattering across the atmosphere. Blue is often linked to", + " calm, depth, and quiet trust. From sapphires to deep oceans, it spans some of nature's most", + " striking sights.", + ], + status: "✻ Brewed for 3s", + working: false + ) + #expect(extractor.extract(lines: rows, agentKind: .claude) == nil) + } + + @Test("a long answer is capped, never folding the whole screen") + func capsAnswerLength() { + let answer = (0..<400).map { "line \($0)" } + // No boundary above the answer: only the cap stops collection. Use Codex, + // whose answer is bullet-less, so the cap (not the answer-top) bounds it. + var rows = answer + rows.append("✢ Forming… (9s)") + let result = extractor.extract(lines: rows, agentKind: .codex) + let lineCount = result?.split(separator: "\n", omittingEmptySubsequences: false).count ?? 0 + #expect(lineCount <= 200) + } +} diff --git a/Packages/Shared/CmuxAgentChat/Tests/CmuxAgentChatTests/ChatConversationStoreTests.swift b/Packages/Shared/CmuxAgentChat/Tests/CmuxAgentChatTests/ChatConversationStoreTests.swift index cc2bbacc1f9f..0b03c413c41d 100644 --- a/Packages/Shared/CmuxAgentChat/Tests/CmuxAgentChatTests/ChatConversationStoreTests.swift +++ b/Packages/Shared/CmuxAgentChat/Tests/CmuxAgentChatTests/ChatConversationStoreTests.swift @@ -177,6 +177,19 @@ struct ChatConversationStoreTests { (0.. ChatMessage { + ChatMessage( + id: "stream:session-1", + seq: Int.max - 1, + role: .agent, + timestamp: baseTime.addingTimeInterval(1000), + kind: .prose(ChatProse(text: text)) + ) + } + private static func makeStore( source: any ChatEventSource, lastReadSeq: Int? = nil, @@ -296,6 +309,64 @@ struct ChatConversationStoreTests { #expect(snaps.first?.message.id == original.id) } + @Test("streaming prose renders as a trailing bubble and clears on nil") + func streamingProseRendersAndClears() async { + let source = FixtureChatEventSource() + let store = Self.makeStore(source: source) + let runTask = Task { await store.run() } + defer { runTask.cancel() } + + #expect(await TestPoller.waitUntil { store.isConnected }) + let preview = Self.streamingMessage(text: "partial answer") + await source.emit(.streamingProse(preview)) + #expect( + await TestPoller.waitUntil { + Self.snapshots(store.rows).contains { $0.message.id == preview.id } + } + ) + await source.emit(.streamingProse(nil)) + #expect( + await TestPoller.waitUntil { + !Self.snapshots(store.rows).contains { $0.message.id == preview.id } + } + ) + } + + @Test("authoritative agent prose supersedes the live preview without a duplicate") + func authoritativeProseSupersedesPreview() async { + let source = FixtureChatEventSource() + let store = Self.makeStore(source: source) + let runTask = Task { await store.run() } + defer { runTask.cancel() } + + #expect(await TestPoller.waitUntil { store.isConnected }) + let preview = Self.streamingMessage(text: "The sky is blue") + await source.emit(.streamingProse(preview)) + #expect( + await TestPoller.waitUntil { + Self.snapshots(store.rows).contains { $0.message.id == preview.id } + } + ) + // The committed transcript line lands; the preview must vanish and only + // the real message remains (no duplicate bubble). + let committed = Self.prose(seq: 0, role: .agent, text: "The sky is blue") + await source.emit(.appended([committed])) + #expect( + await TestPoller.waitUntil { + let snaps = Self.snapshots(store.rows) + return snaps.contains { $0.message.id == committed.id } + && !snaps.contains { $0.message.id == preview.id } + } + ) + let proseTexts = Self.snapshots(store.rows) + .filter { $0.message.role == .agent } + .compactMap { snapshot -> String? in + if case .prose(let prose) = snapshot.message.kind { return prose.text } + return nil + } + #expect(proseTexts == ["The sky is blue"]) + } + @Test("stateChanged event updates agentState") func stateChangedUpdatesAgentState() async { let source = FixtureChatEventSource() diff --git a/Packages/iOS/CmuxMobileShellUI/Sources/CmuxMobileShellUI/CMUXMobileRootView.swift b/Packages/iOS/CmuxMobileShellUI/Sources/CmuxMobileShellUI/CMUXMobileRootView.swift index 45e9e68cbd18..4adcaedf2553 100644 --- a/Packages/iOS/CmuxMobileShellUI/Sources/CmuxMobileShellUI/CMUXMobileRootView.swift +++ b/Packages/iOS/CmuxMobileShellUI/Sources/CmuxMobileShellUI/CMUXMobileRootView.swift @@ -70,6 +70,22 @@ struct CMUXMobileRootView: View { #endif } + private var shouldShowStreamingChatPreview: Bool { + #if os(iOS) && DEBUG + return UITestConfig.streamingChatPreviewEnabled + #else + return false + #endif + } + + @ViewBuilder private var streamingChatPreview: some View { + #if os(iOS) && DEBUG + StreamingChatPreviewView() + #else + EmptyView() + #endif + } + @ViewBuilder private var terminalLayoutPreview: some View { #if os(iOS) && DEBUG TerminalLayoutPreviewView() @@ -188,6 +204,8 @@ struct CMUXMobileRootView: View { terminalLayoutPreview } else if shouldShowWorkspaceListLayoutPreview { workspaceListLayoutPreview + } else if shouldShowStreamingChatPreview { + streamingChatPreview } else if !isAuthenticated { SignInView() } else if store.connectionState != .connected && shouldShowRestoringStoredMac { diff --git a/Packages/iOS/CmuxMobileShellUI/Sources/CmuxMobileShellUI/StreamingChatPreviewView.swift b/Packages/iOS/CmuxMobileShellUI/Sources/CmuxMobileShellUI/StreamingChatPreviewView.swift new file mode 100644 index 000000000000..37b8036b3d59 --- /dev/null +++ b/Packages/iOS/CmuxMobileShellUI/Sources/CmuxMobileShellUI/StreamingChatPreviewView.swift @@ -0,0 +1,106 @@ +#if canImport(UIKit) && DEBUG +import CmuxAgentChat +import CmuxAgentChatUI +import SwiftUI + +/// DEBUG-only self-playing agent chat used to record and verify the live +/// streaming-prose preview on the simulator, with no sign-in or Mac pairing. +/// +/// Mounted by the root view when ``UITestConfig/streamingChatPreviewEnabled`` is +/// set (`CMUX_UITEST_STREAMING_CHAT_PREVIEW=1`). It builds the production +/// ``ChatConversationStore`` + ``ChatScreen`` and drives them with real +/// ``ChatSessionEvent/streamingProse`` events that grow word by word, then clear +/// — exactly the wire traffic the Mac host emits while a turn streams — so the +/// recording exercises the real store/projector/transcript render path. +struct StreamingChatPreviewView: View { + @State private var model: Model? + + var body: some View { + Group { + if let model { + ChatScreen(store: model.store, onOpenTerminal: {}) + } else { + ProgressView() + .task { await start() } + } + } + } + + @MainActor + private func start() async { + guard model == nil else { return } + let descriptor = ChatSessionDescriptor( + id: "streaming-preview", + agentKind: .claude, + title: "Streaming preview" + ) + let prompt = ChatMessage( + id: "preview-user-0", + seq: 0, + role: .user, + timestamp: Date(), + kind: .prose(ChatProse(text: "Reply with three short sentences about the color blue.")) + ) + let source = FixtureChatEventSource(backlog: [prompt]) + let store = ChatConversationStore(descriptor: descriptor, source: source) + let model = Model(store: store, source: source) + model.runTask = Task { await store.run() } + model.driveTask = Task { await Self.drive(source: source) } + self.model = model + } + + /// Loops the streaming lifecycle so any recording window captures a full + /// build: grow the answer word by word via `streamingProse`, hold, then + /// clear with `streamingProse(nil)`, and repeat. + private static func drive(source: FixtureChatEventSource) async { + let full = "The sky looks blue because air scatters short blue wavelengths most. " + + "Blue is widely tied to calm, depth, and quiet focus. " + + "From sapphires to the open ocean, it runs through the natural world." + let words = full.split(separator: " ").map(String.init) + // Let ChatScreen subscribe and load the seeded prompt first. + try? await Task.sleep(for: .milliseconds(1200)) + while !Task.isCancelled { + var accumulated = "" + for word in words { + if Task.isCancelled { return } + accumulated += accumulated.isEmpty ? word : " " + word + await source.emit(.streamingProse(previewMessage(text: accumulated))) + try? await Task.sleep(for: .milliseconds(130)) + } + try? await Task.sleep(for: .milliseconds(1500)) + await source.emit(.streamingProse(nil)) + try? await Task.sleep(for: .milliseconds(800)) + } + } + + private static func previewMessage(text: String) -> ChatMessage { + ChatMessage( + id: "stream:streaming-preview", + seq: Int.max - 1, + role: .agent, + timestamp: Date(), + kind: .prose(ChatProse(text: text)) + ) + } + + /// Retains the store, source, and driver tasks across re-renders; cancels + /// the tasks on teardown. + @MainActor + final class Model { + let store: ChatConversationStore + let source: FixtureChatEventSource + var runTask: Task? + var driveTask: Task? + + init(store: ChatConversationStore, source: FixtureChatEventSource) { + self.store = store + self.source = source + } + + deinit { + runTask?.cancel() + driveTask?.cancel() + } + } +} +#endif diff --git a/Packages/iOS/CmuxMobileSupport/Sources/CmuxMobileSupport/UITestConfig.swift b/Packages/iOS/CmuxMobileSupport/Sources/CmuxMobileSupport/UITestConfig.swift index fcf94e877a55..dcc4206aa43b 100644 --- a/Packages/iOS/CmuxMobileSupport/Sources/CmuxMobileSupport/UITestConfig.swift +++ b/Packages/iOS/CmuxMobileSupport/Sources/CmuxMobileSupport/UITestConfig.swift @@ -94,6 +94,20 @@ public struct UITestConfig { #endif } + /// Whether the standalone streaming-chat preview is enabled. + /// + /// When `CMUX_UITEST_STREAMING_CHAT_PREVIEW=1`, the root view renders a + /// self-playing agent chat that drives the real conversation store with live + /// `streamingProse` events, so the incremental streaming preview can be + /// recorded on the simulator without sign-in or Mac pairing. DEBUG-only. + public static var streamingChatPreviewEnabled: Bool { + #if DEBUG + return ProcessInfo.processInfo.environment["CMUX_UITEST_STREAMING_CHAT_PREVIEW"] == "1" + #else + return false + #endif + } + /// Whether the standalone agent-chat preview is enabled. /// /// When `CMUX_UITEST_AGENT_CHAT_PREVIEW=1`, the root view renders the diff --git a/Sources/Mobile/AgentChat/AgentChatProseStreamer.swift b/Sources/Mobile/AgentChat/AgentChatProseStreamer.swift new file mode 100644 index 000000000000..13abcb2dd0f0 --- /dev/null +++ b/Sources/Mobile/AgentChat/AgentChatProseStreamer.swift @@ -0,0 +1,144 @@ +import CmuxAgentChat +import Foundation + +/// Drives the live agent-prose streaming preview for in-flight turns. +/// +/// The agent CLIs never write token-level prose to their JSONL transcript, so +/// the only token-grained source of a streaming answer is the rendered terminal +/// screen. While a turn is active this polls a screen snapshot for the hosting +/// surface, extracts the in-progress prose with +/// ``AgentChatProseScreenExtractor``, and pushes it to chat clients as a +/// ``ChatSessionEvent/streamingProse`` preview. The preview is cleared the +/// instant the authoritative JSONL line lands (``authoritativeProseArrived``) or +/// the turn ends (``turnEnded``), so it never duplicates a committed message. +/// +/// Gating keeps it off the keystroke hot path: a surface is only snapshotted +/// while the feature flag is on, a chat client is subscribed, and a turn is +/// actively streaming for that session. Idle sessions are never polled. +@MainActor +final class AgentChatProseStreamer { + /// Per-session live-turn bookkeeping. + private struct ActiveTurn { + let surfaceID: UUID + let agentKind: ChatAgentKind + /// Set once the authoritative prose for this turn has landed; the loop + /// stops emitting until the next ``turnStarted``. + var settled: Bool = false + /// Last preview text pushed, so unchanged snapshots don't re-emit. + var lastEmitted: String? + } + + private let extractor = AgentChatProseScreenExtractor() + private var turns: [String: ActiveTurn] = [:] + private var tasks: [String: Task] = [:] + + private let emit: @MainActor (ChatSessionEventFrame) -> Void + private let snapshot: @MainActor (UUID) -> [String]? + private let isEnabled: @MainActor () -> Bool + private let hasSubscribers: @MainActor () -> Bool + private let now: @MainActor () -> Date + private let pollInterval: Duration + private let sleep: @Sendable (Duration) async -> Void + + /// Creates a prose streamer. + /// + /// - Parameters: + /// - emit: Publishes a wire frame to subscribed chat clients. + /// - snapshot: Maps a surface id to its rendered screen rows (top to + /// bottom), or `nil` when the surface is gone. + /// - isEnabled: The live feature-flag read (default-off). + /// - hasSubscribers: Whether any chat client is currently listening. + /// - now: Clock seam for the preview message timestamp. + /// - pollInterval: Snapshot cadence while a turn streams. + /// - sleep: Cancellable sleep seam for the poll loop. + init( + emit: @escaping @MainActor (ChatSessionEventFrame) -> Void, + snapshot: @escaping @MainActor (UUID) -> [String]?, + isEnabled: @escaping @MainActor () -> Bool = { AgentChatProseStreamingFlag.isEnabled }, + hasSubscribers: @escaping @MainActor () -> Bool, + now: @escaping @MainActor () -> Date = { Date() }, + pollInterval: Duration = .milliseconds(150), + sleep: @escaping @Sendable (Duration) async -> Void = { try? await Task.sleep(for: $0) } + ) { + self.emit = emit + self.snapshot = snapshot + self.isEnabled = isEnabled + self.hasSubscribers = hasSubscribers + self.now = now + self.pollInterval = pollInterval + self.sleep = sleep + } + + /// Begins (or re-arms) streaming for a session's in-flight turn. + /// + /// - Parameters: + /// - sessionID: The chat session. + /// - surfaceID: The hosting terminal surface to snapshot. + /// - agentKind: Selects the prose extractor's chrome markers. + func turnStarted(sessionID: String, surfaceID: UUID, agentKind: ChatAgentKind) { + turns[sessionID] = ActiveTurn(surfaceID: surfaceID, agentKind: agentKind) + guard tasks[sessionID] == nil else { return } + tasks[sessionID] = Task { [weak self] in + await self?.runLoop(sessionID: sessionID) + } + } + + /// The authoritative prose for the turn landed; drop the preview and stop + /// emitting until the next turn (kept distinct from ``turnEnded`` so the + /// poll loop is reused across a multi-block turn instead of respawned). + func authoritativeProseArrived(sessionID: String) { + guard turns[sessionID] != nil else { return } + turns[sessionID]?.settled = true + turns[sessionID]?.lastEmitted = nil + clearPreview(sessionID: sessionID) + } + + /// Ends streaming for a session: cancels the loop and clears the preview. + func turnEnded(sessionID: String) { + tasks[sessionID]?.cancel() + tasks[sessionID] = nil + let wasActive = turns[sessionID] != nil + turns[sessionID] = nil + if wasActive { clearPreview(sessionID: sessionID) } + } + + /// Tears down every active stream (app teardown / subscriber loss). + func stopAll() { + for sessionID in Array(tasks.keys) { turnEnded(sessionID: sessionID) } + } + + // MARK: - Internals + + private func runLoop(sessionID: String) async { + while !Task.isCancelled { + emitPreviewIfChanged(sessionID: sessionID) + await sleep(pollInterval) + } + } + + private func emitPreviewIfChanged(sessionID: String) { + guard let turn = turns[sessionID], !turn.settled else { return } + guard isEnabled(), hasSubscribers() else { return } + guard let lines = snapshot(turn.surfaceID) else { return } + guard let prose = extractor.extract(lines: lines, agentKind: turn.agentKind) else { return } + guard prose != turn.lastEmitted else { return } + turns[sessionID]?.lastEmitted = prose + emit(ChatSessionEventFrame(sessionID: sessionID, event: .streamingProse(previewMessage(sessionID: sessionID, text: prose)))) + } + + private func clearPreview(sessionID: String) { + emit(ChatSessionEventFrame(sessionID: sessionID, event: .streamingProse(nil))) + } + + /// Builds the preview message. The id is stable per session so successive + /// previews replace in place; the seq sorts it after the committed window. + private func previewMessage(sessionID: String, text: String) -> ChatMessage { + ChatMessage( + id: "stream:\(sessionID)", + seq: Int.max - 1, + role: .agent, + timestamp: now(), + kind: .prose(ChatProse(text: text)) + ) + } +} diff --git a/Sources/Mobile/AgentChat/AgentChatProseStreamingFlag.swift b/Sources/Mobile/AgentChat/AgentChatProseStreamingFlag.swift new file mode 100644 index 000000000000..2574eaeffe44 --- /dev/null +++ b/Sources/Mobile/AgentChat/AgentChatProseStreamingFlag.swift @@ -0,0 +1,22 @@ +import Foundation + +/// Default-off feature flag gating the live agent-prose streaming preview. +/// +/// Streaming scrapes the rendered terminal screen heuristically and is meant +/// for dogfood before it is trusted on by default, so it is an explicit opt-in +/// rather than a `cmux.json` product setting. Flip it at runtime (no relaunch) +/// with: +/// +/// ``` +/// defaults write CMUXAgentChatProseStreaming -bool YES +/// ``` +/// +/// The value is read live on every poll so toggling takes effect immediately. +enum AgentChatProseStreamingFlag { + static let defaultsKey = "CMUXAgentChatProseStreaming" + + /// Whether the live streaming preview is enabled. Defaults to `false`. + static var isEnabled: Bool { + UserDefaults.standard.bool(forKey: defaultsKey) + } +} diff --git a/Sources/Mobile/AgentChat/AgentChatTranscriptService.swift b/Sources/Mobile/AgentChat/AgentChatTranscriptService.swift index 2d2357e3619a..c292301c11f5 100644 --- a/Sources/Mobile/AgentChat/AgentChatTranscriptService.swift +++ b/Sources/Mobile/AgentChat/AgentChatTranscriptService.swift @@ -1,5 +1,6 @@ import CMUXAgentLaunch import CmuxAgentChat +import CmuxTerminal import Foundation /// Mac-side facade for the agent chat surface: tracks sessions from hook @@ -14,6 +15,8 @@ final class AgentChatTranscriptService { let resolver: AgentChatTranscriptResolver private let coding = ChatWireCoding() private var tailers: [String: AgentChatTranscriptTailer] = [:] + /// Drives the live agent-prose streaming preview (default-off feature flag). + private var proseStreamer: AgentChatProseStreamer! /// Sessions whose transcript could not be resolved; skipped until an /// explicit history request retries, so per-hook-event resolution /// failures don't rescan the filesystem during tool storms. @@ -40,6 +43,21 @@ final class AgentChatTranscriptService { registry.onRecordChanged = { [weak self] record, previous in self?.handleRecordChange(record, previous: previous) } + self.proseStreamer = AgentChatProseStreamer( + emit: { [weak self] frame in self?.emit(frame: frame) }, + snapshot: { surfaceID in Self.screenRows(surfaceID: surfaceID) }, + hasSubscribers: { MobileHostService.hasEventSubscribers(topic: Self.eventTopic) } + ) + } + + /// Rendered screen rows (top to bottom) for a surface, the source the prose + /// streamer scrapes. Mirrors the render-grid observer's surface lookup. + @MainActor + private static func screenRows(surfaceID: UUID) -> [String]? { + guard let surface = GhosttyApp.terminalSurfaceRegistry.terminalSurface(id: surfaceID) else { + return nil + } + return surface.mobileRenderGridFrame(stateSeq: 0, full: true)?.rows } /// A `(session, surface)` resume re-bind cmux authored during session @@ -136,6 +154,25 @@ final class AgentChatTranscriptService { MobileHostService.hasEventSubscribers(topic: Self.eventTopic) { ensureTailer(for: record) } + // Drive the live prose-streaming preview off the turn lifecycle: a + // prompt starts the in-flight turn, Stop ends it. Gated by the + // default-off flag here so a turn loop is never spawned otherwise. + switch event.hookEventName { + case .userPromptSubmit: + if AgentChatProseStreamingFlag.isEnabled, + record.state != .ended, + let surfaceID = record.surfaceID.flatMap(UUID.init(uuidString:)) { + proseStreamer.turnStarted( + sessionID: record.sessionID, + surfaceID: surfaceID, + agentKind: record.agentKind + ) + } + case .stop, .sessionEnd: + proseStreamer.turnEnded(sessionID: record.sessionID) + default: + break + } } /// Lists chat-capable sessions. @@ -315,6 +352,11 @@ final class AgentChatTranscriptService { registry.update(sessionID: sessionID) { $0.title = title } } if !batch.appended.isEmpty { + // The authoritative prose for the turn just landed: settle the live + // preview so the committed message takes over with no duplicate. + if Self.batchContainsAgentProse(batch.appended) { + proseStreamer.authoritativeProseArrived(sessionID: sessionID) + } emit(frame: ChatSessionEventFrame(sessionID: sessionID, event: .appended(batch.appended))) } if !batch.updated.isEmpty { @@ -325,6 +367,16 @@ final class AgentChatTranscriptService { } } + /// Whether a batch carries any committed agent prose, the signal that the + /// streaming preview for the turn should settle. + private static func batchContainsAgentProse(_ messages: [ChatMessage]) -> Bool { + messages.contains { message in + guard message.role == .agent else { return false } + if case .prose = message.kind { return true } + return false + } + } + private static func completedAssistantTurnTimestamp(in messages: [ChatMessage]) -> Date? { guard !messages.isEmpty else { return nil } var completedAt: Date? @@ -347,6 +399,9 @@ final class AgentChatTranscriptService { let stateChanged = previous?.state != record.state let transcriptBecameAvailable = previous?.transcriptPath == nil && record.transcriptPath != nil if stateChanged, record.state == .ended { + // The transcript can no longer grow; stop any live preview loop so + // an agent that exits without a Stop hook doesn't leak the poll task. + proseStreamer.turnEnded(sessionID: record.sessionID) if let tailer = tailers.removeValue(forKey: record.sessionID) { // The transcript can no longer grow; release the file watcher // and cache instead of holding them until app quit. Evicting diff --git a/cmux.xcodeproj/project.pbxproj b/cmux.xcodeproj/project.pbxproj index b5cb4c1f90db..92aba2246a82 100644 --- a/cmux.xcodeproj/project.pbxproj +++ b/cmux.xcodeproj/project.pbxproj @@ -10,6 +10,8 @@ A9E020000000000000000006 /* agent-session-react in Resources */ = {isa = PBXBuildFile; fileRef = A9E010000000000000000006 /* agent-session-react */; }; A9E020000000000000000007 /* agent-session-solid in Resources */ = {isa = PBXBuildFile; fileRef = A9E010000000000000000007 /* agent-session-solid */; }; ACA7C4A70000000000000003 /* AgentChatHookSessionStore.swift in Sources */ = {isa = PBXBuildFile; fileRef = ACA7C4A70000000000000004 /* AgentChatHookSessionStore.swift */; }; + CDFE000000000000000000A1 /* AgentChatProseStreamer.swift in Sources */ = {isa = PBXBuildFile; fileRef = CDFE000000000000000000A2 /* AgentChatProseStreamer.swift */; }; + CDFE000000000000000000B1 /* AgentChatProseStreamingFlag.swift in Sources */ = {isa = PBXBuildFile; fileRef = CDFE000000000000000000B2 /* AgentChatProseStreamingFlag.swift */; }; ACA7C4A70000000000000001 /* AgentChatSessionRecord.swift in Sources */ = {isa = PBXBuildFile; fileRef = ACA7C4A70000000000000002 /* AgentChatSessionRecord.swift */; }; ACA7C4A70000000000000007 /* AgentChatSessionRegistry.swift in Sources */ = {isa = PBXBuildFile; fileRef = ACA7C4A70000000000000008 /* AgentChatSessionRegistry.swift */; }; ACA7C4A70000000000000005 /* AgentChatTranscriptResolver.swift in Sources */ = {isa = PBXBuildFile; fileRef = ACA7C4A70000000000000006 /* AgentChatTranscriptResolver.swift */; }; @@ -1212,6 +1214,8 @@ A9E010000000000000000006 /* agent-session-react */ = {isa = PBXFileReference; lastKnownFileType = folder; path = "agent-session-react"; sourceTree = ""; }; A9E010000000000000000007 /* agent-session-solid */ = {isa = PBXFileReference; lastKnownFileType = folder; path = "agent-session-solid"; sourceTree = ""; }; ACA7C4A70000000000000004 /* AgentChatHookSessionStore.swift */ = {isa = PBXFileReference; includeInIndex = 1; lastKnownFileType = sourcecode.swift; path = "AgentChatHookSessionStore.swift"; sourceTree = ""; }; + CDFE000000000000000000A2 /* AgentChatProseStreamer.swift */ = {isa = PBXFileReference; includeInIndex = 1; lastKnownFileType = sourcecode.swift; path = "AgentChatProseStreamer.swift"; sourceTree = ""; }; + CDFE000000000000000000B2 /* AgentChatProseStreamingFlag.swift */ = {isa = PBXFileReference; includeInIndex = 1; lastKnownFileType = sourcecode.swift; path = "AgentChatProseStreamingFlag.swift"; sourceTree = ""; }; ACA7C4A70000000000000002 /* AgentChatSessionRecord.swift */ = {isa = PBXFileReference; includeInIndex = 1; lastKnownFileType = sourcecode.swift; path = "AgentChatSessionRecord.swift"; sourceTree = ""; }; ACA7C4A70000000000000008 /* AgentChatSessionRegistry.swift */ = {isa = PBXFileReference; includeInIndex = 1; lastKnownFileType = sourcecode.swift; path = "AgentChatSessionRegistry.swift"; sourceTree = ""; }; ACA7C4A70000000000000006 /* AgentChatTranscriptResolver.swift */ = {isa = PBXFileReference; includeInIndex = 1; lastKnownFileType = sourcecode.swift; path = "AgentChatTranscriptResolver.swift"; sourceTree = ""; }; @@ -2460,6 +2464,8 @@ ACA7C4A70000000000000008 /* AgentChatSessionRegistry.swift */, ACA7C4A7000000000000000A /* AgentChatTranscriptTailer.swift */, ACA7C4A7000000000000000C /* AgentChatTranscriptService.swift */, + CDFE000000000000000000A2 /* AgentChatProseStreamer.swift */, + CDFE000000000000000000B2 /* AgentChatProseStreamingFlag.swift */, ); name = AgentChat; path = AgentChat; @@ -3954,6 +3960,8 @@ files = ( ACA7C4A70000000000000003 /* AgentChatHookSessionStore.swift in Sources */, + CDFE000000000000000000A1 /* AgentChatProseStreamer.swift in Sources */, + CDFE000000000000000000B1 /* AgentChatProseStreamingFlag.swift in Sources */, ACA7C4A70000000000000001 /* AgentChatSessionRecord.swift in Sources */, ACA7C4A70000000000000007 /* AgentChatSessionRegistry.swift in Sources */, ACA7C4A70000000000000005 /* AgentChatTranscriptResolver.swift in Sources */, diff --git a/docs/streaming-agent-updates.md b/docs/streaming-agent-updates.md new file mode 100644 index 000000000000..c4e6a3358889 --- /dev/null +++ b/docs/streaming-agent-updates.md @@ -0,0 +1,34 @@ +# iOS streaming agent updates: incremental responses in the chat GUI + +Goal: the iOS chat GUI should show the agent's response building incrementally (partial text / blocks appearing as they generate), not just the final message after the turn completes. + +## Root cause (investigation in this worktree) + +The wire and the iOS renderer already support incremental updates. The pipeline is starved of partial data, not broken. + +1. **Agent JSONL is the source, and it is written per content block, never per token.** Confirmed empirically against the installed versions (Claude Code 2.1.187, Codex 0.142.0) by watching the on-disk file grow during a live turn: a 647-character answer went from absent to fully present in a single 100ms poll interval, not in growing increments. Each Claude `assistant` line is one already-complete content block carrying its `stop_reason`; Codex rollout JSONL contains `agent_message`/`message` items but no `agent_message_delta` (the deltas exist only on Codex's in-process event stream, never persisted). So a single long final-answer block lands all at once when generation finishes. Block-level streaming already works (thinking → tool → text each appear as their line is written); the only gap is intra-block token streaming. (`Packages/Shared/CmuxAgentChat/.../Parsing/ClaudeTranscriptParser.swift:174-189`, `CodexTranscriptParser.swift:135-163` parse complete lines only.) +2. **Host tailer reads complete lines only.** `Sources/Mobile/AgentChat/AgentChatTranscriptTailer.swift:184-254` seeks to the last byte offset, reads only newline-terminated lines (`:224`), and re-reads on file growth (FileWatcher, 200ms throttle). A single final line appears once and is parsed once. +3. **Wire is correct.** `Sources/Mobile/AgentChat/AgentChatTranscriptService.swift:313-329` emits `.appended` for new messages and `.updated` for in-place changes (today only tool results / permission answers), with no coalescing. +4. **iOS render is correct.** `Packages/Shared/CmuxAgentChat/.../Store/ChatConversationStore.swift:457-528` mutates messages on `.updated` and reprojects; `Packages/iOS/CmuxAgentChatUI/.../Transcript/ChatTranscriptListView.swift` renders the final state. It would render streaming updates if they arrived. + +Ranked: (a) agent JSONL has no partials to parse [PRIMARY], (b) parsers can't emit partial prose [follows from a], (c) wire forwards everything correctly [not a cause], (d) iOS renders updates correctly [not a cause]. + +## The real lever: Ghostty's emulated screen grid + +The three candidate sources, now resolved against the installed versions: + +- **B (parser change) is dead.** Confirmed above: neither agent persists token deltas to JSONL. There is no format flag or version that writes partial lines for an interactive session. Nothing to tail. +- **C (stream-json output mode) is incompatible with the interactive UX.** `claude --print --output-format stream-json --include-partial-messages` and `codex exec --json` do emit token deltas, but only in non-interactive print mode, which is mutually exclusive with the interactive TUI the user actually runs. There is no way to get deltas out of an already-running interactive agent except its own stdout (the TUI bytes). Using C would mean cmux owning the agent process in headless mode and rendering it itself, a different product, not a streaming improvement. +- **A (pty tee), but not by parsing raw bytes.** Confirmed empirically by capturing Claude Code 2.1's interactive pty output: it paints the streaming answer as a synchronized-output frame (`ESC[?2026h`/`l`) over a bottom "live" region, placing each word with absolute-column moves (`ESC[G`) and relative cursor moves, repainting the whole region (answer + spinner + input box + footer) every frame. No alt-screen, no `ESC[2J`, but also no linear append: finalized text only commits to scrollback once the turn ends. Recovering the visible prose from these raw bytes means running a terminal emulator (cell grid, cursor, scroll region, synchronized frames). Ghostty already is that emulator. + +So the one viable source is **Ghostty's emulated screen text** (`ghostty_surface_read_text`), polled during an active turn off the typing hot path, with the assistant-prose region isolated from per-agent TUI chrome (Claude's `●`/`✶` bullets and `✢ …ing… (Ns · ↓N tokens)` spinner, the `────` rules, the `❯` input box, the `⏵⏵ auto mode` footer; Codex differs), emitted as a provisional `.prose` message via the existing wire (`.appended` once, then `.updated` repeatedly), and reconciled against the authoritative JSONL line at turn end (replace provisional with parsed-final by message id) so the committed state always matches JSONL and streaming imperfections are transient only. + +This is a large, inherently heuristic subsystem: per-agent chrome stripping that drifts with CLI versions and locale, anchoring "this turn's prose" against the prompt echo / turn-start scrollback position, scrollback inclusion for answers taller than the viewport, JSONL reconciliation, and care on the latency-sensitive terminal path (`MobileTerminalByteTee`, the Ghostty read path) so polling never touches the keystroke hot path. It is also hard to unit-test deterministically (needs recorded TUI byte fixtures replayed through the emulator, or recorded screen snapshots). + +### Active-turn gating + +Poll only while a turn is in flight. The hook lifecycle already marks this: `UserPromptSubmit` starts a turn, `Stop` ends it. Gate the scraper to active turns per surface so idle terminals are never polled (cost + latency), and stop on `Stop` / the authoritative JSONL append. + +## Parity note + +This is consistent with the agent-GUI component map: the streaming source belongs in Layer 1 (host) and the shared model (Layer 2), feeding the existing wire. The iOS view layer (Layer 3) needs no new capability for streaming text beyond what `.updated` already drives. diff --git a/vendor/bonsplit b/vendor/bonsplit index 8ca85b6a6880..c4aa88a559e3 160000 --- a/vendor/bonsplit +++ b/vendor/bonsplit @@ -1 +1 @@ -Subproject commit 8ca85b6a688072acc6e21171dfebba302dc82020 +Subproject commit c4aa88a559e3ad9fb66a70c506d371b5725cf813