| 1 | //! Projection of the existing Runtime owner; no direct provider/tool authority. |
| 2 | use super::*; |
| 3 | use crate::runtime_threads::{ |
| 4 | ExternalApprovalDecision, RuntimeEventRecord, RuntimeTurnStatus, RuntimeTurnStopReason, |
| 5 | StartTurnRequest, TurnRecord, |
| 6 | }; |
| 7 | use std::time::Duration; |
| 8 | struct Permission { |
| 9 | thread_id: String, |
| 10 | turn_id: String, |
| 11 | approval_id: String, |
| 12 | execution_id: String, |
| 13 | } |
| 14 | struct Projection<'a> { |
| 15 | session: &'a str, |
| 16 | thread: &'a str, |
| 17 | turn: &'a str, |
| 18 | cursor: u64, |
| 19 | emit: bool, |
| 20 | calls: HashMap<String, PendingToolCall>, |
| 21 | permissions: HashMap<String, Permission>, |
| 22 | } |
| 23 | /// Drop interrupts only this transport's admitted turn on the current scheduler. |
| 24 | struct TurnClaim { |
| 25 | runtime: Arc<RuntimeThreadManager>, |
| 26 | thread_id: String, |
| 27 | turn_id: String, |
| 28 | armed: bool, |
| 29 | } |
| 30 | impl Drop for TurnClaim { |
| 31 | fn drop(&mut self) { |
| 32 | if !self.armed { |
| 33 | return; |
| 34 | } |
| 35 | let runtime = Arc::clone(&self.runtime); |
| 36 | let thread = self.thread_id.clone(); |
| 37 | let turn = self.turn_id.clone(); |
| 38 | if let Ok(handle) = tokio::runtime::Handle::try_current() { |
| 39 | handle.spawn(async move { |
| 40 | let _ = runtime.interrupt_turn(&thread, &turn).await; |
| 41 | }); |
| 42 | } |
| 43 | } |
| 44 | } |
| 45 | impl AcpServer { |
| 46 | pub(super) async fn drive_prompt<R: AsyncBufRead + Unpin, W: AsyncWrite + Unpin>( |
| 47 | &mut self, |
| 48 | params: Value, |
| 49 | reader: &mut codewhale_app_server::BoundedLines<R>, |
| 50 | writer: &mut W, |
| 51 | ) -> Result<&'static str> { |
| 52 | let (session_id, prompt) = self |
| 53 | .validate_prompt(¶ms) |
| 54 | .map_err(|error| anyhow!(error.message))?; |
| 55 | let binding = self |
| 56 | .sessions |
| 57 | .get(&session_id) |
| 58 | .ok_or_else(|| anyhow!("unknown sessionId"))?; |
| 59 | let thread_id = binding.thread_id.clone(); |
| 60 | // Subscribe before snapshot/admission; replay deduplicates by seq. |
| 61 | let mut receiver = self.runtime.subscribe_events(); |
| 62 | let cursor = self.runtime.get_thread_detail(&thread_id).await?.latest_seq; |
| 63 | let turn = self |
| 64 | .runtime |
| 65 | .start_acp_turn( |
| 66 | &thread_id, |
| 67 | StartTurnRequest { |
| 68 | prompt, |
| 69 | ..Default::default() |
| 70 | }, |
| 71 | ) |
| 72 | .await?; |
| 73 | let mut claim = TurnClaim { |
| 74 | runtime: Arc::clone(&self.runtime), |
| 75 | thread_id: thread_id.clone(), |
| 76 | turn_id: turn.id.clone(), |
| 77 | armed: true, |
| 78 | }; |
| 79 | let mut projection = Projection { |
| 80 | session: &session_id, |
| 81 | thread: &thread_id, |
| 82 | turn: &turn.id, |
| 83 | cursor, |
| 84 | emit: true, |
| 85 | calls: HashMap::new(), |
| 86 | permissions: HashMap::new(), |
| 87 | }; |
| 88 | let terminal = match self |
| 89 | .project_turn(&mut projection, &mut receiver, reader, writer) |
| 90 | .await |
| 91 | { |
| 92 | Ok(terminal) => terminal, |
| 93 | Err(error) => { |
| 94 | self.interrupt_claimed_turn(&thread_id, &turn.id).await?; |
| 95 | let settled = tokio::time::timeout(Duration::from_secs(30), async { |
| 96 | loop { |
| 97 | let detail = self.runtime.get_thread_detail(&thread_id).await?; |
| 98 | if let Some(record) = detail.turns.into_iter().find(|record| record.id == turn.id && !matches!(record.status, RuntimeTurnStatus::Queued | RuntimeTurnStatus::InProgress)) { return Ok::<_, anyhow::Error>(record); } |
| 99 | tokio::time::sleep(Duration::from_millis(25)).await; |
| 100 | } |
| 101 | }).await.map_err(|_| anyhow!("ACP projection failed ({error}); Core cancellation is still unconfirmed, inspect its canonical receipt"))??; |
| 102 | self.save_checkpoint(&thread_id, &session_id).await?; |
| 103 | claim.armed = false; |
| 104 | return Err(anyhow!( |
| 105 | "ACP projection failed ({error}); actual Core terminal is {:?}; partial receipts retained", |
| 106 | settled.status |
| 107 | )); |
| 108 | } |
| 109 | }; |
| 110 | claim.armed = false; |
| 111 | if let Some(binding) = self.sessions.get_mut(&session_id) { |
| 112 | binding.cursor = projection.cursor; |
| 113 | } |
| 114 | self.save_checkpoint(&thread_id, &session_id).await?; |
| 115 | if terminal |
| 116 | .model_request_diagnostics |
| 117 | .as_ref() |
| 118 | .and_then(|facts| facts.stop_reason) |
| 119 | == Some(RuntimeTurnStopReason::StepBudgetExhausted) |
| 120 | { |
| 121 | return Ok("max_turn_requests"); |
| 122 | } |
| 123 | match terminal.status { |
| 124 | RuntimeTurnStatus::Completed => Ok("end_turn"), |
| 125 | RuntimeTurnStatus::Interrupted | RuntimeTurnStatus::Canceled => Ok("cancelled"), |
| 126 | _ => Err(anyhow!( |
| 127 | "{}", |
| 128 | terminal.error.unwrap_or_else( |
| 129 | || "Core turn failed; inspect its canonical receipt".to_string() |
| 130 | ) |
| 131 | )), |
| 132 | } |
| 133 | } |
| 134 | async fn interrupt_claimed_turn(&self, thread: &str, turn: &str) -> Result<()> { |
| 135 | match self.runtime.interrupt_turn(thread, turn).await { |
| 136 | Ok(_) => Ok(()), |
| 137 | Err(error) => { |
| 138 | // Completion may win the race with EOF, cancel, or a failed |
| 139 | // writer. Accept only the actual terminal of this claimed turn. |
| 140 | let detail = self.runtime.get_thread_detail(thread).await.map_err(|read| anyhow!("ACP interruption failed ({error}); Core receipt cannot be read ({read}); inspect without replay"))?; |
| 141 | if detail.turns.iter().any(|record| { |
| 142 | record.id == turn |
| 143 | && !matches!( |
| 144 | record.status, |
| 145 | RuntimeTurnStatus::Queued | RuntimeTurnStatus::InProgress |
| 146 | ) |
| 147 | }) { |
| 148 | Ok(()) |
| 149 | } else { |
| 150 | Err(anyhow!( |
| 151 | "ACP cannot confirm interruption of its claimed Core turn: {error}; inspect without replay" |
| 152 | )) |
| 153 | } |
| 154 | } |
| 155 | } |
| 156 | } |
| 157 | async fn save_checkpoint(&self, thread: &str, session: &str) -> Result<()> { |
| 158 | // The same full Engine snapshot writer as HTTP, never projected text. |
| 159 | let axum::Json(_saved_checkpoint) = crate::runtime_api::sessions::save_session_in_runtime( |
| 160 | &self.runtime, |
| 161 | &self.sessions_dir, |
| 162 | crate::runtime_api::sessions::SaveSessionRequest { |
| 163 | thread_id: Some(thread.to_string()), |
| 164 | session_id: Some(session.to_string()), |
| 165 | }, |
| 166 | ) |
| 167 | .await |
| 168 | .map_err(|error| { |
| 169 | anyhow!( |
| 170 | "Core turn finished but its checkpoint failed; do not replay: {}", |
| 171 | error.message |
| 172 | ) |
| 173 | })?; |
| 174 | Ok(()) |
| 175 | } |
| 176 | async fn project_turn<R: AsyncBufRead + Unpin, W: AsyncWrite + Unpin>( |
| 177 | &self, |
| 178 | projection: &mut Projection<'_>, |
| 179 | receiver: &mut tokio::sync::broadcast::Receiver<RuntimeEventRecord>, |
| 180 | reader: &mut codewhale_app_server::BoundedLines<R>, |
| 181 | writer: &mut W, |
| 182 | ) -> Result<TurnRecord> { |
| 183 | let session = projection.session; |
| 184 | let thread = projection.thread; |
| 185 | let turn = projection.turn; |
| 186 | let mut eof = false; |
| 187 | let mut cancelled = false; |
| 188 | let mut cancel_deadline: Option<tokio::time::Instant> = None; |
| 189 | let mut replay: Option<crate::runtime_threads::RuntimeEventReplay> = None; |
| 190 | let mut replay_batch = VecDeque::new(); |
| 191 | loop { |
| 192 | if cancel_deadline.is_some_and(|deadline| tokio::time::Instant::now() >= deadline) { |
| 193 | return Err(anyhow!( |
| 194 | "Core cancellation remains nonterminal; ACP cannot report a fabricated cancelled receipt" |
| 195 | )); |
| 196 | } |
| 197 | tokio::select! { |
| 198 | () = async { match cancel_deadline { Some(deadline) => tokio::time::sleep_until(deadline).await, None => std::future::pending().await } } => return Err(anyhow!("Core cancellation remains nonterminal; ACP cannot report a fabricated cancelled receipt")), |
| 199 | line = reader.next_line(), if !eof => { |
| 200 | let Some(line) = line? else { |
| 201 | eof = true; cancelled = true; projection.permissions.clear(); self.interrupt_claimed_turn(thread, turn).await?; |
| 202 | cancel_deadline = Some(tokio::time::Instant::now() + Duration::from_secs(30)); continue; |
| 203 | }; |
| 204 | if line.trim().is_empty() { continue; } |
| 205 | let message: Value = match serde_json::from_str(&line) { |
| 206 | Ok(message) => message, |
| 207 | Err(error) => { write_jsonrpc_error(writer, None, -32700, format!("invalid json: {error}")).await?; continue; } |
| 208 | }; |
| 209 | if codewhale_app_server::is_control_input_closed(&message) { |
| 210 | eof=true; cancelled=true; projection.permissions.clear();self.interrupt_claimed_turn(thread,turn).await?; |
| 211 | cancel_deadline.get_or_insert(tokio::time::Instant::now()+Duration::from_secs(30));continue; |
| 212 | } |
| 213 | let response_id = message.get("id").cloned().map(|id| self.response_id_policy.response_id(id)); |
| 214 | if message.get("jsonrpc").and_then(Value::as_str) != Some("2.0") { write_jsonrpc_error(writer, response_id, -32600, "jsonrpc version must be 2.0").await?; continue; } |
| 215 | if is_jsonrpc_response(&message) { |
| 216 | if !cancelled && let Some(id) = message.get("id").and_then(Value::as_str) && let Some(permission) = projection.permissions.remove(id) { self.answer_permission(&permission, &message).await?; } |
| 217 | continue; |
| 218 | } |
| 219 | if message.get("method").and_then(Value::as_str) == Some("session/cancel") { |
| 220 | if message.pointer("/params/sessionId").and_then(Value::as_str) == Some(session) { |
| 221 | cancelled = true; projection.permissions.clear(); self.interrupt_claimed_turn(thread, turn).await?; |
| 222 | cancel_deadline.get_or_insert(tokio::time::Instant::now() + Duration::from_secs(30)); |
| 223 | } |
| 224 | if let Some(id) = response_id { write_jsonrpc_result(writer, id, json!(null)).await?; } |
| 225 | } else if let Some(id) = response_id { write_jsonrpc_error(writer, Some(id), -32603, "a session/prompt turn is already in progress").await?; } |
| 226 | } |
| 227 | record = async { replay_batch.pop_front().expect("nonempty bounded replay batch") }, if !replay_batch.is_empty() => { |
| 228 | projection.emit = !cancelled && !eof; |
| 229 | if let Some(done) = self.project_event(projection, &record, writer).await? { return Ok(done); } |
| 230 | } |
| 231 | batch = async { replay.as_mut().expect("active replay").batches.recv().await }, if replay.is_some() && replay_batch.is_empty() => { |
| 232 | match batch { |
| 233 | Some(batch) => replay_batch.extend(batch.map_err(|error| anyhow!(error))?), |
| 234 | None => replay = None, |
| 235 | } |
| 236 | } |
| 237 | event = receiver.recv(), if replay.is_none() && replay_batch.is_empty() => { |
| 238 | match event { |
| 239 | Ok(record) => { projection.emit = !cancelled && !eof; if let Some(done) = self.project_event(projection, &record, writer).await? { return Ok(done); } }, |
| 240 | Err(tokio::sync::broadcast::error::RecvError::Lagged(_)) => { |
| 241 | let captured = self.runtime.replay_events(thread, Some(projection.cursor), None).await?; |
| 242 | if captured.base_seq > projection.cursor { return Err(anyhow!("ACP replay no longer covers the cursor; projection is incomplete")); } |
| 243 | replay = Some(captured); |
| 244 | } |
| 245 | Err(tokio::sync::broadcast::error::RecvError::Closed) => return Err(anyhow!("Runtime owner closed before its terminal receipt")), |
| 246 | } |
| 247 | } |
| 248 | } |
| 249 | } |
| 250 | } |
| 251 | async fn answer_permission(&self, permission: &Permission, response: &Value) -> Result<()> { |
| 252 | let detail = self |
| 253 | .runtime |
| 254 | .get_thread_detail(&permission.thread_id) |
| 255 | .await?; |
| 256 | if !detail.pending_approvals.iter().any(|pending| { |
| 257 | pending.id == permission.approval_id |
| 258 | && pending.turn_id == permission.turn_id |
| 259 | && pending.tool_call_id.as_deref() == Some(permission.execution_id.as_str()) |
| 260 | }) { |
| 261 | return Ok(()); |
| 262 | } |
| 263 | let allow = response.get("error").is_none() |
| 264 | && response |
| 265 | .pointer("/result/outcome/outcome") |
| 266 | .and_then(Value::as_str) |
| 267 | == Some("selected") |
| 268 | && response |
| 269 | .pointer("/result/outcome/optionId") |
| 270 | .and_then(Value::as_str) |
| 271 | == Some("allow-once"); |
| 272 | self.runtime.deliver_external_approval( |
| 273 | &permission.approval_id, |
| 274 | if allow { |
| 275 | ExternalApprovalDecision::Allow { remember: false } |
| 276 | } else { |
| 277 | ExternalApprovalDecision::Deny { remember: false } |
| 278 | }, |
| 279 | ); |
| 280 | Ok(()) |
| 281 | } |
| 282 | async fn project_event<W: AsyncWrite + Unpin>( |
| 283 | &self, |
| 284 | projection: &mut Projection<'_>, |
| 285 | record: &RuntimeEventRecord, |
| 286 | writer: &mut W, |
| 287 | ) -> Result<Option<TurnRecord>> { |
| 288 | let session = projection.session; |
| 289 | let thread = projection.thread; |
| 290 | let turn = projection.turn; |
| 291 | let cursor = &mut projection.cursor; |
| 292 | let calls = &mut projection.calls; |
| 293 | let permissions = &mut projection.permissions; |
| 294 | let emit = projection.emit; |
| 295 | if record.thread_id != thread || record.seq <= *cursor { |
| 296 | return Ok(None); |
| 297 | } |
| 298 | *cursor = record.seq; |
| 299 | if record.turn_id.as_deref() != Some(turn) { |
| 300 | return Ok(None); |
| 301 | } |
| 302 | let p = &record.payload; |
| 303 | match record.event.as_str() { |
| 304 | "item.delta" |
| 305 | if emit && p.get("kind").and_then(Value::as_str) == Some("agent_message") => |
| 306 | { |
| 307 | if let Some(text) = p.get("delta").and_then(Value::as_str) { |
| 308 | write_session_update(writer, session, text.to_string()).await?; |
| 309 | } |
| 310 | } |
| 311 | "item.started" => { |
| 312 | if let Some(tool) = p.get("tool") { |
| 313 | let id = tool |
| 314 | .get("id") |
| 315 | .and_then(Value::as_str) |
| 316 | .ok_or_else(|| anyhow!("Tool start lacks Core identity"))? |
| 317 | .to_string(); |
| 318 | let call = PendingToolCall { |
| 319 | execution_id: id.clone(), |
| 320 | name: tool |
| 321 | .get("name") |
| 322 | .and_then(Value::as_str) |
| 323 | .unwrap_or_default() |
| 324 | .to_string(), |
| 325 | input: tool.get("input").cloned().unwrap_or(Value::Null), |
| 326 | }; |
| 327 | if emit { |
| 328 | write_tool_call_start(writer, session, &call).await?; |
| 329 | } |
| 330 | calls.insert(id, call); |
| 331 | } |
| 332 | } |
| 333 | "tool.execution_started" if emit => { |
| 334 | let id = p |
| 335 | .get("execution_id") |
| 336 | .and_then(Value::as_str) |
| 337 | .ok_or_else(|| anyhow!("Dispatch lacks Core identity"))?; |
| 338 | let call = calls |
| 339 | .get(id) |
| 340 | .ok_or_else(|| anyhow!("Dispatch projection lacks its admitted start"))?; |
| 341 | write_tool_call_update(writer, session, call, "in_progress", None).await?; |
| 342 | } |
| 343 | "item.completed" | "item.failed" | "item.interrupted" if emit => { |
| 344 | if let Some(item) = p.get("item") |
| 345 | && let Some(id) = item |
| 346 | .pointer("/metadata/execution_id") |
| 347 | .and_then(Value::as_str) |
| 348 | && let Some(call) = calls.remove(id) |
| 349 | { |
| 350 | let blocks: Vec<codewhale_tools::ToolResultContentBlock> = item |
| 351 | .pointer("/metadata/acp_result_content") |
| 352 | .cloned() |
| 353 | .map(serde_json::from_value) |
| 354 | .transpose()? |
| 355 | .unwrap_or_default(); |
| 356 | let status = if item.get("status").and_then(Value::as_str) == Some("completed") |
| 357 | { |
| 358 | "completed" |
| 359 | } else { |
| 360 | "failed" |
| 361 | }; |
| 362 | write_tool_call_update_with_blocks( |
| 363 | writer, |
| 364 | session, |
| 365 | &call, |
| 366 | status, |
| 367 | item.get("detail").and_then(Value::as_str), |
| 368 | &blocks, |
| 369 | ) |
| 370 | .await?; |
| 371 | } |
| 372 | } |
| 373 | "approval.required" if emit => { |
| 374 | let approval = p |
| 375 | .get("approval_id") |
| 376 | .and_then(Value::as_str) |
| 377 | .ok_or_else(|| anyhow!("Approval lacks Runtime identity"))?; |
| 378 | let execution = p |
| 379 | .get("tool_call_id") |
| 380 | .and_then(Value::as_str) |
| 381 | .ok_or_else(|| anyhow!("Approval lacks Core correlator"))?; |
| 382 | let detail = self.runtime.get_thread_detail(thread).await?; |
| 383 | if detail.pending_approvals.iter().any(|pending| { |
| 384 | pending.id == approval |
| 385 | && pending.turn_id == turn |
| 386 | && pending.tool_call_id.as_deref() == Some(execution) |
| 387 | }) { |
| 388 | let call = calls |
| 389 | .get(execution) |
| 390 | .ok_or_else(|| anyhow!("Approval projection lacks admitted tool start"))?; |
| 391 | let id = format!( |
| 392 | "codewhale-permission-{}", |
| 393 | NEXT_ACP_PERMISSION_REQUEST_ID.fetch_add(1, Ordering::Relaxed) |
| 394 | ); |
| 395 | write_json_line(writer, json!({"jsonrpc":"2.0","id":id,"method":"session/request_permission","params":{"sessionId":session,"toolCall":{"toolCallId":execution,"title":tool_call_title(call),"kind":tool_call_kind(call),"status":"pending","rawInput":call.input,"content":[{"type":"content","content":{"type":"text","text":p.get("description")}}]},"options":[{"optionId":"allow-once","name":"Allow once","kind":"allow_once"},{"optionId":"reject-once","name":"Reject","kind":"reject_once"}]}})).await?; |
| 396 | permissions.insert( |
| 397 | id, |
| 398 | Permission { |
| 399 | thread_id: thread.to_string(), |
| 400 | turn_id: turn.to_string(), |
| 401 | approval_id: approval.to_string(), |
| 402 | execution_id: execution.to_string(), |
| 403 | }, |
| 404 | ); |
| 405 | } |
| 406 | } |
| 407 | "approval.withdrawn" | "approval.decided" => { |
| 408 | if let Some(id) = p.get("approval_id").and_then(Value::as_str) { |
| 409 | permissions.retain(|_, permission| permission.approval_id != id); |
| 410 | } |
| 411 | } |
| 412 | crate::runtime_threads::RUNTIME_STORE_FAILURE_EVENT |
| 413 | if p.get("terminal").and_then(Value::as_bool) == Some(true) => |
| 414 | { |
| 415 | return Err(anyhow!( |
| 416 | "Core store failure has no terminal receipt: {}; inspect without replay", |
| 417 | p.get("message") |
| 418 | .and_then(Value::as_str) |
| 419 | .unwrap_or("canonical store unavailable") |
| 420 | )); |
| 421 | } |
| 422 | "turn.completed" => { |
| 423 | permissions.clear(); |
| 424 | let terminal: TurnRecord = serde_json::from_value( |
| 425 | p.get("turn") |
| 426 | .cloned() |
| 427 | .ok_or_else(|| anyhow!("Terminal event lacks its actual record"))?, |
| 428 | )?; |
| 429 | if terminal.id != turn { |
| 430 | return Err(anyhow!("Terminal does not match admitted turn")); |
| 431 | } |
| 432 | return Ok(Some(terminal)); |
| 433 | } |
| 434 | _ => {} |
| 435 | } |
| 436 | Ok(None) |
| 437 | } |
| 438 | } |
| 439 | |
| 440 | #[cfg(test)] |
| 441 | mod tests { |
| 442 | use super::super::tests::{Rig, fixture_config}; |
| 443 | use super::*; |
| 444 | fn event(seq: u64, thread: &str, turn: &str, name: &str, payload: Value) -> RuntimeEventRecord { |
| 445 | serde_json::from_value(json!({"seq":seq,"timestamp":chrono::Utc::now(),"thread_id":thread,"turn_id":turn,"event":name,"payload":payload})).unwrap() |
| 446 | } |
| 447 | #[tokio::test(flavor = "current_thread")] |
| 448 | async fn projection_deduplicates_replay_and_refuses_dispatch_without_core_start() -> Result<()> |
| 449 | { |
| 450 | let mut rig = Rig::new(fixture_config(), vec![])?; |
| 451 | let id = rig.new_session().await; |
| 452 | let thread = rig.server.sessions[&id].thread_id.clone(); |
| 453 | let mut p = Projection { |
| 454 | session: &id, |
| 455 | thread: &thread, |
| 456 | turn: "turn-current", |
| 457 | cursor: 0, |
| 458 | emit: true, |
| 459 | calls: HashMap::new(), |
| 460 | permissions: HashMap::new(), |
| 461 | }; |
| 462 | let mut output = Vec::new(); |
| 463 | let start = event( |
| 464 | 1, |
| 465 | &thread, |
| 466 | p.turn, |
| 467 | "item.started", |
| 468 | json!({"tool":{"id":"core-call","name":"read","input":{"path":"x"}}}), |
| 469 | ); |
| 470 | rig.server |
| 471 | .project_event(&mut p, &start, &mut output) |
| 472 | .await?; |
| 473 | let before = output.len(); |
| 474 | rig.server |
| 475 | .project_event(&mut p, &start, &mut output) |
| 476 | .await?; |
| 477 | assert_eq!(output.len(), before); |
| 478 | let dispatch = event( |
| 479 | 2, |
| 480 | &thread, |
| 481 | p.turn, |
| 482 | "tool.execution_started", |
| 483 | json!({"execution_id":"unknown"}), |
| 484 | ); |
| 485 | assert!( |
| 486 | rig.server |
| 487 | .project_event(&mut p, &dispatch, &mut output) |
| 488 | .await |
| 489 | .is_err() |
| 490 | ); |
| 491 | rig.close().await; |
| 492 | Ok(()) |
| 493 | } |
| 494 | #[tokio::test(flavor = "current_thread")] |
| 495 | async fn projection_filters_foreign_turns_reasoning_and_unminted_approvals() -> Result<()> { |
| 496 | let mut rig = Rig::new(fixture_config(), vec![])?; |
| 497 | let id = rig.new_session().await; |
| 498 | let thread = rig.server.sessions[&id].thread_id.clone(); |
| 499 | let mut p = Projection { |
| 500 | session: &id, |
| 501 | thread: &thread, |
| 502 | turn: "turn-current", |
| 503 | cursor: 0, |
| 504 | emit: true, |
| 505 | calls: HashMap::new(), |
| 506 | permissions: HashMap::new(), |
| 507 | }; |
| 508 | let mut output = Vec::new(); |
| 509 | for record in [ |
| 510 | event( |
| 511 | 1, |
| 512 | "foreign", |
| 513 | p.turn, |
| 514 | "item.delta", |
| 515 | json!({"kind":"agent_message","delta":"foreign"}), |
| 516 | ), |
| 517 | event( |
| 518 | 2, |
| 519 | &thread, |
| 520 | "foreign-turn", |
| 521 | "item.delta", |
| 522 | json!({"kind":"agent_message","delta":"other turn"}), |
| 523 | ), |
| 524 | event( |
| 525 | 3, |
| 526 | &thread, |
| 527 | p.turn, |
| 528 | "item.delta", |
| 529 | json!({"kind":"reasoning","delta":"private reasoning"}), |
| 530 | ), |
| 531 | event( |
| 532 | 4, |
| 533 | &thread, |
| 534 | p.turn, |
| 535 | "approval.required", |
| 536 | json!({"approval_id":"unminted","tool_call_id":"provider-id"}), |
| 537 | ), |
| 538 | ] { |
| 539 | rig.server |
| 540 | .project_event(&mut p, &record, &mut output) |
| 541 | .await?; |
| 542 | } |
| 543 | assert!(output.is_empty()); |
| 544 | assert!(p.permissions.is_empty()); |
| 545 | rig.close().await; |
| 546 | Ok(()) |
| 547 | } |
| 548 | #[tokio::test(flavor = "current_thread")] |
| 549 | async fn terminal_store_fault_refuses_a_fabricated_completion() -> Result<()> { |
| 550 | let mut rig = Rig::new(fixture_config(), vec![])?; |
| 551 | let id = rig.new_session().await; |
| 552 | let thread = rig.server.sessions[&id].thread_id.clone(); |
| 553 | let mut p = Projection { |
| 554 | session: &id, |
| 555 | thread: &thread, |
| 556 | turn: "turn-current", |
| 557 | cursor: 0, |
| 558 | emit: true, |
| 559 | calls: HashMap::new(), |
| 560 | permissions: HashMap::new(), |
| 561 | }; |
| 562 | let mut output = Vec::new(); |
| 563 | let failure = event( |
| 564 | 1, |
| 565 | &thread, |
| 566 | p.turn, |
| 567 | crate::runtime_threads::RUNTIME_STORE_FAILURE_EVENT, |
| 568 | json!({"terminal":true,"message":"fixture store fault"}), |
| 569 | ); |
| 570 | let error = rig |
| 571 | .server |
| 572 | .project_event(&mut p, &failure, &mut output) |
| 573 | .await |
| 574 | .unwrap_err(); |
| 575 | assert!(error.to_string().contains("no terminal receipt")); |
| 576 | assert!(output.is_empty()); |
| 577 | rig.close().await; |
| 578 | Ok(()) |
| 579 | } |
| 580 | struct FailedAnswerWriter(Vec<u8>); |
| 581 | impl AsyncWrite for FailedAnswerWriter { |
| 582 | fn poll_write( |
| 583 | mut self: std::pin::Pin<&mut Self>, |
| 584 | _: &mut std::task::Context<'_>, |
| 585 | bytes: &[u8], |
| 586 | ) -> std::task::Poll<std::io::Result<usize>> { |
| 587 | if String::from_utf8_lossy(bytes).contains("agent_message_chunk") { |
| 588 | return std::task::Poll::Ready(Err(std::io::Error::new( |
| 589 | std::io::ErrorKind::BrokenPipe, |
| 590 | "fixture client disappeared after tool effect", |
| 591 | ))); |
| 592 | } |
| 593 | self.0.extend_from_slice(bytes); |
| 594 | std::task::Poll::Ready(Ok(bytes.len())) |
| 595 | } |
| 596 | fn poll_flush( |
| 597 | self: std::pin::Pin<&mut Self>, |
| 598 | _: &mut std::task::Context<'_>, |
| 599 | ) -> std::task::Poll<std::io::Result<()>> { |
| 600 | std::task::Poll::Ready(Ok(())) |
| 601 | } |
| 602 | fn poll_shutdown( |
| 603 | self: std::pin::Pin<&mut Self>, |
| 604 | _: &mut std::task::Context<'_>, |
| 605 | ) -> std::task::Poll<std::io::Result<()>> { |
| 606 | std::task::Poll::Ready(Ok(())) |
| 607 | } |
| 608 | } |
| 609 | #[tokio::test(flavor = "current_thread")] |
| 610 | async fn failed_transport_after_completed_write_retains_actual_core_terminal_and_checkpoint() |
| 611 | -> Result<()> { |
| 612 | use crate::llm_client::mock::canned; |
| 613 | let mut config = fixture_config(); |
| 614 | config.yolo = Some(true); |
| 615 | let mut rig = Rig::new( |
| 616 | config, |
| 617 | vec![ |
| 618 | canned::tool_call_turn( |
| 619 | "write", |
| 620 | "write", |
| 621 | r#"{"path":"writer-failure.txt","content":"completed"}"#, |
| 622 | ), |
| 623 | canned::simple_text_turn("final answer"), |
| 624 | ], |
| 625 | )?; |
| 626 | let id = rig.new_session().await; |
| 627 | let (client, input) = tokio::io::duplex(1024); |
| 628 | let mut reader = codewhale_app_server::BoundedLines::new(BufReader::new(input)); |
| 629 | let mut writer = FailedAnswerWriter(Vec::new()); |
| 630 | let error = tokio::time::timeout( |
| 631 | Duration::from_secs(20), |
| 632 | rig.server.drive_prompt( |
| 633 | json!({"sessionId":id,"prompt":"write then answer"}), |
| 634 | &mut reader, |
| 635 | &mut writer, |
| 636 | ), |
| 637 | ) |
| 638 | .await? |
| 639 | .unwrap_err(); |
| 640 | drop(client); |
| 641 | assert!(error.to_string().contains("partial receipts retained")); |
| 642 | assert_eq!( |
| 643 | std::fs::read_to_string(rig.workspace.join("writer-failure.txt"))?, |
| 644 | "completed" |
| 645 | ); |
| 646 | let store = crate::session_manager::SessionManager::new(rig.server.sessions_dir.clone())?; |
| 647 | let saved = store.load_session(&id)?; |
| 648 | assert!( |
| 649 | saved |
| 650 | .messages |
| 651 | .iter() |
| 652 | .flat_map(|m| &m.content) |
| 653 | .any(|b| matches!(b, codewhale_models::ContentBlock::ToolResult { .. })) |
| 654 | ); |
| 655 | let thread = &rig.server.sessions[&id].thread_id; |
| 656 | let detail = rig.server.runtime.get_thread_detail(thread).await?; |
| 657 | assert!(!matches!( |
| 658 | detail.turns.last().unwrap().status, |
| 659 | RuntimeTurnStatus::Queued | RuntimeTurnStatus::InProgress |
| 660 | )); |
| 661 | rig.close().await; |
| 662 | Ok(()) |
| 663 | } |
| 664 | } |
| 665 |