diff --git a/src-tauri/src/commands/conversations.rs b/src-tauri/src/commands/conversations.rs index 3cd1c35954..5bdfb08bfd 100644 --- a/src-tauri/src/commands/conversations.rs +++ b/src-tauri/src/commands/conversations.rs @@ -1602,6 +1602,21 @@ fn sig_from_turn_blocks(blocks: &[ContentBlock]) -> Option> Some(sig) } +/// How many USER turns [`apply_in_flight_message_id`]'s walk will compare +/// before giving up. +/// +/// A cost bound, not a correctness one — the ambiguity check inside the walk is +/// what keeps it from stamping an earlier round's prompt. `sig_from_turn_blocks` +/// copies each candidate's text and image bytes, and the walk runs on every +/// detail fetch, so a transcript whose timestamps are all parse instants (see +/// the walk) must not turn that into a scan of every prompt ever sent. +/// +/// Counts USER turns only, deliberately. Counting turns would count the length +/// of the reply, which is unrelated to the hazard and routinely in the +/// hundreds: a real window holds the prompt and whatever the user sent +/// mid-turn, so 32 is never reached by a turn with usable timestamps. +const MAX_IN_FLIGHT_WALK_USER_TURNS: usize = 32; + /// Stamp the persisted in-flight user turn with the broadcast `message_id`. /// /// A cross-client viewer renders the in-flight prompt from two sources that use @@ -1611,12 +1626,23 @@ fn sig_from_turn_blocks(blocks: &[ContentBlock]) -> Option> /// frontend's id-dedup collapse the two into one instead of showing the prompt /// twice. /// -/// The in-flight prompt is located tail-bounded: -/// - the trailing user turn (Claude/Codex write the assistant turn only on -/// completion, so mid-stream the transcript ends exactly at the prompt); or -/// - the user turn immediately before a *single* trailing assistant turn -/// (OpenCode and Gemini persist a partial assistant turn mid-stream, so the -/// transcript tail is `[.., user X, partial assistant Y]`). +/// The in-flight prompt is located by walking back from the transcript tail +/// over the turns this turn produced — those persisted at/after `started_at`, +/// comparing at most [`MAX_IN_FLIGHT_WALK_USER_TURNS`] user turns of them — and +/// keeping the EARLIEST user turn whose content matches. The tail itself is the +/// prompt only for Claude/Codex, which write the assistant turn on completion; +/// OpenCode and Gemini persist a partial reply mid-stream (which a parser may +/// split), and a message the user sends mid-turn is written after the prompt as +/// a user turn of its own. Earliest, because the agent writes the prompt before +/// anything it produces in reply, so a mid-turn message repeating the prompt's +/// own words cannot take the stamp from it. +/// +/// That "earliest" only orders the round's own turns, which presumes the walk +/// stopped at the round's start. It does when the gate below fires. When it +/// never fires — every turn in hand is at/after `started_at`, which a parser +/// stamping parse instants makes routine — the walk saw the whole transcript +/// and earliest means nothing, so a second matching copy leaves it unable to +/// say which round it is in: it then stamps nothing rather than guess. /// /// A recency check then disambiguates: the in-flight prompt was persisted by the /// agent CLI at/after `started_at` (the agent — a local subprocess sharing this @@ -1660,14 +1686,39 @@ fn apply_in_flight_message_id( return None; } let started_at = started_at?; - let target_idx = match turns[n - 1].role { - TurnRole::User => n - 1, - TurnRole::Assistant if n >= 2 && matches!(turns[n - 2].role, TurnRole::User) => n - 2, - _ => return None, - }; - // Recency gate. `started_at` is recorded when the backend broadcasts the - // `UserMessage` event, which happens *before* the agent request is issued - // (see `connection.rs`), so the agent — a local subprocess on this machine's + let want = sig_from_user_message_blocks(&pending.blocks); + + // Walk back over the turns THIS turn produced and keep the EARLIEST user + // turn whose content is the pending prompt's. + // + // A walk, not the tail. The prompt is the last turn only while the agent + // has written nothing else, and what trails it is not bounded to one + // assistant turn: + // + // * a message the user sends MID-TURN (`/steering`) is written into the + // transcript as a USER turn after the prompt, so the tail becomes that + // message, whose content is not the prompt's; + // * OpenCode and Gemini persist the reply as it goes, and a parser that + // splits it leaves two or more assistant turns behind the prompt. + // + // Anchoring on the last one or two turns lost the stamp in the middle of + // exactly those turns, and every consumer reads a missing stamp as "this + // detail is settled, the turn is over": `computeTimelinePrefix` stops + // hiding the persisted half of the reply the live stream is re-showing (so + // the first half renders twice), the runtime store's `detailIsInFlight` + // lets a mid-turn refetch clear `liveMessage` / `localTurns` / + // `optimisticTurns`, and `collectInFlightPersistedToolCalls` stops marking + // the round's unfinished tool calls, which then paint as completed. + // + // EARLIEST, not last: the agent writes the prompt before anything it + // produces in reply, so within one turn the first copy of those words is + // the prompt itself. That is what keeps a mid-turn message repeating the + // prompt's own text ("continue" twice) from taking the stamp off it. + // + // Recency gate, which is both what makes the walk safe and what bounds it. + // `started_at` is recorded when the backend broadcasts the `UserMessage` + // event, which happens *before* the agent request is issued (see + // `connection.rs`), so the agent — a local subprocess on this machine's // clock — necessarily persists the in-flight prompt at a wall-clock instant // at or after `started_at`. A *prior* identical prompt was persisted during // an earlier turn and is therefore strictly older. We allow no backward @@ -1676,29 +1727,85 @@ fn apply_in_flight_message_id( // second), and stamping it would HIDE the genuinely new prompt via the // frontend's keep-first user dedup. Erring the other way only ever yields a // recoverable visible duplicate, so the strict bound is the safe one. - if turns[target_idx].timestamp < started_at { - return None; - } - let want = sig_from_user_message_blocks(&pending.blocks); - if sig_from_turn_blocks(&turns[target_idx].blocks) == Some(want) { - // Never create a duplicate id. The broadcast id is normally disjoint from - // parser `turn-N` ids (and `is_reserved_turn_id` in the manager rejects a - // client id of that shape), but defend the invariant here too: if the id - // already exists on another turn, stamping would make two turns share an - // id and the frontend's id-keyed dedup could hide one. Leave the turn - // under its parser id — a recoverable visible duplicate, never a hidden - // prompt — and report nothing. - let collides = turns - .iter() - .enumerate() - .any(|(i, t)| i != target_idx && t.id == pending.message_id); - if collides { + // + // Turns are in transcript order, so the first one older than the start ends + // the search — the walk never reads past the running turn, and an + // out-of-order timestamp can only cut it short, which stamps nothing. + // + // …unless the timestamps are not the agent's at all. Two shipping parsers + // fall back to the PARSE INSTANT when a record carries no usable time: + // `cline.rs` seeds `last_ts` from `Utc::now()` when the manifest has no + // `started_at` and hands it to every message with `ts` missing or 0, and + // `antigravity.rs` writes `ts.unwrap_or_else(Utc::now)` per step. A parse + // instant is by construction at/after `started_at`, so there the gate never + // fires, the "window" is the whole transcript, and "earliest match" could + // reach an identical prompt from an earlier round — stamping THAT makes + // `visiblePersistedTurns` hide every assistant turn after it. + // + // `gate_fired` is what distinguishes the two. It is false in exactly two + // shapes: the transcript holds nothing older than this turn (a first turn — + // there is no earlier round to confuse anything with), or the timestamps are + // parse instants (there may be). Both are covered by asking whether the + // decision was AMBIGUOUS: one matching user turn in the whole transcript is + // the prompt whatever the clocks say, because the text is then unique to it. + // Two or more, with nothing proving where this turn began, is a guess — and + // guessing wrong here hides turns, so it refuses instead. + // + // NOT a cap on turns walked. The obvious bound — stop after N turns — counts + // the LENGTH OF THE REPLY, which is the one quantity that has nothing to do + // with the hazard: every mid-turn-persisting parser emits one assistant turn + // per assistant record (`claude.rs`, `opencode.rs`, `gemini.rs` all keep + // turns small for virtualization), so a measured Claude round runs to a + // median of 44 and a p90 of 317. Any N small enough to bound the hazard + // voids the stamp part-way through most ordinary rounds, mid-turn, which is + // the exact failure this walk exists to remove. + let mut target_idx: Option = None; + let mut matched = 0usize; + let mut user_turns_seen = 0usize; + let mut gate_fired = false; + for i in (0..n).rev() { + if turns[i].timestamp < started_at { + gate_fired = true; + break; + } + if !matches!(turns[i].role, TurnRole::User) { + continue; + } + // A cost bound, not a correctness one — `sig_from_turn_blocks` copies + // each turn's text and image bytes, and this runs on every detail + // fetch. Only reachable when the gate never fires (a real window holds + // the prompt plus the handful of messages sent mid-turn), and refusing + // there is the same safe direction as the ambiguity check below. + user_turns_seen += 1; + if user_turns_seen > MAX_IN_FLIGHT_WALK_USER_TURNS { return None; } - turns[target_idx].id = pending.message_id.clone(); - return Some(pending.message_id.clone()); + if sig_from_turn_blocks(&turns[i].blocks).as_ref() == Some(&want) { + matched += 1; + target_idx = Some(i); + } + } + if !gate_fired && matched > 1 { + return None; } - None + let target_idx = target_idx?; + + // Never create a duplicate id. The broadcast id is normally disjoint from + // parser `turn-N` ids (and `is_reserved_turn_id` in the manager rejects a + // client id of that shape), but defend the invariant here too: if the id + // already exists on another turn, stamping would make two turns share an + // id and the frontend's id-keyed dedup could hide one. Leave the turn + // under its parser id — a recoverable visible duplicate, never a hidden + // prompt — and report nothing. + let collides = turns + .iter() + .enumerate() + .any(|(i, t)| i != target_idx && t.id == pending.message_id); + if collides { + return None; + } + turns[target_idx].id = pending.message_id.clone(); + Some(pending.message_id.clone()) } /// Resolve the raw `tailTurns` / `fromIndex` request fields into a window @@ -2917,32 +3024,195 @@ mod tests { } #[test] - fn does_not_reach_back_past_the_last_two_turns() { + fn does_not_reach_back_into_an_earlier_round() { // The matching prompt sits buried before another full user/assistant - // round; only the trailing user turn or the user-before-trailing- - // assistant are eligible, so it is never stamped. + // round. The walk is bounded by the recency gate, not by a turn count: + // it stops at the first turn older than this turn's start, which is the + // completed reply above — so the identical prompt behind it is out of + // reach however short the transcript is. let mut turns = vec![ user_text_turn("turn-0", "hello", at(-30)), assistant_text_turn("turn-1", "a", at(-29), true), user_text_turn("turn-2", "ok", at(1)), assistant_text_turn("turn-3", "b", at(2), false), ]; - apply_in_flight_message_id(&mut turns, &pending_text("msg-live", "hello"), Some(turn_started())); + let stamped = + apply_in_flight_message_id(&mut turns, &pending_text("msg-live", "hello"), Some(turn_started())); + assert_eq!(stamped, None, "nothing in the running turn matches"); assert_eq!(turns[0].id, "turn-0"); assert_eq!(turns[2].id, "turn-2", "non-matching tail user turn untouched"); } #[test] - fn does_not_stamp_with_two_trailing_assistant_turns() { - // Bounded to a single trailing assistant: a deeper assistant tail means - // we can't be sure the user prompt is the in-flight one, so bail. + fn stamps_nothing_when_the_clocks_cannot_say_which_continue_this_is() { + // `cline.rs` and `antigravity.rs` fall back to the PARSE INSTANT when a + // record carries no usable time, and a parse instant is by construction + // at or after `started_at` — so for those the recency gate never fires + // and NOTHING here says where this turn began. Two rounds of "continue" + // are then indistinguishable from one round the user steered with the + // same word, and stamping the older one would make + // `visiblePersistedTurns` hide every assistant turn after it. + let parsed_at = at(1); + let mut turns = vec![ + user_text_turn("turn-0", "continue", parsed_at), + assistant_text_turn("turn-1", "done", parsed_at, true), + user_text_turn("turn-2", "continue", parsed_at), + ]; + let stamped = apply_in_flight_message_id( + &mut turns, + &pending_text("msg-live", "continue"), + Some(turn_started()), + ); + assert_eq!(stamped, None, "ambiguous and unprovable → stamp nothing"); + assert_eq!(turns[0].id, "turn-0", "untouched"); + assert_eq!(turns[2].id, "turn-2", "untouched"); + } + + #[test] + fn stamps_an_unambiguous_match_even_with_no_proof_of_where_the_turn_began() { + // Same parse-instant transcript, but the prompt's text occurs once. A + // single candidate in the WHOLE transcript cannot be an earlier round's + // prompt confused with this one — the text is unique to it — so the + // clocks have nothing left to disambiguate and the stamp is safe. + // + // This is also the ordinary first turn of a conversation, where there is + // simply nothing older than `started_at` for the gate to find. + let parsed_at = at(1); + let mut turns = vec![ + user_text_turn("turn-0", "run the tests", parsed_at), + assistant_text_turn("turn-1", "first half", parsed_at, false), + user_text_turn("turn-2", "also lint", parsed_at), + ]; + let stamped = apply_in_flight_message_id( + &mut turns, + &pending_text("msg-live", "run the tests"), + Some(turn_started()), + ); + assert_eq!(stamped.as_deref(), Some("msg-live")); + assert_eq!(turns[0].id, "msg-live"); + } + + #[test] + fn a_long_reply_never_costs_the_stamp() { + // The bound counts USER turns, not turns. Every mid-turn-persisting + // parser emits one assistant turn per assistant record, so a real + // Claude round runs to a median of 44 of them and a p90 of 317 — a + // bound on turns walked would void the stamp part-way through most + // ordinary rounds, mid-turn, which is the failure this walk removes. + let mut turns = vec![ + user_text_turn("old", "earlier", at(-30)), + assistant_text_turn("old-a", "done", at(-29), true), + user_text_turn("turn-0", "hello", at(1)), + ]; + for i in 0..400 { + turns.push(assistant_text_turn( + &format!("a-{i}"), + "step", + at(2), + false, + )); + } + let stamped = + apply_in_flight_message_id(&mut turns, &pending_text("msg-live", "hello"), Some(turn_started())); + assert_eq!(stamped.as_deref(), Some("msg-live")); + assert_eq!(turns[2].id, "msg-live"); + } + + #[test] + fn stops_comparing_prompts_once_the_walk_is_plainly_not_in_one_turn() { + // The cost bound. Only reachable when the gate never fires, since a real + // window holds the prompt plus whatever was sent mid-turn; here every + // turn is a user turn with a parse instant, so the walk would otherwise + // rebuild a content signature for every prompt ever sent, on every + // detail fetch. Refusing is the same safe direction as ambiguity. + let parsed_at = at(1); + let mut turns = vec![user_text_turn("wanted", "hello", parsed_at)]; + for i in 0..MAX_IN_FLIGHT_WALK_USER_TURNS { + turns.push(user_text_turn(&format!("u-{i}"), "filler", parsed_at)); + } + let stamped = + apply_in_flight_message_id(&mut turns, &pending_text("msg-live", "hello"), Some(turn_started())); + assert_eq!(stamped, None, "past the cost bound, stamp nothing"); + assert_eq!(turns[0].id, "wanted", "untouched"); + + // One fewer and the same transcript is inside the bound, so it is the + // bound that refused above and not some other gate. + turns.pop(); + let stamped = + apply_in_flight_message_id(&mut turns, &pending_text("msg-live", "hello"), Some(turn_started())); + assert_eq!(stamped.as_deref(), Some("msg-live")); + } + + #[test] + fn stamps_the_prompt_behind_a_reply_the_parser_split() { + // Previously refused: the rule was "the tail, or the user before a + // SINGLE trailing assistant", so a deeper assistant tail bailed. + // + // That bound predates the recency gate and is redundant beside it. A + // user turn at or after this turn's start was persisted DURING it, and + // its content is the pending prompt's — there is nothing else it could + // be, whatever the agent has written since. Refusing here instead cost + // the stamp for the whole of every OpenCode/Gemini turn whose partial + // reply the parser split in two, which is exactly the shape the + // frontend's partial suppression exists for. let mut turns = vec![ user_text_turn("turn-0", "hello", at(1)), assistant_text_turn("turn-1", "a", at(2), false), assistant_text_turn("turn-2", "b", at(3), false), ]; - apply_in_flight_message_id(&mut turns, &pending_text("msg-live", "hello"), Some(turn_started())); - assert_eq!(turns[0].id, "turn-0", "left untouched"); + let stamped = + apply_in_flight_message_id(&mut turns, &pending_text("msg-live", "hello"), Some(turn_started())); + assert_eq!(stamped.as_deref(), Some("msg-live")); + assert_eq!(turns[0].id, "msg-live"); + } + + #[test] + fn stamps_the_prompt_behind_a_message_sent_mid_turn() { + // The steering shape. The agent writes a message the user sent DURING + // the turn into its own transcript, so the tail is that message and its + // content is not the prompt's. The old tail rule reported nothing here, + // in the middle of the turn, and every consumer reads that as "settled". + let mut turns = vec![ + user_text_turn("turn-0", "hello", at(1)), + assistant_text_turn("turn-1", "first half", at(2), false), + user_text_turn("turn-2", "also check the tests", at(3)), + ]; + let stamped = + apply_in_flight_message_id(&mut turns, &pending_text("msg-live", "hello"), Some(turn_started())); + assert_eq!(stamped.as_deref(), Some("msg-live")); + assert_eq!(turns[0].id, "msg-live", "the prompt keeps the stamp"); + assert_eq!(turns[2].id, "turn-2", "the mid-turn message is untouched"); + } + + #[test] + fn stamps_the_prompt_not_a_mid_turn_message_repeating_its_words() { + // "continue" is the most repeatable thing a person steers with, and the + // prompt itself may have been the same word. Both copies match by + // content and both are inside the turn, so recency cannot separate + // them — ORDER does: the agent writes the prompt before anything it + // produces, so the earliest copy in the turn is the prompt. Stamping + // the later one would move the anchor past the reply's first half and + // leave it beside the live copy of itself. + // + // An earlier round sits in front, which is what lets the recency gate + // fire and prove where this turn begins; without that proof the two + // copies are ambiguous and the walk refuses instead (see below). + let mut turns = vec![ + user_text_turn("old", "hello", at(-30)), + assistant_text_turn("old-a", "hi", at(-29), true), + user_text_turn("turn-0", "continue", at(1)), + assistant_text_turn("turn-1", "working", at(2), false), + user_text_turn("turn-2", "continue", at(3)), + ]; + let stamped = apply_in_flight_message_id( + &mut turns, + &pending_text("msg-live", "continue"), + Some(turn_started()), + ); + assert_eq!(stamped.as_deref(), Some("msg-live")); + assert_eq!(turns[2].id, "msg-live", "the earliest copy is the prompt"); + assert_eq!(turns[4].id, "turn-2", "the mid-turn repeat is untouched"); + assert_eq!(turns[0].id, "old", "the earlier round is untouched"); } #[test] diff --git a/src/contexts/conversation-runtime-context.test.tsx b/src/contexts/conversation-runtime-context.test.tsx index 2c635ba5c5..602c9e324d 100644 --- a/src/contexts/conversation-runtime-context.test.tsx +++ b/src/contexts/conversation-runtime-context.test.tsx @@ -709,6 +709,14 @@ describe("ConversationRuntimeProvider delegation kickoff projection", () => { mockGetFolderConversation.mockImplementation(() => new Promise(() => {})) }) + /** The live message the strip stands in for: one that IS showing the reply. */ + const streamingReply: LiveMessage = { + id: "lm-streaming", + role: "assistant", + content: [{ type: "text", text: "working on it" }], + startedAt: 0, + } + it("synthesizes the kickoff user turn (and strips the persisted reply) while the transcript has no user turn yet", async () => { // DB lags: only a partial assistant turn is persisted, no user turn. mockGetFolderConversation.mockResolvedValueOnce( @@ -721,7 +729,7 @@ describe("ConversationRuntimeProvider delegation kickoff projection", () => { api().setLiveOwnsActiveTurn(99, true, "do the thing") }) act(() => { - api().setLiveMessage(99, LIVE_MSG, true) + api().setLiveMessage(99, streamingReply, true) }) await act(async () => { api().refetchDetail(99, { preserveLive: true }) @@ -744,6 +752,35 @@ describe("ConversationRuntimeProvider delegation kickoff projection", () => { ).toBe(false) }) + it("keeps the persisted reply while the live message is showing nothing", async () => { + // The child's next turn has begun: `status_changed → prompting` put a fresh + // `content: []` live message on the connection and the viewer bridged it, + // but no chunk has arrived. Stripping the reply then leaves the dialog + // showing a prompt with nothing under it. + mockGetFolderConversation.mockResolvedValueOnce( + detailWithTurns([userTurn("u1"), assistantTurn("a1")]) + ) + renderProvider() + const api = () => runtimeHolder.current! + + act(() => { + api().setLiveOwnsActiveTurn(99, true, "do the thing") + }) + act(() => { + api().setLiveMessage(99, LIVE_MSG, true) + }) + await act(async () => { + api().refetchDetail(99, { preserveLive: true }) + await Promise.resolve() + }) + + expect( + api() + .getTimelineTurns(99) + .map((t) => t.turn.id) + ).toEqual(["u1", "a1"]) + }) + it("uses the real persisted user turn instead of synthesizing once it has landed", async () => { mockGetFolderConversation.mockResolvedValueOnce( detailWithTurns([userTurn("u1"), assistantTurn("a1")]) diff --git a/src/stores/conversation-runtime-store.ts b/src/stores/conversation-runtime-store.ts index 8333981d76..114796e130 100644 --- a/src/stores/conversation-runtime-store.ts +++ b/src/stores/conversation-runtime-store.ts @@ -2826,8 +2826,12 @@ interface TimelinePrefixDeps { optimisticTurns: MessageTurn[] liveOwnsActiveTurn: boolean delegationKickoffText: string | null - hasLiveMessage: boolean + liveShowsReply: boolean liveStartedAt: number | null + /** The steered-copy set, flattened for the `===` comparison below. `""` in + * an ordinary turn — see `computeTimelinePrefix`. */ + steeredCopyKey: string + liveStreamedRoundStart: boolean } interface TimelinePrefixEntry { deps: TimelinePrefixDeps @@ -3092,8 +3096,10 @@ function timelinePrefixDepsEqual( a.optimisticTurns === b.optimisticTurns && a.liveOwnsActiveTurn === b.liveOwnsActiveTurn && a.delegationKickoffText === b.delegationKickoffText && - a.hasLiveMessage === b.hasLiveMessage && - a.liveStartedAt === b.liveStartedAt + a.liveShowsReply === b.liveShowsReply && + a.liveStartedAt === b.liveStartedAt && + a.steeredCopyKey === b.steeredCopyKey && + a.liveStreamedRoundStart === b.liveStreamedRoundStart ) } @@ -3141,9 +3147,40 @@ function collectInFlightPersistedToolCalls( return out } +/** + * @param liveShowsReply whether the live message this session holds actually + * produced an assistant turn — see [`computeTimeline`], which derives it from + * the same build the streaming tail is made of. Both suppressions below hide a + * persisted assistant turn *because the live stream is showing that reply*, so + * they must key off what the live message RENDERS, never off the existence of a + * live message object. Those two differ: a live message with nothing renderable + * in it is an ordinary state (`STATUS_CHANGED` → `prompting` installs + * `content: []` at the start of every turn, and the runtime mirror never writes + * a null back over it), and keying off the object hid a reply with nothing put + * in its place — a blank agent turn. + * @param steeredCopyIds the detail's own copies of this turn's mid-turn + * messages (see `collectSteeredPersistedCopyIds`). A steered message is a USER + * turn that the agent writes into the MIDDLE of a round, so both round anchors + * below — the viewer's persisted-tail strip and the in-flight partial + * suppression — have to look straight through it. Anchored on it instead, they + * read the round as having ended at the interruption and leave the reply's + * already-persisted first half beside the live copy of the same text, which + * `mergeConsecutiveAssistantTurns` then glues into one run-on bubble. + * @param liveStreamedRoundStart whether this session's live message carries the + * reply from BEFORE the first mid-turn message, i.e. it holds the whole round + * and can stand in for every persisted turn of it. Both anchor adjustments are + * gated on it: they hide persisted assistant turns, so they may only run where + * the live stream is provably re-showing them. False for a session that first + * saw this turn from a snapshot (which carries no steering block at all, so no + * copies are identified either) and for one whose live message opens on the + * interruption itself. + */ function computeTimelinePrefix( session: ConversationRuntimeSession, - conversationId: number + conversationId: number, + liveShowsReply: boolean, + steeredCopyIds: ReadonlySet | null, + liveStreamedRoundStart: boolean ): TimelinePrefixEntry { const detail = session.detail // Everything Phases 1–3 read, snapshotted for the `===` validity check. @@ -3157,8 +3194,12 @@ function computeTimelinePrefix( optimisticTurns: session.optimisticTurns, liveOwnsActiveTurn: session.liveOwnsActiveTurn, delegationKickoffText: session.delegationKickoffText, - hasLiveMessage: session.liveMessage !== null, + liveShowsReply, liveStartedAt: session.liveMessage?.startedAt ?? null, + // Content, not identity: the set is rebuilt on every streaming batch, and + // an ordinary turn's `""` keeps the cache hitting exactly as before. + steeredCopyKey: steeredCopyIds ? [...steeredCopyIds].join("") : "", + liveStreamedRoundStart, } if (detail) { const cached = timelinePrefixCache.get(detail) @@ -3182,14 +3223,24 @@ function computeTimelinePrefix( // anchor, so fall back to the first assistant turn: the only assistant // content that can exist is the reply being streamed. const rawPersistedTurns = session.detail?.turns ?? [] + // A persisted copy of a message the user sent mid-turn is NOT a round + // boundary: it sits inside the round it interrupted, so the anchors below + // step over it. Only while the live stream provably holds the round from + // before that interruption — stepping over one hides the persisted turns + // between it and the real prompt, which is only sound where the live copy is + // showing them. + const roundAnchorSkipIds = + liveStreamedRoundStart && steeredCopyIds ? steeredCopyIds : null + const isRoundBoundaryUserTurn = (turn: MessageTurn): boolean => + turn.role === "user" && !roundAnchorSkipIds?.has(turn.id) const hasLiveOrLocalReply = session.liveOwnsActiveTurn && - (session.liveMessage !== null || session.localTurns.length > 0) + (liveShowsReply || session.localTurns.length > 0) let stripFrom = -1 if (hasLiveOrLocalReply) { let lastUserIdx = -1 for (let i = rawPersistedTurns.length - 1; i >= 0; i--) { - if (rawPersistedTurns[i]!.role === "user") { + if (isRoundBoundaryUserTurn(rawPersistedTurns[i]!)) { lastUserIdx = i break } @@ -3216,13 +3267,12 @@ function computeTimelinePrefix( // into `detail` it sits beside the live reply (a separate assistant turn // under a `live-…` id), and `mergeConsecutiveAssistantTurns` concatenates // the two — so the already-persisted head (e.g. the first reasoning block) - // renders twice. Hide that persisted partial, but ONLY while `liveMessage` - // is in hand: the live stream carries the full reply (the attach snapshot is - // built atomically and includes it), so this only ever hides from render - // what the live stream is concurrently showing — never dropping a reply we - // can't re-show. The moment the turn ends, `liveMessage` clears and the - // persisted copy (now complete) renders normally; the brief promote→refetch - // grace window can show a transient visible duplicate, never a hidden turn. + // renders twice. Hide that persisted partial, but ONLY while the live message + // is actually SHOWING a reply (`liveShowsReply`): that is what makes this a + // choice between two renderings of one reply rather than a deletion. The + // moment the turn ends, `liveMessage` clears and the persisted copy (now + // complete) renders normally; the brief promote→refetch grace window can show + // a transient visible duplicate, never a hidden turn. // // The in-flight prompt is identified authoritatively by the backend, which // reports the id of the persisted user turn it stamped as the in-flight one @@ -3232,15 +3282,47 @@ function computeTimelinePrefix( // client clock on the streaming path — neither can locate the prompt across // machines. When the new prompt isn't persisted yet the backend reports no // id, so an earlier completed round's reply is never mistaken for a partial. + // + // A mid-turn message used to take that id away, and the fallback below is + // what is left of that. `apply_in_flight_message_id` matched the transcript + // TAIL, so once the agent had written a steered message the tail was that + // message, whose content is not the prompt's, and nothing was stamped — in + // the middle of the round the suppression exists for. That is fixed at the + // source now: the backend walks back over the turns of the running turn and + // stamps the earliest copy of the prompt, which is the only reading that also + // repairs `detailIsInFlight` (a mid-turn refetch clearing the live buffers) + // and `collectInFlightPersistedToolCalls` (unfinished tool calls painting as + // completed) — neither of which this file can reach. + // + // Kept as a BACKSTOP, and for grok as the only thing there is: `grok.rs` + // parses its per-line `timestamp` as whole UNIX SECONDS, so a prompt the + // agent persisted a few hundred milliseconds into the turn floors to an + // instant BEFORE the (millisecond-precision) turn start and the backend's + // recency gate refuses it. The stamp is therefore unreachable there however + // the backend locates the prompt, and this is grok's round anchor. Fall back + // to the + // last persisted turn this client can still prove opened the round: the + // newest user turn that is not one of those copies. Reachable only once a + // copy is actually in `detail.turns`, and only while the live stream holds + // the round from before the interruption. const inFlightPromptId = session.detail?.in_flight_user_turn_id ?? null - const inFlightPromptIdx = - !hasLiveOrLocalReply && - session.liveMessage !== null && - inFlightPromptId !== null - ? persistedTurns.findIndex( - (t) => t.role === "user" && t.id === inFlightPromptId - ) - : -1 + // `liveShowsReply`, not `liveMessage !== null`: this hides a persisted + // assistant turn because the live stream is showing that reply, so a live + // message rendering nothing may not stand in for it. + const canSuppressInFlightPartial = !hasLiveOrLocalReply && liveShowsReply + let inFlightPromptIdx = -1 + if (canSuppressInFlightPartial && inFlightPromptId !== null) { + inFlightPromptIdx = persistedTurns.findIndex( + (t) => t.role === "user" && t.id === inFlightPromptId + ) + } else if (canSuppressInFlightPartial && roundAnchorSkipIds) { + for (let i = persistedTurns.length - 1; i >= 0; i--) { + if (isRoundBoundaryUserTurn(persistedTurns[i]!)) { + inFlightPromptIdx = i + break + } + } + } const visiblePersistedTurns = inFlightPromptIdx === -1 ? persistedTurns @@ -3381,27 +3463,25 @@ function computeTimelinePrefix( } /** - * Hide the persisted copy of a message the user sent mid-turn, when the live - * stream is already showing it. + * Ids of the DETAIL's own copies of the messages the user sent mid-turn, i.e. + * the persisted user turns that the live stream is already showing as steered + * messages. `null` when this turn steered nothing (the overwhelmingly common + * case), so every caller below is free in an ordinary turn. * * The agent writes a steered message into its own transcript, so a detail * fetch that lands DURING the turn brings it back as an ordinary user turn — - * under a parser id, which no id-keyed dedup can match to the live copy. Both - * would render. - * - * The live copy is the one to keep: it sits between the two halves of the - * reply, where the message was actually sent, while the persisted copy is - * appended after the in-flight prompt with the reply's first half suppressed - * around it (see `visiblePersistedTurns`), which would put the interruption - * before the text it interrupted. + * under a parser id, which no id-keyed dedup can match to the live copy. Three + * separate rules need to know which persisted turns those are: + * `suppressPersistedSteeredPrompts` hides them, and the two round anchors in + * `computeTimelinePrefix` must not mistake one for the start of a new round. * * Matched on CONTENT, the same way `APPEND_VIEWER_USER_TURN` reconciles the two * id namespaces of one prompt — but content ALONE cannot say which message it * matched. Steered text is short and repeatable ("continue", "stop", "not - * done"), so a bare content match reaches back and hides the identical prompt - * the user sent three rounds ago, for as long as the turn runs. Suppressing a - * user turn is the one failure that hides a message rather than duplicating - * it, so the match is bounded by WHEN: + * done"), so a bare content match reaches back and finds the identical prompt + * the user sent three rounds ago. Suppressing a user turn is the one failure + * that hides a message rather than duplicating it, so the match is bounded by + * WHEN: * * - each `steering` block carries the note's `created_at`, taken on the * agent's machine BEFORE the backend handed it the text (an invariant of @@ -3411,23 +3491,16 @@ function computeTimelinePrefix( * including this round's own prompt, which the agent wrote before the user * steered. * - * Candidates are further limited to turns the DETAIL projected, so every - * timestamp compared comes from the agent's own clock; a promoted `localTurns` - * copy (client clock, and kept across a mid-turn refetch by `preserveLive`) is - * never a candidate. Anything unreadable — no parseable instant on either side - * — suppresses nothing, leaving the two copies to coexist: a visible duplicate, + * Candidates are limited to turns the DETAIL projected, so every timestamp + * compared comes from the agent's own clock; a promoted `localTurns` copy + * (client clock, and kept across a mid-turn refetch by `preserveLive`) is never + * a candidate. Anything unreadable — no parseable instant on either side — + * matches nothing, leaving the two copies to coexist: a visible duplicate, * never a hidden message. - * - * Deliberately NOT anchored on `detail.in_flight_user_turn_id`: the backend - * stamps that by matching the pending prompt against the transcript TAIL (see - * `apply_in_flight_message_id`), and once the agent has written the steered - * message the tail is that message, not the prompt — so the stamp is gone in - * exactly the shape this function exists for. */ -function suppressPersistedSteeredPrompts( - prefix: ConversationTimelineTurn[], +function collectSteeredPersistedCopyIds( session: ConversationRuntimeSession -): ConversationTimelineTurn[] { +): Set | null { // Content key → the earliest instant a copy of it could have been written. // Read from the blocks rather than from the built turns: a block with no // readable stamp shows under the turn's start time, and treating THAT as the @@ -3444,27 +3517,83 @@ function suppressPersistedSteeredPrompts( if (known === undefined || at < known) steeredAt.set(key, at) if (at < earliestSteerAt) earliestSteerAt = at } - if (!steeredAt) return prefix + if (!steeredAt) return null const detailTurns = session.detail?.turns - if (!detailTurns) return prefix - // Ids are unique across the timeline's phases (a same-id copy in another - // phase is the same turn — see `dedupeTimeline`), so membership alone tells - // a detail-projected turn from a locally promoted one. - const detailUserIds = new Set() + if (!detailTurns) return null + let ids: Set | null = null for (const turn of detailTurns) { - if (turn.role === "user") detailUserIds.add(turn.id) - } - const filtered = prefix.filter((item) => { - if (item.phase !== "persisted" || item.turn.role !== "user") return true - if (!detailUserIds.has(item.turn.id)) return true + if (turn.role !== "user") continue // Cheap gate first: everything written before the earliest steer is out, // so history never reaches the content key (which serializes full text and - // full image data, and this runs on every streaming batch). - const at = Date.parse(item.turn.timestamp) - if (!Number.isFinite(at) || at < earliestSteerAt) return true - const steered = steeredAt.get(userTurnContentKey(item.turn)) - return steered === undefined || at < steered - }) + // full image data, and this runs whenever the prefix is rebuilt). + const at = Date.parse(turn.timestamp) + if (!Number.isFinite(at) || at < earliestSteerAt) continue + const steered = steeredAt.get(userTurnContentKey(turn)) + if (steered === undefined || at < steered) continue + ids ??= new Set() + ids.add(turn.id) + } + return ids +} + +/** + * Whether this session's live message opens on the reply rather than on the + * interruption: it holds at least one block from BEFORE the first mid-turn + * message, so it is showing the round from its start and can stand in for + * every persisted turn of it. + * + * The gate on moving a round anchor past a steered copy, because doing that + * hides the persisted turns between the copy and the real prompt. A live + * message that begins at the interruption is not evidence for them: a session + * that adopts a snapshot mid-turn starts from the backend's live message, which + * carries no `steering` block at all (see `snapshot-denormalize`), so a steer + * arriving afterwards can be the first thing this client ever saw of the turn. + * Hiding the reply's persisted first half there would put it nowhere. + */ +function liveMessageOpensBeforeFirstSteer( + liveMessage: LiveMessage | null +): boolean { + for (const block of liveMessage?.content ?? []) { + if (block.type === "steering") return false + // Parented subagent output never reaches the main thread (see + // `buildStreamingTurnsFromLiveMessage`), so it is not evidence that this + // client holds the reply either. + if ( + (block.type === "text" || block.type === "thinking") && + block.parentToolUseId + ) { + continue + } + return true + } + return false +} + +/** + * Hide the persisted copies of the messages the user sent mid-turn, since the + * live stream is already showing them. + * + * The live copy is the one to keep: it sits between the two halves of the + * reply, where the message was actually sent, while the persisted copy is + * appended after the in-flight prompt with the reply's first half suppressed + * around it (see `visiblePersistedTurns`), which would put the interruption + * before the text it interrupted. + * + * Ids are unique across the timeline's phases (a same-id copy in another phase + * is the same turn — see `dedupeTimeline`), so id membership alone tells a + * detail-projected turn from a locally promoted one. + */ +function suppressPersistedSteeredPrompts( + prefix: ConversationTimelineTurn[], + steeredCopyIds: ReadonlySet | null +): ConversationTimelineTurn[] { + if (!steeredCopyIds) return prefix + const filtered = prefix.filter( + (item) => + item.phase !== "persisted" || + item.turn.role !== "user" || + !steeredCopyIds.has(item.turn.id) + ) return filtered.length === prefix.length ? prefix : filtered } @@ -3478,14 +3607,40 @@ function computeTimeline( const cached = timelineCache.get(session) if (cached) return cached - // Phases 1–3 (already deduped), reused across streaming batches. - const { prefix, prefixKeys } = computeTimelinePrefix(session, conversationId) - - // Phase 4: Streaming turns (live agent response, split into rounds) + // Phase 4 first: Phases 1–3 hide the persisted copy of the reply this build + // is showing, so they need its verdict, and deriving that from the same build + // is what keeps the two from disagreeing. A live message can hold nothing + // renderable — `content: []` from the turn's own `prompting` transition, or + // only blocks this build drops — and a check for the message OBJECT then hid + // a persisted reply that nothing replaced. const streamingMessage = session.liveMessage const built = streamingMessage ? buildStreamingTurnsFromLiveMessage(conversationId, streamingMessage) : null + // A `user` turn here is a message the user sent mid-turn (native steering), + // not a rendering of the reply — a live message that produced only those is + // showing no reply and must suppress nothing. + const liveShowsReply = + built?.turns.some((turn) => turn.role === "assistant") ?? false + + // The detail's own copies of this turn's mid-turn messages, and whether the + // live stream can stand in for the round they interrupted. Derived once and + // shared: the prefix's two round anchors and the suppression below must + // agree on which persisted user turns are steered copies, or one of them + // hides a turn another is still anchoring on. + const steeredCopyIds = collectSteeredPersistedCopyIds(session) + const liveStreamedRoundStart = steeredCopyIds + ? liveMessageOpensBeforeFirstSteer(session.liveMessage) + : false + + // Phases 1–3 (already deduped), reused across streaming batches. + const { prefix, prefixKeys } = computeTimelinePrefix( + session, + conversationId, + liveShowsReply, + steeredCopyIds, + liveStreamedRoundStart + ) let deduped: ConversationTimelineTurn[] if (!built || built.turns.length === 0) { @@ -3516,7 +3671,7 @@ function computeTimeline( } seenTailKeys?.add(key) } - const head = suppressPersistedSteeredPrompts(prefix, session) + const head = suppressPersistedSteeredPrompts(prefix, steeredCopyIds) deduped = collides ? dedupeTimeline(head.concat(tail)) : head.concat(tail) } diff --git a/src/stores/runtime-empty-live-message.test.ts b/src/stores/runtime-empty-live-message.test.ts new file mode 100644 index 0000000000..a61d593865 --- /dev/null +++ b/src/stores/runtime-empty-live-message.test.ts @@ -0,0 +1,207 @@ +/** + * Two timeline rules hide a persisted assistant turn while a reply streams: the + * `liveOwnsActiveTurn` tail strip (delegation-child dialog) and the + * `in_flight_user_turn_id` partial suppression (cross-client viewer). Both are + * only sound because the live stream is showing that same reply — so both have + * to key off what the live message RENDERS, not off a live message existing. + * + * Those differ, and routinely. `status_changed → prompting` installs a fresh + * `content: []` live message at the start of every turn and mirrors it into + * this store (see "fires with isLive=true and a fresh non-null liveMessage when + * a turn starts" in acp-connections-context.test.tsx); the mirror never writes a + * null back over it, so the same object stays in hand for any part of a turn + * that produces nothing this build renders. Keyed on the object, the persisted + * reply was hidden with nothing put in its place: a blank agent turn. + */ + +import { afterEach, describe, expect, it } from "vitest" + +import type { LiveMessage } from "@/contexts/acp-connections-context" +import type { DbConversationDetail, MessageTurn, TurnRole } from "@/lib/types" +import { + getTimelineTurns, + resetConversationRuntimeStore, + useConversationRuntimeStore, +} from "@/stores/conversation-runtime-store" + +const CID = 77 +const TS = "2026-09-06T00:00:00.000Z" + +function turn(id: string, role: TurnRole): MessageTurn { + return { id, role, blocks: [{ type: "text", text: id }], timestamp: TS } +} + +/** What the turn's own `prompting` transition installs, before any content. */ +const promptingLiveMessage: LiveMessage = { + id: "m1", + role: "assistant", + content: [], + startedAt: Date.parse(TS), +} + +/** A block Phase 2 drops, so this message renders exactly as much as `[]`. */ +const emptyTextLiveMessage: LiveMessage = { + ...promptingLiveMessage, + content: [{ type: "text", text: "" }], +} + +const replyLiveMessage: LiveMessage = { + ...promptingLiveMessage, + content: [{ type: "text", text: "streaming…" }], +} + +/** A message the user sent mid-turn, with no reply to it yet. */ +const steeringOnlyLiveMessage: LiveMessage = { + ...promptingLiveMessage, + content: [ + { + type: "steering", + id: "note-1", + text: "also check the tests", + createdAt: TS, + }, + ], +} + +function seed( + turns: MessageTurn[], + overrides: { + liveMessage?: LiveMessage | null + liveOwnsActiveTurn?: boolean + localTurns?: MessageTurn[] + inFlightUserTurnId?: string | null + } +) { + const detail: DbConversationDetail = { + summary: { + id: CID, + folder_id: 1, + title: "t", + title_locked: false, + agent_type: "claude_code", + status: "in_progress", + kind: "regular", + model: null, + git_branch: null, + external_id: null, + message_count: turns.length, + child_count: 0, + created_at: TS, + updated_at: TS, + pinned_at: null, + }, + turns, + in_flight_user_turn_id: overrides.inFlightUserTurnId ?? null, + } + const next = new Map(useConversationRuntimeStore.getState().byConversationId) + next.set(CID, { + conversationId: CID, + externalId: null, + dbConversationId: null, + detail, + detailLoading: false, + detailError: null, + acpLoadError: null, + localTurns: overrides.localTurns ?? [], + backgroundTurns: [], + pendingBackgroundSettlements: [], + optimisticTurns: [], + liveMessage: overrides.liveMessage ?? null, + syncState: "idle" as const, + activeTurnToken: null, + lastTurnOwned: false, + liveOwnsActiveTurn: overrides.liveOwnsActiveTurn ?? false, + delegationKickoffText: null, + sessionStats: null, + historyAssistantBaseline: null, + batchBoundaryIndex: null, + batchBoundaryPrefixHash: null, + loadingOlderTurns: false, + olderTurnsPrependEpoch: 0, + pendingOutOfTurnContent: false, + pendingCleanup: false, + }) + useConversationRuntimeStore.setState({ byConversationId: next }) +} + +const timelineIds = () => getTimelineTurns(CID).map((t) => t.turn.id) + +afterEach(() => { + resetConversationRuntimeStore() +}) + +describe("persisted-tail strip vs. a live message that renders nothing", () => { + it("keeps the child's reply while the new turn has produced nothing yet", () => { + seed([turn("u1", "user"), turn("a1", "assistant")], { + liveOwnsActiveTurn: true, + liveMessage: promptingLiveMessage, + }) + expect(timelineIds()).toEqual(["u1", "a1"]) + }) + + it("keeps the child's reply when the live message holds only an empty block", () => { + seed([turn("u1", "user"), turn("a1", "assistant")], { + liveOwnsActiveTurn: true, + liveMessage: emptyTextLiveMessage, + }) + expect(timelineIds()).toEqual(["u1", "a1"]) + }) + + it("still strips the persisted copy once the live message shows the reply", () => { + seed([turn("u1", "user"), turn("a1", "assistant")], { + liveOwnsActiveTurn: true, + liveMessage: replyLiveMessage, + }) + expect(timelineIds()).toEqual(["u1", `live-${CID}-m1`]) + }) + + it("still strips for a promoted reply, which renders on its own", () => { + seed([turn("u1", "user"), turn("a1", "assistant")], { + liveOwnsActiveTurn: true, + liveMessage: promptingLiveMessage, + localTurns: [turn("promoted", "assistant")], + }) + expect(timelineIds()).toEqual(["u1", "promoted"]) + }) +}) + +describe("in-flight partial suppression vs. a live message that renders nothing", () => { + it("keeps the persisted partial while the turn has produced nothing yet", () => { + seed([turn("u1", "user"), turn("a1", "assistant")], { + liveMessage: promptingLiveMessage, + inFlightUserTurnId: "u1", + }) + expect(timelineIds()).toEqual(["u1", "a1"]) + }) + + it("keeps every persisted reply of the round, not just the newest", () => { + seed( + [ + turn("u1", "user"), + turn("a1", "assistant"), + turn("a2", "assistant"), + turn("a3", "assistant"), + ], + { liveMessage: emptyTextLiveMessage, inFlightUserTurnId: "u1" } + ) + expect(timelineIds()).toEqual(["u1", "a1", "a2", "a3"]) + }) + + it("still hides the persisted partial once the live message shows the reply", () => { + seed([turn("u1", "user"), turn("a1", "assistant")], { + liveMessage: replyLiveMessage, + inFlightUserTurnId: "u1", + }) + expect(timelineIds()).toEqual(["u1", `live-${CID}-m1`]) + }) + + it("keeps the persisted reply when the live message carries only a steer", () => { + // A mid-turn message is the user's, not a rendering of the reply, so it + // cannot stand in for the persisted copy it would otherwise hide. + seed([turn("u1", "user"), turn("a1", "assistant")], { + liveMessage: steeringOnlyLiveMessage, + inFlightUserTurnId: "u1", + }) + expect(timelineIds()).toEqual(["u1", "a1", `live-${CID}-m1`]) + }) +}) diff --git a/src/stores/runtime-steered-round-anchor.test.ts b/src/stores/runtime-steered-round-anchor.test.ts new file mode 100644 index 0000000000..a11c617c07 --- /dev/null +++ b/src/stores/runtime-steered-round-anchor.test.ts @@ -0,0 +1,406 @@ +/** + * A message the user sends mid-turn is a USER turn the agent writes into the + * MIDDLE of a round. Two timeline rules locate the running round by its newest + * persisted user turn, and both read that message as the start of a new one: + * + * - the viewer's persisted-tail strip anchors on the last persisted user turn + * (`liveOwnsActiveTurn`), so the anchor jumps forward past the reply's + * already-persisted first half and stops stripping it; + * - the in-flight partial suppression anchors on + * `detail.in_flight_user_turn_id`, which the backend stamps by matching the + * pending prompt against the transcript TAIL — and once the steered message + * IS the tail, the content no longer matches and nothing is stamped at all. + * + * Either way the persisted first half of the reply lands beside the live copy + * of the same text, and `mergeConsecutiveAssistantTurns` glues them into one + * bubble: the run-on paragraph the mid-turn split exists to prevent, back again + * for every turn that is actually steered. + * + * Both anchors now step over the detail's own copies of this turn's mid-turn + * messages — but only while the live message provably holds the round from + * before the interruption, since stepping over one hides the persisted turns + * behind it. + */ + +import { afterEach, describe, expect, it } from "vitest" + +import type { + LiveContentBlock, + LiveMessage, +} from "@/contexts/acp-connections-context" +import type { DbConversationDetail, MessageTurn, TurnRole } from "@/lib/types" +import { + getTimelineTurns, + resetConversationRuntimeStore, + useConversationRuntimeStore, +} from "@/stores/conversation-runtime-store" + +const CID = 91 +/** Older than every steer below, so history is provably a different message. */ +const BEFORE = "2026-09-05T00:00:00.000Z" +/** When the backend injected the message, stamped where the agent runs. */ +const STEER_AT = "2026-09-05T00:05:00.000Z" +/** A turn the agent wrote after that injection, i.e. its own copy of it. */ +const AFTER_STEER = "2026-09-05T00:05:01.000Z" + +const STEER_TEXT = "use the other API" + +function turn( + id: string, + role: TurnRole, + text: string, + timestamp = BEFORE +): MessageTurn { + return { id, role, blocks: [{ type: "text", text }], timestamp } +} + +function steerBlock(text = STEER_TEXT, id = "n1"): LiveContentBlock { + return { type: "steering", id, text, createdAt: STEER_AT, blocks: null } +} + +function live(content: LiveContentBlock[]): LiveMessage { + return { + id: "lm", + role: "assistant", + content, + startedAt: Date.parse(BEFORE), + } +} + +/** The whole round in hand: the reply opened, the user cut in, it went on. */ +const splitReply = live([ + { type: "text", text: "half one" }, + steerBlock(), + { type: "text", text: "half two" }, +]) + +function seed( + turns: MessageTurn[], + o: { + liveMessage?: LiveMessage | null + liveOwnsActiveTurn?: boolean + inFlightUserTurnId?: string | null + } +) { + const detail: DbConversationDetail = { + summary: { + id: CID, + folder_id: 1, + title: "t", + title_locked: false, + agent_type: "claude_code", + status: "in_progress", + kind: "regular", + model: null, + git_branch: null, + external_id: null, + message_count: turns.length, + child_count: 0, + created_at: BEFORE, + updated_at: BEFORE, + pinned_at: null, + }, + turns, + in_flight_user_turn_id: o.inFlightUserTurnId ?? null, + } + const next = new Map(useConversationRuntimeStore.getState().byConversationId) + next.set(CID, { + conversationId: CID, + externalId: null, + dbConversationId: null, + detail, + detailLoading: false, + detailError: null, + acpLoadError: null, + localTurns: [], + backgroundTurns: [], + pendingBackgroundSettlements: [], + optimisticTurns: [], + liveMessage: o.liveMessage ?? null, + syncState: "idle" as const, + activeTurnToken: null, + lastTurnOwned: false, + liveOwnsActiveTurn: o.liveOwnsActiveTurn ?? false, + delegationKickoffText: null, + sessionStats: null, + historyAssistantBaseline: null, + batchBoundaryIndex: null, + batchBoundaryPrefixHash: null, + loadingOlderTurns: false, + olderTurnsPrependEpoch: 0, + pendingOutOfTurnContent: false, + pendingCleanup: false, + }) + useConversationRuntimeStore.setState({ byConversationId: next }) +} + +/** What the thread actually reads, top to bottom. */ +const timelineTexts = () => + getTimelineTurns(CID).map((t) => + t.turn.blocks[0]?.type === "text" ? t.turn.blocks[0].text : "?" + ) + +/** The transcript mid-turn: the prompt, the reply so far, then the message the + * user cut in with — which the agent records as an ordinary user turn. */ +const steeredTranscript = [ + turn("u1", "user", "the original prompt"), + turn("a1", "assistant", "half one"), + turn("u2", "user", STEER_TEXT, AFTER_STEER), +] + +afterEach(() => { + resetConversationRuntimeStore() +}) + +describe("persisted-tail strip vs. a message sent mid-turn", () => { + it("strips the reply's first half, which the live stream is re-showing", () => { + seed(steeredTranscript, { + liveOwnsActiveTurn: true, + liveMessage: splitReply, + }) + expect(timelineTexts()).toEqual([ + "the original prompt", + "half one", + STEER_TEXT, + "half two", + ]) + }) + + it("keeps the first half when the live message opens on the interruption", () => { + // A session that adopted a snapshot mid-turn starts from the backend's live + // message, which carries no steering block (the wire has no such kind), so a + // steer arriving afterwards can be the first thing it holds of this turn. + // The persisted first half is then rendered NOWHERE else — keep it, and let + // the mid-turn message be the only thing the live copy stands in for. + seed(steeredTranscript, { + liveOwnsActiveTurn: true, + liveMessage: live([steerBlock(), { type: "text", text: "half two" }]), + }) + expect(timelineTexts()).toEqual([ + "the original prompt", + "half one", + STEER_TEXT, + "half two", + ]) + }) + + it("leaves an earlier round's identical message in history", () => { + // Steered text is short and repeatable. Reaching past the round's own + // prompt would erase the same words the user sent three rounds ago for as + // long as the turn runs. + seed( + [ + turn("u0", "user", STEER_TEXT), + turn("a0", "assistant", "sure"), + ...steeredTranscript, + ], + { liveOwnsActiveTurn: true, liveMessage: splitReply } + ) + expect(timelineTexts()).toEqual([ + STEER_TEXT, + "sure", + "the original prompt", + "half one", + STEER_TEXT, + "half two", + ]) + }) + + it("still anchors on the newest prompt when nothing was steered", () => { + seed( + [ + turn("u1", "user", "the original prompt"), + turn("a1", "assistant", "half one"), + ], + { + liveOwnsActiveTurn: true, + liveMessage: live([{ type: "text", text: "half one" }]), + } + ) + expect(timelineTexts()).toEqual(["the original prompt", "half one"]) + }) +}) + +describe("in-flight partial suppression vs. a message sent mid-turn", () => { + it("hides the persisted partial once the steer has taken the backend's stamp away", () => { + // `apply_in_flight_message_id` matches the pending prompt against the + // transcript tail; the tail is now the steered message, whose content is + // not the prompt's, so the detail comes back with no stamp at all. + seed(steeredTranscript, { + liveMessage: splitReply, + inFlightUserTurnId: null, + }) + expect(timelineTexts()).toEqual([ + "the original prompt", + "half one", + STEER_TEXT, + "half two", + ]) + }) + + it("keeps the persisted partial when the live message opens on the interruption", () => { + seed(steeredTranscript, { + liveMessage: live([steerBlock(), { type: "text", text: "half two" }]), + inFlightUserTurnId: null, + }) + expect(timelineTexts()).toEqual([ + "the original prompt", + "half one", + STEER_TEXT, + "half two", + ]) + }) + + it("keeps using the backend's stamp while it is there", () => { + seed(steeredTranscript, { + liveMessage: splitReply, + inFlightUserTurnId: "u1", + }) + expect(timelineTexts()).toEqual([ + "the original prompt", + "half one", + STEER_TEXT, + "half two", + ]) + }) + + it("suppresses nothing when no steered copy has landed yet", () => { + // The agent hasn't written its copy, so there is no evidence of which round + // is running and no stamp to read: both persisted turns stand. + seed( + [ + turn("u1", "user", "the original prompt"), + turn("a1", "assistant", "half one"), + ], + { + liveMessage: splitReply, + inFlightUserTurnId: null, + } + ) + expect(timelineTexts()).toEqual([ + "the original prompt", + "half one", + "half one", + STEER_TEXT, + "half two", + ]) + }) + + it("handles the same words sent twice in one turn", () => { + // The shape content matching is worst at: two mid-turn messages the round + // cannot tell apart by text. Each still has to render exactly once, in + // place, with the reply between them. + seed( + [ + turn("u1", "user", "the original prompt"), + turn("a1", "assistant", "half one"), + turn("u2", "user", "continue", AFTER_STEER), + turn("a2", "assistant", "half two"), + turn("u3", "user", "continue", "2026-09-05T00:06:01.000Z"), + ], + { + liveMessage: live([ + { type: "text", text: "half one" }, + steerBlock("continue"), + { type: "text", text: "half two" }, + { + type: "steering", + id: "n2", + text: "continue", + createdAt: "2026-09-05T00:06:00.000Z", + blocks: null, + }, + { type: "text", text: "half three" }, + ]), + inFlightUserTurnId: null, + } + ) + expect(timelineTexts()).toEqual([ + "the original prompt", + "half one", + "continue", + "half two", + "continue", + "half three", + ]) + }) + + it("matches a message that carried an attachment", () => { + // A steer with a draft attached is keyed on the blocks it actually sent, + // which is what the agent writes to its transcript — the display text says + // something else entirely. + const image = { + type: "image" as const, + mime_type: "image/png", + data: "AAAA", + uri: "file:///shot.png", + } + seed( + [ + turn("u1", "user", "the original prompt"), + turn("a1", "assistant", "half one"), + { + id: "u2", + role: "user", + blocks: [image, { type: "text", text: "look at this" }], + timestamp: AFTER_STEER, + }, + ], + { + liveMessage: live([ + { type: "text", text: "half one" }, + { + type: "steering", + id: "n1", + text: "look at this [image]", + createdAt: STEER_AT, + blocks: [ + { type: "image", mime_type: "image/png", data: "AAAA" }, + { type: "text", text: "look at this" }, + ], + }, + { type: "text", text: "half two" }, + ]), + inFlightUserTurnId: null, + } + ) + expect(getTimelineTurns(CID).map((t) => `${t.phase}:${t.turn.id}`)).toEqual( + [ + "persisted:u1", + `streaming:live-${CID}-lm`, + `streaming:live-${CID}-lm-1`, + `streaming:live-${CID}-lm-2`, + ] + ) + }) + + it("handles two messages sent in the same turn", () => { + seed( + [ + turn("u1", "user", "the original prompt"), + turn("a1", "assistant", "half one"), + turn("u2", "user", STEER_TEXT, AFTER_STEER), + turn("a2", "assistant", "half two"), + turn("u3", "user", "and the other thing", AFTER_STEER), + ], + { + liveMessage: live([ + { type: "text", text: "half one" }, + steerBlock(), + { type: "text", text: "half two" }, + steerBlock("and the other thing", "n2"), + { type: "text", text: "half three" }, + ]), + inFlightUserTurnId: null, + } + ) + expect(timelineTexts()).toEqual([ + "the original prompt", + "half one", + STEER_TEXT, + "half two", + "and the other thing", + "half three", + ]) + }) +})