diff --git a/CodegiOS/Features/SessionDetail/LiveTurn.swift b/CodegiOS/Features/SessionDetail/LiveTurn.swift index fd446bc..6bf8fdc 100644 --- a/CodegiOS/Features/SessionDetail/LiveTurn.swift +++ b/CodegiOS/Features/SessionDetail/LiveTurn.swift @@ -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)" } } } @@ -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 = [] /// 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] = [] @@ -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 @@ -209,6 +239,20 @@ 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. @@ -216,7 +260,7 @@ final class LiveTurn: Identifiable { for segment in segments { switch segment { case .text(let run), .thinking(let run): run.flushNow() - case .tool: break + case .tool, .userMessage: break } } } @@ -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): @@ -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 } } diff --git a/CodegiOS/Features/SessionDetail/Rendering/ToolCallVM.swift b/CodegiOS/Features/SessionDetail/Rendering/ToolCallVM.swift index 3ac86f8..820a477 100644 --- a/CodegiOS/Features/SessionDetail/Rendering/ToolCallVM.swift +++ b/CodegiOS/Features/SessionDetail/Rendering/ToolCallVM.swift @@ -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. @@ -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 diff --git a/CodegiOS/Features/SessionDetail/SessionDetailViewModel.swift b/CodegiOS/Features/SessionDetail/SessionDetailViewModel.swift index 1a23eb3..70d4dfd 100644 --- a/CodegiOS/Features/SessionDetail/SessionDetailViewModel.swift +++ b/CodegiOS/Features/SessionDetail/SessionDetailViewModel.swift @@ -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) @@ -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() } diff --git a/CodegiOS/Features/SessionDetail/Timeline/TimelineNode.swift b/CodegiOS/Features/SessionDetail/Timeline/TimelineNode.swift index 859ce07..436d773 100644 --- a/CodegiOS/Features/SessionDetail/Timeline/TimelineNode.swift +++ b/CodegiOS/Features/SessionDetail/Timeline/TimelineNode.swift @@ -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)) } @@ -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-"), 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): diff --git a/CodegiOS/Models/AcpEvent.swift b/CodegiOS/Models/AcpEvent.swift index da36bb1..b6b5b65 100644 --- a/CodegiOS/Models/AcpEvent.swift +++ b/CodegiOS/Models/AcpEvent.swift @@ -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. @@ -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 @@ -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 } @@ -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",