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