返回 CodeWhale
events.rs
根目录 / crates / tui / src / conformance / events.rs
1 //! Golden-events family: one full turn through the one turn loop.
2 //!
3 //! Each case drives `Engine::run` (and so `Engine::run_turn`) with a scripted
4 //! provider — a queue of normalized stream events per model request, written
5 //! in the [`super::stream_json`] format — and a small host driver that
6 //! answers approvals or cancels on cue. The golden is the protocol projection
7 //! (`protocol_parity::event_to_protocol`, the `EventMsg` wire shape every
8 //! frontend consumes), normalized, plus one trailing `harness_summary` line
9 //! with the provider request count and the workspace after the turn.
10 //!
11 //! Normalization (also stated in the fixture README): the per-session
12 //! `thread_id`/`session_id` routing envelope is dropped; liveness heartbeats
13 //! are dropped; UUIDs, timestamps, temp paths and durations are masked; the
14 //! system prompt and tool catalog bodies are replaced by a marker because the
15 //! prompt family owns those bytes; an uninterrupted run of completion events
16 //! is put in a canonical order because parallel completions race and their
17 //! pairs interleave: `operation_activity_completed` by span id, then
18 //! `tool_call_complete` by tool call id. The activity events of one announced
19 //! parallel batch (the events after the `Executing N ... parallel chunk(s)`
20 //! status) are ordered the same way, starts first, because a fast tool can
21 //! finish before its sibling has started. Only a causally valid batch is
22 //! reordered; no error or other event is crossed.
23
24 use std::collections::{BTreeMap, VecDeque};
25 use std::path::Path;
26 use std::sync::Mutex;
27 use std::sync::atomic::{AtomicUsize, Ordering};
28 use std::time::Duration;
29
30 use anyhow::{Result, anyhow};
31 use serde_json::{Value, json};
32
33 use codewhale_config::AppMode;
34 use codewhale_execpolicy::ApprovalMode;
35 use codewhale_models::{MessageRequest, MessageResponse, StreamEvent};
36 use codewhale_protocol::ids::{SessionId, ThreadId};
37
38 use super::golden::{self, Failures, Masker, Sandbox};
39 use super::stream_json;
40 use crate::compaction::CompactionConfig;
41 use crate::config::Config;
42 use crate::core::engine::{Engine, EngineConfig};
43 use crate::core::events::Event;
44 use crate::core::ops::{Op, TurnSpec, UserInputProvenance};
45 use crate::core::protocol_parity::{ProtocolIds, event_to_protocol};
46 use crate::llm_client::{LlmClient, StreamEventBox};
47
48 const FAMILY: &str = "events";
49 pub(super) const PREFIX_OWNED_BY_PROMPT_FAMILY: &str = "<pinned by the prompt family>";
50
51 /// Keys whose values are clocks or durations.
52 const VOLATILE_KEYS: &[&str] = &[
53 "created_at",
54 "duration_ms",
55 "first_token_ms",
56 "request_ms",
57 "elapsed_ms",
58 "pinned_combined_hash",
59 ];
60
61 fn event_deadline() -> Duration {
62 if cfg!(windows) {
63 Duration::from_secs(60)
64 } else {
65 Duration::from_secs(20)
66 }
67 }
68
69 enum Step {
70 Stream {
71 events: Vec<StreamEvent>,
72 hang: bool,
73 },
74 Refuse(String),
75 }
76
77 /// The scripted fake provider: one queued step per model request. Unlike the
78 /// general-purpose `MockLlmClient` it appends nothing — a script without
79 /// `message_stop` is a truncated stream, exactly as written.
80 pub(super) struct ScriptedProvider {
81 model: String,
82 steps: Mutex<VecDeque<Step>>,
83 stream_requests: AtomicUsize,
84 other_requests: AtomicUsize,
85 captured: Mutex<Vec<MessageRequest>>,
86 }
87
88 impl ScriptedProvider {
89 pub(super) fn from_script(model: &str, script: &Value) -> Self {
90 assert!(
91 script.as_array().is_some_and(|steps| !steps.is_empty()),
92 "provider_script must contain at least one response"
93 );
94 let steps = script
95 .as_array()
96 .expect("provider_script is an array")
97 .iter()
98 .map(|step| {
99 if let Some(message) = step.get("request_error").and_then(Value::as_str) {
100 return Step::Refuse(message.to_string());
101 }
102 let events = step["events"]
103 .as_array()
104 .expect("provider_script step has events")
105 .iter()
106 .map(|event| {
107 stream_json::from_json(event)
108 .unwrap_or_else(|error| panic!("script event {event}: {error}"))
109 })
110 .collect();
111 Step::Stream {
112 events,
113 hang: step.get("then").and_then(Value::as_str) == Some("hang"),
114 }
115 })
116 .collect();
117 Self {
118 model: model.to_string(),
119 steps: Mutex::new(steps),
120 stream_requests: AtomicUsize::new(0),
121 other_requests: AtomicUsize::new(0),
122 captured: Mutex::new(Vec::new()),
123 }
124 }
125
126 pub(super) fn stream_requests(&self) -> usize {
127 self.stream_requests.load(Ordering::SeqCst)
128 }
129
130 pub(super) fn other_requests(&self) -> usize {
131 self.other_requests.load(Ordering::SeqCst)
132 }
133
134 pub(super) fn captured(&self) -> Vec<MessageRequest> {
135 self.captured.lock().expect("captured requests").clone()
136 }
137 }
138
139 impl LlmClient for ScriptedProvider {
140 fn provider_name(&self) -> &'static str {
141 "conformance"
142 }
143
144 fn model(&self) -> &str {
145 &self.model
146 }
147
148 async fn create_message(&self, _request: MessageRequest) -> Result<MessageResponse> {
149 self.other_requests.fetch_add(1, Ordering::SeqCst);
150 Err(anyhow!(
151 "conformance provider serves streaming turns only; a non-streaming request is not scripted"
152 ))
153 }
154
155 async fn create_message_stream(&self, request: MessageRequest) -> Result<StreamEventBox> {
156 let call = self.stream_requests.fetch_add(1, Ordering::SeqCst) + 1;
157 self.captured
158 .lock()
159 .expect("captured requests")
160 .push(request);
161 let step = self.steps.lock().expect("script").pop_front();
162 match step {
163 None => Err(anyhow!(
164 "conformance script exhausted: model request #{call} has no scripted response"
165 )),
166 Some(Step::Refuse(message)) => Err(anyhow!(message)),
167 Some(Step::Stream { events, hang }) => {
168 let stream: StreamEventBox = Box::pin(async_stream::stream! {
169 for event in events {
170 yield Ok::<StreamEvent, anyhow::Error>(event);
171 }
172 if hang {
173 std::future::pending::<()>().await;
174 }
175 });
176 Ok(stream)
177 }
178 }
179 }
180 }
181
182 fn app_mode(case: &Value) -> AppMode {
183 match case["mode"].as_str().unwrap_or("agent") {
184 "agent" => AppMode::Agent,
185 "plan" => AppMode::Plan,
186 "operate" => AppMode::Operate,
187 other => panic!("unknown mode `{other}`"),
188 }
189 }
190
191 fn approval_mode(case: &Value) -> ApprovalMode {
192 match case["approval_mode"].as_str().unwrap_or("suggest") {
193 "suggest" => ApprovalMode::Suggest,
194 "auto" => ApprovalMode::Auto,
195 "bypass" => ApprovalMode::Bypass,
196 "never" => ApprovalMode::Never,
197 other => panic!("unknown approval_mode `{other}`"),
198 }
199 }
200
201 pub(super) fn send_message_op(case: &Value, config: &Config) -> Op {
202 let model = case["model"].as_str().expect("case.model");
203 let route = crate::route_runtime::resolve_runtime_route(
204 config,
205 config.active_provider_identity().unwrap().provider,
206 Some(model),
207 )
208 .expect("resolve conformance route");
209 Op::SendMessage(TurnSpec {
210 max_output_tokens: None,
211 submission_id: None,
212 content: case["user_message"]
213 .as_str()
214 .expect("case.user_message")
215 .to_string(),
216 images: Vec::new(),
217 mode: app_mode(case),
218 route: Box::new(route),
219 compaction: Box::new(CompactionConfig::default()),
220 initial_routed_usage: Box::default(),
221 goal_objective: None,
222 goal_token_budget: None,
223 goal_status: crate::tools::goal::GoalStatus::Active,
224 reasoning_effort: None,
225 reasoning_effort_auto: false,
226 auto_model: false,
227 allow_shell: case["allow_shell"].as_bool().unwrap_or(true),
228 trust_mode: false,
229 auto_approve: case["auto_approve"].as_bool().unwrap_or(false),
230 approval_mode: approval_mode(case),
231 translation_enabled: false,
232 allowed_tools: None,
233 dynamic_tools: Vec::new(),
234 hook_executor: None,
235 verbosity: None,
236 provenance: UserInputProvenance::ExternalUser,
237 })
238 }
239
240 /// Relative path → sha256 prefix for every file left in the workspace.
241 fn workspace_listing(workspace: &Path) -> BTreeMap<String, String> {
242 fn walk(root: &Path, dir: &Path, out: &mut BTreeMap<String, String>) {
243 let Ok(entries) = std::fs::read_dir(dir) else {
244 return;
245 };
246 for entry in entries.flatten() {
247 let path = entry.path();
248 if path.is_dir() {
249 walk(root, &path, out);
250 } else if let Ok(bytes) = std::fs::read(&path) {
251 let relative = path
252 .strip_prefix(root)
253 .expect("inside workspace")
254 .to_string_lossy()
255 .replace('\\', "/");
256 let digest = crate::hashing::sha256_hex(&bytes);
257 out.insert(relative, digest[..16].to_string());
258 }
259 }
260 }
261 let mut out = BTreeMap::new();
262 walk(workspace, workspace, &mut out);
263 out
264 }
265
266 pub(super) struct TurnRecord {
267 pub(super) events: Vec<Event>,
268 pub(super) failure: Option<String>,
269 }
270
271 /// Drive one scripted turn inside `sandbox` on a fresh current-thread
272 /// runtime, dropping the runtime (and its blocking post-turn work) before
273 /// returning so nothing outlives the sandboxed environment.
274 pub(super) fn run_scripted_turn(
275 sandbox: &Sandbox,
276 case: &Value,
277 script: &Value,
278 ) -> (TurnRecord, std::sync::Arc<ScriptedProvider>) {
279 run_scripted_turn_with_deadline(sandbox, case, script, event_deadline())
280 }
281
282 fn run_scripted_turn_with_deadline(
283 sandbox: &Sandbox,
284 case: &Value,
285 script: &Value,
286 deadline: Duration,
287 ) -> (TurnRecord, std::sync::Arc<ScriptedProvider>) {
288 let model = case["model"].as_str().expect("case.model");
289 let provider = std::sync::Arc::new(ScriptedProvider::from_script(model, script));
290 let driver = case["driver"].as_array().cloned().unwrap_or_default();
291 let runtime = tokio::runtime::Builder::new_current_thread()
292 .enable_all()
293 .build()
294 .expect("runtime");
295 let (os, shell) = recorded_environment(case);
296 // Held across the whole turn: the engine builds and refreshes its system
297 // prompt on this thread's current-thread runtime.
298 let _environment = crate::prompts::pin_recorded_environment(&os, &shell);
299 let _host_tools = crate::dependencies::pin_recorded_host_tools(recorded_host_tools(case));
300 let record = runtime.block_on(async {
301 let config = Config::default();
302 let engine_config = EngineConfig {
303 workspace: sandbox.workspace.clone(),
304 snapshots_enabled: false,
305 subagents_enabled: false,
306 session_id: Some("conformance-session".to_string()),
307 ..EngineConfig::default()
308 };
309 let client: crate::core::model_client::SharedModelClient = provider.clone();
310 let (mut engine, handle) = Engine::new_with_model_client(engine_config, &config, client);
311 let (enforcement, no_new_privs_active) = recorded_platform(case);
312 engine.pin_recorded_platform_posture(enforcement, no_new_privs_active);
313 let op = send_message_op(case, &config);
314 drive_turn(engine, handle, op, &driver, deadline).await
315 });
316 drop(runtime);
317 (record, provider)
318 }
319
320 /// The execution boundary a case's golden was recorded under.
321 ///
322 /// The Engine names its sandbox posture to the model in `<turn_meta>`, and a
323 /// live probe makes that line a fact about the runner: macOS applies a local
324 /// OS sandbox, a Linux runner without bubblewrap is policy-only, and Linux
325 /// startup may relax no-new-privs. Every scripted case therefore declares the
326 /// recorded facts and the harness replays exactly those, so a golden states
327 /// what the Engine does on that platform rather than which machine ran it.
328 /// A case without them fails loud instead of silently probing the host. The
329 /// label for each platform is owned and tested by `sandbox::policy`.
330 fn recorded_platform(case: &Value) -> (crate::sandbox::policy::SandboxEnforcement, Option<bool>) {
331 use crate::sandbox::policy::SandboxEnforcement;
332 let platform = case
333 .get("recorded_platform")
334 .expect("scripted conformance case must declare `recorded_platform`");
335 let enforcement = match platform.get("sandbox_enforcement").and_then(Value::as_str) {
336 Some("local_os") => SandboxEnforcement::LocalOs,
337 Some("unavailable") => SandboxEnforcement::Unavailable,
338 Some("external_backend") => SandboxEnforcement::ExternalBackend,
339 other => panic!("unknown recorded_platform.sandbox_enforcement: {other:?}"),
340 };
341 let no_new_privs_active = match platform.get("no_new_privs_active") {
342 Some(Value::Null) => None,
343 Some(Value::Bool(active)) => Some(*active),
344 other => panic!("recorded_platform.no_new_privs_active must be null or a bool: {other:?}"),
345 };
346 (enforcement, no_new_privs_active)
347 }
348
349 /// The OS and shell the golden's `## Environment` block was recorded with.
350 /// They enter the frozen prompt prefix, and so its hash in
351 /// `prefix_cache_change`; a case without them fails loud.
352 fn recorded_environment(case: &Value) -> (String, String) {
353 let platform = case
354 .get("recorded_platform")
355 .expect("scripted conformance case must declare `recorded_platform`");
356 let field = |name: &str| {
357 platform
358 .get(name)
359 .and_then(Value::as_str)
360 .filter(|value| !value.trim().is_empty())
361 .unwrap_or_else(|| panic!("recorded_platform.{name} must be a non-empty string"))
362 .to_string()
363 };
364 (field("os"), field("shell"))
365 }
366
367 /// The optional host-backed tools the recording machine had (see
368 /// `dependencies::host_tool_available`). They change the registry the golden
369 /// snapshot counts; a case without the list fails loud.
370 fn recorded_host_tools(case: &Value) -> Vec<String> {
371 case.get("recorded_platform")
372 .and_then(|platform| platform.get("host_tools"))
373 .and_then(Value::as_array)
374 .expect("recorded_platform.host_tools must list the recording host's optional tools")
375 .iter()
376 .map(|tool| {
377 tool.as_str()
378 .expect("recorded_platform.host_tools entries are tool names")
379 .to_string()
380 })
381 .collect()
382 }
383
384 /// Run one turn and return every engine event up to and including
385 /// `TurnComplete`, then everything the engine emits while it shuts down.
386 async fn drive_turn(
387 engine: Engine,
388 handle: crate::core::engine::EngineHandle,
389 op: Op,
390 driver: &[Value],
391 deadline: Duration,
392 ) -> TurnRecord {
393 let mut run = tokio::spawn(engine.run());
394 handle.send(op).await.expect("send conformance turn");
395 let rx = handle.rx_event.clone();
396 let ids = protocol_ids();
397 let mut events = Vec::new();
398 let mut fired = vec![false; driver.len()];
399 let mut failure = None;
400 let mut completed = false;
401 let expires = tokio::time::Instant::now() + deadline;
402 loop {
403 if tokio::time::Instant::now() >= expires {
404 failure = Some(format!(
405 "harness timeout: engine turn did not complete within {deadline:?}"
406 ));
407 break;
408 }
409 let next = golden::complete_within(
410 "engine turn",
411 expires.saturating_duration_since(tokio::time::Instant::now()),
412 async { rx.write().await.recv().await },
413 )
414 .await;
415 let event = match next {
416 Ok(Some(event)) => event,
417 Ok(None) => break,
418 Err(error) => {
419 failure = Some(error);
420 break;
421 }
422 };
423 let projected = serde_json::to_value(event_to_protocol(&event, &ids)).expect("EventMsg");
424 for (index, rule) in driver.iter().enumerate() {
425 if fired[index] || !rule_matches(&rule["when"], &projected) {
426 continue;
427 }
428 fired[index] = true;
429 match rule["action"].as_str().expect("driver action") {
430 "cancel" => handle.cancel(),
431 action @ ("approve" | "deny") => {
432 let Event::ApprovalRequired { id, .. } = &event else {
433 panic!("driver `{action}` must match an approval_required event");
434 };
435 let answered = if action == "approve" {
436 handle.approve_tool_call(id.clone()).await
437 } else {
438 handle.deny_tool_call(id.clone()).await
439 };
440 answered.expect("answer approval");
441 }
442 other => panic!("unknown driver action `{other}`"),
443 }
444 }
445 let terminal = matches!(event, Event::TurnComplete { .. });
446 events.push(event);
447 if terminal {
448 completed = true;
449 break;
450 }
451 }
452 if !completed && failure.is_none() {
453 failure = Some("engine closed before TurnComplete".to_string());
454 }
455 handle.cancel();
456 match golden::complete_within(
457 "engine shutdown request",
458 deadline,
459 handle.send(Op::Shutdown),
460 )
461 .await
462 {
463 Ok(Ok(())) => {}
464 Ok(Err(error)) => {
465 failure.get_or_insert_with(|| format!("engine shutdown request failed: {error}"));
466 }
467 Err(error) => {
468 failure.get_or_insert(error);
469 }
470 }
471 match golden::complete_within("engine shutdown", deadline, &mut run).await {
472 Ok(Ok(_)) => {}
473 Ok(Err(error)) => {
474 failure.get_or_insert_with(|| format!("engine task failed: {error}"));
475 }
476 Err(error) => {
477 failure.get_or_insert(error);
478 run.abort();
479 let _ = run.await;
480 }
481 }
482 let mut rx = rx.write().await;
483 while let Ok(event) = rx.try_recv() {
484 events.push(event);
485 }
486 TurnRecord { events, failure }
487 }
488
489 fn rule_matches(when: &Value, projected: &Value) -> bool {
490 if when["event"].as_str() != projected["event"].as_str() {
491 return false;
492 }
493 if let Some(needle) = when.get("delta_contains").and_then(Value::as_str) {
494 return projected["delta"]
495 .as_str()
496 .is_some_and(|delta| delta.contains(needle));
497 }
498 true
499 }
500
501 fn protocol_ids() -> ProtocolIds {
502 ProtocolIds {
503 thread_id: ThreadId::from_string("conformance-thread"),
504 session_id: SessionId::from_string("conformance-session"),
505 }
506 }
507
508 /// Project, drop, mask and order engine events into golden lines.
509 pub(super) fn normalize_events(events: &[Event], masker: &mut Masker) -> Vec<Value> {
510 let ids = protocol_ids();
511 let mut lines: Vec<Value> = Vec::new();
512 for event in events {
513 let mut value = serde_json::to_value(event_to_protocol(event, &ids)).expect("EventMsg");
514 let kind = value["event"].as_str().unwrap_or_default().to_string();
515 if kind == "tool_call_heartbeat" {
516 continue;
517 }
518 if let Some(map) = value.as_object_mut() {
519 map.remove("thread_id");
520 map.remove("session_id");
521 for owned in ["tool_catalog", "system_prompt"] {
522 if map.get(owned).is_some_and(|value| !value.is_null()) {
523 map.insert(owned.to_string(), json!(PREFIX_OWNED_BY_PROMPT_FAMILY));
524 }
525 }
526 }
527 masker.value(&mut value);
528 lines.push(golden::canonical(&value));
529 }
530 order_parallel_batches(&mut lines);
531 order_parallel_completions(&mut lines);
532 lines
533 }
534
535 /// The activity and tool events a parallel batch emits.
536 fn batch_rank(line: &Value) -> Option<(u8, &'static str)> {
537 match line["event"].as_str() {
538 Some("operation_activity_started") => Some((0, "span_id")),
539 Some("operation_activity_completed") => Some((1, "span_id")),
540 Some("tool_call_complete") => Some((2, "tool_call_id")),
541 _ => None,
542 }
543 }
544
545 /// Whether `run` is a legal interleaving of parallel tools: every span starts
546 /// once and completes at most once after its start, and a tool's completion
547 /// follows its own activity completion. A run that fails this is left in
548 /// its recorded order so the golden reports the violation.
549 fn batch_is_causal(run: &[Value]) -> bool {
550 let mut started = std::collections::BTreeSet::new();
551 let mut completed = std::collections::BTreeSet::new();
552 let mut finished_calls = std::collections::BTreeSet::new();
553 for line in run {
554 let span = line["span_id"].as_str().unwrap_or_default();
555 let call = span.split('#').next().unwrap_or_default();
556 match line["event"].as_str() {
557 Some("operation_activity_started") => {
558 if !started.insert(span) {
559 return false;
560 }
561 }
562 Some("operation_activity_completed") => {
563 if !started.contains(span) || !completed.insert(span) {
564 return false;
565 }
566 }
567 Some("tool_call_complete") => {
568 finished_calls.insert(line["tool_call_id"].as_str().unwrap_or_default());
569 }
570 _ => {}
571 }
572 // A tool completion may not precede its own activity completion.
573 if line["event"].as_str() == Some("operation_activity_completed")
574 && finished_calls.contains(call)
575 {
576 return false;
577 }
578 }
579 true
580 }
581
582 /// Parallel tools in one announced batch run concurrently, so one tool can
583 /// start, finish and report before its sibling has started (a blocking-pool
584 /// read that completes before its first poll). The run of activity and tool
585 /// events after the batch's `Executing N ... parallel chunk(s)` status is put
586 /// in one canonical order: starts by span, activity completions by span, tool
587 /// completions by call id. Only a causally valid run is reordered, and every
588 /// event, outcome and boundary is kept.
589 fn order_parallel_batches(lines: &mut [Value]) {
590 let mut at = 0;
591 while at < lines.len() {
592 let announces_batch = lines[at]["event"].as_str() == Some("status")
593 && lines[at]["message"].as_str().is_some_and(|message| {
594 message.starts_with("Executing ") && message.ends_with(" parallel chunk(s)")
595 });
596 at += 1;
597 if !announces_batch {
598 continue;
599 }
600 let start = at;
601 let mut end = start;
602 while end < lines.len() && batch_rank(&lines[end]).is_some() {
603 end += 1;
604 }
605 if batch_is_causal(&lines[start..end]) {
606 lines[start..end].sort_by(|left, right| {
607 let (left_rank, left_key) = batch_rank(left).unwrap_or((0, "span_id"));
608 let (right_rank, right_key) = batch_rank(right).unwrap_or((0, "span_id"));
609 left_rank
610 .cmp(&right_rank)
611 .then_with(|| left[left_key].as_str().cmp(&right[right_key].as_str()))
612 });
613 }
614 at = end;
615 }
616 }
617
618 /// Parallel tools race to report: each emits its activity completion and then
619 /// its tool completion, and under load the two tools' pairs interleave. One
620 /// uninterrupted run of completion events is therefore put in a canonical
621 /// order: activity completions by span, then tool completions by call id.
622 /// Every event, its outcome and the run's boundaries are kept; only the order
623 /// inside the run is normalized.
624 fn order_parallel_completions(lines: &mut [Value]) {
625 fn completion(line: &Value) -> Option<(u8, &'static str)> {
626 match line["event"].as_str() {
627 Some("operation_activity_completed") => Some((0, "span_id")),
628 Some("tool_call_complete") => Some((1, "tool_call_id")),
629 _ => None,
630 }
631 }
632 let mut start = 0;
633 while start < lines.len() {
634 if completion(&lines[start]).is_none() {
635 start += 1;
636 continue;
637 }
638 let mut end = start;
639 while end < lines.len() && completion(&lines[end]).is_some() {
640 end += 1;
641 }
642 lines[start..end].sort_by(|left, right| {
643 let (left_kind, left_key) = completion(left).unwrap_or((0, "span_id"));
644 let (right_kind, right_key) = completion(right).unwrap_or((0, "span_id"));
645 left_kind
646 .cmp(&right_kind)
647 .then_with(|| left[left_key].as_str().cmp(&right[right_key].as_str()))
648 });
649 start = end;
650 }
651 }
652
653 #[test]
654 fn parallel_completion_projection_keeps_outcomes_spans_and_causal_boundaries() {
655 let first = json!({"event": "operation_activity_completed", "span_id": "call_a#1", "outcome": "succeeded"});
656 let second = json!({"event": "operation_activity_completed", "span_id": "call_b#2", "outcome": "failed"});
657 let boundary = json!({"event": "error", "message": "exact error bytes"});
658 let start = json!({"event": "operation_activity_started", "span_id": "call_a#1"});
659 let original = vec![
660 start.clone(),
661 second.clone(),
662 first.clone(),
663 boundary.clone(),
664 second.clone(),
665 start.clone(),
666 first.clone(),
667 ];
668 let expected = vec![
669 start.clone(),
670 first.clone(),
671 second.clone(),
672 boundary,
673 second,
674 start,
675 first,
676 ];
677 let mut normalized = original.clone();
678 order_parallel_completions(&mut normalized);
679 assert_eq!(normalized, expected);
680 assert_eq!(normalized.len(), original.len());
681
682 // Two tools' completion pairs interleaved under load normalize to the
683 // same order as the uninterleaved run.
684 let tool = |id: &str| json!({"event": "tool_call_complete", "tool_call_id": id});
685 let activity = |span: &str| json!({"event": "operation_activity_completed", "span_id": span, "outcome": "succeeded"});
686 let mut interleaved = vec![
687 activity("call_b#2"),
688 tool("call_b"),
689 activity("call_a#1"),
690 tool("call_a"),
691 ];
692 order_parallel_completions(&mut interleaved);
693 assert_eq!(
694 interleaved,
695 vec![
696 activity("call_a#1"),
697 activity("call_b#2"),
698 tool("call_a"),
699 tool("call_b")
700 ]
701 );
702
703 let mut wrong_outcome = original.clone();
704 wrong_outcome[1]["outcome"] = json!("succeeded");
705 order_parallel_completions(&mut wrong_outcome);
706 assert_ne!(wrong_outcome, expected, "wrong outcome was normalized away");
707
708 let mut wrong_span = original;
709 wrong_span[1]["span_id"] = json!("call_other#2");
710 order_parallel_completions(&mut wrong_span);
711 assert_ne!(
712 wrong_span, expected,
713 "wrong span relationship was normalized away"
714 );
715 }
716
717 #[test]
718 fn parallel_batch_projection_orders_legal_interleavings_and_keeps_violations() {
719 let status =
720 json!({"event": "status", "message": "Executing 2 read-only tools in 1 parallel chunk(s)"});
721 let start = |span: &str| json!({"event": "operation_activity_started", "span_id": span});
722 let done = |span: &str, outcome: &str| json!({"event": "operation_activity_completed", "span_id": span, "outcome": outcome});
723 let tool = |id: &str| json!({"event": "tool_call_complete", "tool_call_id": id});
724 let tail = json!({"event": "session_updated"});
725 let canonical = vec![
726 status.clone(),
727 start("a#1"),
728 start("b#2"),
729 done("a#1", "succeeded"),
730 done("b#2", "succeeded"),
731 tool("a"),
732 tool("b"),
733 tail.clone(),
734 ];
735 let with_tail = |run: Vec<Value>| {
736 let mut lines = vec![status.clone()];
737 lines.extend(run);
738 lines.push(tail.clone());
739 lines
740 };
741 // Recorded: both start, then both finish. Hosted CI also saw one tool
742 // finish before its sibling started, and both pairs fully interleaved.
743 for run in [
744 vec![
745 start("a#1"),
746 start("b#2"),
747 done("a#1", "succeeded"),
748 done("b#2", "succeeded"),
749 tool("a"),
750 tool("b"),
751 ],
752 vec![
753 start("a#1"),
754 done("a#1", "succeeded"),
755 start("b#2"),
756 done("b#2", "succeeded"),
757 tool("a"),
758 tool("b"),
759 ],
760 vec![
761 start("a#1"),
762 done("a#1", "succeeded"),
763 tool("a"),
764 start("b#2"),
765 done("b#2", "succeeded"),
766 tool("b"),
767 ],
768 vec![
769 start("b#2"),
770 start("a#1"),
771 done("b#2", "succeeded"),
772 tool("b"),
773 done("a#1", "succeeded"),
774 tool("a"),
775 ],
776 ] {
777 let mut lines = with_tail(run);
778 order_parallel_batches(&mut lines);
779 order_parallel_completions(&mut lines);
780 assert_eq!(lines, canonical);
781 }
782
783 // Outcomes survive the reordering.
784 let mut lines = with_tail(vec![
785 start("a#1"),
786 done("a#1", "failed"),
787 start("b#2"),
788 done("b#2", "succeeded"),
789 tool("a"),
790 tool("b"),
791 ]);
792 order_parallel_batches(&mut lines);
793 assert_eq!(lines[3], done("a#1", "failed"));
794
795 // An illegal run is left exactly as recorded: a completion with no start,
796 // a repeated start, a repeated completion, a tool completion before its
797 // activity completion.
798 for run in [
799 vec![done("a#1", "succeeded"), start("a#1"), tool("a")],
800 vec![
801 start("a#1"),
802 start("a#1"),
803 done("a#1", "succeeded"),
804 tool("a"),
805 ],
806 vec![
807 start("a#1"),
808 done("a#1", "succeeded"),
809 done("a#1", "succeeded"),
810 tool("a"),
811 ],
812 vec![start("a#1"), tool("a"), done("a#1", "succeeded")],
813 ] {
814 let recorded = with_tail(run);
815 let mut lines = recorded.clone();
816 order_parallel_batches(&mut lines);
817 assert_eq!(lines, recorded, "an illegal run was normalized away");
818 }
819
820 // Without an announced parallel batch nothing is reordered: a serial run
821 // that starts its second tool early is a real change.
822 let serial = vec![start("b#2"), done("b#2", "succeeded"), start("a#1")];
823 let mut lines = serial.clone();
824 order_parallel_batches(&mut lines);
825 assert_eq!(lines, serial);
826 }
827
828 fn check_invariants(case: &Value, workspace: &Path, provider: &ScriptedProvider) -> Vec<String> {
829 let mut violations = Vec::new();
830 for invariant in case["invariants"].as_array().into_iter().flatten() {
831 if let Some(file) = invariant
832 .get("workspace_file_absent")
833 .and_then(Value::as_str)
834 && workspace.join(file).exists()
835 {
836 violations.push(format!(
837 "`{file}` exists: a tool side effect ran that the scenario forbids"
838 ));
839 }
840 if let Some(file) = invariant
841 .get("workspace_file_present")
842 .and_then(Value::as_str)
843 && !workspace.join(file).exists()
844 {
845 violations.push(format!(
846 "`{file}` is gone: a tool side effect ran that the scenario forbids"
847 ));
848 }
849 if let Some(max) = invariant.get("max_model_requests").and_then(Value::as_u64)
850 && provider.stream_requests() as u64 > max
851 {
852 violations.push(format!(
853 "{} model requests were issued; at most {max} allowed",
854 provider.stream_requests()
855 ));
856 }
857 }
858 violations
859 }
860
861 fn run_case(name: &str, case: &Value, failures: &mut Failures) {
862 run_case_with_deadline(name, case, failures, event_deadline());
863 }
864
865 fn run_case_with_deadline(name: &str, case: &Value, failures: &mut Failures, deadline: Duration) {
866 let sandbox = Sandbox::new(case);
867 let workspace = sandbox.workspace.clone();
868 let (record, provider) =
869 run_scripted_turn_with_deadline(&sandbox, case, &case["provider_script"], deadline);
870 if let Some(error) = record.failure {
871 failures.push(name, error);
872 return;
873 }
874 let mut masker = sandbox.masker(VOLATILE_KEYS);
875 let mut lines = normalize_events(&record.events, &mut masker);
876 let summary = json!({
877 "model_requests": provider.stream_requests(),
878 "non_streaming_requests": provider.other_requests(),
879 "workspace_after": workspace_listing(&workspace),
880 });
881 lines.push(json!({ "harness_summary": summary }));
882
883 let violations = check_invariants(case, &workspace, &provider);
884 if !violations.is_empty() {
885 failures.push(name, format!("invariant broken: {}", violations.join("; ")));
886 return;
887 }
888
889 failures.record(
890 name,
891 golden::check_golden(
892 &golden::family_dir(FAMILY).join(format!("{name}.golden.jsonl")),
893 &golden::jsonl(&lines),
894 ),
895 );
896 }
897
898 #[test]
899 fn harness_timeout_rejects_a_real_stalled_turn() {
900 let mut case = golden::read_case(FAMILY, "plain_answer");
901 case["provider_script"][0]["then"] = json!("hang");
902 case["provider_script"][0]["events"] = json!([]);
903 let mut failures = Failures::default();
904 // This is the same recording path as normal cases. In update mode it must
905 // still fail and must not create this deliberately absent golden.
906 let name = "harness_timeout_control";
907 let path = golden::family_dir(FAMILY).join(format!("{name}.golden.jsonl"));
908 assert!(!path.exists());
909 run_case_with_deadline(name, &case, &mut failures, Duration::from_secs(1));
910 assert!(failures.contains("harness timeout"));
911 assert!(
912 !path.exists(),
913 "an incomplete turn was recorded as a golden"
914 );
915 }
916
917 #[test]
918 fn golden_turn_events_match() {
919 // The goldens carry the span numbers of one pass over the cases in an
920 // otherwise idle process (`one_tool_call` #1, `parallel_tool_calls` #2
921 // and #3); replay that numbering whatever else this process runs.
922 let _spans = crate::core::engine::pin_replay_span_sequence();
923 let names = golden::case_names(FAMILY);
924 let mut failures = Failures::default();
925 for name in &names {
926 let case = golden::read_case(FAMILY, name);
927 run_case(name, &case, &mut failures);
928 }
929 failures.finish(FAMILY, names.len());
930 }
931
931 lines RUST