| 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 |