返回 CodeWhale
runtime.rs
根目录 / crates / tui / src / acp_server / runtime.rs
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(&params)
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
665 lines RUST