1use 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
25pub 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#[derive(Debug, Default, Clone)]
41pub struct TurnContext(Arc<Mutex<TurnContextState>>);
42
43impl TurnContext {
44 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 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 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 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 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 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}