Skip to main content

mj_transcript/
turn_context.rs

1//! Process-local classification evidence using the shared transcript summary.
2use agent_client_protocol::schema::v1::{ContentBlock, SessionUpdate};
3use mj_core::activity::ActivityFacts;
4use mj_core::activity::verdict::*;
5use mj_core::config::HarnessKind;
6use std::sync::{Arc, Mutex};
7
8#[derive(Debug, Default)]
9struct TurnContextState {
10    generation: u64,
11    decision_log: Option<mj_core::jev::DecisionLog>,
12    decision: Option<(u64, String)>,
13    user_prompt_tail: String,
14    message_id: Option<String>,
15    assistant_text_tail: String,
16    last_completed_message: String,
17    summary: crate::summary::TranscriptSummary,
18    background_commands: usize,
19    queued_commands: usize,
20    background_inventory: Vec<mj_core::relay::BackgroundCommand>,
21    native_agent_ids: Vec<String>,
22    session_id: String,
23}
24
25/// The prompt that became visible to the harness at this durable observation.
26/// A queued or unconfirmed steering prompt has not been delivered yet.
27pub fn delivered_prompt_command_id(observation: &mj_core::relay::RelayObservation) -> Option<&str> {
28    use mj_core::relay::{RelayCommandOutcome, RelayObservation};
29    match observation {
30        RelayObservation::CommandStarted { command_id, .. } => Some(command_id),
31        RelayObservation::CommandCompleted {
32            outcome: RelayCommandOutcome::Steered { queued_command_id },
33            ..
34        } => Some(queued_command_id),
35        _ => None,
36    }
37}
38
39/// The relay and runtime share this small process-local evidence accumulator.
40#[derive(Debug, Default, Clone)]
41pub struct TurnContext(Arc<Mutex<TurnContextState>>);
42
43impl TurnContext {
44    pub fn set_decision_log(&self, log: mj_core::jev::DecisionLog) {
45        self.0
46            .lock()
47            .expect("turn context lock poisoned")
48            .decision_log = Some(log);
49    }
50    pub fn decision_log(&self) -> Option<mj_core::jev::DecisionLog> {
51        self.0
52            .lock()
53            .expect("turn context lock poisoned")
54            .decision_log
55            .clone()
56    }
57    pub fn set_decision(&self, generation: u64, id: String) {
58        self.0.lock().expect("turn context lock poisoned").decision = Some((generation, id));
59    }
60    pub fn decision(&self) -> Option<String> {
61        let state = self.0.lock().expect("turn context lock poisoned");
62        state
63            .decision
64            .as_ref()
65            .filter(|(generation, _)| *generation == state.generation)
66            .map(|(_, id)| id.clone())
67    }
68    pub fn reset(&self, prompt: &str) {
69        let mut state = self.0.lock().expect("turn context lock poisoned");
70        let generation = state.generation.wrapping_add(1);
71        *state = TurnContextState {
72            generation,
73            decision_log: state.decision_log.clone(),
74            user_prompt_tail: tail(prompt, USER_PROMPT_BYTES),
75            background_commands: state.background_commands,
76            queued_commands: state.queued_commands,
77            background_inventory: std::mem::take(&mut state.background_inventory),
78            native_agent_ids: std::mem::take(&mut state.native_agent_ids),
79            session_id: std::mem::take(&mut state.session_id),
80            summary: std::mem::take(&mut state.summary),
81            ..Default::default()
82        };
83    }
84
85    /// Record durable transcript observations through the same path during live work and replay.
86    pub fn observe_relay(
87        &self,
88        observation: &mj_core::relay::RelayObservation,
89        prompt: Option<&str>,
90    ) {
91        use mj_core::relay::{RelayCommandOutcome, RelayObservation};
92        match observation {
93            RelayObservation::CommandStarted { .. }
94            | RelayObservation::CommandCompleted {
95                outcome: RelayCommandOutcome::Steered { .. },
96                ..
97            } => {
98                if let Some(prompt) = prompt {
99                    self.reset(prompt);
100                    self.0
101                        .lock()
102                        .expect("turn context lock poisoned")
103                        .summary
104                        .push_user(prompt);
105                }
106            }
107            RelayObservation::SessionUpdate { update } => self.observe(update),
108            RelayObservation::TerminalOutput {
109                terminal_id,
110                output,
111                truncated,
112                exit_code,
113                signal,
114            } => {
115                self.observe_terminal(&mj_core::transcript::TerminalOutputRecord {
116                    terminal_id: terminal_id.clone(),
117                    output: output.clone(),
118                    truncated: *truncated,
119                    exit_code: *exit_code,
120                    signal: signal.clone(),
121                });
122            }
123            RelayObservation::CommandCompleted {
124                outcome: RelayCommandOutcome::ContextCleared { .. },
125                ..
126            } => self.clear_history(),
127            _ => {}
128        }
129    }
130
131    /// Invalidate a pending verdict at a lifecycle boundary without discarding evidence.
132    pub fn invalidate(&self) {
133        let mut state = self.0.lock().expect("turn context lock poisoned");
134        state.generation = state.generation.wrapping_add(1);
135    }
136
137    pub fn generation(&self) -> u64 {
138        self.0
139            .lock()
140            .expect("turn context lock poisoned")
141            .generation
142    }
143
144    pub fn counts(&self) -> (usize, usize) {
145        let state = self.0.lock().expect("turn context lock poisoned");
146        (state.background_commands, state.queued_commands)
147    }
148
149    pub fn set_counts(&self, background_commands: usize, queued_commands: usize) {
150        let mut state = self.0.lock().expect("turn context lock poisoned");
151        if (state.background_commands, state.queued_commands)
152            != (background_commands, queued_commands)
153        {
154            state.generation = state.generation.wrapping_add(1);
155            state.background_commands = background_commands;
156            state.queued_commands = queued_commands;
157        }
158    }
159
160    pub fn set_session_id(&self, session_id: &str) {
161        self.0
162            .lock()
163            .expect("turn context lock poisoned")
164            .session_id = session_id.into();
165    }
166
167    pub fn session_id(&self) -> String {
168        self.0
169            .lock()
170            .expect("turn context lock poisoned")
171            .session_id
172            .clone()
173    }
174
175    /// Compare identities as well as counts; identical level reports are not activity.
176    pub fn set_background_inventory(
177        &self,
178        mut commands: Vec<mj_core::relay::BackgroundCommand>,
179        mut native_agent_ids: Vec<String>,
180    ) {
181        commands.sort_by(|a, b| a.id.cmp(&b.id));
182        native_agent_ids.sort();
183        let mut state = self.0.lock().expect("turn context lock poisoned");
184        if state.background_inventory != commands || state.native_agent_ids != native_agent_ids {
185            state.background_inventory = commands;
186            state.native_agent_ids = native_agent_ids;
187            state.generation = state.generation.wrapping_add(1);
188        }
189    }
190
191    pub fn observe(&self, update: &SessionUpdate) {
192        let mut state = self.0.lock().expect("turn context lock poisoned");
193        if matches!(
194            update,
195            SessionUpdate::AgentMessageChunk(_)
196                | SessionUpdate::AgentThoughtChunk(_)
197                | SessionUpdate::ToolCall(_)
198                | SessionUpdate::ToolCallUpdate(_)
199        ) {
200            state.generation = state.generation.wrapping_add(1);
201        }
202        state.summary.observe(update);
203        match update {
204            SessionUpdate::AgentMessageChunk(chunk) => {
205                let id = chunk.message_id.as_ref().map(ToString::to_string);
206                if id != state.message_id {
207                    complete_message(&mut state);
208                    state.message_id = id;
209                }
210                if let ContentBlock::Text(content) = &chunk.content {
211                    // Trim each chunk before appending, so a huge chunk never grows retained state.
212                    state
213                        .assistant_text_tail
214                        .push_str(&tail(&content.text, ASSISTANT_TEXT_BYTES));
215                    state.assistant_text_tail =
216                        tail(&state.assistant_text_tail, ASSISTANT_TEXT_BYTES);
217                }
218            }
219            SessionUpdate::ToolCall(_) => {
220                complete_message(&mut state);
221            }
222            SessionUpdate::ToolCallUpdate(_) | SessionUpdate::AgentThoughtChunk(_) => {
223                complete_message(&mut state)
224            }
225            _ => {}
226        }
227    }
228
229    pub fn observe_terminal(&self, record: &mj_core::transcript::TerminalOutputRecord) {
230        let mut state = self.0.lock().expect("turn context lock poisoned");
231        state.summary.observe_terminal(record);
232        state.generation = state.generation.wrapping_add(1);
233    }
234
235    pub fn mark_earlier_history_omitted(&self) {
236        self.0
237            .lock()
238            .expect("turn context lock poisoned")
239            .summary
240            .mark_earlier_history_omitted();
241    }
242
243    pub fn clear_history(&self) {
244        let mut state = self.0.lock().expect("turn context lock poisoned");
245        state.summary = Default::default();
246        state.user_prompt_tail.clear();
247        state.assistant_text_tail.clear();
248        state.last_completed_message.clear();
249        state.message_id = None;
250        state.generation = state.generation.wrapping_add(1);
251    }
252
253    pub fn evidence(
254        &self,
255        harness: HarnessKind,
256        phase: TurnPhase,
257        facts: &ActivityFacts,
258        now_ms: i64,
259    ) -> TurnEvidence {
260        let state = self.0.lock().expect("turn context lock poisoned");
261        let summary = state.summary.latest_user_messages();
262        let mut evidence = TurnEvidence {
263            harness,
264            phase,
265            silent_for_s: facts
266                .last_acp_activity_at_ms
267                .map_or(0, |last| now_ms.saturating_sub(last).max(0) as u64 / 1000),
268            tools_in_flight: facts
269                .tools_in_flight
270                .iter()
271                .take(IN_FLIGHT_TOOLS)
272                .map(|tool| ToolEvidence {
273                    title: tail(
274                        &state
275                            .summary
276                            .entries
277                            .iter()
278                            .find(|e| e.id == tool.tool_call_id)
279                            .map(|e| e.text.clone())
280                            .unwrap_or_else(|| "unknown tool".into()),
281                        TOOL_TITLE_BYTES,
282                    ),
283                    running_s: now_ms.saturating_sub(tool.started_at_ms).max(0) as u64 / 1000,
284                })
285                .collect(),
286            transcript_summary: summary.render(48 * 1024),
287            background_commands: facts.background_commands,
288            queued_commands: facts.queued_commands,
289            user_prompt_tail: state.user_prompt_tail.clone(),
290            assistant_text_tail: if state.assistant_text_tail.is_empty() {
291                state.last_completed_message.clone()
292            } else {
293                state.assistant_text_tail.clone()
294            },
295        };
296        let mut limit = 48 * 1024;
297        while serde_json::to_vec(&evidence)
298            .expect("serialize evidence")
299            .len()
300            > 60 * 1024
301        {
302            limit /= 2;
303            evidence.transcript_summary = summary.render(limit);
304        }
305        evidence
306    }
307}
308
309fn complete_message(state: &mut TurnContextState) {
310    if !state.assistant_text_tail.is_empty() {
311        state.last_completed_message = std::mem::take(&mut state.assistant_text_tail);
312    }
313    state.message_id = None;
314}
315
316fn tail(text: &str, maximum_bytes: usize) -> String {
317    let start = text.floor_char_boundary(text.len().saturating_sub(maximum_bytes));
318    let mut value = text[start..].to_owned();
319    mj_core::transcript::truncate_string_start(&mut value, maximum_bytes);
320    value
321}
322
323#[cfg(test)]
324mod tests {
325    use super::*;
326    use mj_core::activity::InFlightToolCall;
327    use serde_json::json;
328    fn message(id: &str, text: &str) -> SessionUpdate {
329        serde_json::from_value(json!({"sessionUpdate":"agent_message_chunk","messageId":id,"content":{"type":"text","text":text}})).unwrap()
330    }
331
332    fn tool(title: &str) -> SessionUpdate {
333        serde_json::from_value(json!({"sessionUpdate":"tool_call","toolCallId":title,"title":title,"status":"in_progress"})).unwrap()
334    }
335
336    #[test]
337    fn diagnostic_updates_do_not_change_evidence_or_generation_and_new_input_clears_attribution() {
338        let context = TurnContext::default();
339        context.reset("Implement the parser");
340        let generation = context.generation();
341        let before = serde_json::to_value(context.evidence(
342            HarnessKind::Codex,
343            TurnPhase::Running,
344            &Default::default(),
345            0,
346        ))
347        .unwrap();
348        let dir = tempfile::tempdir().unwrap();
349        let log = mj_core::jev::DecisionLog::open(dir.path().into()).unwrap();
350        context.set_decision_log(log.clone());
351        let attempt = log.start("s", "activity", "Who acts?", "Current request");
352        context.set_decision(generation, attempt.id());
353        attempt.finish("unchanged", "Kept runtime facts");
354        assert_eq!(context.generation(), generation);
355        assert_eq!(
356            serde_json::to_value(context.evidence(
357                HarnessKind::Codex,
358                TurnPhase::Running,
359                &Default::default(),
360                0
361            ))
362            .unwrap(),
363            before
364        );
365        assert!(context.decision().is_some());
366        context.reset("New work");
367        assert!(context.decision().is_none());
368        assert!(context.decision_log().is_some());
369    }
370
371    #[test]
372    fn verdict_uses_latest_delivered_user_without_tool_history_but_keeps_live_facts() {
373        use mj_core::relay::RelayObservation;
374        let context = TurnContext::default();
375        let start = |id: &str| RelayObservation::CommandStarted {
376            command_id: id.into(),
377            started_at_ms: 0,
378        };
379        context.observe_relay(&start("old"), Some("OLD REQUEST"));
380        context.observe(&message("old", "OLD ANSWER"));
381        context.observe(&tool("active-build"));
382        context.observe_relay(&start("new"), Some("CURRENT REQUEST"));
383        let noisy: SessionUpdate = serde_json::from_value(json!({
384            "sessionUpdate":"tool_call", "toolCallId":"noisy", "title":"noisy", "status":"completed",
385            "rawInput":{"command":"TOOL_BODY".repeat(20_000)}
386        })).unwrap();
387        context.observe(&noisy);
388        context.observe(&message("new", "CURRENT ANSWER"));
389        let facts = ActivityFacts {
390            background_commands: 2,
391            queued_commands: 3,
392            tools_in_flight: vec![InFlightToolCall {
393                tool_call_id: "active-build".into(),
394                title: Some("active-build".into()),
395                status: agent_client_protocol::schema::v1::ToolCallStatus::InProgress,
396                started_at_ms: 1_000,
397            }],
398            ..Default::default()
399        };
400        let evidence = context.evidence(HarnessKind::Codex, TurnPhase::Running, &facts, 5_000);
401        assert!(evidence.transcript_summary.contains("CURRENT REQUEST"));
402        assert!(evidence.transcript_summary.contains("CURRENT ANSWER"));
403        for excluded in [
404            "OLD REQUEST",
405            "OLD ANSWER",
406            "TOOL_BODY",
407            "<tool",
408            "bytes omitted",
409        ] {
410            assert!(
411                !evidence.transcript_summary.contains(excluded),
412                "{excluded}"
413            );
414        }
415        assert_eq!(evidence.user_prompt_tail, "CURRENT REQUEST");
416        assert_eq!(evidence.assistant_text_tail, "CURRENT ANSWER");
417        assert_eq!(evidence.background_commands, 2);
418        assert_eq!(evidence.queued_commands, 3);
419        assert_eq!(evidence.tools_in_flight.len(), 1);
420        assert!(evidence.tools_in_flight[0].title.contains("active-build"));
421        assert_eq!(evidence.tools_in_flight[0].running_s, 4);
422        let full = context.0.lock().unwrap().summary.render(256 * 1024);
423        assert!(full.contains("OLD REQUEST") && full.contains("TOOL_BODY"));
424    }
425
426    #[test]
427    fn evidence_caps_utf8_text_tools_and_durations() {
428        let context = TurnContext::default();
429        context.reset(&"é".repeat(20_000));
430        context.observe(&message("first", &"🦀".repeat(20_000)));
431        let first = context.evidence(
432            HarnessKind::Claude,
433            TurnPhase::Running,
434            &ActivityFacts::default(),
435            0,
436        );
437        assert_eq!(first.user_prompt_tail.len(), USER_PROMPT_BYTES);
438        assert_eq!(first.assistant_text_tail.len(), ASSISTANT_TEXT_BYTES);
439        for n in 0..20 {
440            context.observe(&tool(&format!("{n}{}", "é".repeat(300))));
441        }
442        let facts = ActivityFacts {
443            last_acp_activity_at_ms: Some(1_000),
444            background_commands: 3,
445            queued_commands: 2,
446            tools_in_flight: (0..100)
447                .map(|n| InFlightToolCall {
448                    tool_call_id: n.to_string(),
449                    title: Some("🦀".repeat(200)),
450                    status: agent_client_protocol::schema::v1::ToolCallStatus::InProgress,
451                    started_at_ms: 2_000,
452                })
453                .collect(),
454            ..Default::default()
455        };
456        let evidence = context.evidence(HarnessKind::Claude, TurnPhase::Running, &facts, 95_000);
457        assert_eq!(evidence.tools_in_flight.len(), IN_FLIGHT_TOOLS);
458        assert!(
459            evidence
460                .tools_in_flight
461                .iter()
462                .all(|tool| tool.title.len() <= TOOL_TITLE_BYTES && tool.running_s == 93)
463        );
464        assert!(evidence.transcript_summary.len() <= 48 * 1024);
465        assert_eq!(evidence.assistant_text_tail, first.assistant_text_tail);
466        assert_eq!(evidence.silent_for_s, 94);
467        let json = serde_json::to_value(evidence).unwrap();
468        assert_eq!(json["phase"], "running");
469        assert_eq!(json["harness"], "claude");
470        assert_eq!(json["background_commands"], 3);
471        assert_eq!(json["queued_commands"], 2);
472        assert!(serde_json::to_vec(&json).unwrap().len() < 64 * 1024);
473    }
474
475    #[test]
476    fn message_boundaries_and_prompt_reset_do_not_mix_answers() {
477        let context = TurnContext::default();
478        context.reset("First prompt");
479        context.observe(&message("one", "Old answer"));
480        context.observe(&message("two", "New "));
481        context.observe(&message("two", "answer?"));
482        let evidence = context.evidence(
483            HarnessKind::Claude,
484            TurnPhase::Replied,
485            &ActivityFacts::default(),
486            0,
487        );
488        assert_eq!(evidence.assistant_text_tail, "New answer?");
489        let generation = context.generation();
490        context.set_counts(1, 2);
491        assert_ne!(context.generation(), generation);
492        let generation = context.generation();
493        context.set_counts(1, 2);
494        assert_eq!(context.generation(), generation);
495        assert_eq!(context.counts(), (1, 2));
496        context.invalidate();
497        assert_ne!(context.generation(), generation);
498        assert_eq!(context.counts(), (1, 2));
499        let generation = context.generation();
500        context.reset("Second prompt");
501        assert_eq!(context.counts(), (1, 2));
502        assert_ne!(context.generation(), generation);
503        let evidence = context.evidence(
504            HarnessKind::Claude,
505            TurnPhase::Running,
506            &ActivityFacts::default(),
507            0,
508        );
509        assert_eq!(evidence.user_prompt_tail, "Second prompt");
510        assert!(evidence.assistant_text_tail.is_empty());
511        assert!(!evidence.transcript_summary.contains("<user>"));
512    }
513}