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    parent_activity: Option<std::time::Instant>,
12    decision_log: Option<mj_core::jev::DecisionLog>,
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    /// Process-local parent clock; child traffic and inventory updates never mark it.
45    pub fn mark_parent_activity(&self) {
46        self.0
47            .lock()
48            .expect("turn context lock poisoned")
49            .parent_activity = Some(std::time::Instant::now());
50    }
51
52    pub fn parent_activity(&self) -> Option<std::time::Instant> {
53        self.0
54            .lock()
55            .expect("turn context lock poisoned")
56            .parent_activity
57    }
58
59    pub fn set_decision_log(&self, log: mj_core::jev::DecisionLog) {
60        self.0
61            .lock()
62            .expect("turn context lock poisoned")
63            .decision_log = Some(log);
64    }
65    pub fn decision_log(&self) -> Option<mj_core::jev::DecisionLog> {
66        self.0
67            .lock()
68            .expect("turn context lock poisoned")
69            .decision_log
70            .clone()
71    }
72    pub fn reset(&self, prompt: &str) {
73        let mut state = self.0.lock().expect("turn context lock poisoned");
74        let generation = state.generation.wrapping_add(1);
75        *state = TurnContextState {
76            generation,
77            parent_activity: Some(std::time::Instant::now()),
78            decision_log: state.decision_log.clone(),
79            user_prompt_tail: tail(prompt, USER_PROMPT_BYTES),
80            background_commands: state.background_commands,
81            queued_commands: state.queued_commands,
82            background_inventory: std::mem::take(&mut state.background_inventory),
83            native_agent_ids: std::mem::take(&mut state.native_agent_ids),
84            session_id: std::mem::take(&mut state.session_id),
85            summary: std::mem::take(&mut state.summary),
86            ..Default::default()
87        };
88    }
89
90    /// Record durable transcript observations through the same path during live work and replay.
91    pub fn observe_relay(
92        &self,
93        observation: &mj_core::relay::RelayObservation,
94        prompt: Option<&str>,
95    ) {
96        use mj_core::relay::{RelayCommandOutcome, RelayObservation};
97        match observation {
98            RelayObservation::CommandStarted { .. }
99            | RelayObservation::CommandCompleted {
100                outcome: RelayCommandOutcome::Steered { .. },
101                ..
102            } => {
103                if let Some(prompt) = prompt {
104                    self.reset(prompt);
105                    self.0
106                        .lock()
107                        .expect("turn context lock poisoned")
108                        .summary
109                        .push_user(prompt);
110                }
111            }
112            RelayObservation::SessionUpdate { update } => self.observe(update),
113            RelayObservation::TerminalOutput {
114                terminal_id,
115                output,
116                truncated,
117                exit_code,
118                signal,
119            } => {
120                self.observe_terminal(&mj_core::transcript::TerminalOutputRecord {
121                    terminal_id: terminal_id.clone(),
122                    output: output.clone(),
123                    truncated: *truncated,
124                    exit_code: *exit_code,
125                    signal: signal.clone(),
126                });
127            }
128            RelayObservation::CommandCompleted {
129                outcome: RelayCommandOutcome::ContextCleared { .. },
130                ..
131            } => self.clear_history(),
132            _ => {}
133        }
134    }
135
136    /// Invalidate a pending verdict at a lifecycle boundary without discarding evidence.
137    pub fn invalidate(&self) {
138        let mut state = self.0.lock().expect("turn context lock poisoned");
139        state.generation = state.generation.wrapping_add(1);
140    }
141
142    pub fn generation(&self) -> u64 {
143        self.0
144            .lock()
145            .expect("turn context lock poisoned")
146            .generation
147    }
148
149    pub fn counts(&self) -> (usize, usize) {
150        let state = self.0.lock().expect("turn context lock poisoned");
151        (state.background_commands, state.queued_commands)
152    }
153
154    pub fn set_counts(&self, background_commands: usize, queued_commands: usize) {
155        let mut state = self.0.lock().expect("turn context lock poisoned");
156        if (state.background_commands, state.queued_commands)
157            != (background_commands, queued_commands)
158        {
159            state.generation = state.generation.wrapping_add(1);
160            state.background_commands = background_commands;
161            state.queued_commands = queued_commands;
162        }
163    }
164
165    pub fn set_session_id(&self, session_id: &str) {
166        self.0
167            .lock()
168            .expect("turn context lock poisoned")
169            .session_id = session_id.into();
170    }
171
172    pub fn session_id(&self) -> String {
173        self.0
174            .lock()
175            .expect("turn context lock poisoned")
176            .session_id
177            .clone()
178    }
179
180    /// Compare identities as well as counts; identical level reports are not activity.
181    pub fn set_background_inventory(
182        &self,
183        mut commands: Vec<mj_core::relay::BackgroundCommand>,
184        mut native_agent_ids: Vec<String>,
185    ) {
186        commands.sort_by(|a, b| a.id.cmp(&b.id));
187        native_agent_ids.sort();
188        let mut state = self.0.lock().expect("turn context lock poisoned");
189        if state.background_inventory != commands || state.native_agent_ids != native_agent_ids {
190            state.background_inventory = commands;
191            state.native_agent_ids = native_agent_ids;
192            state.generation = state.generation.wrapping_add(1);
193        }
194    }
195
196    pub fn observe(&self, update: &SessionUpdate) {
197        let mut state = self.0.lock().expect("turn context lock poisoned");
198        state.parent_activity = Some(std::time::Instant::now());
199        if matches!(
200            update,
201            SessionUpdate::AgentMessageChunk(_)
202                | SessionUpdate::AgentThoughtChunk(_)
203                | SessionUpdate::ToolCall(_)
204                | SessionUpdate::ToolCallUpdate(_)
205        ) {
206            state.generation = state.generation.wrapping_add(1);
207        }
208        state.summary.observe(update);
209        match update {
210            SessionUpdate::AgentMessageChunk(chunk) => {
211                let id = chunk.message_id.as_ref().map(ToString::to_string);
212                if id != state.message_id {
213                    complete_message(&mut state);
214                    state.message_id = id;
215                }
216                if let ContentBlock::Text(content) = &chunk.content {
217                    // Trim each chunk before appending, so a huge chunk never grows retained state.
218                    state
219                        .assistant_text_tail
220                        .push_str(&tail(&content.text, ASSISTANT_TEXT_BYTES));
221                    state.assistant_text_tail =
222                        tail(&state.assistant_text_tail, ASSISTANT_TEXT_BYTES);
223                }
224            }
225            SessionUpdate::ToolCall(_) => {
226                complete_message(&mut state);
227            }
228            SessionUpdate::ToolCallUpdate(_) | SessionUpdate::AgentThoughtChunk(_) => {
229                complete_message(&mut state)
230            }
231            _ => {}
232        }
233    }
234
235    pub fn observe_terminal(&self, record: &mj_core::transcript::TerminalOutputRecord) {
236        let mut state = self.0.lock().expect("turn context lock poisoned");
237        state.summary.observe_terminal(record);
238        state.generation = state.generation.wrapping_add(1);
239    }
240
241    pub fn mark_earlier_history_omitted(&self) {
242        self.0
243            .lock()
244            .expect("turn context lock poisoned")
245            .summary
246            .mark_earlier_history_omitted();
247    }
248
249    pub fn clear_history(&self) {
250        let mut state = self.0.lock().expect("turn context lock poisoned");
251        state.summary = Default::default();
252        state.user_prompt_tail.clear();
253        state.assistant_text_tail.clear();
254        state.last_completed_message.clear();
255        state.message_id = None;
256        state.generation = state.generation.wrapping_add(1);
257    }
258
259    pub fn evidence(
260        &self,
261        harness: HarnessKind,
262        phase: TurnPhase,
263        facts: &ActivityFacts,
264        now_ms: i64,
265    ) -> TurnEvidence {
266        let state = self.0.lock().expect("turn context lock poisoned");
267        let summary = state.summary.latest_user_messages();
268        let last_assistant = state
269            .summary
270            .entries
271            .iter()
272            .rposition(|e| e.role == crate::summary::SummaryRole::Assistant);
273        let final_tool_calls = state
274            .summary
275            .entries
276            .iter()
277            .skip(last_assistant.map_or(0, |i| i + 1))
278            .filter(|e| e.role == crate::summary::SummaryRole::Tool)
279            .take(IN_FLIGHT_TOOLS)
280            .map(|e| mj_core::activity::verdict::ToolOutcome {
281                name: tail(
282                    e.tool
283                        .as_ref()
284                        .and_then(|t| t.get("name"))
285                        .and_then(|n| n.as_str())
286                        .unwrap_or(&e.text),
287                    TOOL_TITLE_BYTES,
288                ),
289                status: e
290                    .tool
291                    .as_ref()
292                    .and_then(|t| t.get("status"))
293                    .and_then(|s| s.as_str())
294                    .unwrap_or("unknown")
295                    .to_owned(),
296            })
297            .collect();
298        let mut evidence = TurnEvidence {
299            authorization: None,
300            final_tool_calls,
301            background: Vec::new(),
302            harness,
303            phase,
304            silent_for_s: facts
305                .last_acp_activity_at_ms
306                .map_or(0, |last| now_ms.saturating_sub(last).max(0) as u64 / 1000),
307            tools_in_flight: facts
308                .tools_in_flight
309                .iter()
310                .take(IN_FLIGHT_TOOLS)
311                .map(|tool| ToolEvidence {
312                    title: tail(
313                        &state
314                            .summary
315                            .entries
316                            .iter()
317                            .find(|e| e.id == tool.tool_call_id)
318                            .map(|e| e.text.clone())
319                            .unwrap_or_else(|| "unknown tool".into()),
320                        TOOL_TITLE_BYTES,
321                    ),
322                    running_s: now_ms.saturating_sub(tool.started_at_ms).max(0) as u64 / 1000,
323                })
324                .collect(),
325            transcript_summary: summary.render(48 * 1024),
326            background_commands: facts.background_commands,
327            queued_commands: facts.queued_commands,
328            user_prompt_tail: state.user_prompt_tail.clone(),
329            assistant_text_tail: if state.assistant_text_tail.is_empty() {
330                state.last_completed_message.clone()
331            } else {
332                state.assistant_text_tail.clone()
333            },
334            completion: None,
335        };
336        let mut limit = 48 * 1024;
337        while serde_json::to_vec(&evidence)
338            .expect("serialize evidence")
339            .len()
340            > 60 * 1024
341        {
342            limit /= 2;
343            evidence.transcript_summary = summary.render(limit);
344        }
345        evidence
346    }
347}
348
349fn complete_message(state: &mut TurnContextState) {
350    if !state.assistant_text_tail.is_empty() {
351        state.last_completed_message = std::mem::take(&mut state.assistant_text_tail);
352    }
353    state.message_id = None;
354}
355
356fn tail(text: &str, maximum_bytes: usize) -> String {
357    let start = text.floor_char_boundary(text.len().saturating_sub(maximum_bytes));
358    let mut value = text[start..].to_owned();
359    mj_core::transcript::truncate_string_start(&mut value, maximum_bytes);
360    value
361}
362
363#[cfg(test)]
364mod tests {
365    use super::*;
366    use mj_core::activity::InFlightToolCall;
367    use serde_json::json;
368    fn message(id: &str, text: &str) -> SessionUpdate {
369        serde_json::from_value(json!({"sessionUpdate":"agent_message_chunk","messageId":id,"content":{"type":"text","text":text}})).unwrap()
370    }
371
372    fn tool(title: &str) -> SessionUpdate {
373        serde_json::from_value(json!({"sessionUpdate":"tool_call","toolCallId":title,"title":title,"status":"in_progress"})).unwrap()
374    }
375
376    #[test]
377    fn diagnostic_logging_does_not_change_evidence_or_generation_and_survives_reset() {
378        let context = TurnContext::default();
379        context.reset("Implement the parser");
380        let generation = context.generation();
381        let before = serde_json::to_value(context.evidence(
382            HarnessKind::Codex,
383            TurnPhase::Running,
384            &Default::default(),
385            0,
386        ))
387        .unwrap();
388        let dir = tempfile::tempdir().unwrap();
389        let log = mj_core::jev::DecisionLog::open(dir.path().into()).unwrap();
390        context.set_decision_log(log.clone());
391        let attempt = log.start("s", "activity", "Who acts?", "Current request");
392        attempt.finish("unchanged", "Kept runtime facts");
393        assert_eq!(context.generation(), generation);
394        assert_eq!(
395            serde_json::to_value(context.evidence(
396                HarnessKind::Codex,
397                TurnPhase::Running,
398                &Default::default(),
399                0
400            ))
401            .unwrap(),
402            before
403        );
404        context.reset("New work");
405        assert!(context.decision_log().is_some());
406    }
407
408    #[test]
409    fn verdict_uses_latest_delivered_user_without_tool_history_but_keeps_live_facts() {
410        use mj_core::relay::RelayObservation;
411        let context = TurnContext::default();
412        let start = |id: &str| RelayObservation::CommandStarted {
413            command_id: id.into(),
414            started_at_ms: 0,
415        };
416        context.observe_relay(&start("old"), Some("OLD REQUEST"));
417        context.observe(&message("old", "OLD ANSWER"));
418        context.observe(&tool("active-build"));
419        context.observe_relay(&start("new"), Some("CURRENT REQUEST"));
420        let noisy: SessionUpdate = serde_json::from_value(json!({
421            "sessionUpdate":"tool_call", "toolCallId":"noisy", "title":"noisy", "status":"completed",
422            "rawInput":{"command":"TOOL_BODY".repeat(20_000)}
423        })).unwrap();
424        context.observe(&noisy);
425        context.observe(&message("new", "CURRENT ANSWER"));
426        let facts = ActivityFacts {
427            background_commands: 2,
428            queued_commands: 3,
429            tools_in_flight: vec![InFlightToolCall {
430                tool_call_id: "active-build".into(),
431                title: Some("active-build".into()),
432                status: agent_client_protocol::schema::v1::ToolCallStatus::InProgress,
433                started_at_ms: 1_000,
434            }],
435            ..Default::default()
436        };
437        let evidence = context.evidence(HarnessKind::Codex, TurnPhase::Running, &facts, 5_000);
438        assert!(evidence.transcript_summary.contains("CURRENT REQUEST"));
439        assert!(evidence.transcript_summary.contains("CURRENT ANSWER"));
440        for excluded in [
441            "OLD REQUEST",
442            "OLD ANSWER",
443            "TOOL_BODY",
444            "<tool",
445            "bytes omitted",
446        ] {
447            assert!(
448                !evidence.transcript_summary.contains(excluded),
449                "{excluded}"
450            );
451        }
452        assert_eq!(evidence.user_prompt_tail, "CURRENT REQUEST");
453        assert_eq!(evidence.assistant_text_tail, "CURRENT ANSWER");
454        assert_eq!(evidence.background_commands, 2);
455        assert_eq!(evidence.queued_commands, 3);
456        assert_eq!(evidence.tools_in_flight.len(), 1);
457        assert!(evidence.tools_in_flight[0].title.contains("active-build"));
458        assert_eq!(evidence.tools_in_flight[0].running_s, 4);
459        let full = context.0.lock().unwrap().summary.render(256 * 1024);
460        assert!(full.contains("OLD REQUEST") && full.contains("TOOL_BODY"));
461    }
462
463    #[test]
464    fn tool_calls_after_the_last_assistant_text_are_named_in_the_evidence() {
465        use mj_core::relay::RelayObservation;
466        let context = TurnContext::default();
467        context.observe_relay(
468            &RelayObservation::CommandStarted {
469                command_id: "audit".into(),
470                started_at_ms: 0,
471            },
472            Some("Audit the cache and hand back."),
473        );
474        context.observe(&message("m1", "I'm checking the update paths now."));
475        let handback: SessionUpdate = serde_json::from_value(json!({
476            "sessionUpdate":"tool_call", "toolCallId":"hb", "title":"mcp__mj-agents__handback",
477            "status":"completed", "rawInput":{"report":"done"}
478        }))
479        .unwrap();
480        context.observe(&handback);
481        let facts = ActivityFacts::default();
482        let evidence = context.evidence(HarnessKind::Codex, TurnPhase::Replied, &facts, 5_000);
483        assert_eq!(evidence.final_tool_calls.len(), 1);
484        assert!(evidence.final_tool_calls[0].name.contains("handback"));
485        assert_eq!(evidence.final_tool_calls[0].status, "completed");
486        // Text after the call closes the list again.
487        context.observe(&message("m2", "Handed back."));
488        let evidence = context.evidence(HarnessKind::Codex, TurnPhase::Replied, &facts, 5_000);
489        assert!(evidence.final_tool_calls.is_empty());
490    }
491
492    #[test]
493    fn evidence_caps_utf8_text_tools_and_durations() {
494        let context = TurnContext::default();
495        context.reset(&"é".repeat(20_000));
496        context.observe(&message("first", &"🦀".repeat(20_000)));
497        let first = context.evidence(
498            HarnessKind::Claude,
499            TurnPhase::Running,
500            &ActivityFacts::default(),
501            0,
502        );
503        assert_eq!(first.user_prompt_tail.len(), USER_PROMPT_BYTES);
504        assert_eq!(first.assistant_text_tail.len(), ASSISTANT_TEXT_BYTES);
505        for n in 0..20 {
506            context.observe(&tool(&format!("{n}{}", "é".repeat(300))));
507        }
508        let facts = ActivityFacts {
509            last_acp_activity_at_ms: Some(1_000),
510            background_commands: 3,
511            queued_commands: 2,
512            tools_in_flight: (0..100)
513                .map(|n| InFlightToolCall {
514                    tool_call_id: n.to_string(),
515                    title: Some("🦀".repeat(200)),
516                    status: agent_client_protocol::schema::v1::ToolCallStatus::InProgress,
517                    started_at_ms: 2_000,
518                })
519                .collect(),
520            ..Default::default()
521        };
522        let evidence = context.evidence(HarnessKind::Claude, TurnPhase::Running, &facts, 95_000);
523        assert_eq!(evidence.tools_in_flight.len(), IN_FLIGHT_TOOLS);
524        assert!(
525            evidence
526                .tools_in_flight
527                .iter()
528                .all(|tool| tool.title.len() <= TOOL_TITLE_BYTES && tool.running_s == 93)
529        );
530        assert!(evidence.transcript_summary.len() <= 48 * 1024);
531        assert_eq!(evidence.assistant_text_tail, first.assistant_text_tail);
532        assert_eq!(evidence.silent_for_s, 94);
533        let json = serde_json::to_value(evidence).unwrap();
534        assert_eq!(json["phase"], "running");
535        assert_eq!(json["harness"], "claude");
536        assert_eq!(json["background_commands"], 3);
537        assert_eq!(json["queued_commands"], 2);
538        assert!(serde_json::to_vec(&json).unwrap().len() < 64 * 1024);
539    }
540
541    #[test]
542    fn message_boundaries_and_prompt_reset_do_not_mix_answers() {
543        let context = TurnContext::default();
544        context.reset("First prompt");
545        context.observe(&message("one", "Old answer"));
546        context.observe(&message("two", "New "));
547        context.observe(&message("two", "answer?"));
548        let evidence = context.evidence(
549            HarnessKind::Claude,
550            TurnPhase::Replied,
551            &ActivityFacts::default(),
552            0,
553        );
554        assert_eq!(evidence.assistant_text_tail, "New answer?");
555        let generation = context.generation();
556        context.set_counts(1, 2);
557        assert_ne!(context.generation(), generation);
558        let generation = context.generation();
559        context.set_counts(1, 2);
560        assert_eq!(context.generation(), generation);
561        assert_eq!(context.counts(), (1, 2));
562        context.invalidate();
563        assert_ne!(context.generation(), generation);
564        assert_eq!(context.counts(), (1, 2));
565        let generation = context.generation();
566        context.reset("Second prompt");
567        assert_eq!(context.counts(), (1, 2));
568        assert_ne!(context.generation(), generation);
569        let evidence = context.evidence(
570            HarnessKind::Claude,
571            TurnPhase::Running,
572            &ActivityFacts::default(),
573            0,
574        );
575        assert_eq!(evidence.user_prompt_tail, "Second prompt");
576        assert!(evidence.assistant_text_tail.is_empty());
577        assert!(!evidence.transcript_summary.contains("<user>"));
578    }
579}