Skip to main content

vtcode_core/utils/
transcript.rs

1use once_cell::sync::Lazy;
2use parking_lot::RwLock;
3use std::cell::RefCell;
4use std::collections::VecDeque;
5use std::sync::Arc;
6
7use crate::ui::{InlineHandle, InlineMessageKind, InlineSegment, InlineTextStyle};
8pub use crate::utils::message_style::MessageStyle;
9
10const MAX_LINES: usize = 4000;
11const MAX_QUEUE_SIZE: usize = 100;
12
13static TRANSCRIPT: Lazy<RwLock<Vec<String>>> = Lazy::new(|| RwLock::new(Vec::new()));
14static INLINE_HANDLE: Lazy<RwLock<Option<Arc<InlineHandle>>>> = Lazy::new(|| RwLock::new(None));
15/// Session-scoped replaceable tracker transcript block. Shared by every tracker
16/// writer (pipeline + plan-approval handoff) so one surface owns replace/dedupe.
17static REPLACEABLE_TRACKER_BLOCK: Lazy<RwLock<Option<Vec<String>>>> = Lazy::new(|| RwLock::new(None));
18/// UI line count last written to `InlineHandle` by the tracker transcript
19/// helper. Used instead of deriving a UI `replace_last` count from TRANSCRIPT
20/// (the two stores can diverge when other writers hit only one side).
21static REPLACEABLE_TRACKER_UI_LEN: Lazy<RwLock<Option<usize>>> = Lazy::new(|| RwLock::new(None));
22
23#[derive(Clone, Copy, Debug, PartialEq, Eq)]
24enum TranscriptMode {
25    #[expect(
26        dead_code,
27        reason = "Intentional compatibility, platform, test, or API-shape suppression."
28    )]
29    Normal,
30    Suppressed,
31}
32
33thread_local! {
34    static MODE_STACK: RefCell<Vec<TranscriptMode>> = const { RefCell::new(Vec::new()) };
35}
36
37fn is_suppressed() -> bool {
38    MODE_STACK.with(|stack| matches!(stack.borrow().last(), Some(TranscriptMode::Suppressed)))
39}
40
41struct SuspensionGuard {
42    active: bool,
43}
44
45impl SuspensionGuard {
46    fn new() -> Self {
47        MODE_STACK.with(|stack| stack.borrow_mut().push(TranscriptMode::Suppressed));
48        Self { active: true }
49    }
50}
51
52impl Drop for SuspensionGuard {
53    fn drop(&mut self) {
54        if !self.active {
55            return;
56        }
57        MODE_STACK.with(|stack| {
58            let mut stack = stack.borrow_mut();
59            match stack.pop() {
60                Some(TranscriptMode::Suppressed) | None => {}
61                Some(TranscriptMode::Normal) => {
62                    debug_assert!(false, "transcript suspension stack corrupted: expected Suppressed");
63                }
64            };
65        });
66        self.active = false;
67    }
68}
69
70fn suspend() -> SuspensionGuard {
71    SuspensionGuard::new()
72}
73
74pub fn with_suppressed<F, R>(operation: F) -> R
75where
76    F: FnOnce() -> R,
77{
78    let guard = suspend();
79    let result = operation();
80    drop(guard);
81    result
82}
83
84/// Structured message with metadata for queuing
85#[derive(Clone, Debug)]
86struct QueuedMessage {
87    text: String,
88    kind: InlineMessageKind,
89    style: InlineTextStyle,
90}
91
92static MESSAGE_QUEUE: Lazy<RwLock<VecDeque<QueuedMessage>>> = Lazy::new(|| RwLock::new(VecDeque::new()));
93/// Messages enqueued while no inline handle was attached. Drained FIFO on
94/// `set_inline_handle` so none are lost. Bounded like `MESSAGE_QUEUE`.
95static PENDING_QUEUE: Lazy<RwLock<VecDeque<QueuedMessage>>> = Lazy::new(|| RwLock::new(VecDeque::new()));
96
97pub fn append(line: &str) {
98    if is_suppressed() || line.trim().is_empty() {
99        return;
100    }
101    let mut log = TRANSCRIPT.write();
102    if log.len() == MAX_LINES {
103        let drop_count = MAX_LINES / 5;
104        log.drain(0..drop_count);
105    }
106    if log.last().is_some_and(|last| last == line) {
107        return;
108    }
109    log.push(line.to_string());
110}
111
112pub fn replace_last(count: usize, lines: &[String]) {
113    if is_suppressed() {
114        return;
115    }
116    let mut log = TRANSCRIPT.write();
117    let new_len = log.len().saturating_sub(count);
118    log.truncate(new_len);
119    for line in lines {
120        if log.len() == MAX_LINES {
121            let drop_count = MAX_LINES / 5;
122            log.drain(0..drop_count);
123        }
124        log.push(line.clone());
125    }
126}
127
128pub fn tail_matches(lines: &[String]) -> bool {
129    if lines.is_empty() {
130        return false;
131    }
132
133    let log = TRANSCRIPT.read();
134    if lines.len() > log.len() {
135        return false;
136    }
137
138    log[log.len() - lines.len()..]
139        .iter()
140        .zip(lines.iter())
141        .all(|(left, right)| left == right)
142}
143
144/// Remember the current user-facing tracker transcript block.
145///
146/// `ui_line_count` is how many UI lines the tracker helper last wrote to
147/// `InlineHandle` for this block (used for tail-safe UI replace).
148pub fn remember_tracker_block_with_ui_len(lines: Vec<String>, ui_line_count: usize) {
149    let has_lines = !lines.is_empty();
150    *REPLACEABLE_TRACKER_BLOCK.write() = has_lines.then_some(lines);
151    *REPLACEABLE_TRACKER_UI_LEN.write() = has_lines.then_some(ui_line_count);
152}
153
154/// Remember the current user-facing tracker transcript block (UI length = line count).
155pub fn remember_tracker_block(lines: Vec<String>) {
156    let ui_len = lines.len();
157    remember_tracker_block_with_ui_len(lines, ui_len);
158}
159
160/// Line count of the remembered tracker transcript block, if any.
161pub fn tracker_block_len() -> Option<usize> {
162    REPLACEABLE_TRACKER_BLOCK.read().as_ref().map(|lines| lines.len())
163}
164
165/// UI write length last recorded for the tracker block, if any.
166pub fn tracker_ui_write_len() -> Option<usize> {
167    *REPLACEABLE_TRACKER_UI_LEN.read()
168}
169
170/// Whether `lines` match the remembered tracker transcript block exactly.
171pub fn tracker_block_matches(lines: &[String]) -> bool {
172    REPLACEABLE_TRACKER_BLOCK.read().as_deref() == Some(lines)
173}
174
175/// Length of the remembered tracker block **only if** it is still the transcript tail.
176pub fn tracker_block_len_if_at_tail() -> Option<usize> {
177    let remembered = REPLACEABLE_TRACKER_BLOCK.read().clone()?;
178    if tail_matches(&remembered) {
179        Some(remembered.len())
180    } else {
181        None
182    }
183}
184
185/// Clear remembered tracker replace state (tests / session teardown).
186pub fn clear_tracker_block() {
187    *REPLACEABLE_TRACKER_BLOCK.write() = None;
188    *REPLACEABLE_TRACKER_UI_LEN.write() = None;
189}
190
191/// Last non-empty transcript line, if any.
192pub fn last_line() -> Option<String> {
193    TRANSCRIPT.read().iter().rev().find(|line| !line.trim().is_empty()).cloned()
194}
195
196pub fn snapshot() -> Vec<String> {
197    TRANSCRIPT.read().clone()
198}
199
200pub fn len() -> usize {
201    TRANSCRIPT.read().len()
202}
203
204pub fn clear() {
205    TRANSCRIPT.write().clear();
206    clear_tracker_block();
207}
208
209/// Set the inline handle for immediate message display.
210///
211/// Any messages enqueued while no handle was attached are replayed FIFO
212/// exactly once, in order, so none are lost. The handle is captured before
213/// the drain so a concurrent `clear_inline_handle` cannot drop already-taken
214/// pending messages.
215pub fn set_inline_handle(handle: Arc<InlineHandle>) {
216    *INLINE_HANDLE.write() = Some(handle.clone());
217    let pending: Vec<QueuedMessage> = PENDING_QUEUE.write().drain(..).collect();
218    for msg in pending {
219        display_message_to(&handle, &msg.text, msg.kind, &msg.style);
220    }
221}
222
223/// Remove the inline handle
224pub fn clear_inline_handle() {
225    *INLINE_HANDLE.write() = None;
226}
227
228/// Map MessageStyle to InlineMessageKind
229fn message_kind(style: MessageStyle) -> InlineMessageKind {
230    style.message_kind()
231}
232
233/// Enqueue a message with a specific style and display it immediately
234pub fn enqueue_message(message: &str, style: MessageStyle) {
235    enqueue_message_with_kind(message, message_kind(style), InlineTextStyle::default())
236}
237
238/// Enqueue a message with a specific kind and display it immediately
239pub fn enqueue_message_with_kind(message: &str, kind: InlineMessageKind, text_style: InlineTextStyle) {
240    if message.trim().is_empty() {
241        return;
242    }
243
244    let queued = QueuedMessage { text: message.to_string(), kind, style: text_style };
245
246    // Record history (bounded FIFO, drops oldest on overflow).
247    {
248        let mut queue = MESSAGE_QUEUE.write();
249        if queue.len() >= MAX_QUEUE_SIZE {
250            queue.pop_front();
251        }
252        queue.push_back(queued.clone());
253    }
254
255    // Retain FIFO pending first, then drain for display when a handle is
256    // attached. Never re-read `.back()`: under concurrency that shows the
257    // last writer's message twice while dropping earlier ones. Pushing
258    // before the handle check closes the lost-wakeup race where `set_inline_handle`
259    // drains between the check and the push; the atomic drain below guarantees
260    // each pending message displays exactly once across both paths.
261    {
262        let mut pending = PENDING_QUEUE.write();
263        if pending.len() >= MAX_QUEUE_SIZE {
264            pending.pop_front();
265        }
266        pending.push_back(queued);
267    }
268    if let Some(handle) = INLINE_HANDLE.read().clone() {
269        let pending: Vec<QueuedMessage> = PENDING_QUEUE.write().drain(..).collect();
270        for msg in pending {
271            display_message_to(&handle, &msg.text, msg.kind, &msg.style);
272        }
273    }
274
275    // Also add to transcript for persistence (plain text)
276    append(message);
277}
278
279/// Display an error message in the transcript instead of the input field
280#[cold]
281pub fn display_error(message: &str) {
282    enqueue_message(message, MessageStyle::Error);
283}
284
285/// Display an info message in the transcript instead of the input field
286pub fn display_info(message: &str) {
287    enqueue_message(message, MessageStyle::Info);
288}
289
290/// Display a message immediately without queuing (low-level function)
291fn display_message_now(text: &str, kind: InlineMessageKind, style: &InlineTextStyle) {
292    if let Some(handle) = INLINE_HANDLE.read().as_ref() {
293        display_message_to(handle, text, kind, style);
294    }
295}
296
297/// Append one message to a specific handle. Used by drain paths that already
298/// captured the handle so a concurrent detach cannot swallow the message.
299fn display_message_to(handle: &InlineHandle, text: &str, kind: InlineMessageKind, style: &InlineTextStyle) {
300    handle.append_line(
301        kind,
302        vec![InlineSegment {
303            text: text.to_string(),
304            style: Arc::new(style.clone()),
305        }],
306    );
307}
308
309/// Enqueue a message and display it immediately (defaults to Output/Pty style)
310pub fn enqueue(message: &str) {
311    enqueue_message(message, MessageStyle::Output);
312}
313
314/// Display a message immediately without enqueueing or adding to transcript
315pub fn display_immediate(message: &str, style: MessageStyle) {
316    display_message_now(message, message_kind(style), &InlineTextStyle::default());
317}
318
319/// Get all queued messages as plain text
320pub fn get_queued_messages() -> Vec<String> {
321    MESSAGE_QUEUE.read().iter().map(|m| m.text.clone()).collect()
322}
323
324/// Get all queued messages with their metadata
325pub fn get_queued_messages_with_metadata() -> Vec<(String, InlineMessageKind)> {
326    MESSAGE_QUEUE.read().iter().map(|m| (m.text.clone(), m.kind)).collect()
327}
328
329/// Clear the message queue (history and undisplayed pending).
330pub fn clear_queue() {
331    MESSAGE_QUEUE.write().clear();
332    PENDING_QUEUE.write().clear();
333}
334
335/// Get queue length (history length).
336pub fn queue_len() -> usize {
337    MESSAGE_QUEUE.read().len()
338}
339
340/// Get pending (undisplayed) queue length.
341pub fn pending_queue_len() -> usize {
342    PENDING_QUEUE.read().len()
343}
344
345/// Atomically take every queued message FIFO and clear the history queue.
346/// Prefer this over `get_queued_messages` + `clear_queue`: the snapshot
347/// pattern races and callers that only use `.last()` drop all but the
348/// newest message.
349pub fn drain_queued_messages() -> Vec<String> {
350    PENDING_QUEUE.write().clear();
351    MESSAGE_QUEUE.write().drain(..).map(|m| m.text).collect()
352}
353
354/// Atomically take every queued message with metadata FIFO and clear.
355pub fn drain_queued_messages_with_metadata() -> Vec<(String, InlineMessageKind)> {
356    PENDING_QUEUE.write().clear();
357    MESSAGE_QUEUE.write().drain(..).map(|m| (m.text, m.kind)).collect()
358}
359
360/// Replay all queued messages to the current inline handle (useful for recovery)
361pub fn replay_queued_messages() {
362    let messages: Vec<QueuedMessage> = MESSAGE_QUEUE.read().iter().cloned().collect();
363    for msg in messages {
364        display_message_now(&msg.text, msg.kind, &msg.style);
365    }
366}
367
368#[cfg(test)]
369mod tests {
370    use super::*;
371
372    #[test]
373    #[serial_test::serial(transcript_state)]
374    fn append_and_snapshot_store_lines() {
375        clear();
376        append("first");
377        append("second");
378        assert_eq!(len(), 2);
379        let snap = snapshot();
380        assert_eq!(snap, vec!["first".to_owned(), "second".to_owned()]);
381        clear();
382    }
383
384    #[test]
385    #[serial_test::serial(transcript_state)]
386    fn append_skips_adjacent_duplicate_lines() {
387        clear();
388        append("duplicate");
389        append("duplicate");
390        append("different");
391        append("duplicate");
392        let snap = snapshot();
393        assert_eq!(snap, vec!["duplicate".to_owned(), "different".to_owned(), "duplicate".to_owned()]);
394        clear();
395    }
396
397    #[test]
398    #[serial_test::serial(transcript_state)]
399    fn transcript_drops_oldest_chunk_when_full() {
400        clear();
401        for idx in 0..MAX_LINES {
402            append(&format!("line {idx}"));
403        }
404        assert_eq!(len(), MAX_LINES);
405        for extra in 0..10 {
406            append(&format!("extra {extra}"));
407        }
408        assert_eq!(len(), MAX_LINES - (MAX_LINES / 5) + 10);
409        let snap = snapshot();
410        assert_eq!(snap.first().unwrap(), &format!("line {}", MAX_LINES / 5).to_owned());
411        clear();
412    }
413
414    #[test]
415    #[serial_test::serial(transcript_state)]
416    fn message_queue_enqueue_and_retrieve() {
417        clear_queue();
418        assert_eq!(queue_len(), 0);
419
420        enqueue("first message");
421        enqueue_message("second message", MessageStyle::Info);
422
423        assert_eq!(queue_len(), 2);
424        let messages = get_queued_messages();
425        assert_eq!(messages, vec!["first message", "second message"]);
426
427        clear_queue();
428        assert_eq!(queue_len(), 0);
429    }
430
431    #[test]
432    #[serial_test::serial(transcript_state)]
433    fn message_queue_preserves_metadata() {
434        clear_queue();
435
436        enqueue_message("info message", MessageStyle::Info);
437        enqueue_message("error message", MessageStyle::Error);
438        enqueue_message("user message", MessageStyle::User);
439
440        let messages = get_queued_messages_with_metadata();
441        assert_eq!(messages.len(), 3);
442        assert_eq!(messages[0].1, InlineMessageKind::Info);
443        assert_eq!(messages[1].1, InlineMessageKind::Error);
444        assert_eq!(messages[2].1, InlineMessageKind::User);
445
446        clear_queue();
447    }
448
449    #[test]
450    #[serial_test::serial(transcript_state)]
451    fn message_queue_size_limit() {
452        clear_queue();
453
454        // Fill queue to max capacity
455        for i in 0..MAX_QUEUE_SIZE {
456            enqueue(&format!("message {i}"));
457        }
458        assert_eq!(queue_len(), MAX_QUEUE_SIZE);
459
460        // Add one more message - should drop oldest
461        enqueue("overflow message");
462        assert_eq!(queue_len(), MAX_QUEUE_SIZE);
463
464        let messages = get_queued_messages();
465        assert_eq!(messages.first().unwrap(), "message 1"); // First message should be dropped
466        assert_eq!(messages.last().unwrap(), "overflow message");
467
468        clear_queue();
469    }
470
471    #[test]
472    #[serial_test::serial(transcript_state)]
473    fn suppressed_scope_skips_transcript_entries() {
474        clear();
475        with_suppressed(|| {
476            append("hidden");
477        });
478        append("visible");
479        let snap = snapshot();
480        assert_eq!(snap, vec!["visible".to_owned()]);
481        clear();
482    }
483
484    #[test]
485    #[serial_test::serial(transcript_state)]
486    fn tracker_block_store_replaces_and_clears() {
487        clear();
488        assert_eq!(tracker_block_len(), None);
489        remember_tracker_block(vec!["• Plan 0/1".to_string()]);
490        assert_eq!(tracker_block_len(), Some(1));
491        assert!(tracker_block_matches(&["• Plan 0/1".to_string()]));
492        remember_tracker_block(vec!["• Plan 1/1".to_string()]);
493        assert!(!tracker_block_matches(&["• Plan 0/1".to_string()]));
494        clear();
495        assert_eq!(tracker_block_len(), None);
496    }
497
498    #[test]
499    #[serial_test::serial(transcript_state)]
500    fn drain_queued_messages_returns_all_fifo_and_clears() {
501        clear();
502        clear_queue();
503        clear_inline_handle();
504        assert_eq!(pending_queue_len(), 0);
505
506        enqueue("alpha");
507        enqueue("beta");
508        enqueue("gamma");
509
510        assert_eq!(queue_len(), 3);
511        assert_eq!(pending_queue_len(), 3);
512        assert_eq!(get_queued_messages(), vec!["alpha".to_owned(), "beta".to_owned(), "gamma".to_owned()]);
513
514        let drained = drain_queued_messages();
515        assert_eq!(drained, vec!["alpha".to_owned(), "beta".to_owned(), "gamma".to_owned()]);
516        assert_eq!(queue_len(), 0);
517        assert_eq!(pending_queue_len(), 0);
518        assert!(get_queued_messages().is_empty());
519
520        clear();
521        clear_queue();
522        clear_inline_handle();
523    }
524
525    #[test]
526    #[serial_test::serial(transcript_state)]
527    fn drain_queued_messages_with_metadata_preserves_order() {
528        clear();
529        clear_queue();
530        clear_inline_handle();
531
532        enqueue_message("first", MessageStyle::Info);
533        enqueue_message("second", MessageStyle::Error);
534
535        let drained = drain_queued_messages_with_metadata();
536        assert_eq!(drained.len(), 2);
537        assert_eq!(drained[0].0, "first");
538        assert_eq!(drained[0].1, InlineMessageKind::Info);
539        assert_eq!(drained[1].0, "second");
540        assert_eq!(drained[1].1, InlineMessageKind::Error);
541        assert_eq!(queue_len(), 0);
542
543        clear();
544        clear_queue();
545        clear_inline_handle();
546    }
547
548    #[test]
549    #[serial_test::serial(transcript_state)]
550    fn pending_queue_replays_fifo_once_on_handle_set() {
551        use crate::ui::InlineCommand;
552
553        clear();
554        clear_queue();
555        clear_inline_handle();
556
557        enqueue("first pending");
558        enqueue("second pending");
559        assert_eq!(pending_queue_len(), 2);
560
561        let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
562        let handle = InlineHandle::new_for_tests(tx);
563        set_inline_handle(Arc::new(handle));
564        assert_eq!(pending_queue_len(), 0);
565        // History is retained for inspection; pending is what was replayed.
566        assert_eq!(queue_len(), 2);
567
568        let mut texts = Vec::new();
569        while let Ok(cmd) = rx.try_recv() {
570            if let InlineCommand::AppendLine { segments, .. } = cmd {
571                texts.extend(segments.into_iter().map(|s| s.text));
572            }
573        }
574        assert_eq!(texts, vec!["first pending".to_owned(), "second pending".to_owned()]);
575
576        // Second handle attach must not duplicate replay.
577        let (tx2, mut rx2) = tokio::sync::mpsc::unbounded_channel();
578        let handle2 = InlineHandle::new_for_tests(tx2);
579        set_inline_handle(Arc::new(handle2));
580        assert!(rx2.try_recv().is_err());
581
582        clear();
583        clear_queue();
584        clear_inline_handle();
585    }
586
587    #[test]
588    #[serial_test::serial(transcript_state)]
589    fn immediate_display_uses_each_message_not_last_only() {
590        use crate::ui::InlineCommand;
591
592        clear();
593        clear_queue();
594        clear_inline_handle();
595
596        let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
597        let handle = InlineHandle::new_for_tests(tx);
598        set_inline_handle(Arc::new(handle));
599
600        enqueue("one");
601        enqueue("two");
602        assert_eq!(pending_queue_len(), 0);
603
604        let mut texts = Vec::new();
605        while let Ok(cmd) = rx.try_recv() {
606            if let InlineCommand::AppendLine { segments, .. } = cmd {
607                texts.extend(segments.into_iter().map(|s| s.text));
608            }
609        }
610        assert_eq!(texts, vec!["one".to_owned(), "two".to_owned()]);
611
612        clear();
613        clear_queue();
614        clear_inline_handle();
615    }
616}