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    /// A completed tool call whose turn was interrupted while it ran. The
423    /// harness kept the command going and it ended later on its own, so its
424    /// completion is not the agent's finished work.
425    #[serde(default, skip_serializing_if = "is_false")]
426    pub ended_after_interrupt: bool,
427    /// The materialized transcript item this entry was derived from, when it
428    /// came from the controller's projection. Provenance only, so it is
429    /// neither serialized nor part of the entry's value.
430    #[serde(skip)]
431    pub source: TranscriptSource,
432}
433
434/// Handle on the transcript item an entry was derived from. Unchanged items
435/// keep the same `Arc` from one projection to the next, so a pointer
436/// comparison replaces re-reading the item and re-parsing its JSON.
437///
438/// The handle records where an entry came from, not what it says, so two
439/// entries with equal content are equal whatever they were derived from.
440#[derive(Debug, Clone, Default)]
441pub struct TranscriptSource(pub Option<Arc<TranscriptItem>>);
442
443impl TranscriptSource {
444    pub fn is(&self, item: &Arc<TranscriptItem>) -> bool {
445        self.0
446            .as_ref()
447            .is_some_and(|source| Arc::ptr_eq(source, item))
448    }
449}
450
451impl PartialEq for TranscriptSource {
452    fn eq(&self, _other: &Self) -> bool {
453        true
454    }
455}
456
457impl Eq for TranscriptSource {}
458
459impl ChatEntry {
460    /// Whether this entry is the durable marker emitted when a session's
461    /// control plane restarts. The source identity is authoritative for
462    /// materialized entries; the role/text check also covers entries built
463    /// from older worker snapshots that have no materialized source handle.
464    pub fn is_session_restart(&self) -> bool {
465        self.source
466            .0
467            .as_ref()
468            .is_some_and(|item| item.is_session_restart())
469            || (self.role == ChatRole::System && self.text == SESSION_RESTART_TEXT)
470    }
471
472    pub fn plan(seq: u64, plan: Vec<PlanLine>) -> Self {
473        Self {
474            start_seq: seq,
475            seq,
476            role: ChatRole::Plan,
477            text: String::new(),
478            recorded_at_ms: None,
479            revision: 0,
480            message_id: None,
481            tool_call_id: None,
482            tool_status: None,
483            tool_summary: None,
484            tool_presentation: None,
485            tool_content: Vec::new(),
486            tool_diffstats: Vec::new(),
487            tool_locations: Vec::new(),
488            plan,
489            leading_omitted: false,
490            raw_only: false,
491            ended_after_interrupt: false,
492            source: TranscriptSource::default(),
493        }
494    }
495
496    pub fn touch(&mut self, seq: u64) {
497        self.seq = seq;
498        self.revision = self.revision.wrapping_add(1);
499    }
500
501    /// Bound one entry to the sizes the dashboard's summary tolerates.
502    ///
503    /// Compiled unconditionally and hidden from the documentation because the
504    /// chat crate's tests need it, and a `#[cfg(test)]` item is invisible to
505    /// another crate.
506    #[doc(hidden)]
507    pub fn bounded_for_dashboard(mut self) -> Self {
508        self.bound_dashboard_content();
509        self
510    }
511
512    fn bound_dashboard_content(&mut self) {
513        const TEXT_BYTES: usize = 64 * 1024;
514        const DETAIL_BYTES: usize = 2 * 1024;
515        const DETAIL_COUNT: usize = 8;
516
517        self.leading_omitted |= truncate_string_start(&mut self.text, TEXT_BYTES);
518        for values in [
519            &mut self.tool_content,
520            &mut self.tool_diffstats,
521            &mut self.tool_locations,
522        ] {
523            values.truncate(DETAIL_COUNT);
524            for value in values {
525                truncate_string_start(value, DETAIL_BYTES);
526            }
527        }
528        if let Some(summary) = &mut self.tool_summary {
529            truncate_string_start(summary, DETAIL_BYTES);
530        }
531        if let Some(presentation) = &mut self.tool_presentation {
532            truncate_string_start(&mut presentation.summary, DETAIL_BYTES);
533            truncate_string_start(&mut presentation.source, TEXT_BYTES);
534        }
535        self.plan.truncate(DETAIL_COUNT);
536        for line in &mut self.plan {
537            truncate_string_start(&mut line.text, DETAIL_BYTES);
538        }
539    }
540
541    pub fn with_recorded_at(mut self, recorded_at_ms: Option<i64>) -> Self {
542        self.recorded_at_ms = recorded_at_ms;
543        self
544    }
545}
546
547/// Constructors that sanitize the text they are given, so terminal escape
548/// sequences from a harness never reach a transcript entry.
549impl ChatEntry {
550    pub fn plain(seq: u64, role: ChatRole, text: impl Into<String>) -> Self {
551        Self {
552            start_seq: seq,
553            seq,
554            role,
555            text: sanitize_terminal_text(&text.into()),
556            recorded_at_ms: None,
557            revision: 0,
558            message_id: None,
559            tool_call_id: None,
560            tool_status: None,
561            tool_summary: None,
562            tool_presentation: None,
563            tool_content: Vec::new(),
564            tool_diffstats: Vec::new(),
565            tool_locations: Vec::new(),
566            plan: Vec::new(),
567            leading_omitted: false,
568            raw_only: false,
569            ended_after_interrupt: false,
570            source: TranscriptSource::default(),
571        }
572    }
573
574    pub fn tool(
575        seq: u64,
576        title: impl Into<String>,
577        tool_call_id: Option<String>,
578        tool_status: ToolStatus,
579    ) -> Self {
580        Self {
581            start_seq: seq,
582            seq,
583            role: ChatRole::Tool,
584            text: sanitize_terminal_text(&title.into()),
585            recorded_at_ms: None,
586            revision: 0,
587            message_id: None,
588            tool_call_id,
589            tool_status: Some(tool_status),
590            tool_summary: None,
591            tool_presentation: None,
592            tool_content: Vec::new(),
593            tool_diffstats: Vec::new(),
594            tool_locations: Vec::new(),
595            plan: Vec::new(),
596            leading_omitted: false,
597            raw_only: false,
598            ended_after_interrupt: false,
599            source: TranscriptSource::default(),
600        }
601    }
602}
603
604pub(crate) fn is_false(value: &bool) -> bool {
605    !*value
606}
607
608pub fn plan_status(status: &PlanEntryStatus) -> PlanStatus {
609    match status {
610        PlanEntryStatus::InProgress => PlanStatus::Running,
611        PlanEntryStatus::Completed => PlanStatus::Completed,
612        _ => PlanStatus::Pending,
613    }
614}
615
616/// Remove terminal controls while preserving user-visible whitespace.
617pub fn sanitize_terminal_text(text: &str) -> String {
618    let mut sanitized = String::with_capacity(text.len());
619    let mut chars = text.chars().peekable();
620    while let Some(ch) = chars.next() {
621        if ch == '\x1b' {
622            // One escape can end at the ESC introducing the next one, so keep
623            // consuming rather than recursing: transcript text is untrusted and
624            // may nest these arbitrarily deep.
625            while consume_escape_body(&mut chars) {}
626        } else if ch == '\r' {
627            if chars.peek() != Some(&'\n') {
628                sanitized.push('\n');
629            }
630        } else if matches!(ch, '\n' | '\t') || !ch.is_control() {
631            sanitized.push(ch);
632        }
633    }
634    sanitized
635}
636
637/// Consume one escape sequence's body, after its introducing ESC. Returns
638/// whether the body ended at another ESC, which introduces the next sequence.
639///
640/// Dropping the ESC alone is not enough: an OSC payload (a build tool setting
641/// the window title) or the second byte of a charset selection would otherwise
642/// reach the transcript as visible text.
643fn consume_escape_body(chars: &mut std::iter::Peekable<std::str::Chars<'_>>) -> bool {
644    match chars.next() {
645        // CSI: parameter and intermediate bytes up to a final byte.
646        Some('[') => {
647            let _ = chars.find(|ch| ('@'..='~').contains(ch));
648            false
649        }
650        // OSC, DCS, SOS, PM, and APC all carry a string payload.
651        Some(']' | 'P' | 'X' | '^' | '_') => consume_string_body(chars),
652        // Two-byte sequences: charset selection (ESC ( B), ESC # 8, ESC SP F.
653        Some('(' | ')' | '*' | '+' | '-' | '.' | '/' | '#' | '%' | ' ') => {
654            chars.next();
655            false
656        }
657        // Everything else is a complete one-byte escape: ESC 7, ESC 8, ESC M,
658        // ESC =, and a trailing ESC with nothing after it.
659        _ => false,
660    }
661}
662
663/// Consume a string payload, which ends at BEL or at ST (ESC \). A line break
664/// or a cancel control aborts it instead, so one malformed OSC cannot swallow
665/// the rest of a transcript.
666fn consume_string_body(chars: &mut std::iter::Peekable<std::str::Chars<'_>>) -> bool {
667    while let Some(&ch) = chars.peek() {
668        match ch {
669            '\n' | '\r' | '\x18' | '\x1a' => return false,
670            '\x07' => {
671                chars.next();
672                return false;
673            }
674            '\x1b' => {
675                chars.next();
676                return true;
677            }
678            _ => {
679                chars.next();
680            }
681        }
682    }
683    false
684}
685
686pub fn materialized_content_text(content: &[serde_json::Value]) -> String {
687    let text = content
688        .iter()
689        .map(materialized_value_text)
690        .filter(|text| !text.is_empty())
691        .collect::<Vec<_>>()
692        .join("\n");
693    crate::relay::strip_hidden_prompt_context(&text).to_owned()
694}
695
696/// What produced a transcript item, as a stable wire name.
697///
698/// The chat view has its own role enum shaped around how it renders; this is
699/// the name the HTTP API publishes, so it changes only when the transcript
700/// model does.
701#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
702#[serde(rename_all = "snake_case")]
703pub enum TranscriptRole {
704    User,
705    Agent,
706    Thought,
707    Tool,
708    Terminal,
709    Plan,
710    PlanProposal,
711    System,
712}
713
714impl TranscriptRole {
715    pub fn as_str(self) -> &'static str {
716        match self {
717            Self::User => "user",
718            Self::Agent => "agent",
719            Self::Thought => "thought",
720            Self::Tool => "tool",
721            Self::Terminal => "terminal",
722            Self::Plan => "plan",
723            Self::PlanProposal => "plan_proposal",
724            Self::System => "system",
725        }
726    }
727    pub fn storage_kind(self) -> &'static str {
728        match self {
729            Self::Terminal => "terminal_output",
730            other => other.as_str(),
731        }
732    }
733}
734
735pub fn transcript_item_role(body: &TranscriptBody) -> &'static str {
736    let role = match body {
737        TranscriptBody::User { .. } => TranscriptRole::User,
738        TranscriptBody::Agent { .. } => TranscriptRole::Agent,
739        TranscriptBody::Thought { .. } => TranscriptRole::Thought,
740        TranscriptBody::Tool { .. } => TranscriptRole::Tool,
741        TranscriptBody::TerminalOutput { .. } => TranscriptRole::Terminal,
742        TranscriptBody::Plan { .. } => TranscriptRole::Plan,
743        TranscriptBody::PlanProposal { .. } => TranscriptRole::PlanProposal,
744        TranscriptBody::System { .. } => TranscriptRole::System,
745    };
746    role.as_str()
747}
748
749pub fn materialized_chunks_text(chunks: &[serde_json::Value]) -> String {
750    chunks
751        .iter()
752        .filter_map(|value| match ContentChunk::deserialize(value) {
753            Ok(chunk) => Some(chunk),
754            Err(error) => {
755                tracing::warn!(%error, "could not decode a stored content chunk");
756                None
757            }
758        })
759        .filter_map(|chunk| content_block_text(&chunk.content))
760        .map(|text| sanitize_terminal_text(&text))
761        .collect::<Vec<_>>()
762        .join("")
763}
764
765fn materialized_value_text(value: &serde_json::Value) -> String {
766    if let Ok(block) = ContentBlock::deserialize(value)
767        && let Some(text) = content_block_text(&block)
768    {
769        return sanitize_terminal_text(&text);
770    }
771    if let Some(text) = value.as_str() {
772        return sanitize_terminal_text(text);
773    }
774    sanitize_terminal_text(&serde_json::to_string(value).unwrap_or_else(|_| "[content]".into()))
775}
776
777pub fn tool_status(status: &ToolCallStatus) -> ToolStatus {
778    match status {
779        ToolCallStatus::InProgress => ToolStatus::Running,
780        ToolCallStatus::Completed => ToolStatus::Completed,
781        ToolCallStatus::Failed => ToolStatus::Failed,
782        _ => ToolStatus::Pending,
783    }
784}
785
786pub fn content_block_text(content: &ContentBlock) -> Option<String> {
787    match content {
788        ContentBlock::Text(text) => Some(text.text.clone()),
789        ContentBlock::Image(_) => Some("[image]".into()),
790        ContentBlock::Audio(_) => Some("[audio]".into()),
791        ContentBlock::ResourceLink(link) => Some(format!("[{}]({})", link.name, link.uri)),
792        ContentBlock::Resource(resource) => Some(match &resource.resource {
793            EmbeddedResourceResource::TextResourceContents(resource) => resource.text.clone(),
794            EmbeddedResourceResource::BlobResourceContents(resource) => {
795                format!("[embedded resource: {}]", resource.uri)
796            }
797            _ => "[embedded resource]".into(),
798        }),
799        _ => None,
800    }
801}
802
803pub fn truncate_string_start(value: &mut String, maximum_bytes: usize) -> bool {
804    if value.len() <= maximum_bytes {
805        return false;
806    }
807    let mut start = value.len() - maximum_bytes;
808    while !value.is_char_boundary(start) {
809        start += 1;
810    }
811    value.drain(..start);
812    true
813}
814
815/// The ACP tool states needed to keep a compact tool block visually useful.
816#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
817pub enum ToolStatus {
818    Pending,
819    Running,
820    Completed,
821    Failed,
822}
823
824#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
825pub enum PlanStatus {
826    Pending,
827    Running,
828    Completed,
829}
830
831#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
832pub struct PlanLine {
833    pub text: String,
834    pub status: PlanStatus,
835}
836
837#[cfg(test)]
838mod tests {
839    use super::*;
840    use serde_json::json;
841
842    fn text_chunk(text: &str, message_id: Option<&str>) -> serde_json::Value {
843        match message_id {
844            Some(id) => json!({"content": {"type": "text", "text": text}, "messageId": id}),
845            None => json!({"content": {"type": "text", "text": text}}),
846        }
847    }
848
849    #[test]
850    fn push_content_chunk_merges_adjacent_text_for_the_same_message_id() {
851        let mut chunks = vec![text_chunk("The", Some("m1"))];
852        push_content_chunk(&mut chunks, text_chunk(" quick", Some("m1")));
853        push_content_chunk(&mut chunks, text_chunk(" fox", Some("m1")));
854        assert_eq!(chunks, vec![text_chunk("The quick fox", Some("m1"))]);
855    }
856
857    #[test]
858    fn push_content_chunk_merges_adjacent_text_without_message_ids() {
859        let mut chunks = vec![text_chunk("one", None)];
860        push_content_chunk(&mut chunks, text_chunk(" two", None));
861        assert_eq!(chunks, vec![text_chunk("one two", None)]);
862    }
863
864    #[test]
865    fn push_content_chunk_keeps_chunks_from_different_message_ids_apart() {
866        let mut chunks = vec![text_chunk("first", Some("m1"))];
867        push_content_chunk(&mut chunks, text_chunk("second", Some("m2")));
868        push_content_chunk(&mut chunks, text_chunk("third", None));
869        assert_eq!(
870            chunks,
871            vec![
872                text_chunk("first", Some("m1")),
873                text_chunk("second", Some("m2")),
874                text_chunk("third", None),
875            ]
876        );
877    }
878
879    #[test]
880    fn push_content_chunk_keeps_non_text_content_separate() {
881        let image = json!({"content": {"type": "image", "data": "abc", "mimeType": "image/png"}});
882        let mut chunks = vec![text_chunk("before", None)];
883        push_content_chunk(&mut chunks, image.clone());
884        push_content_chunk(&mut chunks, image.clone());
885        push_content_chunk(&mut chunks, text_chunk("after", None));
886        assert_eq!(
887            chunks,
888            vec![
889                text_chunk("before", None),
890                image.clone(),
891                image,
892                text_chunk("after", None),
893            ]
894        );
895    }
896
897    #[test]
898    fn push_content_chunk_keeps_chunks_with_differing_metadata_apart() {
899        let mut chunks =
900            vec![json!({"content": {"type": "text", "text": "a"}, "meta": {"source": "one"}})];
901        push_content_chunk(
902            &mut chunks,
903            json!({"content": {"type": "text", "text": "b"}, "meta": {"source": "two"}}),
904        );
905        push_content_chunk(
906            &mut chunks,
907            json!({"content": {"type": "text", "text": "c"}, "meta": {"source": "two"}}),
908        );
909        assert_eq!(
910            chunks,
911            vec![
912                json!({"content": {"type": "text", "text": "a"}, "meta": {"source": "one"}}),
913                json!({"content": {"type": "text", "text": "bc"}, "meta": {"source": "two"}}),
914            ]
915        );
916    }
917
918    #[test]
919    fn push_content_chunk_keeps_chunks_with_differing_annotations_apart() {
920        let mut chunks = vec![
921            json!({"content": {"type": "text", "text": "a", "annotations": {"audience": ["user"]}}}),
922        ];
923        push_content_chunk(
924            &mut chunks,
925            json!({"content": {"type": "text", "text": "b"}}),
926        );
927        assert_eq!(
928            chunks,
929            vec![
930                json!({"content": {"type": "text", "text": "a", "annotations": {"audience": ["user"]}}}),
931                json!({"content": {"type": "text", "text": "b"}}),
932            ]
933        );
934    }
935
936    #[test]
937    fn coalesce_content_chunks_collapses_runs_and_keeps_segment_boundaries() {
938        let mut chunks = vec![
939            text_chunk("He", Some("m1")),
940            text_chunk("llo", Some("m1")),
941            text_chunk("!", Some("m1")),
942            text_chunk("next", Some("m2")),
943            text_chunk(" turn", Some("m2")),
944        ];
945        coalesce_content_chunks(&mut chunks);
946        assert_eq!(
947            chunks,
948            vec![
949                text_chunk("Hello!", Some("m1")),
950                text_chunk("next turn", Some("m2")),
951            ]
952        );
953    }
954
955    #[test]
956    fn coalesce_content_chunks_leaves_unmergeable_chunks_alone() {
957        let original = vec![
958            text_chunk("a", Some("m1")),
959            text_chunk("b", Some("m2")),
960            json!({"content": {"type": "image", "data": "x", "mimeType": "image/png"}}),
961        ];
962        let mut chunks = original.clone();
963        coalesce_content_chunks(&mut chunks);
964        assert_eq!(chunks, original);
965    }
966}