Skip to main content

mj_core/
transcript.rs

1//! Transcript data shared by the worker, controller state, and the chat UI.
2//!
3//! [`ChatEntry`] is what a worker snapshot carries in its transcript tail and
4//! what the chat view renders, and [`TranscriptItem`] is the materialized
5//! form controller state persists, so both live below the modules that use
6//! them rather than inside any one of them.
7//!
8//! The text helpers that read one of those shapes live here for the same
9//! reason: the database, the projection, controller state, the compactor and
10//! the review host all need the plain text of a stored message, and none of
11//! them should have to reach up into the chat view to get it.
12
13use std::sync::Arc;
14
15use agent_client_protocol::schema::v1::{
16    ContentBlock, ContentChunk, EmbeddedResourceResource, PlanEntryStatus, ToolCallStatus, ToolKind,
17};
18use anyhow::{Result, bail};
19use serde::{Deserialize, Serialize};
20
21pub const SESSION_RESTART_TEXT: &str = "[session restarted]";
22pub const SESSION_RESTART_ITEM_PREFIX: &str = "system:session-restarted:";
23pub const WORK_INTERRUPTED_ITEM_PREFIX: &str = "system:work-interrupted:";
24/// Marks where a turn the user cancelled stopped.
25pub const TURN_INTERRUPTED_TEXT: &str = "Interrupted";
26pub const TURN_INTERRUPTED_ITEM_PREFIX: &str = "system:turn-interrupted:";
27/// Marks the point where the harness resumed work with no prompt in flight.
28pub const HARNESS_TURN_TEXT: &str = "Agent continued on its own";
29pub const HARNESS_TURN_ITEM_PREFIX: &str = "harness-turn:";
30
31/// Where a tool's compact presentation source came from.
32#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
33#[serde(rename_all = "snake_case")]
34pub enum ToolSummarySourceKind {
35    RawInput,
36    RawOutput,
37    Title,
38}
39
40/// Bounded presentation metadata derived from an ACP tool call.
41#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
42pub struct ToolCallPresentation {
43    pub summary: String,
44    pub source: String,
45    pub source_kind: ToolSummarySourceKind,
46    pub tool_kind: ToolKind,
47    /// Version of the parser rules that produced `summary`. A missing value
48    /// identifies presentation metadata written before parser versioning.
49    #[serde(default)]
50    pub summary_version: u8,
51}
52
53/// The current value of one logical transcript item. ACP structures whose
54/// schemas can grow are kept as JSON values, while logical item identity and
55/// lifecycle remain controller-owned and stable.
56#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
57#[serde(tag = "kind", rename_all = "snake_case")]
58pub enum TranscriptBody {
59    User {
60        content: Vec<serde_json::Value>,
61    },
62    Agent {
63        /// Complete ACP `ContentChunk` values, including message IDs, content
64        /// metadata, and non-text content blocks.
65        chunks: Vec<serde_json::Value>,
66        streaming: bool,
67    },
68    Thought {
69        /// Complete ACP `ContentChunk` values, including message IDs, content
70        /// metadata, and non-text content blocks.
71        chunks: Vec<serde_json::Value>,
72        streaming: bool,
73    },
74    Tool {
75        /// Complete current ACP `ToolCall`, updated field-for-field as
76        /// `ToolCallUpdate` notifications arrive.
77        call: serde_json::Value,
78        /// Output of the terminals this call's content refers to. It is a
79        /// sibling of `call` rather than part of it because `ToolCall::update`
80        /// replaces `content` wholesale, which would discard anything injected
81        /// into the stored ACP value.
82        #[serde(default, skip_serializing_if = "Vec::is_empty")]
83        terminal_outputs: Vec<TerminalOutputRecord>,
84        /// Every terminal this call has ever referred to. Agents that replace
85        /// `content` wholesale can drop a terminal reference before the
86        /// terminal is reaped, so the current call is not enough to decide
87        /// where a terminal's output belongs.
88        #[serde(default, skip_serializing_if = "Vec::is_empty")]
89        terminal_refs: Vec<String>,
90        /// Cached label data used by Rich and browser projections. Older
91        /// transcript items omit this and derive it from `call` when read.
92        #[serde(default, skip_serializing_if = "Option::is_none")]
93        presentation: Option<Box<ToolCallPresentation>>,
94    },
95    /// Terminal output that no tool call refers to yet. It becomes a
96    /// `Tool` item's `terminal_outputs` entry as soon as a call naming the
97    /// terminal arrives, and stays here permanently otherwise so output is
98    /// never dropped.
99    TerminalOutput {
100        record: TerminalOutputRecord,
101    },
102    Plan {
103        /// Complete current ACP `Plan`, including entry priorities and all
104        /// plan- and entry-level metadata.
105        plan: serde_json::Value,
106    },
107    /// A plan the harness asked the user to approve, captured where the
108    /// decision happened so it renders inline and survives restart and export.
109    ///
110    /// It is a record of the proposal, not conversation input: Hel never
111    /// replays it to a model as a user or agent message.
112    PlanProposal {
113        /// Identity of the plan review that carried this proposal.
114        proposal_id: String,
115        /// Exact proposal text the harness sent.
116        plan: String,
117    },
118    System {
119        text: String,
120    },
121}
122
123/// Append one ACP `ContentChunk` value to a transcript item's chunk list,
124/// merging it into the previous chunk when the two are the same text stream.
125///
126/// Agents stream text one token at a time, so a long turn arrives as
127/// thousands of chunks that differ only in `content.text`. Stored separately
128/// they cost far more memory than the text they carry: each chunk is its own
129/// pair of nested `serde_json::Value` maps. Merging them keeps one chunk per
130/// run of text, which is what every reader of `chunks` already reconstructs.
131///
132/// Two chunks merge only when nothing but the text differs: both are objects
133/// whose `content.type` is `"text"` with a string `content.text`, their
134/// `messageId` values are equal (both absent counts as equal), and every
135/// other top-level key (such as `meta`) and every other `content` key (such
136/// as `annotations`) is identical. Anything else is pushed as its own chunk.
137pub fn push_content_chunk(chunks: &mut Vec<serde_json::Value>, chunk: serde_json::Value) {
138    if chunks
139        .last()
140        .is_some_and(|last| text_chunks_mergeable(last, &chunk))
141    {
142        let addition = chunk
143            .get("content")
144            .and_then(|content| content.get("text"))
145            .and_then(serde_json::Value::as_str)
146            .unwrap_or_default()
147            .to_owned();
148        if let Some(serde_json::Value::Object(last)) = chunks.last_mut()
149            && let Some(serde_json::Value::Object(content)) = last.get_mut("content")
150            && let Some(serde_json::Value::String(text)) = content.get_mut("text")
151        {
152            text.push_str(&addition);
153            return;
154        }
155    }
156    chunks.push(chunk);
157}
158
159/// Collapse runs of per-token text chunks that were stored before they were
160/// merged on the way in. Rebuilds the list through [`push_content_chunk`], so
161/// it applies exactly the same merge rule and leaves chunk boundaries that
162/// carry real differences (a new message ID, non-text content, differing
163/// metadata) where they are.
164pub fn coalesce_content_chunks(chunks: &mut Vec<serde_json::Value>) {
165    if chunks.len() < 2 {
166        return;
167    }
168    let mut merged = Vec::with_capacity(chunks.len());
169    for chunk in std::mem::take(chunks) {
170        push_content_chunk(&mut merged, chunk);
171    }
172    merged.shrink_to_fit();
173    *chunks = merged;
174}
175
176/// Whether `next` carries only more text for the same stream as `last`, so
177/// the two can share one chunk. The single definition of the merge rule used
178/// by both [`push_content_chunk`] and [`coalesce_content_chunks`].
179fn text_chunks_mergeable(last: &serde_json::Value, next: &serde_json::Value) -> bool {
180    let (serde_json::Value::Object(last), serde_json::Value::Object(next)) = (last, next) else {
181        return false;
182    };
183    let (
184        Some(serde_json::Value::Object(last_content)),
185        Some(serde_json::Value::Object(next_content)),
186    ) = (last.get("content"), next.get("content"))
187    else {
188        return false;
189    };
190    let is_text = |content: &serde_json::Map<String, serde_json::Value>| {
191        content.get("type").and_then(serde_json::Value::as_str) == Some("text")
192            && content
193                .get("text")
194                .is_some_and(serde_json::Value::is_string)
195    };
196    if !is_text(last_content) || !is_text(next_content) {
197        return false;
198    }
199    // Every other top-level key, `messageId` and `meta` included, must match.
200    if last.len() != next.len()
201        || !last
202            .iter()
203            .all(|(key, value)| key == "content" || next.get(key) == Some(value))
204    {
205        return false;
206    }
207    // ... as must every other content key, such as `annotations`.
208    last_content.len() == next_content.len()
209        && last_content
210            .iter()
211            .all(|(key, value)| key == "text" || next_content.get(key) == Some(value))
212}
213
214/// What one client-run terminal produced, as hel recorded it when the child
215/// was reaped. `exit_code` and `signal` mirror ACP `TerminalExitStatus`; both
216/// are `None` when the terminal was released before a status was observed.
217#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
218pub struct TerminalOutputRecord {
219    pub terminal_id: String,
220    pub output: String,
221    #[serde(default, skip_serializing_if = "is_false")]
222    pub truncated: bool,
223    #[serde(default, skip_serializing_if = "Option::is_none")]
224    pub exit_code: Option<u32>,
225    #[serde(default, skip_serializing_if = "Option::is_none")]
226    pub signal: Option<String>,
227}
228
229impl TerminalOutputRecord {
230    /// Whether the command ended the way a caller asked for: exit status zero
231    /// and no signal. Anything else — a nonzero exit, a signal, or no status at
232    /// all because the terminal was released before one was observed — is
233    /// abnormal, and stays visible in every render mode.
234    pub fn exited_cleanly(&self) -> bool {
235        self.exit_code == Some(0) && self.signal.is_none()
236    }
237
238    /// Whether a completed ACP tool's provider-specific raw result is this
239    /// child result. Kimi reports shell output as a byte array beside its exit
240    /// status but omits the ACP terminal reference, so the exact result is the
241    /// only ownership information it publishes.
242    pub fn matches_tool_raw_result(&self, call: &serde_json::Value) -> bool {
243        if !matches!(
244            call.get("status").and_then(serde_json::Value::as_str),
245            Some("completed" | "failed")
246        ) {
247            return false;
248        }
249        let Some(raw) = call.get("rawOutput") else {
250            return false;
251        };
252        let Some(exit_code) = raw
253            .get("exit_code")
254            .and_then(serde_json::Value::as_u64)
255            .and_then(|code| u32::try_from(code).ok())
256        else {
257            return false;
258        };
259        if self.exit_code != Some(exit_code) || self.signal.is_some() {
260            return false;
261        }
262        match raw.get("output") {
263            Some(serde_json::Value::Array(bytes)) => {
264                bytes.len() == self.output.len()
265                    && bytes
266                        .iter()
267                        .zip(self.output.as_bytes())
268                        .all(|(value, byte)| value.as_u64() == Some(u64::from(*byte)))
269            }
270            Some(serde_json::Value::String(output)) => output == &self.output,
271            _ => false,
272        }
273    }
274}
275
276#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
277#[serde(deny_unknown_fields)]
278pub struct TranscriptItem {
279    pub stable_id: String,
280    /// Ordinal of the relay event that first created this logical item.
281    pub position: u64,
282    /// Ordinal of the most recent content chunk for an agent message. This is
283    /// `None` for every other logical item.
284    pub latest_content_event_ordinal: Option<u64>,
285    pub created_at_ms: i64,
286    pub last_changed_at_ms: i64,
287    pub body: TranscriptBody,
288}
289
290impl TranscriptItem {
291    pub fn is_work_interruption(&self) -> bool {
292        self.stable_id.starts_with(WORK_INTERRUPTED_ITEM_PREFIX)
293    }
294
295    pub fn is_session_restart(&self) -> bool {
296        self.stable_id.starts_with(SESSION_RESTART_ITEM_PREFIX)
297    }
298
299    /// The relay ordinal a reader pages by.
300    ///
301    /// An agent message is rewritten as its content streams in, and its
302    /// `latest_content_event_ordinal` is where that stopped, so paging by it
303    /// hands a caller the finished message once instead of the partial one it
304    /// was created with. Every other body is created once, so its position is
305    /// its sequence.
306    pub fn seq(&self) -> u64 {
307        self.latest_content_event_ordinal.unwrap_or(self.position)
308    }
309
310    /// Whether this item begins a turn: a user message, or the marker for a
311    /// turn the harness started on its own. The recovery boundary and the
312    /// scope of a plan update both key on the newest of these.
313    pub fn is_turn_start(&self) -> bool {
314        matches!(self.body, TranscriptBody::User { .. })
315            || self.stable_id.starts_with(HARNESS_TURN_ITEM_PREFIX)
316    }
317
318    pub fn is_nonempty_agent_message(&self) -> bool {
319        let TranscriptBody::Agent { chunks, .. } = &self.body else {
320            return false;
321        };
322        chunks.iter().any(|chunk| {
323            let Some(content) = chunk.get("content") else {
324                return false;
325            };
326            match content.get("type").and_then(serde_json::Value::as_str) {
327                Some("text") => content
328                    .get("text")
329                    .and_then(serde_json::Value::as_str)
330                    .is_some_and(|text| !text.trim().is_empty()),
331                Some(_) => true,
332                None => false,
333            }
334        })
335    }
336
337    pub fn validate(&self, through: u64) -> Result<()> {
338        if self.stable_id.trim().is_empty() {
339            bail!("materialized transcript item has an empty stable id");
340        }
341        if self.position == 0 || self.position > through {
342            bail!(
343                "materialized transcript item {:?} has invalid position {} at frontier {through}",
344                self.stable_id,
345                self.position
346            );
347        }
348        match (&self.body, self.latest_content_event_ordinal) {
349            (TranscriptBody::Agent { .. }, Some(ordinal))
350                if ordinal >= self.position && ordinal <= through => {}
351            (TranscriptBody::Agent { .. }, Some(ordinal)) => bail!(
352                "materialized agent message {:?} has invalid latest content ordinal {ordinal} at position {} and frontier {through}",
353                self.stable_id,
354                self.position
355            ),
356            (TranscriptBody::Agent { .. }, None) => bail!(
357                "materialized agent message {:?} has no latest content ordinal",
358                self.stable_id
359            ),
360            (_, Some(ordinal)) => bail!(
361                "non-agent transcript item {:?} has latest content ordinal {ordinal}",
362                self.stable_id
363            ),
364            (_, None) => {}
365        }
366        if self.last_changed_at_ms < self.created_at_ms {
367            bail!(
368                "materialized transcript item {:?} changed before it was created",
369                self.stable_id
370            );
371        }
372        Ok(())
373    }
374}
375
376#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
377pub enum ChatRole {
378    User,
379    Agent,
380    /// Agent reasoning stream, rendered dimmed.
381    Thought,
382    /// Tool invocation titles.
383    Tool,
384    /// Current agent plan.
385    Plan,
386    /// A plan proposal awaiting, or already given, a decision.
387    PlanProposal,
388    System,
389}
390
391#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
392pub struct ChatEntry {
393    #[serde(default)]
394    pub start_seq: u64,
395    pub seq: u64,
396    pub role: ChatRole,
397    pub text: String,
398    pub recorded_at_ms: Option<i64>,
399    pub revision: u64,
400    pub message_id: Option<String>,
401    pub tool_call_id: Option<String>,
402    pub tool_status: Option<ToolStatus>,
403    /// Compact label used by Rich and browser projections. `text` remains the
404    /// original provider title for Raw mode.
405    #[serde(default, skip_serializing_if = "Option::is_none")]
406    pub tool_summary: Option<String>,
407    /// Selected bounded source retained so partial ACP updates can preserve a
408    /// raw-command-derived summary without retaining arbitrary raw JSON.
409    #[serde(default, skip_serializing_if = "Option::is_none")]
410    pub tool_presentation: Option<ToolCallPresentation>,
411    pub tool_content: Vec<String>,
412    pub tool_diffstats: Vec<String>,
413    pub tool_locations: Vec<String>,
414    pub plan: Vec<PlanLine>,
415    #[serde(default, skip_serializing_if = "is_false")]
416    pub leading_omitted: bool,
417    /// Detail the decluttered feed leaves out: the entry renders only in the
418    /// raw transcript mode. Set once, when the entry is built, because Alt-T
419    /// switches render mode without rebuilding entries.
420    #[serde(default, skip_serializing_if = "is_false")]
421    pub raw_only: bool,
422    /// The materialized transcript item this entry was derived from, when it
423    /// came from the controller's projection. Provenance only, so it is
424    /// neither serialized nor part of the entry's value.
425    #[serde(skip)]
426    pub source: TranscriptSource,
427}
428
429/// Handle on the transcript item an entry was derived from. Unchanged items
430/// keep the same `Arc` from one projection to the next, so a pointer
431/// comparison replaces re-reading the item and re-parsing its JSON.
432///
433/// The handle records where an entry came from, not what it says, so two
434/// entries with equal content are equal whatever they were derived from.
435#[derive(Debug, Clone, Default)]
436pub struct TranscriptSource(pub Option<Arc<TranscriptItem>>);
437
438impl TranscriptSource {
439    pub fn is(&self, item: &Arc<TranscriptItem>) -> bool {
440        self.0
441            .as_ref()
442            .is_some_and(|source| Arc::ptr_eq(source, item))
443    }
444}
445
446impl PartialEq for TranscriptSource {
447    fn eq(&self, _other: &Self) -> bool {
448        true
449    }
450}
451
452impl Eq for TranscriptSource {}
453
454impl ChatEntry {
455    /// Whether this entry is the durable marker emitted when a session's
456    /// control plane restarts. The source identity is authoritative for
457    /// materialized entries; the role/text check also covers entries built
458    /// from older worker snapshots that have no materialized source handle.
459    pub fn is_session_restart(&self) -> bool {
460        self.source
461            .0
462            .as_ref()
463            .is_some_and(|item| item.is_session_restart())
464            || (self.role == ChatRole::System && self.text == SESSION_RESTART_TEXT)
465    }
466
467    pub fn plan(seq: u64, plan: Vec<PlanLine>) -> Self {
468        Self {
469            start_seq: seq,
470            seq,
471            role: ChatRole::Plan,
472            text: String::new(),
473            recorded_at_ms: None,
474            revision: 0,
475            message_id: None,
476            tool_call_id: None,
477            tool_status: None,
478            tool_summary: None,
479            tool_presentation: None,
480            tool_content: Vec::new(),
481            tool_diffstats: Vec::new(),
482            tool_locations: Vec::new(),
483            plan,
484            leading_omitted: false,
485            raw_only: false,
486            source: TranscriptSource::default(),
487        }
488    }
489
490    pub fn touch(&mut self, seq: u64) {
491        self.seq = seq;
492        self.revision = self.revision.wrapping_add(1);
493    }
494
495    /// Bound one entry to the sizes the dashboard's summary tolerates.
496    ///
497    /// Compiled unconditionally and hidden from the documentation because the
498    /// chat crate's tests need it, and a `#[cfg(test)]` item is invisible to
499    /// another crate.
500    #[doc(hidden)]
501    pub fn bounded_for_dashboard(mut self) -> Self {
502        self.bound_dashboard_content();
503        self
504    }
505
506    fn bound_dashboard_content(&mut self) {
507        const TEXT_BYTES: usize = 64 * 1024;
508        const DETAIL_BYTES: usize = 2 * 1024;
509        const DETAIL_COUNT: usize = 8;
510
511        self.leading_omitted |= truncate_string_start(&mut self.text, TEXT_BYTES);
512        for values in [
513            &mut self.tool_content,
514            &mut self.tool_diffstats,
515            &mut self.tool_locations,
516        ] {
517            values.truncate(DETAIL_COUNT);
518            for value in values {
519                truncate_string_start(value, DETAIL_BYTES);
520            }
521        }
522        if let Some(summary) = &mut self.tool_summary {
523            truncate_string_start(summary, DETAIL_BYTES);
524        }
525        if let Some(presentation) = &mut self.tool_presentation {
526            truncate_string_start(&mut presentation.summary, DETAIL_BYTES);
527            truncate_string_start(&mut presentation.source, TEXT_BYTES);
528        }
529        self.plan.truncate(DETAIL_COUNT);
530        for line in &mut self.plan {
531            truncate_string_start(&mut line.text, DETAIL_BYTES);
532        }
533    }
534
535    pub fn with_recorded_at(mut self, recorded_at_ms: Option<i64>) -> Self {
536        self.recorded_at_ms = recorded_at_ms;
537        self
538    }
539}
540
541/// Constructors that sanitize the text they are given, so terminal escape
542/// sequences from a harness never reach a transcript entry.
543impl ChatEntry {
544    pub fn plain(seq: u64, role: ChatRole, text: impl Into<String>) -> Self {
545        Self {
546            start_seq: seq,
547            seq,
548            role,
549            text: sanitize_terminal_text(&text.into()),
550            recorded_at_ms: None,
551            revision: 0,
552            message_id: None,
553            tool_call_id: None,
554            tool_status: None,
555            tool_summary: None,
556            tool_presentation: None,
557            tool_content: Vec::new(),
558            tool_diffstats: Vec::new(),
559            tool_locations: Vec::new(),
560            plan: Vec::new(),
561            leading_omitted: false,
562            raw_only: false,
563            source: TranscriptSource::default(),
564        }
565    }
566
567    pub fn tool(
568        seq: u64,
569        title: impl Into<String>,
570        tool_call_id: Option<String>,
571        tool_status: ToolStatus,
572    ) -> Self {
573        Self {
574            start_seq: seq,
575            seq,
576            role: ChatRole::Tool,
577            text: sanitize_terminal_text(&title.into()),
578            recorded_at_ms: None,
579            revision: 0,
580            message_id: None,
581            tool_call_id,
582            tool_status: Some(tool_status),
583            tool_summary: None,
584            tool_presentation: None,
585            tool_content: Vec::new(),
586            tool_diffstats: Vec::new(),
587            tool_locations: Vec::new(),
588            plan: Vec::new(),
589            leading_omitted: false,
590            raw_only: false,
591            source: TranscriptSource::default(),
592        }
593    }
594}
595
596pub(crate) fn is_false(value: &bool) -> bool {
597    !*value
598}
599
600pub fn plan_status(status: &PlanEntryStatus) -> PlanStatus {
601    match status {
602        PlanEntryStatus::InProgress => PlanStatus::Running,
603        PlanEntryStatus::Completed => PlanStatus::Completed,
604        _ => PlanStatus::Pending,
605    }
606}
607
608/// Remove terminal controls while preserving user-visible whitespace.
609pub fn sanitize_terminal_text(text: &str) -> String {
610    let mut sanitized = String::with_capacity(text.len());
611    let mut chars = text.chars().peekable();
612    while let Some(ch) = chars.next() {
613        if ch == '\x1b' {
614            // One escape can end at the ESC introducing the next one, so keep
615            // consuming rather than recursing: transcript text is untrusted and
616            // may nest these arbitrarily deep.
617            while consume_escape_body(&mut chars) {}
618        } else if ch == '\r' {
619            if chars.peek() != Some(&'\n') {
620                sanitized.push('\n');
621            }
622        } else if matches!(ch, '\n' | '\t') || !ch.is_control() {
623            sanitized.push(ch);
624        }
625    }
626    sanitized
627}
628
629/// Consume one escape sequence's body, after its introducing ESC. Returns
630/// whether the body ended at another ESC, which introduces the next sequence.
631///
632/// Dropping the ESC alone is not enough: an OSC payload (a build tool setting
633/// the window title) or the second byte of a charset selection would otherwise
634/// reach the transcript as visible text.
635fn consume_escape_body(chars: &mut std::iter::Peekable<std::str::Chars<'_>>) -> bool {
636    match chars.next() {
637        // CSI: parameter and intermediate bytes up to a final byte.
638        Some('[') => {
639            let _ = chars.find(|ch| ('@'..='~').contains(ch));
640            false
641        }
642        // OSC, DCS, SOS, PM, and APC all carry a string payload.
643        Some(']' | 'P' | 'X' | '^' | '_') => consume_string_body(chars),
644        // Two-byte sequences: charset selection (ESC ( B), ESC # 8, ESC SP F.
645        Some('(' | ')' | '*' | '+' | '-' | '.' | '/' | '#' | '%' | ' ') => {
646            chars.next();
647            false
648        }
649        // Everything else is a complete one-byte escape: ESC 7, ESC 8, ESC M,
650        // ESC =, and a trailing ESC with nothing after it.
651        _ => false,
652    }
653}
654
655/// Consume a string payload, which ends at BEL or at ST (ESC \). A line break
656/// or a cancel control aborts it instead, so one malformed OSC cannot swallow
657/// the rest of a transcript.
658fn consume_string_body(chars: &mut std::iter::Peekable<std::str::Chars<'_>>) -> bool {
659    while let Some(&ch) = chars.peek() {
660        match ch {
661            '\n' | '\r' | '\x18' | '\x1a' => return false,
662            '\x07' => {
663                chars.next();
664                return false;
665            }
666            '\x1b' => {
667                chars.next();
668                return true;
669            }
670            _ => {
671                chars.next();
672            }
673        }
674    }
675    false
676}
677
678pub fn materialized_content_text(content: &[serde_json::Value]) -> String {
679    let text = content
680        .iter()
681        .map(materialized_value_text)
682        .filter(|text| !text.is_empty())
683        .collect::<Vec<_>>()
684        .join("\n");
685    crate::relay::strip_hidden_prompt_context(&text).to_owned()
686}
687
688/// What produced a transcript item, as a stable wire name.
689///
690/// The chat view has its own role enum shaped around how it renders; this is
691/// the name the HTTP API publishes, so it changes only when the transcript
692/// model does.
693#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
694#[serde(rename_all = "snake_case")]
695pub enum TranscriptRole {
696    User,
697    Agent,
698    Thought,
699    Tool,
700    Terminal,
701    Plan,
702    PlanProposal,
703    System,
704}
705
706impl TranscriptRole {
707    pub fn as_str(self) -> &'static str {
708        match self {
709            Self::User => "user",
710            Self::Agent => "agent",
711            Self::Thought => "thought",
712            Self::Tool => "tool",
713            Self::Terminal => "terminal",
714            Self::Plan => "plan",
715            Self::PlanProposal => "plan_proposal",
716            Self::System => "system",
717        }
718    }
719    pub fn storage_kind(self) -> &'static str {
720        match self {
721            Self::Terminal => "terminal_output",
722            other => other.as_str(),
723        }
724    }
725}
726
727pub fn transcript_item_role(body: &TranscriptBody) -> &'static str {
728    let role = match body {
729        TranscriptBody::User { .. } => TranscriptRole::User,
730        TranscriptBody::Agent { .. } => TranscriptRole::Agent,
731        TranscriptBody::Thought { .. } => TranscriptRole::Thought,
732        TranscriptBody::Tool { .. } => TranscriptRole::Tool,
733        TranscriptBody::TerminalOutput { .. } => TranscriptRole::Terminal,
734        TranscriptBody::Plan { .. } => TranscriptRole::Plan,
735        TranscriptBody::PlanProposal { .. } => TranscriptRole::PlanProposal,
736        TranscriptBody::System { .. } => TranscriptRole::System,
737    };
738    role.as_str()
739}
740
741pub fn materialized_chunks_text(chunks: &[serde_json::Value]) -> String {
742    chunks
743        .iter()
744        .filter_map(|value| match ContentChunk::deserialize(value) {
745            Ok(chunk) => Some(chunk),
746            Err(error) => {
747                tracing::warn!(%error, "could not decode a stored content chunk");
748                None
749            }
750        })
751        .filter_map(|chunk| content_block_text(&chunk.content))
752        .map(|text| sanitize_terminal_text(&text))
753        .collect::<Vec<_>>()
754        .join("")
755}
756
757fn materialized_value_text(value: &serde_json::Value) -> String {
758    if let Ok(block) = ContentBlock::deserialize(value)
759        && let Some(text) = content_block_text(&block)
760    {
761        return sanitize_terminal_text(&text);
762    }
763    if let Some(text) = value.as_str() {
764        return sanitize_terminal_text(text);
765    }
766    sanitize_terminal_text(&serde_json::to_string(value).unwrap_or_else(|_| "[content]".into()))
767}
768
769pub fn tool_status(status: &ToolCallStatus) -> ToolStatus {
770    match status {
771        ToolCallStatus::InProgress => ToolStatus::Running,
772        ToolCallStatus::Completed => ToolStatus::Completed,
773        ToolCallStatus::Failed => ToolStatus::Failed,
774        _ => ToolStatus::Pending,
775    }
776}
777
778pub fn content_block_text(content: &ContentBlock) -> Option<String> {
779    match content {
780        ContentBlock::Text(text) => Some(text.text.clone()),
781        ContentBlock::Image(_) => Some("[image]".into()),
782        ContentBlock::Audio(_) => Some("[audio]".into()),
783        ContentBlock::ResourceLink(link) => Some(format!("[{}]({})", link.name, link.uri)),
784        ContentBlock::Resource(resource) => Some(match &resource.resource {
785            EmbeddedResourceResource::TextResourceContents(resource) => resource.text.clone(),
786            EmbeddedResourceResource::BlobResourceContents(resource) => {
787                format!("[embedded resource: {}]", resource.uri)
788            }
789            _ => "[embedded resource]".into(),
790        }),
791        _ => None,
792    }
793}
794
795pub fn truncate_string_start(value: &mut String, maximum_bytes: usize) -> bool {
796    if value.len() <= maximum_bytes {
797        return false;
798    }
799    let mut start = value.len() - maximum_bytes;
800    while !value.is_char_boundary(start) {
801        start += 1;
802    }
803    value.drain(..start);
804    true
805}
806
807/// The ACP tool states needed to keep a compact tool block visually useful.
808#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
809pub enum ToolStatus {
810    Pending,
811    Running,
812    Completed,
813    Failed,
814}
815
816#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
817pub enum PlanStatus {
818    Pending,
819    Running,
820    Completed,
821}
822
823#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
824pub struct PlanLine {
825    pub text: String,
826    pub status: PlanStatus,
827}
828
829#[cfg(test)]
830mod tests {
831    use super::*;
832    use serde_json::json;
833
834    fn text_chunk(text: &str, message_id: Option<&str>) -> serde_json::Value {
835        match message_id {
836            Some(id) => json!({"content": {"type": "text", "text": text}, "messageId": id}),
837            None => json!({"content": {"type": "text", "text": text}}),
838        }
839    }
840
841    #[test]
842    fn push_content_chunk_merges_adjacent_text_for_the_same_message_id() {
843        let mut chunks = vec![text_chunk("The", Some("m1"))];
844        push_content_chunk(&mut chunks, text_chunk(" quick", Some("m1")));
845        push_content_chunk(&mut chunks, text_chunk(" fox", Some("m1")));
846        assert_eq!(chunks, vec![text_chunk("The quick fox", Some("m1"))]);
847    }
848
849    #[test]
850    fn push_content_chunk_merges_adjacent_text_without_message_ids() {
851        let mut chunks = vec![text_chunk("one", None)];
852        push_content_chunk(&mut chunks, text_chunk(" two", None));
853        assert_eq!(chunks, vec![text_chunk("one two", None)]);
854    }
855
856    #[test]
857    fn push_content_chunk_keeps_chunks_from_different_message_ids_apart() {
858        let mut chunks = vec![text_chunk("first", Some("m1"))];
859        push_content_chunk(&mut chunks, text_chunk("second", Some("m2")));
860        push_content_chunk(&mut chunks, text_chunk("third", None));
861        assert_eq!(
862            chunks,
863            vec![
864                text_chunk("first", Some("m1")),
865                text_chunk("second", Some("m2")),
866                text_chunk("third", None),
867            ]
868        );
869    }
870
871    #[test]
872    fn push_content_chunk_keeps_non_text_content_separate() {
873        let image = json!({"content": {"type": "image", "data": "abc", "mimeType": "image/png"}});
874        let mut chunks = vec![text_chunk("before", None)];
875        push_content_chunk(&mut chunks, image.clone());
876        push_content_chunk(&mut chunks, image.clone());
877        push_content_chunk(&mut chunks, text_chunk("after", None));
878        assert_eq!(
879            chunks,
880            vec![
881                text_chunk("before", None),
882                image.clone(),
883                image,
884                text_chunk("after", None),
885            ]
886        );
887    }
888
889    #[test]
890    fn push_content_chunk_keeps_chunks_with_differing_metadata_apart() {
891        let mut chunks =
892            vec![json!({"content": {"type": "text", "text": "a"}, "meta": {"source": "one"}})];
893        push_content_chunk(
894            &mut chunks,
895            json!({"content": {"type": "text", "text": "b"}, "meta": {"source": "two"}}),
896        );
897        push_content_chunk(
898            &mut chunks,
899            json!({"content": {"type": "text", "text": "c"}, "meta": {"source": "two"}}),
900        );
901        assert_eq!(
902            chunks,
903            vec![
904                json!({"content": {"type": "text", "text": "a"}, "meta": {"source": "one"}}),
905                json!({"content": {"type": "text", "text": "bc"}, "meta": {"source": "two"}}),
906            ]
907        );
908    }
909
910    #[test]
911    fn push_content_chunk_keeps_chunks_with_differing_annotations_apart() {
912        let mut chunks = vec![
913            json!({"content": {"type": "text", "text": "a", "annotations": {"audience": ["user"]}}}),
914        ];
915        push_content_chunk(
916            &mut chunks,
917            json!({"content": {"type": "text", "text": "b"}}),
918        );
919        assert_eq!(
920            chunks,
921            vec![
922                json!({"content": {"type": "text", "text": "a", "annotations": {"audience": ["user"]}}}),
923                json!({"content": {"type": "text", "text": "b"}}),
924            ]
925        );
926    }
927
928    #[test]
929    fn coalesce_content_chunks_collapses_runs_and_keeps_segment_boundaries() {
930        let mut chunks = vec![
931            text_chunk("He", Some("m1")),
932            text_chunk("llo", Some("m1")),
933            text_chunk("!", Some("m1")),
934            text_chunk("next", Some("m2")),
935            text_chunk(" turn", Some("m2")),
936        ];
937        coalesce_content_chunks(&mut chunks);
938        assert_eq!(
939            chunks,
940            vec![
941                text_chunk("Hello!", Some("m1")),
942                text_chunk("next turn", Some("m2")),
943            ]
944        );
945    }
946
947    #[test]
948    fn coalesce_content_chunks_leaves_unmergeable_chunks_alone() {
949        let original = vec![
950            text_chunk("a", Some("m1")),
951            text_chunk("b", Some("m2")),
952            json!({"content": {"type": "image", "data": "x", "mimeType": "image/png"}}),
953        ];
954        let mut chunks = original.clone();
955        coalesce_content_chunks(&mut chunks);
956        assert_eq!(chunks, original);
957    }
958}