Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
12 changes: 4 additions & 8 deletions Sources/Mobile/AgentChat/AgentChatProseStreamer.swift
Original file line number Diff line number Diff line change
Expand Up @@ -12,9 +12,9 @@ import Foundation
/// 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.
/// It stays off the keystroke hot path: a surface is only snapshotted while 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.
Expand All @@ -34,7 +34,6 @@ final class AgentChatProseStreamer {

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
Expand All @@ -46,23 +45,20 @@ final class AgentChatProseStreamer {
/// - 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
Expand Down Expand Up @@ -118,7 +114,7 @@ final class AgentChatProseStreamer {

private func emitPreviewIfChanged(sessionID: String) {
guard let turn = turns[sessionID], !turn.settled else { return }
guard isEnabled(), hasSubscribers() else { return }
guard 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 }
Expand Down
22 changes: 0 additions & 22 deletions Sources/Mobile/AgentChat/AgentChatProseStreamingFlag.swift

This file was deleted.

8 changes: 3 additions & 5 deletions Sources/Mobile/AgentChat/AgentChatTranscriptService.swift
Original file line number Diff line number Diff line change
Expand Up @@ -15,7 +15,7 @@ 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).
/// Drives the live agent-prose streaming preview.
private var proseStreamer: AgentChatProseStreamer!
/// Sessions whose transcript could not be resolved; skipped until an
/// explicit history request retries, so per-hook-event resolution
Expand Down Expand Up @@ -155,12 +155,10 @@ final class AgentChatTranscriptService {
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.
// prompt starts the in-flight turn, Stop ends it.
switch event.hookEventName {
case .userPromptSubmit:
if AgentChatProseStreamingFlag.isEnabled,
record.state != .ended,
if record.state != .ended,
let surfaceID = record.surfaceID.flatMap(UUID.init(uuidString:)) {
proseStreamer.turnStarted(
sessionID: record.sessionID,
Expand Down
8 changes: 4 additions & 4 deletions cmux.xcodeproj/project.pbxproj
Original file line number Diff line number Diff line change
Expand Up @@ -11,7 +11,7 @@
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 */; };
CDFE000000000000000000C1 /* AgentChatProseStreamerTests.swift in Sources */ = {isa = PBXBuildFile; fileRef = CDFE000000000000000000C2 /* AgentChatProseStreamerTests.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 */; };
Expand Down Expand Up @@ -1277,7 +1277,7 @@
A9E010000000000000000007 /* agent-session-solid */ = {isa = PBXFileReference; lastKnownFileType = folder; path = "agent-session-solid"; sourceTree = "<group>"; };
ACA7C4A70000000000000004 /* AgentChatHookSessionStore.swift */ = {isa = PBXFileReference; includeInIndex = 1; lastKnownFileType = sourcecode.swift; path = "AgentChatHookSessionStore.swift"; sourceTree = "<group>"; };
CDFE000000000000000000A2 /* AgentChatProseStreamer.swift */ = {isa = PBXFileReference; includeInIndex = 1; lastKnownFileType = sourcecode.swift; path = "AgentChatProseStreamer.swift"; sourceTree = "<group>"; };
CDFE000000000000000000B2 /* AgentChatProseStreamingFlag.swift */ = {isa = PBXFileReference; includeInIndex = 1; lastKnownFileType = sourcecode.swift; path = "AgentChatProseStreamingFlag.swift"; sourceTree = "<group>"; };
CDFE000000000000000000C2 /* AgentChatProseStreamerTests.swift */ = {isa = PBXFileReference; lastKnownFileType = sourcecode.swift; path = AgentChatProseStreamerTests.swift; sourceTree = "<group>"; };
ACA7C4A70000000000000002 /* AgentChatSessionRecord.swift */ = {isa = PBXFileReference; includeInIndex = 1; lastKnownFileType = sourcecode.swift; path = "AgentChatSessionRecord.swift"; sourceTree = "<group>"; };
ACA7C4A70000000000000008 /* AgentChatSessionRegistry.swift */ = {isa = PBXFileReference; includeInIndex = 1; lastKnownFileType = sourcecode.swift; path = "AgentChatSessionRegistry.swift"; sourceTree = "<group>"; };
ACA7C4A70000000000000006 /* AgentChatTranscriptResolver.swift */ = {isa = PBXFileReference; includeInIndex = 1; lastKnownFileType = sourcecode.swift; path = "AgentChatTranscriptResolver.swift"; sourceTree = "<group>"; };
Expand Down Expand Up @@ -2590,7 +2590,6 @@
ACA7C4A7000000000000000A /* AgentChatTranscriptTailer.swift */,
ACA7C4A7000000000000000C /* AgentChatTranscriptService.swift */,
CDFE000000000000000000A2 /* AgentChatProseStreamer.swift */,
CDFE000000000000000000B2 /* AgentChatProseStreamingFlag.swift */,
);
name = AgentChat;
path = AgentChat;
Expand Down Expand Up @@ -3445,6 +3444,7 @@
F6572003A1B2C3D4E5F60718 /* SessionPersistenceResumeBindingTests.swift */,
F5000001A1B2C3D4E5F60718 /* SessionPersistenceTests.swift */,
B35750000000000000000009 /* PiVaultAgentPersistenceTests.swift */,
CDFE000000000000000000C2 /* AgentChatProseStreamerTests.swift */,
D3610B010000000000000002 /* AgentSessionAutoResumeSettingsTests.swift */,
D3610B020000000000000002 /* AgentSessionAutoResumeSwiftTests.swift */,
A9E010000000000000000005 /* AgentExecutableResolverTests.swift */,
Expand Down Expand Up @@ -4149,7 +4149,6 @@

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 */,
Expand Down Expand Up @@ -4978,6 +4977,7 @@
isa = PBXSourcesBuildPhase;
buildActionMask = 2147483647;
files = (
CDFE000000000000000000C1 /* AgentChatProseStreamerTests.swift in Sources */,
A9E020000000000000000005 /* AgentExecutableResolverTests.swift in Sources */,
D36A00020000000000000001 /* AgentHibernationTests.swift in Sources */,
D3610B010000000000000001 /* AgentSessionAutoResumeSettingsTests.swift in Sources */,
Expand Down
88 changes: 88 additions & 0 deletions cmuxTests/AgentChatProseStreamerTests.swift
Original file line number Diff line number Diff line change
@@ -0,0 +1,88 @@
import CmuxAgentChat
import Foundation
import XCTest

#if canImport(cmux_DEV)
@testable import cmux_DEV
#elseif canImport(cmux)
@testable import cmux
#endif

@MainActor
final class AgentChatProseStreamerTests: XCTestCase {
private actor SleepGate {
private var continuation: CheckedContinuation<Void, Never>?

func wait() async {
await withCheckedContinuation { continuation in
self.continuation = continuation
}
}

func resume() {
continuation?.resume()
continuation = nil
}
}

func testStreamsWhenLegacyUserDefaultsFlagIsFalse() async throws {
let legacyDefaultsKey = "CMUXAgentChatProseStreaming"
let previousLegacyValue = UserDefaults.standard.object(forKey: legacyDefaultsKey)
UserDefaults.standard.set(false, forKey: legacyDefaultsKey)
defer {
if let previousLegacyValue {
UserDefaults.standard.set(previousLegacyValue, forKey: legacyDefaultsKey)
} else {
UserDefaults.standard.removeObject(forKey: legacyDefaultsKey)
}
}

let surfaceID = UUID()
let sessionID = "session-with-always-on-streaming"
let expectedText = "The sky is blue."
let screenRows = [
"> Reply with one short sentence about blue.",
"",
expectedText,
"",
"Working (3s Esc to interrupt)",
"> ",
]

let emittedFrame = expectation(description: "streaming prose frame emitted")
let sleepGate = SleepGate()
var emittedFrames: [ChatSessionEventFrame] = []
let streamer = AgentChatProseStreamer(
emit: { frame in
if emittedFrames.isEmpty {
emittedFrame.fulfill()
}
emittedFrames.append(frame)
},
snapshot: { requestedSurfaceID in
requestedSurfaceID == surfaceID ? screenRows : nil
},
hasSubscribers: { true },
now: { Date(timeIntervalSince1970: 1_711_111_111) },
pollInterval: .seconds(60),
sleep: { _ in await sleepGate.wait() }
)

streamer.turnStarted(sessionID: sessionID, surfaceID: surfaceID, agentKind: .codex)
await fulfillment(of: [emittedFrame], timeout: 1.0)
streamer.turnEnded(sessionID: sessionID)
await sleepGate.resume()

let frame = try XCTUnwrap(emittedFrames.first)
XCTAssertEqual(frame.sessionID, sessionID)
guard case .streamingProse(let message?) = frame.event else {
return XCTFail("Expected a streaming prose preview frame")
}
XCTAssertEqual(message.id, "stream:\(sessionID)")
XCTAssertEqual(message.role, .agent)
guard case .prose(let prose) = message.kind else {
return XCTFail("Expected prose preview content")
}
XCTAssertEqual(prose.text, expectedText)
}
}
Loading