Skip to content
Open
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
101 changes: 88 additions & 13 deletions CodegiOS/Features/SessionDetail/LiveTurn.swift
Original file line number Diff line number Diff line change
Expand Up @@ -64,20 +64,38 @@ final class LiveToolCall: Identifiable {
}
}

/// A message the user sent WHILE the turn was running, delivered to the agent
/// over the live-steering channel. Not agent output: it marks the point in the
/// stream where the user interrupted, so the transcript can show it as a user
/// message there and start the reply to it as a fresh run.
struct LiveUserMessage {
/// The feedback note's id. Also the dedup key — the server broadcasts a
/// submission to every attached client, including the one that sent it.
let id: String
let text: String
/// When the note was created (its `created_at`), so the message is stamped
/// with the moment it was sent rather than the start of the reply it
/// interrupted.
let sentAt: Date
}

/// One ordered piece of an in-flight assistant turn. The live turn is a sequence
/// of these so text, reasoning, and tool calls render in arrival order, exactly
/// as a finalized transcript would. The associated reference types expose a
/// `nonisolated let id`, so the enum's `id` needs no actor hop.
/// of these so text, reasoning, tool calls, and any message the user sent
/// mid-turn render in arrival order, exactly as a finalized transcript would.
/// The associated reference types expose a `nonisolated let id`, so the enum's
/// `id` needs no actor hop.
enum LiveSegment: Identifiable {
case text(LiveTextRun)
case thinking(LiveTextRun)
case tool(LiveToolCall)
case userMessage(LiveUserMessage)

var id: String {
switch self {
case .text(let run): return "t-\(run.id)"
case .thinking(let run): return "k-\(run.id)"
case .tool(let call): return "x-\(call.id)"
case .userMessage(let message): return "u-\(message.id)"
}
}
}
Expand Down Expand Up @@ -167,6 +185,9 @@ final class LiveTurn: Identifiable {
private(set) var segments: [LiveSegment] = []
/// Tool-call lookup so `tool_call_update` finds its target in O(1).
private var toolIndex: [String: LiveToolCall] = [:]
/// Feedback-note ids already spliced in as mid-turn user messages, so the
/// submit broadcast can be applied idempotently.
private var userMessageIDs: Set<String> = []
/// The agent's live plan/TODO list (`plan_update`). Replaced wholesale per
/// event (each carries the full list); rendered as a checklist above the turn.
var livePlan: [PlanEntry] = []
Expand All @@ -184,6 +205,15 @@ final class LiveTurn: Identifiable {
/// True before any content has streamed — used to show a "waiting" shimmer.
var isEmpty: Bool { segments.isEmpty && errorMessage == nil && livePlan.isEmpty }

/// True while there is nothing from the agent to show yet: before the turn's
/// first content, and again right after a mid-turn user message until the
/// reply to it starts. Drives the "waiting" shimmer at the tail of the turn,
/// so interrupting the agent doesn't leave the transcript looking stalled.
var isAwaitingOutput: Bool {
if case .userMessage? = segments.last { return true }
return isEmpty
}

/// Replace the live plan from a `plan_update` event (or snapshot `plan` block).
func updatePlan(_ entries: [PlanEntry]) {
livePlan = entries
Expand All @@ -209,14 +239,28 @@ final class LiveTurn: Identifiable {
}
}

/// Splice a message the user sent mid-turn into the stream at the point the
/// agent received it. It ends the run the agent was writing, so the reply to
/// it starts a fresh one instead of continuing the same paragraph — the same
/// split the persisted transcript makes for a mid-turn `user_message_chunk`,
/// so the live view and a reload agree.
///
/// Idempotent by note id: the submit is broadcast to every attached client,
/// and one of them may be the sender.
func appendUserMessage(id: String, text: String, sentAt: Date) {
guard !text.isEmpty, !userMessageIDs.contains(id) else { return }
userMessageIDs.insert(id)
segments.append(.userMessage(LiveUserMessage(id: id, text: text, sentAt: sentAt)))
}

/// Flush every text/reasoning run's pending coalesce immediately. Called when
/// the turn ends so the finalized render sees complete text (no trailing
/// reflow) and a snapshot taken right after is whole.
func flushAllText() {
for segment in segments {
switch segment {
case .text(let run), .thinking(let run): run.flushNow()
case .tool: break
case .tool, .userMessage: break
}
}
}
Expand Down Expand Up @@ -287,15 +331,37 @@ final class LiveTurn: Identifiable {
return nil
}

/// Snapshot this (finalized) turn as an immutable assistant `MessageTurn`, so a
/// reply that finished streaming but isn't yet in the server transcript can be
/// folded into the persisted list — surviving a subsequent send — until a
/// reconcile replaces it with the authoritative copy. Segment→block mapping
/// mirrors `MessageRender.adaptLive` so the folded copy renders the same; a
/// tool's diff `content` rides in the result preview (its only persisted slot)
/// when there's no separate output text.
func snapshotAsMessageTurn() -> MessageTurn {
/// Snapshot this (finalized) turn as immutable `MessageTurn`s, so a reply that
/// finished streaming but isn't yet in the server transcript can be folded into
/// the persisted list — surviving a subsequent send — until a reconcile replaces
/// it with the authoritative copy. Segment→block mapping mirrors
/// `MessageRender.adaptLive` so the folded copy renders the same; a tool's diff
/// `content` rides in the result preview (its only persisted slot) when there's
/// no separate output text.
///
/// Usually one assistant turn. A message the user sent mid-turn splits it, the
/// way the server's own projection does, so the folded copy keeps the message
/// in place and the two halves of the reply don't run together. Turns with no
/// renderable blocks are omitted: a turn can be non-empty *solely* because of an
/// inline error / "Cancelled." message, which `ContentBlock` has no case for,
/// and folding in that zero-block turn would render as "No content".
func snapshotAsMessageTurns() -> [MessageTurn] {
var out: [MessageTurn] = []
var blocks: [ContentBlock] = []

func flushAssistant() {
guard !blocks.isEmpty else { return }
// The first turn keeps the live turn's own id so a `List` row that was
// rendering the stream keeps its identity across the handoff.
out.append(MessageTurn(
id: out.isEmpty ? self.id : "\(self.id)#\(out.count)",
role: .assistant,
blocks: blocks,
timestamp: Date()
))
blocks = []
}

for segment in segments {
switch segment {
case .text(let run):
Expand All @@ -310,8 +376,17 @@ final class LiveTurn: Identifiable {
blocks.append(.toolUse(id: call.id, name: call.title, inputPreview: call.rawInput, meta: call.meta))
let output = call.rawOutput.isEmpty ? call.content : call.rawOutput
blocks.append(.toolResult(id: call.id, outputPreview: output, isError: call.isError))
case .userMessage(let message):
flushAssistant()
out.append(MessageTurn(
id: "u-\(message.id)",
role: .user,
blocks: [.text(message.text)],
timestamp: message.sentAt
))
}
}
return MessageTurn(id: id, role: .assistant, blocks: blocks, timestamp: Date())
flushAssistant()
return out
}
}
11 changes: 11 additions & 0 deletions CodegiOS/Features/SessionDetail/Rendering/ToolCallVM.swift
Original file line number Diff line number Diff line change
Expand Up @@ -93,6 +93,10 @@ enum RenderPart {
case reasoning(text: String) // finalized reasoning
case liveText(LiveTextRun, streaming: Bool) // streaming prose (plain, @Bindable)
case liveReasoning(LiveTextRun, streaming: Bool)
/// A message the user sent while the reply was streaming, carried as the user
/// turn it renders as. It breaks the grouping runs below, so the answer to it
/// never folds into what the agent was saying before it arrived.
case userMessage(MessageTurn)
case tool(ToolCallVM)
case toolGroup(items: [ToolCallVM], streaming: Bool)
/// A run of consecutive `get_delegation_status` polls, merged into one card.
Expand Down Expand Up @@ -185,6 +189,13 @@ enum MessageRender {
parts.append(.liveText(run, streaming: turn.isStreaming && isLast))
case .thinking(let run):
parts.append(.liveReasoning(run, streaming: turn.isStreaming && isLast))
case .userMessage(let message):
parts.append(.userMessage(MessageTurn(
id: "u-\(message.id)",
role: .user,
blocks: [.text(message.text)],
timestamp: message.sentAt
)))
case .tool(let call):
let state: ToolCallState = call.isFinished ? (call.isError ? .error : .done) : .running
// Same boundary-marker treatment as the persisted path: codex emits
Expand Down
33 changes: 24 additions & 9 deletions CodegiOS/Features/SessionDetail/SessionDetailViewModel.swift
Original file line number Diff line number Diff line change
Expand Up @@ -1275,6 +1275,22 @@ final class SessionDetailViewModel {
// The server echoes our own prompt; we already showed it optimistically.
break

case .feedbackSubmitted(let note):
// A note that is already `delivered` on submission was pushed INTO the
// running turn over the native steering channel, so the agent has the
// text as a user message — show it as one, where it interrupted. A
// `pending` note is the cooperative `check_user_feedback` pull channel:
// the agent reads it as a tool result, never as a message, and it gets
// no user turn on reload either.
//
// `isTurnActive` is the "there is a turn to splice into" guard. The
// submission is recorded ungated, so a note can land just after the
// turn settled, and appending there would graft it onto a finished
// reply. The agent recorded it either way, so a reload still shows it.
guard isTurnActive, note.status == "delivered" else { break }
live.appendUserMessage(id: note.id, text: note.text, sentAt: note.createdAt ?? Date())
requestScrollToBottom()

case .turnComplete(let stopReason):
finalize(live: live, stopReason: stopReason)

Expand Down Expand Up @@ -1675,16 +1691,15 @@ final class SessionDetailViewModel {
guard liveTurn === live else { return }
turns.append(contentsOf: pendingUserTurns)
pendingUserTurns.removeAll()
// Only fold in an assistant turn that actually has renderable content. A
// Usually one assistant turn; a message the user sent mid-turn splits it
// into the same user/assistant sequence the server's projection produces.
// Turns with no renderable content are dropped by the snapshot itself: a
// finalized turn can be non-empty *solely* because of an inline error /
// "Cancelled." message (which `snapshotAsMessageTurn` can't represent as a
// persisted block, since ContentBlock has no error case) — appending its
// zero-block snapshot would render as "No content". Such a transient error
// placeholder is simply dropped on the next send; the user turn is kept.
let snapshot = live.snapshotAsMessageTurn()
if !snapshot.blocks.isEmpty {
turns.append(snapshot)
}
// "Cancelled." message (which it can't represent as a persisted block,
// since ContentBlock has no error case) and appending that zero-block turn
// would render as "No content". Such a transient error placeholder is
// simply dropped on the next send; the user turn is kept.
turns.append(contentsOf: live.snapshotAsMessageTurns())
liveTurn = nil
requestScrollToBottom()
}
Expand Down
15 changes: 12 additions & 3 deletions CodegiOS/Features/SessionDetail/Timeline/TimelineNode.swift
Original file line number Diff line number Diff line change
Expand Up @@ -187,12 +187,15 @@ enum TranscriptTimeline {
if !live.livePlan.isEmpty {
out.append(TimelineNode(id: "\(live.id)#plan", content: .plan(live.livePlan, streaming: live.isStreaming), agent: agent))
}
if live.isEmpty && live.isStreaming {
out.append(TimelineNode(id: "\(live.id)#thinking", content: .thinking, agent: agent))
}
out.append(contentsOf: MessageRender.adaptLive(live).enumerated().map { idx, part in
node(for: part, ownerID: live.id, index: idx, agent: agent)
})
// The "waiting" shimmer sits at the TAIL, after whatever has arrived so
// far: nothing has, or the user interrupted and the agent hasn't answered
// yet. (With no parts it lands exactly where it always did.)
if live.isStreaming && live.isAwaitingOutput {
out.append(TimelineNode(id: "\(live.id)#thinking", content: .thinking, agent: agent))
}
if let error = live.errorMessage {
out.append(TimelineNode(id: "\(live.id)#error", content: .error(error), agent: agent))
}
Expand All @@ -218,6 +221,12 @@ enum TranscriptTimeline {
return TimelineNode(id: "t-\(run.id)", content: .liveText(run, streaming: streaming), agent: agent)
case .liveReasoning(let run, let streaming):
return TimelineNode(id: "k-\(run.id)", content: .liveReasoning(run, streaming: streaming), agent: agent)
case .userMessage(let turn):
// Keyed by the feedback note (`turn.id` is "u-<note id>"), not the
// part index, so the row keeps its identity as the reply grows past
// it. `startsGroup` matches a persisted user turn: interrupting the
// agent begins a new exchange.
return TimelineNode(id: turn.id, content: .user(turn), agent: agent, startsGroup: true)
case .tool(let vm):
return TimelineNode(id: "tool-\(vm.id)", content: .tool(vm), agent: agent)
case .toolGroup(let items, let streaming):
Expand Down
42 changes: 41 additions & 1 deletion CodegiOS/Models/AcpEvent.swift
Original file line number Diff line number Diff line change
Expand Up @@ -40,6 +40,34 @@ enum UserMessageBlock: Hashable, Sendable, Decodable {
}
}

/// A live-feedback note (Rust `FeedbackItem`) — a short message the user sent
/// while the agent was mid-turn.
///
/// `status` decides what it *is*. A note that is already `delivered` when it is
/// submitted was pushed into the running turn over the native steering channel,
/// so the agent has the text as a user message. A `pending` one is waiting for
/// the agent to pull it with `check_user_feedback`, which returns it as a tool
/// result and never as a message. Kept as the raw string so an unknown future
/// status decodes instead of throwing.
struct FeedbackNote: Hashable, Sendable, Decodable {
let id: String
let text: String
let status: String
/// When the note was created, stamped on the machine the agent runs on.
/// Optional: an unparseable instant must not cost us the message.
let createdAt: Date?

private enum CodingKeys: String, CodingKey { case id, text, status, createdAt }

init(from decoder: Decoder) throws {
let c = try decoder.container(keyedBy: CodingKeys.self)
id = try c.decodeIfPresent(String.self, forKey: .id) ?? ""
text = try c.decodeIfPresent(String.self, forKey: .text) ?? ""
status = try c.decodeIfPresent(String.self, forKey: .status) ?? ""
createdAt = (try? c.decodeIfPresent(Date.self, forKey: .createdAt)) ?? nil
}
}

/// A backend→client ACP event (Rust `AcpEvent`, internally tagged `type`).
/// Decode-only. Unknown event types decode to `.unknown` so a new server event
/// never breaks the stream.
Expand All @@ -56,6 +84,10 @@ enum AcpEvent: Hashable, Sendable, Decodable {
case usageUpdate(used: UInt64, size: UInt64)
case userMessage(messageId: String, blocks: [UserMessageBlock])
case userPromptSent(textPreview: String)
/// The user submitted a live-feedback note while the agent is mid-turn.
/// Broadcast to every client watching the conversation, including the one
/// that sent it.
case feedbackSubmitted(note: FeedbackNote)
case error(message: String, code: String?)
/// Agent asks the user to approve a tool call before it runs. Also carries
/// ExitPlanMode — the proposed plan rides inside `toolCall`. Resolve via
Expand Down Expand Up @@ -84,7 +116,7 @@ enum AcpEvent: Hashable, Sendable, Decodable {
case type, text, title, kind, status, content, meta
case toolCallId, rawInput, rawOutput, rawOutputAppend
case stopReason, sessionId, conversationId, folderId
case used, size, messageId, blocks, message, code, textPreview
case used, size, messageId, blocks, message, code, textPreview, item
case requestId, toolCall, options, questionId, questions, entries
case approvalId, planMarkdown
}
Expand Down Expand Up @@ -152,6 +184,14 @@ enum AcpEvent: Hashable, Sendable, Decodable {
)
case "user_prompt_sent":
self = .userPromptSent(textPreview: try c.decodeIfPresent(String.self, forKey: .textPreview) ?? "")
case "feedback_submitted":
// No sensible stand-in for a missing note, so an event without one
// degrades to `.unknown` rather than inventing an empty message.
if let note = try c.decodeIfPresent(FeedbackNote.self, forKey: .item) {
self = .feedbackSubmitted(note: note)
} else {
self = .unknown(type: type)
}
case "error":
self = .error(
message: try c.decodeIfPresent(String.self, forKey: .message) ?? "Unknown error",
Expand Down