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