1use std::path::{Path, PathBuf};
8use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
9use std::sync::{Arc, Mutex};
10
11use a3s_effect::{
12 answer_fact, confirm_fact, ingest_coding, message_fact, resume_coding, ActorError, CodingPhase,
13 CodingServices, CodingView, Compactor, Completion, CompletionRequest, Exit, FileLog,
14 HarnessConfig, HarnessGraph, LogStore, ModelDecision, NewFact, ToolCall, ToolRunner, ToolSpec,
15};
16use anyhow::Result;
17
18use crate::ask_user::{self, AskUserError};
19use crate::llm::{LlmClient, LlmResponse, Message, ToolDefinition};
20use crate::permissions::{PermissionDecision, PermissionPolicy};
21use crate::queue::SessionLane;
22use crate::tools::ToolExecutor;
23
24const THREAD: &str = "thread-1";
25
26pub fn thread_for_session(session_id: &str) -> String {
29 let valid = !session_id.is_empty()
30 && session_id.len() <= 128
31 && session_id
32 .chars()
33 .all(|ch| ch.is_ascii_alphanumeric() || matches!(ch, '_' | '.' | ':' | '-'));
34 if valid {
35 return session_id.to_string();
36 }
37 let hex: String = session_id.bytes().fold(String::new(), |mut hex, byte| {
38 hex.push_str(&format!("{byte:02x}"));
39 hex
40 });
41 let mut thread = format!("s-{hex}");
42 thread.truncate(128);
43 thread
44}
45
46pub fn confirmation_required(policy: &PermissionPolicy, tool_name: &str) -> bool {
51 policy.check(tool_name, &serde_json::json!({})) == PermissionDecision::Ask
52}
53
54pub struct FactRun {
55 dir: PathBuf,
56 thread: String,
57 actor: a3s_effect::Actor<CodingServices, CodingView>,
58 services: Arc<CodingServices>,
59 limit: u32,
60}
61
62#[derive(Default)]
64pub struct FactRunHostMounts<'a> {
65 pub compose: Option<&'a crate::meta_harness::HarnessComposeOptions>,
66 pub registry: Option<&'a dyn crate::meta_harness::HostHarnessRegistry>,
67 pub assembler: Option<&'a dyn crate::meta_harness::HostHarnessAssembler>,
68}
69
70struct CappedCompletion {
71 inner: Arc<dyn Completion>,
72 run_id: String,
73}
74
75impl Completion for CappedCompletion {
76 fn complete(
77 &self,
78 request: CompletionRequest,
79 ) -> a3s_effect::coding::BoxFuture<Result<ModelDecision, ActorError>> {
80 let inner = Arc::clone(&self.inner);
81 let run_id = self.run_id.clone();
82 Box::pin(async move {
83 let decision = inner.complete(request).await?;
84 if let ModelDecision::Question { .. } = &decision {
85 if let Err(error) = ask_user::begin(&run_id, "q", "cap-check", &[], false) {
86 if matches!(error, AskUserError::CapExceeded) {
87 return Err(ActorError::Defect("ask_user question cap exceeded".into()));
88 }
89 }
90 }
91 Ok(decision)
92 })
93 }
94}
95
96#[derive(Clone)]
101pub(crate) struct SessionSurface {
102 pub(crate) agent: crate::agent::AgentLoop,
103 pub(crate) session_id: String,
104 pub(crate) checkpoint: Option<CheckpointEmit>,
105 pub(crate) events: Option<tokio::sync::mpsc::Sender<crate::agent::AgentEvent>>,
106 pub(crate) cancel: tokio_util::sync::CancellationToken,
107 pub(crate) transcript: Arc<Mutex<Vec<Message>>>,
108 pub(crate) usage: Arc<Mutex<crate::llm::TokenUsage>>,
109 pub(crate) confirmation: Option<Arc<dyn crate::hitl::ConfirmationProvider>>,
110 pub(crate) run_store: Option<Arc<crate::run::InMemoryRunStore>>,
111 pub(crate) run_id: Option<String>,
112 pub(crate) ledger: Arc<Mutex<crate::harness_loop::MutationLedger>>,
113 pub(crate) reports: Arc<Mutex<Vec<crate::verification::VerificationReport>>>,
114 pub(crate) run_control: Option<Arc<crate::run_control::RunControlInbox>>,
115 pub(crate) harness: Option<crate::meta_harness::HarnessComposeOptions>,
116 pub(crate) host_harness_registry:
117 Option<std::sync::Arc<dyn crate::meta_harness::HostHarnessRegistry>>,
118 pub(crate) host_harness_assembler:
119 Option<std::sync::Arc<dyn crate::meta_harness::HostHarnessAssembler>>,
120}
121
122#[derive(Clone)]
124pub(crate) struct CheckpointEmit {
125 pub(crate) sink: Arc<dyn crate::loop_checkpoint::LoopCheckpointSink>,
126 pub(crate) run_id: String,
127 pub(crate) session_id: String,
128 pub(crate) capability_binding: Option<crate::capability::RunCapabilityBindingV1>,
129}
130
131pub struct LiveCompletion {
132 client: Arc<dyn LlmClient>,
133 calls: Arc<AtomicUsize>,
134 policy: PermissionPolicy,
135 catalog: Vec<ToolDefinition>,
136 surface: Option<SessionSurface>,
137 context_ready: Arc<AtomicBool>,
138 cached_system: Arc<Mutex<Option<String>>>,
139 cached_prompt: Arc<Mutex<Option<String>>>,
140}
141
142impl LiveCompletion {
143 pub fn new(
144 client: Arc<dyn LlmClient>,
145 policy: PermissionPolicy,
146 catalog: Vec<ToolDefinition>,
147 ) -> (Self, Arc<AtomicUsize>) {
148 let calls = Arc::new(AtomicUsize::new(0));
149 (
150 Self {
151 client,
152 calls: Arc::clone(&calls),
153 policy,
154 catalog,
155 surface: None,
156 context_ready: Arc::new(AtomicBool::new(false)),
157 cached_system: Arc::new(Mutex::new(None)),
158 cached_prompt: Arc::new(Mutex::new(None)),
159 },
160 calls,
161 )
162 }
163
164 pub(crate) fn with_surface(mut self, surface: SessionSurface) -> Self {
165 self.surface = Some(surface);
166 self
167 }
168}
169
170fn interrupted_response() -> crate::llm::LlmResponse {
171 crate::llm::LlmResponse {
172 message: Message::assistant("(Response interrupted by the user.)"),
173 usage: crate::llm::TokenUsage::default(),
174 stop_reason: Some("cancelled".to_string()),
175 token_logprobs: Vec::new(),
176 meta: None,
177 }
178}
179
180async fn complete_detached(
181 client: Arc<dyn LlmClient>,
182 messages: &[Message],
183 system: Option<&str>,
184 tools: &[ToolDefinition],
185) -> Result<crate::llm::LlmResponse, ActorError> {
186 let cancel = tokio_util::sync::CancellationToken::new();
187 match client
188 .complete_streaming(messages, system, tools, cancel)
189 .await
190 {
191 Ok(mut events) => {
192 let mut done = None;
193 while let Some(event) = events.recv().await {
194 if let crate::llm::StreamEvent::Done(response) = event {
195 done = Some(response);
196 }
197 }
198 if let Some(response) = done {
199 return Ok(response);
200 }
201 client
202 .complete(messages, system, tools)
203 .await
204 .map_err(model_error)
205 }
206 Err(_) => client
207 .complete(messages, system, tools)
208 .await
209 .map_err(model_error),
210 }
211}
212
213fn model_error(error: anyhow::Error) -> ActorError {
214 if let Some(message) = crate::llm::non_retryable_llm_error_message(&error) {
215 return ActorError::Defect(message.to_string());
216 }
217 ActorError::Handler {
218 key: "infer".into(),
219 message: error.to_string(),
220 }
221}
222
223async fn model_response(
224 client: Arc<dyn LlmClient>,
225 messages: Vec<Message>,
226 system: Option<String>,
227 tools: Vec<ToolDefinition>,
228 surface: Option<SessionSurface>,
229) -> Result<crate::llm::LlmResponse, ActorError> {
230 let Some(surface) = surface else {
231 return complete_detached(client, &messages, system.as_deref(), &tools).await;
232 };
233 surface
234 .agent
235 .fact_budget_gate(&surface.session_id, &surface.cancel)
236 .await
237 .map_err(|error| ActorError::Defect(error.to_string()))?;
238 if surface.cancel.is_cancelled() {
239 return Ok(interrupted_response());
240 }
241 let timeout = surface.agent.llm_api_timeout();
246 let cancel = surface.cancel.clone();
247 let attempt = surface.cancel.child_token();
248 let attempt_for_task = attempt.clone();
249 let mut handle = tokio::spawn(async move {
250 read_model_response(
251 client.as_ref(),
252 &messages,
253 system.as_deref(),
254 &tools,
255 &surface,
256 &attempt_for_task,
257 )
258 .await
259 });
260 let Some(timeout) = timeout else {
261 return await_model_task(handle, &cancel).await;
262 };
263 tokio::select! {
264 biased;
265 _ = cancel.cancelled() => {
266 handle.abort();
267 let _ = handle.await;
268 Ok(interrupted_response())
269 }
270 joined = tokio::time::timeout(timeout, &mut handle) => match joined {
271 Ok(joined) => flatten_model_task(joined),
272 Err(_) => {
273 handle.abort();
274 attempt.cancel();
275 Err(ActorError::Handler {
276 key: "infer".into(),
277 message: format!("LLM call timed out after {} ms", timeout.as_millis()),
278 })
279 }
280 },
281 }
282}
283
284async fn await_model_task(
285 mut handle: tokio::task::JoinHandle<Result<crate::llm::LlmResponse, ActorError>>,
286 cancel: &tokio_util::sync::CancellationToken,
287) -> Result<crate::llm::LlmResponse, ActorError> {
288 tokio::select! {
289 biased;
290 _ = cancel.cancelled() => {
291 handle.abort();
292 let _ = handle.await;
293 Ok(interrupted_response())
294 }
295 joined = &mut handle => flatten_model_task(joined),
296 }
297}
298
299fn flatten_model_task(
300 joined: Result<Result<crate::llm::LlmResponse, ActorError>, tokio::task::JoinError>,
301) -> Result<crate::llm::LlmResponse, ActorError> {
302 match joined {
303 Ok(result) => result,
304 Err(error) if error.is_cancelled() => Err(ActorError::Handler {
305 key: "infer".into(),
306 message: "LLM call timed out".into(),
307 }),
308 Err(error) => Err(ActorError::Defect(format!("model task failed: {error}"))),
309 }
310}
311
312async fn read_model_response(
313 client: &dyn LlmClient,
314 messages: &[Message],
315 system: Option<&str>,
316 tools: &[ToolDefinition],
317 surface: &SessionSurface,
318 attempt: &tokio_util::sync::CancellationToken,
319) -> Result<crate::llm::LlmResponse, ActorError> {
320 let (evidence_events, usage_binding) = surface
321 .agent
322 .fact_model_evidence(messages, system, tools)
323 .await;
324 for event in evidence_events {
325 record_run_event(surface, event).await;
326 }
327 match client
328 .complete_streaming(messages, system, tools, attempt.clone())
329 .await
330 {
331 Ok(mut events) => {
332 let mut done = None;
333 loop {
334 let event = tokio::select! {
335 biased;
336 _ = surface.cancel.cancelled() => return Ok(interrupted_response()),
339 _ = attempt.cancelled() => {
340 return Err(ActorError::Handler {
341 key: "infer".into(),
342 message: "LLM call timed out".into(),
343 });
344 }
345 event = events.recv() => event,
346 };
347 let Some(event) = event else {
348 break;
349 };
350 match event {
351 crate::llm::StreamEvent::TextDelta(text) => {
352 let event = crate::agent::AgentEvent::TextDelta { text };
353 if let (Some(store), Some(run_id)) = (&surface.run_store, &surface.run_id) {
354 store.record_event(run_id, event.clone()).await;
355 }
356 if let Some(sender) = &surface.events {
357 let _ = sender.send(event).await;
358 }
359 }
360 crate::llm::StreamEvent::Done(response) => done = Some(response),
361 _ => {}
362 }
363 }
364 if let Some(response) = done {
365 record_usage(surface, &response);
366 record_model_usage_event(surface, usage_binding.as_ref(), &response).await;
367 return Ok(response);
368 }
369 if surface.cancel.is_cancelled() || attempt.is_cancelled() {
370 return Ok(interrupted_response());
371 }
372 Err(ActorError::Handler {
373 key: "infer".into(),
374 message: "stream ended before a response".into(),
375 })
376 }
377 Err(error) => {
378 if surface.cancel.is_cancelled() {
379 return Ok(interrupted_response());
380 }
381 if crate::llm::non_retryable_llm_error_message(&error).is_some() {
382 return Err(model_error(error));
383 }
384 let response = client
385 .complete(messages, system, tools)
386 .await
387 .map_err(model_error)?;
388 record_usage(surface, &response);
389 record_model_usage_event(surface, usage_binding.as_ref(), &response).await;
390 Ok(response)
391 }
392 }
393}
394
395async fn record_run_event(surface: &SessionSurface, event: crate::agent::AgentEvent) {
396 if let (Some(store), Some(run_id)) = (&surface.run_store, &surface.run_id) {
397 store.record_event(run_id, event.clone()).await;
398 }
399 if let Some(sender) = &surface.events {
400 let _ = sender.send(event).await;
401 }
402}
403
404async fn record_model_usage_event(
405 surface: &SessionSurface,
406 binding: Option<&crate::harness_evidence::ModelUsageBinding>,
407 response: &crate::llm::LlmResponse,
408) {
409 let Some(binding) = binding else {
410 return;
411 };
412 let Some(event) = crate::agent::AgentLoop::fact_model_usage_event(binding, &response.usage)
413 else {
414 return;
415 };
416 record_run_event(surface, event).await;
417}
418
419fn record_usage(surface: &SessionSurface, response: &crate::llm::LlmResponse) {
420 let mut usage = response.usage.clone();
421 if usage.total_tokens == 0 {
424 let reply = response.text();
425 if !reply.is_empty() {
426 let estimated = reply.len().div_ceil(4).max(1);
427 usage.completion_tokens = usage.completion_tokens.max(estimated);
428 usage.total_tokens = estimated;
429 }
430 }
431 surface
432 .usage
433 .lock()
434 .unwrap_or_else(std::sync::PoisonError::into_inner)
435 .accumulate(&usage);
436}
437
438async fn model_context(
442 surface: &SessionSurface,
443 prompt: &str,
444 message_count: usize,
445) -> Result<(String, Option<String>), ActorError> {
446 let work = surface
447 .agent
448 .fact_model_context(prompt, &surface.session_id, message_count);
449 let Some(limit) = surface.agent.llm_api_timeout() else {
450 return work
451 .await
452 .map_err(|error| ActorError::Defect(error.to_string()));
453 };
454 let limit = limit.min(std::time::Duration::from_secs(20));
455 match tokio::time::timeout(limit, work).await {
456 Ok(Ok(value)) => Ok(value),
457 Ok(Err(error)) => Err(ActorError::Defect(error.to_string())),
458 Err(_) => Ok((prompt.to_string(), None)),
459 }
460}
461
462async fn apply_run_controls(surface: &SessionSurface, messages: &mut Vec<Message>) {
463 let Some(control) = &surface.run_control else {
464 return;
465 };
466 let snapshot = control.snapshot().await;
467 let pending = control.drain().await;
468 if pending.is_empty() {
469 return;
470 }
471 let now_ms = std::time::SystemTime::now()
472 .duration_since(std::time::UNIX_EPOCH)
473 .map(|elapsed| u64::try_from(elapsed.as_millis()).unwrap_or(u64::MAX))
474 .unwrap_or(0);
475 for item in pending {
476 let (input, reason) = match &item.request.command {
477 crate::run_control::RunControlCommand::Steer { input } => (Some(input.clone()), None),
478 crate::run_control::RunControlCommand::Interrupt { reason, .. } => {
479 (None, reason.clone())
480 }
481 };
482 if let Some(text) = &input {
483 append_steer_fact(surface, &item.receipt.request_id, text);
484 messages.push(Message::user(text));
485 }
486 let receipt = control
487 .mark_applied(
488 &item,
489 snapshot.turn_id.clone(),
490 snapshot.turn_revision,
491 now_ms,
492 )
493 .await;
494 if receipt.state != crate::run_control::RunControlReceiptState::Applied {
495 continue;
496 }
497 record_run_event(
498 surface,
499 crate::agent::AgentEvent::RunControlApplied {
500 request_id: receipt.request_id,
501 operation: receipt.operation,
502 turn_id: receipt.turn_id,
503 turn_revision: receipt.turn_revision,
504 input,
505 reason,
506 },
507 )
508 .await;
509 }
510}
511
512fn append_steer_fact(surface: &SessionSurface, request_id: &str, text: &str) {
513 let workspace = &surface.agent.tool_context_handle().workspace;
514 let Ok(log) = FileLog::open(log_dir(workspace)) else {
515 return;
516 };
517 let _ = log.append(
518 &thread_for_session(&surface.session_id),
519 &[message_fact(format!("steer:{request_id}"), text)],
520 None,
521 );
522}
523
524fn definitions_for(catalog: &[ToolDefinition], specs: &[ToolSpec]) -> Vec<ToolDefinition> {
525 specs
526 .iter()
527 .map(|spec| {
528 catalog
529 .iter()
530 .find(|tool| tool.name == spec.name)
531 .cloned()
532 .unwrap_or(ToolDefinition {
533 name: spec.name.clone(),
534 description: spec.description.clone(),
535 parameters: serde_json::json!({
536 "type": "object",
537 "additionalProperties": true
538 }),
539 })
540 })
541 .collect()
542}
543
544fn verified_turn_text(surface: &SessionSurface, messages: &[String]) -> Option<String> {
548 let last_is_tool = messages
549 .last()
550 .is_some_and(|line| split_folded_message(line).0 == "tool");
551 if !last_is_tool {
552 return None;
553 }
554 let ledger = surface
555 .ledger
556 .lock()
557 .unwrap_or_else(std::sync::PoisonError::into_inner)
558 .clone();
559 let reports = surface
560 .reports
561 .lock()
562 .unwrap_or_else(std::sync::PoisonError::into_inner)
563 .clone();
564 match surface.agent.fact_completion_gate(&ledger, &reports) {
565 crate::harness_loop::CompletionGate::Allow(
566 crate::harness_loop::CompletionTerminal::Verified { .. }
567 | crate::harness_loop::CompletionTerminal::Waived { .. },
568 ) => Some("completed".into()),
569 _ => None,
570 }
571}
572
573fn split_folded_message(text: &str) -> (&str, String) {
574 if let Some(body) = text.strip_prefix("user\n") {
575 ("user", body.to_string())
576 } else if let Some(body) = text.strip_prefix("tool\n") {
577 ("tool", body.to_string())
578 } else {
579 ("user", text.to_string())
580 }
581}
582
583fn question_from_call(call: &crate::llm::ToolCall) -> ModelDecision {
584 let question = call
585 .args
586 .get("question")
587 .and_then(serde_json::Value::as_str)
588 .unwrap_or("")
589 .to_string();
590 let options = call
591 .args
592 .get("options")
593 .and_then(serde_json::Value::as_array)
594 .map(|items| {
595 items
596 .iter()
597 .filter_map(|item| item.as_str().map(str::to_string))
598 .collect()
599 })
600 .unwrap_or_default();
601 let allow_free_text = call
602 .args
603 .get("allow_free_text")
604 .and_then(serde_json::Value::as_bool)
605 .unwrap_or(false);
606 ModelDecision::Question {
607 question_id: if call.id.is_empty() {
608 "ask".into()
609 } else {
610 call.id.clone()
611 },
612 question,
613 allow_free_text,
614 options,
615 }
616}
617
618fn execution_permission(
624 policy: &PermissionPolicy,
625 checker: Option<&dyn crate::permissions::PermissionChecker>,
626 name: &str,
627 args: &serde_json::Value,
628) -> PermissionDecision {
629 let policy_decision = policy.check(name, args);
630 if policy_decision == PermissionDecision::Deny {
631 return PermissionDecision::Deny;
632 }
633 let Some(checker) = checker else {
634 return policy_decision;
635 };
636 if policy_decision == PermissionDecision::Allow {
637 return PermissionDecision::Allow;
638 }
639 checker.check(name, args)
640}
641
642fn skill_denial(agent: &crate::agent::AgentLoop, name: &str) -> Option<(String, String)> {
643 agent.skill_restriction_denial(name)
644}
645
646async fn confirmation_parked(
647 policy: &PermissionPolicy,
648 surface: Option<&SessionSurface>,
649 name: &str,
650 args: &serde_json::Value,
651) -> bool {
652 if let Some(surface) = surface {
653 if skill_denial(&surface.agent, name).is_some() {
654 return false;
655 }
656 }
657 let checker = surface.and_then(|surface| surface.agent.permission_checker());
658 let checker = checker.as_deref();
659 if execution_permission(policy, checker, name, args) != PermissionDecision::Ask {
660 return false;
661 }
662 match surface.and_then(|surface| surface.confirmation.as_ref()) {
663 Some(manager) => manager.requires_confirmation_for(name, args).await,
664 None => surface.is_none(),
667 }
668}
669
670fn leaked_model_calls(response: &LlmResponse) -> Vec<(String, String, serde_json::Value)> {
674 let mut calls = crate::llm::recover_leaked_tool_calls(&response.text());
675 if calls.is_empty() {
676 if let Some(reasoning) = response
677 .message
678 .reasoning_content
679 .as_deref()
680 .filter(|text| !text.is_empty())
681 {
682 calls = crate::llm::recover_leaked_tool_calls(reasoning);
683 }
684 }
685 calls
686}
687
688async fn decision_from_response(
689 response: &LlmResponse,
690 policy: &PermissionPolicy,
691 surface: Option<&SessionSurface>,
692) -> ModelDecision {
693 let structured: Vec<(String, String, serde_json::Value)> = response
694 .tool_calls()
695 .into_iter()
696 .map(|call| (call.id, call.name, call.args))
697 .collect();
698 let calls = if structured.is_empty() {
699 leaked_model_calls(response)
700 } else {
701 structured
702 };
703 if let Some((id, name, args)) = calls.iter().find(|call| call.1 == "ask_user") {
704 return question_from_call(&crate::llm::ToolCall {
705 id: id.clone(),
706 name: name.clone(),
707 args: args.clone(),
708 });
709 }
710 if let Some((id, name, args)) = calls.into_iter().next() {
711 let needs_confirmation = confirmation_parked(policy, surface, &name, &args).await;
712 let text = crate::llm::strip_leaked_tool_protocol(&response.text());
713 let text = if text.trim().is_empty() {
714 None
715 } else {
716 Some(text)
717 };
718 let reasoning = response
719 .message
720 .reasoning_content
721 .as_ref()
722 .map(|text| crate::llm::strip_leaked_tool_protocol(text))
723 .filter(|text| !text.trim().is_empty());
724 return ModelDecision::Tool {
725 call: ToolCall {
726 id: if id.is_empty() { "tool".into() } else { id },
727 needs_confirmation,
728 name,
729 args,
730 text,
731 reasoning,
732 },
733 };
734 }
735 ModelDecision::Text {
736 text: response.text(),
737 }
738}
739
740impl Completion for LiveCompletion {
741 fn complete(
742 &self,
743 request: CompletionRequest,
744 ) -> a3s_effect::coding::BoxFuture<Result<ModelDecision, ActorError>> {
745 self.calls.fetch_add(1, Ordering::SeqCst);
746 let client = Arc::clone(&self.client);
747 let policy = self.policy.clone();
748 let catalog = self.catalog.clone();
749 let surface = self.surface.clone();
750 let ready_flag = Arc::clone(&self.context_ready);
751 let system_slot = Arc::clone(&self.cached_system);
752 let prompt_slot = Arc::clone(&self.cached_prompt);
753 Box::pin(async move {
754 if request.messages.len() == 1 {
755 if let Some(surface) = &surface {
756 if let Some(args) = surface
757 .agent
758 .fact_auto_delegation_args(&split_folded_message(&request.messages[0]).1)
759 {
760 let needs_confirmation =
761 confirmation_parked(&policy, Some(surface), "task", &args).await;
762 return Ok(ModelDecision::Tool {
763 call: ToolCall {
764 id: "auto-task".into(),
765 name: "task".into(),
766 args,
767 needs_confirmation,
768 text: None,
769 reasoning: None,
770 },
771 });
772 }
773 }
774 }
775 if let Some(surface) = &surface {
776 if let Some(text) = verified_turn_text(surface, &request.messages) {
777 return Ok(ModelDecision::Text { text });
778 }
779 }
780 let mut system = if request.system.is_empty() {
781 None
782 } else {
783 Some(request.system.join("\n"))
784 };
785 let mut prompt_override = None;
786 if let Some(surface) = &surface {
787 if ready_flag
788 .compare_exchange(false, true, Ordering::SeqCst, Ordering::SeqCst)
789 .is_ok()
790 {
791 let prompt = request
792 .messages
793 .first()
794 .map(|text| split_folded_message(text).1)
795 .unwrap_or_default();
796 let (effective, augmented) =
797 model_context(surface, &prompt, request.messages.len()).await?;
798 prompt_override = Some(effective);
799 if let Some(augmented) = augmented.filter(|text| !text.is_empty()) {
800 system = Some(match system {
801 Some(base) => format!("{base}\n{augmented}"),
802 None => augmented,
803 });
804 }
805 *prompt_slot
806 .lock()
807 .unwrap_or_else(std::sync::PoisonError::into_inner) =
808 prompt_override.clone();
809 *system_slot
810 .lock()
811 .unwrap_or_else(std::sync::PoisonError::into_inner) = system.clone();
812 } else {
813 prompt_override = prompt_slot
814 .lock()
815 .unwrap_or_else(std::sync::PoisonError::into_inner)
816 .clone();
817 if let Some(augmented) = system_slot
818 .lock()
819 .unwrap_or_else(std::sync::PoisonError::into_inner)
820 .clone()
821 {
822 system = Some(augmented);
823 }
824 }
825 }
826 if !request.summary.is_empty() {
827 if let Some(surface) = &surface {
828 if let Some(sender) = &surface.events {
829 let _ = sender
830 .send(crate::agent::AgentEvent::ContextCompacted {
831 session_id: surface.session_id.clone(),
832 before_messages: request.messages.len().saturating_add(1),
833 after_messages: 1,
834 percent_before: 1.0,
835 summary: Some(request.summary.clone()),
836 })
837 .await;
838 }
839 }
840 }
841 let mut folded = Vec::new();
842 if !request.summary.is_empty() {
843 folded.push(format!("user\n{}", request.summary));
844 }
845 folded.extend(request.messages.iter().cloned());
846 let prompt_index = usize::from(!request.summary.is_empty());
847 let mut recorded_calls = surface
848 .as_ref()
849 .map(|surface| {
850 recorded_tool_calls(
851 &surface.agent.tool_context_handle().workspace,
852 &thread_for_session(&surface.session_id),
853 )
854 })
855 .unwrap_or_default()
856 .into_iter();
857 let mut messages: Vec<Message> = folded
858 .iter()
859 .enumerate()
860 .flat_map(|(index, text)| {
861 let (role, body) = split_folded_message(text);
862 let body = if index == prompt_index {
863 prompt_override.clone().unwrap_or(body)
864 } else {
865 body
866 };
867 if role == "tool" {
868 let call = recorded_calls.next().unwrap_or_else(|| RecordedCall {
869 id: format!("fact-tool-{index}"),
870 name: "tool".into(),
871 args: serde_json::json!({}),
872 text: String::new(),
873 reasoning: None,
874 });
875 let id = call.id.clone();
876 vec![
877 assistant_tool_call(&call),
878 Message::tool_result(&id, &body, false),
879 ]
880 } else {
881 vec![Message::user(&body)]
882 }
883 })
884 .collect();
885 if let Some(surface) = &surface {
886 apply_run_controls(surface, &mut messages).await;
887 if surface.cancel.is_cancelled() {
888 return Ok(ModelDecision::Text {
889 text: "(Response interrupted by the user.)".into(),
890 });
891 }
892 }
893 let tools = definitions_for(&catalog, &request.tools);
894 let response = model_response(
895 Arc::clone(&client),
896 messages.clone(),
897 system.clone(),
898 tools,
899 surface.clone(),
900 )
901 .await?;
902 if let Some(surface) = &surface {
903 if surface.cancel.is_cancelled() {
904 apply_run_controls(surface, &mut messages).await;
905 let text = "(Response interrupted by the user.)";
906 messages.push(Message::assistant(text));
907 *surface
908 .transcript
909 .lock()
910 .unwrap_or_else(std::sync::PoisonError::into_inner) = messages;
911 return Ok(ModelDecision::Text {
912 text: text.to_string(),
913 });
914 }
915 let prompt = request
916 .messages
917 .first()
918 .map(|text| split_folded_message(text).1)
919 .unwrap_or_default();
920 let prompt = prompt.as_str();
921 surface
922 .agent
923 .fact_observe_model(&surface.session_id, prompt, &response)
924 .await;
925 let mut transcript = messages.clone();
926 if response.tool_calls().is_empty() {
927 transcript.push(Message::assistant(&response.text()));
928 }
929 *surface
930 .transcript
931 .lock()
932 .unwrap_or_else(std::sync::PoisonError::into_inner) = transcript;
933 if response.tool_calls().is_empty() {
934 surface
935 .agent
936 .fact_post_response(
937 &surface.session_id,
938 &response.text(),
939 0,
940 &response.usage,
941 )
942 .await;
943 }
944 }
945 Ok(decision_from_response(&response, &policy, surface.as_ref()).await)
946 })
947 }
948}
949
950struct ExecutorTools {
951 executor: Arc<ToolExecutor>,
952 calls: Arc<AtomicUsize>,
953 checkpoint: Option<CheckpointState>,
954 events: Option<tokio::sync::mpsc::Sender<crate::agent::AgentEvent>>,
955 context: crate::tools::ToolContext,
956 permission: PermissionPolicy,
957 checker: Option<Arc<dyn crate::permissions::PermissionChecker>>,
958 agent: Option<crate::agent::AgentLoop>,
959 run_store: Option<Arc<crate::run::InMemoryRunStore>>,
960 run_id: Option<String>,
961 ledger: Option<Arc<Mutex<crate::harness_loop::MutationLedger>>>,
962 reports: Option<Arc<Mutex<Vec<crate::verification::VerificationReport>>>>,
963}
964
965#[derive(Clone)]
966struct CheckpointState {
967 sink: Arc<dyn crate::loop_checkpoint::LoopCheckpointSink>,
968 run_id: String,
969 session_id: String,
970 capability_binding: Option<crate::capability::RunCapabilityBindingV1>,
971 turns: Arc<AtomicUsize>,
972}
973
974impl ExecutorTools {
975 async fn emit(&self, event: crate::agent::AgentEvent) {
976 if let (Some(store), Some(run_id)) = (&self.run_store, &self.run_id) {
977 store.record_event(run_id, event.clone()).await;
978 }
979 if let Some(events) = &self.events {
980 let _ = events.send(event).await;
981 }
982 }
983
984 async fn forward_host_events(
985 &self,
986 bridged: &mut Option<tokio::sync::broadcast::Receiver<crate::agent::AgentEvent>>,
987 ) {
988 let Some(receiver) = bridged.as_mut() else {
989 return;
990 };
991 loop {
992 match receiver.try_recv() {
993 Ok(event) if host_bridge_event(&event) => self.emit(event).await,
994 Ok(_) => {}
995 Err(tokio::sync::broadcast::error::TryRecvError::Lagged(_)) => {}
996 Err(_) => break,
997 }
998 }
999 }
1000
1001 async fn save_checkpoint(&self, call: &ToolCall, output: &str, is_error: bool) {
1002 let Some(checkpoint) = &self.checkpoint else {
1003 return;
1004 };
1005 let turn = checkpoint.turns.fetch_add(1, Ordering::SeqCst) + 1;
1006 let checkpoint_ms = std::time::SystemTime::now()
1007 .duration_since(std::time::UNIX_EPOCH)
1008 .map(|elapsed| u64::try_from(elapsed.as_millis()).unwrap_or(u64::MAX))
1009 .unwrap_or(0);
1010 checkpoint
1011 .sink
1012 .save_checkpoint(&crate::loop_checkpoint::LoopCheckpoint {
1013 schema_version: crate::loop_checkpoint::LOOP_CHECKPOINT_SCHEMA_VERSION,
1014 run_id: checkpoint.run_id.clone(),
1015 session_id: checkpoint.session_id.clone(),
1016 capability_binding: checkpoint.capability_binding.clone(),
1017 turn,
1018 messages: vec![
1019 Message {
1020 role: "assistant".into(),
1021 content: vec![crate::llm::ContentBlock::ToolUse {
1022 id: call.id.clone(),
1023 name: call.name.clone(),
1024 input: call.args.clone(),
1025 }],
1026 reasoning_content: None,
1027 transcript_text: None,
1028 transcript_visibility: Default::default(),
1029 },
1030 Message::tool_result(&call.id, output, is_error),
1031 ],
1032 total_usage: crate::llm::TokenUsage::default(),
1033 tool_calls_count: turn,
1034 verification_reports: Vec::new(),
1035 convergence: crate::loop_checkpoint::LoopConvergenceState::default(),
1036 checkpoint_ms,
1037 })
1038 .await;
1039 }
1040}
1041
1042fn host_bridge_event(event: &crate::agent::AgentEvent) -> bool {
1043 matches!(
1044 event,
1045 crate::agent::AgentEvent::SubagentStart { .. }
1046 | crate::agent::AgentEvent::SubagentProgress { .. }
1047 | crate::agent::AgentEvent::SubagentEnd { .. }
1048 | crate::agent::AgentEvent::TaskUpdated { .. }
1049 | crate::agent::AgentEvent::ConfirmationRequired { .. }
1050 | crate::agent::AgentEvent::ConfirmationReceived { .. }
1051 | crate::agent::AgentEvent::ConfirmationTimeout { .. }
1052 | crate::agent::AgentEvent::UserQuestion { .. }
1053 )
1054}
1055
1056impl ToolRunner for ExecutorTools {
1057 fn run(
1058 &self,
1059 call: ToolCall,
1060 ) -> a3s_effect::coding::BoxFuture<Result<serde_json::Value, ActorError>> {
1061 self.calls.fetch_add(1, Ordering::SeqCst);
1062 let executor = Arc::clone(&self.executor);
1063 let context = self.context.clone();
1064 let events = self.events.clone();
1065 let permission = self.permission.clone();
1066 let checker = self.checker.clone();
1067 let agent = self.agent.clone();
1068 let checkpoint = self.checkpoint.clone();
1069 let run_store = self.run_store.clone();
1070 let run_id = self.run_id.clone();
1071 let ledger = self.ledger.clone();
1072 let reports = self.reports.clone();
1073 Box::pin(async move {
1074 let tools = ExecutorTools {
1075 executor,
1076 calls: Arc::new(AtomicUsize::new(0)),
1077 checkpoint,
1078 events,
1079 context,
1080 permission,
1081 checker,
1082 agent,
1083 run_store,
1084 run_id,
1085 ledger,
1086 reports,
1087 };
1088 emit_tool_request_bound(&tools, &call).await;
1089 if let Some(agent) = &tools.agent {
1090 if let Some((output, reason)) = skill_denial(agent, &call.name) {
1091 tools
1092 .emit(crate::agent::AgentEvent::PermissionDenied {
1093 tool_id: call.id.clone(),
1094 tool_name: call.name.clone(),
1095 args: call.args.clone(),
1096 reason,
1097 })
1098 .await;
1099 tools.save_checkpoint(&call, &output, true).await;
1100 return Ok(serde_json::Value::String(output));
1101 }
1102 }
1103 if execution_permission(
1104 &tools.permission,
1105 tools.checker.as_deref(),
1106 &call.name,
1107 &call.args,
1108 ) == PermissionDecision::Deny
1109 {
1110 let output = format!(
1111 "Permission denied: Tool '{}' is blocked by permission policy.",
1112 call.name
1113 );
1114 tools
1115 .emit(crate::agent::AgentEvent::PermissionDenied {
1116 tool_id: call.id.clone(),
1117 tool_name: call.name.clone(),
1118 args: call.args.clone(),
1119 reason: "Blocked by deny rule in permission policy".into(),
1120 })
1121 .await;
1122 tools.save_checkpoint(&call, &output, true).await;
1123 return Ok(serde_json::Value::String(output));
1124 }
1125 tools
1126 .emit(crate::agent::AgentEvent::ToolStart {
1127 id: call.id.clone(),
1128 name: call.name.clone(),
1129 })
1130 .await;
1131 let mut args = call.args.clone();
1132 if let Some(agent) = &tools.agent {
1133 let session_id = tools.context.session_id.clone().unwrap_or_default();
1134 let decision = agent
1135 .fire_pre_tool_use(&session_id, &call.name, &args, Vec::new())
1136 .await;
1137 if let Some(denial) = decision.denial {
1138 let output = format!("Hook denied {}: {}", call.name, denial.reason);
1139 tools
1140 .emit(crate::agent::AgentEvent::PermissionDenied {
1141 tool_id: call.id.clone(),
1142 tool_name: call.name.clone(),
1143 args: args.clone(),
1144 reason: denial.reason,
1145 })
1146 .await;
1147 return Ok(serde_json::Value::String(output));
1148 }
1149 if let Some(updated) = decision.updated_args {
1150 args = updated;
1151 }
1152 }
1153 tools
1154 .emit(crate::agent::AgentEvent::ToolExecutionStart {
1155 id: call.id.clone(),
1156 name: call.name.clone(),
1157 args: args.clone(),
1158 })
1159 .await;
1160 let mut bridged = tools
1161 .context
1162 .agent_event_tx
1163 .as_ref()
1164 .map(|sender| sender.subscribe());
1165 let cancel = tools.context.cancellation_token();
1166 let started = std::time::Instant::now();
1167 let executed = tokio::select! {
1168 biased;
1169 _ = cancel.cancelled() => {
1170 let message = "cancelled".to_string();
1171 tools
1172 .emit(crate::agent::AgentEvent::ToolEnd {
1173 id: call.id.clone(),
1174 name: call.name.clone(),
1175 args: Some(args.clone()),
1176 output: message.clone(),
1177 exit_code: 1,
1178 metadata: None,
1179 error_kind: None,
1180 })
1181 .await;
1182 return Ok(serde_json::Value::String(message));
1183 }
1184 result = tools
1185 .executor
1186 .execute_with_context(&call.name, &args, &tools.context) => result,
1187 };
1188 tools.forward_host_events(&mut bridged).await;
1189 let result = match executed {
1190 Ok(result) => result,
1191 Err(error) => {
1192 let message = visible_tool_output(tools.agent.as_ref(), error.to_string());
1193 tools
1194 .emit(crate::agent::AgentEvent::ToolEnd {
1195 id: call.id.clone(),
1196 name: call.name.clone(),
1197 args: Some(call.args.clone()),
1198 output: message.clone(),
1199 exit_code: 1,
1200 metadata: None,
1201 error_kind: None,
1202 })
1203 .await;
1204 if workspace_boundary_error(&message) {
1205 return Ok(serde_json::Value::String(message));
1206 }
1207 return Err(ActorError::Handler {
1208 key: call.name.clone(),
1209 message,
1210 });
1211 }
1212 };
1213 let mut output = visible_tool_output(tools.agent.as_ref(), result.output);
1214 if let Some(agent) = &tools.agent {
1215 let session_id = tools.context.session_id.clone().unwrap_or_default();
1216 agent
1217 .fire_post_tool_use(
1218 &session_id,
1219 &call.name,
1220 &args,
1221 &output,
1222 result.exit_code == 0,
1223 u64::try_from(started.elapsed().as_millis()).unwrap_or(u64::MAX),
1224 )
1225 .await;
1226 }
1227 let mut metadata = result.metadata;
1228 if let Some(agent) = &tools.agent {
1229 let mut bound = tools
1230 .reports
1231 .as_ref()
1232 .map(|reports| {
1233 reports
1234 .lock()
1235 .unwrap_or_else(std::sync::PoisonError::into_inner)
1236 .clone()
1237 })
1238 .unwrap_or_default();
1239 if let Some(ledger) = &tools.ledger {
1240 let mut ledger_value = ledger
1241 .lock()
1242 .unwrap_or_else(std::sync::PoisonError::into_inner)
1243 .clone();
1244 agent
1245 .record_fact_tool_effect(
1246 &call.name,
1247 &call.args,
1248 result.exit_code,
1249 &mut output,
1250 &mut metadata,
1251 &mut ledger_value,
1252 &mut bound,
1253 )
1254 .await;
1255 *ledger
1256 .lock()
1257 .unwrap_or_else(std::sync::PoisonError::into_inner) = ledger_value;
1258 }
1259 if let Some(reports) = &tools.reports {
1260 *reports
1261 .lock()
1262 .unwrap_or_else(std::sync::PoisonError::into_inner) = bound;
1263 }
1264 } else if let Some(ledger) = &tools.ledger {
1265 ledger
1266 .lock()
1267 .unwrap_or_else(std::sync::PoisonError::into_inner)
1268 .observe_tool(&call.name, result.exit_code, metadata.as_ref());
1269 }
1270 tools
1271 .emit(crate::agent::AgentEvent::ToolEnd {
1272 id: call.id.clone(),
1273 name: call.name.clone(),
1274 args: Some(call.args.clone()),
1275 output: output.clone(),
1276 exit_code: result.exit_code,
1277 metadata: metadata.clone(),
1278 error_kind: result.error_kind.clone(),
1279 })
1280 .await;
1281 tools
1282 .save_checkpoint(&call, &output, result.exit_code != 0)
1283 .await;
1284 Ok(serde_json::Value::String(output))
1285 })
1286 }
1287}
1288
1289struct NoopCompact;
1290
1291impl Compactor for NoopCompact {
1292 fn compact(
1293 &self,
1294 messages: &[String],
1295 ) -> a3s_effect::coding::BoxFuture<Result<String, ActorError>> {
1296 let summary = messages.join(" ");
1297 Box::pin(async move { Ok(summary) })
1298 }
1299}
1300
1301struct RecordedCall {
1302 id: String,
1303 name: String,
1304 args: serde_json::Value,
1305 text: String,
1306 reasoning: Option<String>,
1307}
1308
1309fn seeded_user_text(message: &Message) -> String {
1310 let tool_text = message
1311 .content
1312 .iter()
1313 .filter_map(|block| {
1314 if let crate::llm::ContentBlock::ToolResult { content, .. } = block {
1315 Some(content.as_text())
1316 } else {
1317 None
1318 }
1319 })
1320 .collect::<Vec<_>>()
1321 .join("");
1322 if !tool_text.is_empty() {
1323 return tool_text;
1324 }
1325 message.text()
1326}
1327
1328fn assistant_tool_call(call: &RecordedCall) -> Message {
1329 let mut content = Vec::new();
1330 if !call.text.is_empty() {
1331 content.push(crate::llm::ContentBlock::Text {
1332 text: call.text.clone(),
1333 });
1334 }
1335 content.push(crate::llm::ContentBlock::ToolUse {
1336 id: call.id.clone(),
1337 name: call.name.clone(),
1338 input: call.args.clone(),
1339 });
1340 Message {
1341 role: "assistant".into(),
1342 content,
1343 reasoning_content: call.reasoning.clone(),
1344 transcript_text: None,
1345 transcript_visibility: crate::llm::TranscriptVisibility::Wire,
1346 }
1347}
1348
1349fn recorded_tool_calls(workspace: &Path, thread: &str) -> Vec<RecordedCall> {
1356 let Ok(facts) = read_workspace_facts(workspace, thread) else {
1357 return Vec::new();
1358 };
1359 let start = facts
1360 .iter()
1361 .rposition(|fact| fact.kind == "compaction.done")
1362 .map(|index| index + 1)
1363 .unwrap_or(0);
1364 let mut pending = None;
1365 let mut calls = Vec::new();
1366 for fact in facts.into_iter().skip(start) {
1367 if fact.kind == "model.turn"
1368 && fact.payload.get("kind").and_then(|kind| kind.as_str()) == Some("tool")
1369 {
1370 let call = fact.payload.get("call");
1371 let id = call
1372 .and_then(|value| value.get("id"))
1373 .and_then(|value| value.as_str())
1374 .unwrap_or("tool");
1375 let name = call
1376 .and_then(|value| value.get("name"))
1377 .and_then(|value| value.as_str())
1378 .unwrap_or("tool");
1379 let args = call
1380 .and_then(|value| value.get("args"))
1381 .cloned()
1382 .unwrap_or_else(|| serde_json::json!({}));
1383 let text = call
1384 .and_then(|value| value.get("text"))
1385 .and_then(|value| value.as_str())
1386 .unwrap_or("")
1387 .to_string();
1388 let reasoning = call
1389 .and_then(|value| value.get("reasoning"))
1390 .and_then(|value| value.as_str())
1391 .filter(|text| !text.is_empty())
1392 .map(str::to_string);
1393 pending = Some(RecordedCall {
1394 id: id.to_string(),
1395 name: name.to_string(),
1396 args,
1397 text,
1398 reasoning,
1399 });
1400 } else if fact.kind == "tool.result" {
1401 if let Some(call) = pending.take() {
1402 calls.push(call);
1403 }
1404 }
1405 }
1406 calls
1407}
1408
1409fn workspace_boundary_error(message: &str) -> bool {
1410 message.contains("escapes workspace") || message.contains("Workspace boundary")
1411}
1412
1413fn visible_tool_output(agent: Option<&crate::agent::AgentLoop>, output: String) -> String {
1414 match agent {
1415 Some(agent) => agent.sanitize_tool_output(&output),
1416 None => output,
1417 }
1418}
1419
1420async fn emit_tool_request_bound(tools: &ExecutorTools, call: &a3s_effect::ToolCall) {
1421 let Ok(snapshot) = crate::harness_evidence::ToolRequestSnapshotV1::capture(
1422 &call.id,
1423 &call.name,
1424 &call.args,
1425 crate::harness_evidence::ToolRequestOriginV1::Agent,
1426 ) else {
1427 return;
1428 };
1429 tools
1430 .emit(crate::agent::AgentEvent::ToolRequestBound {
1431 tool_id: call.id.clone(),
1432 tool_name: call.name.clone(),
1433 snapshot,
1434 })
1435 .await;
1436}
1437
1438fn log_dir(workspace: &Path) -> PathBuf {
1439 workspace.join(".a3s").join("effect-log")
1440}
1441
1442#[cfg(test)]
1444pub(crate) fn reset_session_fact_log(workspace: impl AsRef<Path>) {
1445 let _ = std::fs::remove_dir_all(log_dir(workspace.as_ref()));
1446}
1447
1448impl FactRun {
1449 pub fn open(
1450 dir: impl Into<PathBuf>,
1451 completion: Arc<dyn Completion>,
1452 tools: Arc<dyn ToolRunner>,
1453 config: HarnessConfig,
1454 run_id: impl Into<String>,
1455 ) -> Result<Self> {
1456 Self::open_composed(dir, completion, tools, config, run_id, None)
1457 }
1458
1459 pub fn open_composed(
1461 dir: impl Into<PathBuf>,
1462 completion: Arc<dyn Completion>,
1463 tools: Arc<dyn ToolRunner>,
1464 config: HarnessConfig,
1465 run_id: impl Into<String>,
1466 compose: Option<&crate::meta_harness::HarnessComposeOptions>,
1467 ) -> Result<Self> {
1468 Self::open_composed_with_hosts(
1469 dir,
1470 completion,
1471 tools,
1472 config,
1473 run_id,
1474 FactRunHostMounts {
1475 compose,
1476 registry: None,
1477 assembler: None,
1478 },
1479 )
1480 }
1481
1482 pub fn open_composed_with_hosts(
1484 dir: impl Into<PathBuf>,
1485 completion: Arc<dyn Completion>,
1486 tools: Arc<dyn ToolRunner>,
1487 config: HarnessConfig,
1488 run_id: impl Into<String>,
1489 hosts: FactRunHostMounts<'_>,
1490 ) -> Result<Self> {
1491 let step_limit = config.step_limit();
1492 let graph = if let Some(assembler) = hosts.assembler {
1493 let graph = assembler.assemble(config)?;
1494 let _kernel = crate::meta_harness::KernelPolicy::default().admit();
1495 graph
1496 } else {
1497 let (graph, _kernel) = crate::meta_harness::admit_from_compose_with_registry(
1498 hosts.compose,
1499 hosts.registry,
1500 config,
1501 )?;
1502 graph
1503 };
1504 Self::open_with_graph(dir, completion, tools, step_limit, run_id, graph)
1505 }
1506
1507 pub fn open_with_graph(
1510 dir: impl Into<PathBuf>,
1511 completion: Arc<dyn Completion>,
1512 tools: Arc<dyn ToolRunner>,
1513 step_limit: u32,
1514 run_id: impl Into<String>,
1515 graph: HarnessGraph,
1516 ) -> Result<Self> {
1517 let _kernel = crate::meta_harness::KernelPolicy::default().admit();
1518 let dir = dir.into();
1519 let capped = Arc::new(CappedCompletion {
1520 inner: completion,
1521 run_id: run_id.into(),
1522 });
1523 let services = Arc::new(CodingServices {
1524 completion: capped,
1525 tools,
1526 compactor: Arc::new(NoopCompact),
1527 });
1528 Ok(Self {
1529 dir,
1530 thread: THREAD.to_string(),
1531 actor: graph.into_actor(),
1532 services,
1533 limit: step_limit,
1534 })
1535 }
1536
1537 pub fn open_session(
1538 workspace: &Path,
1539 session_id: &str,
1540 client: Arc<dyn LlmClient>,
1541 executor: Arc<ToolExecutor>,
1542 permission: PermissionPolicy,
1543 yolo_lanes: &[SessionLane],
1544 max_tool_rounds: usize,
1545 ) -> Result<Self> {
1546 let catalog = executor.definitions();
1547 let specs = catalog
1548 .iter()
1549 .map(|tool| ToolSpec {
1550 name: tool.name.clone(),
1551 description: tool.description.clone(),
1552 })
1553 .collect();
1554 let permission = permission.allow_yolo_lanes(yolo_lanes.iter().copied());
1555 let (completion, _) = LiveCompletion::new(client, permission.clone(), catalog);
1556 let cap = u32::try_from(max_tool_rounds).unwrap_or(u32::MAX);
1557 let config = HarnessConfig::new(8, 1_000_000, 32, 2, Vec::new(), specs)
1558 .map_err(|error| anyhow::anyhow!(error))?
1559 .with_tool_round_cap(cap);
1560 let context = executor.registry().context();
1561 let mut run = Self::open(
1562 log_dir(workspace),
1563 Arc::new(completion),
1564 Arc::new(ExecutorTools {
1565 executor,
1566 calls: Arc::new(AtomicUsize::new(0)),
1567 checkpoint: None,
1568 events: None,
1569 context,
1570 permission: permission.clone(),
1571 checker: None,
1572 agent: None,
1573 run_store: None,
1574 run_id: None,
1575 ledger: None,
1576 reports: None,
1577 }),
1578 config,
1579 session_id,
1580 )?;
1581 run.thread = thread_for_session(session_id);
1582 Ok(run)
1583 }
1584
1585 #[allow(clippy::too_many_arguments)]
1586 pub(crate) fn open_projected(
1587 workspace: &Path,
1588 session_id: &str,
1589 client: Arc<dyn LlmClient>,
1590 executor: Arc<ToolExecutor>,
1591 permission: PermissionPolicy,
1592 yolo_lanes: &[SessionLane],
1593 max_tool_rounds: usize,
1594 surface: SessionSurface,
1595 ) -> Result<Self> {
1596 let catalog = executor.definitions();
1597 let specs = catalog
1598 .iter()
1599 .map(|tool| ToolSpec {
1600 name: tool.name.clone(),
1601 description: tool.description.clone(),
1602 })
1603 .collect();
1604 let permission = permission.allow_yolo_lanes(yolo_lanes.iter().copied());
1605 let (completion, _) = LiveCompletion::new(client, permission.clone(), catalog);
1606 let events = surface.events.clone();
1607 let run_store = surface.run_store.clone();
1608 let run_id = surface.run_id.clone();
1609 let ledger = Some(Arc::clone(&surface.ledger));
1610 let reports = Some(Arc::clone(&surface.reports));
1611 let checker = surface.agent.permission_checker();
1612 let agent = surface.agent.clone();
1613 let context = surface
1614 .agent
1615 .tool_context_handle()
1616 .with_cancellation(surface.cancel.clone());
1617 let completion = completion.with_surface(surface.clone());
1618 let cap = u32::try_from(max_tool_rounds).unwrap_or(u32::MAX);
1619 let config = HarnessConfig::new(
1620 8,
1621 agent.fact_compact_after_chars(),
1622 32,
1623 2,
1624 Vec::new(),
1625 specs,
1626 )
1627 .map_err(|error| anyhow::anyhow!(error))?
1628 .with_tool_round_cap(cap);
1629 let harness = surface.harness.clone();
1630 let host_registry = surface.host_harness_registry.clone();
1631 let host_assembler = surface.host_harness_assembler.clone();
1632 let checkpoint = surface.checkpoint.map(|checkpoint| CheckpointState {
1633 sink: checkpoint.sink,
1634 run_id: checkpoint.run_id,
1635 session_id: checkpoint.session_id,
1636 capability_binding: checkpoint.capability_binding,
1637 turns: Arc::new(AtomicUsize::new(0)),
1638 });
1639 let mut run = Self::open_composed_with_hosts(
1640 log_dir(workspace),
1641 Arc::new(completion),
1642 Arc::new(ExecutorTools {
1643 executor,
1644 calls: Arc::new(AtomicUsize::new(0)),
1645 checkpoint,
1646 events,
1647 context,
1648 permission,
1649 checker,
1650 agent: Some(agent),
1651 run_store,
1652 run_id,
1653 ledger,
1654 reports,
1655 }),
1656 config,
1657 session_id,
1658 FactRunHostMounts {
1659 compose: harness.as_ref(),
1660 registry: host_registry.as_deref(),
1661 assembler: host_assembler.as_deref(),
1662 },
1663 )?;
1664 run.thread = thread_for_session(session_id);
1665 Ok(run)
1666 }
1667
1668 fn open_log(&self) -> Result<FileLog> {
1669 FileLog::open(&self.dir).map_err(|error| anyhow::anyhow!(error))
1670 }
1671
1672 pub fn read_facts(&self) -> Result<Vec<a3s_effect::Fact>> {
1673 let log = self.open_log()?;
1674 log.read(&self.thread)
1675 .map_err(|error| anyhow::anyhow!(error))
1676 }
1677
1678 pub fn reset_thread_log(&self) -> Result<()> {
1680 let path = self.dir.join(format!("{}.jsonl", self.thread));
1681 if path.exists() {
1682 std::fs::remove_file(&path).map_err(|error| anyhow::anyhow!(error))?;
1683 }
1684 Ok(())
1685 }
1686
1687 pub fn seed_history(&self, history: &[Message]) -> Result<()> {
1692 let log = self.open_log()?;
1693 if !log
1694 .read(&self.thread)
1695 .map_err(|error| anyhow::anyhow!(error))?
1696 .is_empty()
1697 {
1698 return Ok(());
1699 }
1700 let mut turn = 0u64;
1701 let mut cycle = 0u64;
1702 for message in history {
1703 match message.role.as_str() {
1704 "user" => {
1705 turn += 1;
1706 cycle = 0;
1707 log.append(
1708 &self.thread,
1709 &[message_fact(
1710 format!("m-hist-{turn}"),
1711 seeded_user_text(message),
1712 )],
1713 None,
1714 )
1715 .map_err(|error| anyhow::anyhow!(error))?;
1716 }
1717 "assistant" => {
1718 let payload = serde_json::to_value(ModelDecision::Text {
1719 text: message.text(),
1720 })
1721 .map_err(|error| anyhow::anyhow!(error))?;
1722 log.append(
1723 &self.thread,
1724 &[NewFact {
1725 kind: "model.turn".into(),
1726 key: format!("model:{turn}:{cycle}"),
1727 payload,
1728 }],
1729 Some(&format!("infer:{turn}:{cycle}")),
1730 )
1731 .map_err(|error| anyhow::anyhow!(error))?;
1732 cycle += 1;
1733 }
1734 _ => {}
1735 }
1736 }
1737 Ok(())
1738 }
1739
1740 pub async fn user_text(&self, text: &str) -> Result<a3s_effect::Settlement<CodingView>> {
1741 let log = self.open_log()?;
1742 let key = format!("m-{}", log.read(&self.thread).unwrap_or_default().len());
1743 ingest_coding(
1744 &self.actor,
1745 &log,
1746 Arc::clone(&self.services),
1747 &self.thread,
1748 message_fact(key, text),
1749 self.limit,
1750 )
1751 .await
1752 .map_err(exit_to_error)
1753 }
1754
1755 pub async fn steer(&self, text: &str) -> Result<a3s_effect::Settlement<CodingView>> {
1756 self.user_text(text).await
1757 }
1758
1759 pub async fn resume_limit(&self, limit: u32) -> Result<a3s_effect::Settlement<CodingView>> {
1760 let log = self.open_log()?;
1761 resume_coding(
1762 &self.actor,
1763 &log,
1764 Arc::clone(&self.services),
1765 &self.thread,
1766 limit,
1767 )
1768 .await
1769 .map_err(exit_to_error)
1770 }
1771
1772 pub async fn confirm(
1773 &self,
1774 tool_call_id: &str,
1775 approved: bool,
1776 ) -> Result<a3s_effect::Settlement<CodingView>> {
1777 let log = self.open_log()?;
1778 let key = format!("c-{}", log.read(&self.thread).unwrap_or_default().len());
1779 ingest_coding(
1780 &self.actor,
1781 &log,
1782 Arc::clone(&self.services),
1783 &self.thread,
1784 confirm_fact(key, tool_call_id, approved),
1785 self.limit,
1786 )
1787 .await
1788 .map_err(exit_to_error)
1789 }
1790
1791 pub async fn confirm_if_pending(&self, tool_call_id: &str, approved: bool) -> Result<bool> {
1793 let pending = match self.resume_limit(0).await {
1794 Ok(settled) => {
1795 settled.view.phase == CodingPhase::Confirm
1796 && settled
1797 .view
1798 .pending_confirmation
1799 .as_ref()
1800 .is_some_and(|pending| pending.tool_call_id == tool_call_id)
1801 }
1802 Err(error) => {
1803 let rendered = format!("{error:#}");
1804 if rendered.contains("StepLimit") {
1805 false
1806 } else {
1807 return Err(error);
1808 }
1809 }
1810 };
1811 if !pending {
1812 return Ok(false);
1813 }
1814 let before = self.read_facts()?.len();
1815 let settled = self.confirm(tool_call_id, approved).await?;
1816 Ok(settled.log.len() > before)
1817 }
1818
1819 pub async fn answer(&self, text: &str) -> Result<a3s_effect::Settlement<CodingView>> {
1820 let log = self.open_log()?;
1821 let key = format!("a-{}", log.read(&self.thread).unwrap_or_default().len());
1822 ingest_coding(
1823 &self.actor,
1824 &log,
1825 Arc::clone(&self.services),
1826 &self.thread,
1827 answer_fact(key, text),
1828 self.limit,
1829 )
1830 .await
1831 .map_err(exit_to_error)
1832 }
1833
1834 pub async fn append_model_turn(&self, payload: serde_json::Value) -> Result<()> {
1835 let log = self.open_log()?;
1836 log.append(
1837 &self.thread,
1838 &[NewFact {
1839 kind: "model.turn".into(),
1840 key: "model:bad".into(),
1841 payload,
1842 }],
1843 None,
1844 )
1845 .map_err(|error| anyhow::anyhow!(error))?;
1846 Ok(())
1847 }
1848}
1849
1850pub fn read_workspace_facts(workspace: &Path, thread: &str) -> Result<Vec<a3s_effect::Fact>> {
1851 let log = FileLog::open(log_dir(workspace)).map_err(|error| anyhow::anyhow!(error))?;
1852 log.read(thread).map_err(|error| anyhow::anyhow!(error))
1853}
1854
1855fn exit_to_error(error: Exit<ActorError>) -> anyhow::Error {
1856 let message = match error {
1857 Exit::Die(message) => message,
1858 Exit::Interrupt => "cancelled".to_string(),
1859 Exit::Fail(ActorError::Handler { message, .. }) => message,
1860 Exit::Fail(ActorError::Defect(message)) => message,
1861 Exit::Fail(error @ ActorError::StepLimit { .. }) => format!("StepLimit: {error}"),
1862 Exit::Fail(other) => other.to_string(),
1863 };
1864 anyhow::anyhow!(message)
1865}
1866
1867pub fn phase_name(phase: CodingPhase) -> &'static str {
1868 match phase {
1869 CodingPhase::Idle => "idle",
1870 CodingPhase::Infer => "infer",
1871 CodingPhase::Compact => "compact",
1872 CodingPhase::Confirm => "confirm",
1873 CodingPhase::Question => "question",
1874 CodingPhase::Tool => "tool",
1875 CodingPhase::Deny => "deny",
1876 CodingPhase::Done => "done",
1877 }
1878}
1879
1880pub fn project_needs_confirmation(policy: &PermissionPolicy, tool_name: &str) -> bool {
1882 confirmation_required(policy, tool_name)
1883}
1884
1885#[cfg(test)]
1886mod tests {
1887 use super::*;
1888 use std::sync::Mutex;
1889
1890 struct ScriptModel {
1891 decisions: Mutex<Vec<Result<ModelDecision, ActorError>>>,
1892 calls: Arc<AtomicUsize>,
1893 tool_counts: Mutex<Vec<usize>>,
1894 messages: Mutex<Vec<Vec<String>>>,
1895 }
1896
1897 impl Completion for ScriptModel {
1898 fn complete(
1899 &self,
1900 request: CompletionRequest,
1901 ) -> a3s_effect::coding::BoxFuture<Result<ModelDecision, ActorError>> {
1902 self.calls.fetch_add(1, Ordering::SeqCst);
1903 self.tool_counts.lock().unwrap().push(request.tools.len());
1904 self.messages.lock().unwrap().push(request.messages.clone());
1905 let decision =
1906 self.decisions
1907 .lock()
1908 .unwrap()
1909 .pop()
1910 .unwrap_or(Ok(ModelDecision::Text {
1911 text: "empty".into(),
1912 }));
1913 Box::pin(async move { decision })
1914 }
1915 }
1916
1917 struct ScriptTools {
1918 calls: Arc<AtomicUsize>,
1919 fail_first: AtomicUsize,
1920 }
1921
1922 impl ToolRunner for ScriptTools {
1923 fn run(
1924 &self,
1925 call: ToolCall,
1926 ) -> a3s_effect::coding::BoxFuture<Result<serde_json::Value, ActorError>> {
1927 let n = self.calls.fetch_add(1, Ordering::SeqCst);
1928 let fail = self.fail_first.load(Ordering::SeqCst) > 0 && n == 0;
1929 let name = call.name;
1930 Box::pin(async move {
1931 if fail {
1932 Err(ActorError::Handler {
1933 key: name,
1934 message: "down".into(),
1935 })
1936 } else {
1937 Ok(serde_json::json!(format!("ran {name}")))
1938 }
1939 })
1940 }
1941 }
1942
1943 fn config(cap: Option<u32>, budget: u32) -> HarnessConfig {
1944 let config = HarnessConfig::new(
1945 budget,
1946 1_000_000,
1947 16,
1948 2,
1949 vec!["system".into()],
1950 vec![ToolSpec {
1951 name: "read".into(),
1952 description: "Read".into(),
1953 }],
1954 )
1955 .unwrap();
1956 match cap {
1957 Some(cap) => config.with_tool_round_cap(cap),
1958 None => config,
1959 }
1960 }
1961
1962 fn run_with(
1963 decisions: Vec<Result<ModelDecision, ActorError>>,
1964 cap: Option<u32>,
1965 budget: u32,
1966 fail_first: bool,
1967 ) -> (
1968 tempfile::TempDir,
1969 FactRun,
1970 Arc<AtomicUsize>,
1971 Arc<AtomicUsize>,
1972 Arc<ScriptModel>,
1973 ) {
1974 let dir = tempfile::tempdir().unwrap();
1975 let calls = Arc::new(AtomicUsize::new(0));
1976 let tool_calls = Arc::new(AtomicUsize::new(0));
1977 let model = Arc::new(ScriptModel {
1978 decisions: Mutex::new(decisions.into_iter().rev().collect()),
1979 calls: Arc::clone(&calls),
1980 tool_counts: Mutex::new(Vec::new()),
1981 messages: Mutex::new(Vec::new()),
1982 });
1983 let fact = FactRun::open(
1984 dir.path(),
1985 model.clone(),
1986 Arc::new(ScriptTools {
1987 calls: Arc::clone(&tool_calls),
1988 fail_first: AtomicUsize::new(usize::from(fail_first)),
1989 }),
1990 config(cap, budget),
1991 {
1992 static FACT_RUNS: AtomicUsize = AtomicUsize::new(0);
1993 format!("run-{}", FACT_RUNS.fetch_add(1, Ordering::Relaxed))
1994 },
1995 )
1996 .unwrap();
1997 (dir, fact, calls, tool_calls, model)
1998 }
1999
2000 #[tokio::test]
2001 async fn fact_log_second_resume_does_not_call_the_model() {
2002 let (_dir, fact, calls, tools, _) = run_with(
2003 vec![Ok(ModelDecision::Text { text: "hi".into() })],
2004 None,
2005 4,
2006 false,
2007 );
2008 let settled = fact.user_text("hello").await.unwrap();
2009 assert_eq!(settled.view.phase, CodingPhase::Done);
2010 assert_eq!(calls.load(Ordering::SeqCst), 1);
2011 let again = fact.resume_limit(8).await.unwrap();
2012 assert_eq!(again.steps, 0);
2013 assert_eq!(calls.load(Ordering::SeqCst), 1);
2014 assert_eq!(tools.load(Ordering::SeqCst), 0);
2015 let facts = fact.read_facts().unwrap();
2016 assert!(facts.iter().any(|fact| fact.kind == "model.turn"));
2017 }
2018
2019 #[tokio::test]
2020 async fn fact_log_non_decision_sets_schema_error_without_a_tool_or_parse_retry() {
2021 let (_dir, fact, calls, tools, _) = run_with(vec![], None, 4, false);
2022 fact.append_model_turn(serde_json::json!({"kind": "nope"}))
2023 .await
2024 .unwrap();
2025 let settled = fact.resume_limit(8).await.unwrap();
2026 assert_eq!(settled.view.phase, CodingPhase::Done);
2027 assert!(settled.view.schema_error.is_some());
2028 assert_eq!(calls.load(Ordering::SeqCst), 0);
2029 assert_eq!(tools.load(Ordering::SeqCst), 0);
2030 }
2031
2032 #[tokio::test]
2033 async fn fact_log_denial_does_not_run_the_tool() {
2034 let (_dir, fact, _, tools, _) = run_with(
2035 vec![Ok(ModelDecision::Tool {
2036 call: ToolCall {
2037 id: "t1".into(),
2038 name: "read".into(),
2039 args: serde_json::json!({}),
2040 needs_confirmation: true,
2041 text: None,
2042 reasoning: None,
2043 },
2044 })],
2045 None,
2046 4,
2047 false,
2048 );
2049 let parked = fact.user_text("read it").await.unwrap();
2050 assert_eq!(parked.view.phase, CodingPhase::Confirm);
2051 assert_eq!(tools.load(Ordering::SeqCst), 0);
2052 let denied = fact.confirm("t1", false).await.unwrap();
2053 assert_eq!(denied.view.assistant.as_deref(), Some("denied"));
2054 assert_eq!(tools.load(Ordering::SeqCst), 0);
2055 }
2056
2057 #[tokio::test]
2058 async fn fact_log_question_options_and_allow_free_text_survive_reopen() {
2059 let (dir, fact, calls, tools, _) = run_with(
2060 vec![Ok(ModelDecision::Question {
2061 question_id: "q1".into(),
2062 question: "Which?".into(),
2063 allow_free_text: true,
2064 options: vec!["left".into(), "right".into()],
2065 })],
2066 None,
2067 4,
2068 false,
2069 );
2070 let parked = fact.user_text("ask").await.unwrap();
2071 assert_eq!(parked.view.phase, CodingPhase::Question);
2072 drop(parked);
2073 let reopened = FactRun::open(
2074 dir.path(),
2075 Arc::new(ScriptModel {
2076 decisions: Mutex::new(vec![]),
2077 calls: Arc::clone(&calls),
2078 tool_counts: Mutex::new(Vec::new()),
2079 messages: Mutex::new(Vec::new()),
2080 }),
2081 Arc::new(ScriptTools {
2082 calls: Arc::new(AtomicUsize::new(0)),
2083 fail_first: AtomicUsize::new(0),
2084 }),
2085 config(None, 4),
2086 "run-fact-2",
2087 )
2088 .unwrap();
2089 let still = reopened.resume_limit(8).await.unwrap();
2090 assert_eq!(still.steps, 0);
2091 let question = still.view.pending_question.expect("question");
2092 assert!(question.allow_free_text);
2093 assert_eq!(
2094 question.options,
2095 vec!["left".to_string(), "right".to_string()]
2096 );
2097 assert_eq!(tools.load(Ordering::SeqCst), 0);
2098 }
2099
2100 #[tokio::test]
2101 async fn fact_log_question_stays_parked_until_the_answer_fact() {
2102 let (_dir, fact, calls, _, _) = run_with(
2103 vec![
2104 Ok(ModelDecision::Question {
2105 question_id: "q1".into(),
2106 question: "Which?".into(),
2107 allow_free_text: true,
2108 options: vec!["left".into(), "right".into()],
2109 }),
2110 Ok(ModelDecision::Text {
2111 text: "left".into(),
2112 }),
2113 ],
2114 None,
2115 4,
2116 false,
2117 );
2118 let parked = fact.user_text("ask").await.unwrap();
2119 assert_eq!(parked.view.phase, CodingPhase::Question);
2120 let calls_while_parked = calls.load(Ordering::SeqCst);
2121 let resumed = fact.resume_limit(4).await.unwrap();
2122 assert_eq!(resumed.steps, 0);
2123 assert_eq!(resumed.view.phase, CodingPhase::Question);
2124 assert_eq!(calls.load(Ordering::SeqCst), calls_while_parked);
2125 assert!(resumed
2126 .log
2127 .iter()
2128 .all(|fact| fact.kind != "question.answered"));
2129 let answered = fact.answer("left").await.unwrap();
2130 assert!(answered
2131 .log
2132 .iter()
2133 .any(|fact| fact.kind == "question.answered"));
2134 assert_eq!(answered.view.phase, CodingPhase::Done);
2135 assert!(calls.load(Ordering::SeqCst) > calls_while_parked);
2136 }
2137
2138 #[tokio::test]
2139 async fn fact_log_missing_tool_result_runs_once_and_does_not_ask_the_model_again() {
2140 let (_dir, fact, calls, tools, _) = run_with(
2141 vec![Ok(ModelDecision::Tool {
2142 call: ToolCall {
2143 id: "t1".into(),
2144 name: "read".into(),
2145 args: serde_json::json!({}),
2146 needs_confirmation: false,
2147 text: None,
2148 reasoning: None,
2149 },
2150 })],
2151 None,
2152 4,
2153 true,
2154 );
2155 let failed = fact.user_text("go").await;
2156 assert!(failed.is_err());
2157 assert_eq!(calls.load(Ordering::SeqCst), 1);
2158 assert_eq!(tools.load(Ordering::SeqCst), 1);
2159 let again = fact.resume_limit(1).await;
2160 assert!(again.is_err(), "the step limit stops the following infer");
2161 assert_eq!(tools.load(Ordering::SeqCst), 2);
2162 assert_eq!(calls.load(Ordering::SeqCst), 1);
2163 }
2164
2165 #[tokio::test]
2166 async fn fact_log_step_limit_does_not_start_a_pending_transition() {
2167 let (_dir, fact, calls, _, _) = run_with(vec![], None, 4, false);
2168 let log = fact.open_log().unwrap();
2169 log.append(THREAD, &[message_fact("m1", "hello")], None)
2170 .unwrap();
2171 drop(log);
2172 let error = match fact.resume_limit(0).await {
2173 Ok(_) => panic!("a pending transition at limit 0 must not run"),
2174 Err(error) => error,
2175 };
2176 let rendered = format!("{error:#}");
2177 assert!(rendered.contains("StepLimit"), "{rendered}");
2178 assert_eq!(calls.load(Ordering::SeqCst), 0);
2179 }
2180
2181 #[tokio::test]
2182 async fn fact_log_budget_does_not_call_the_tool_body() {
2183 let (_dir, fact, _, tools, _) = run_with(
2184 vec![
2185 Ok(ModelDecision::Tool {
2186 call: ToolCall {
2187 id: "t1".into(),
2188 name: "read".into(),
2189 args: serde_json::json!({}),
2190 needs_confirmation: false,
2191 text: None,
2192 reasoning: None,
2193 },
2194 }),
2195 Ok(ModelDecision::Tool {
2196 call: ToolCall {
2197 id: "t2".into(),
2198 name: "read".into(),
2199 args: serde_json::json!({}),
2200 needs_confirmation: false,
2201 text: None,
2202 reasoning: None,
2203 },
2204 }),
2205 ],
2206 None,
2207 1,
2208 false,
2209 );
2210 let settled = fact.user_text("two tools").await.unwrap();
2211 assert_eq!(tools.load(Ordering::SeqCst), 1);
2212 assert!(settled.log.iter().any(|fact| fact.kind == "budget.denied"));
2213 }
2214
2215 #[tokio::test]
2216 async fn fact_log_tool_round_cap_uses_an_empty_tool_list() {
2217 let (_dir, fact, calls, _, model) = run_with(
2218 vec![Ok(ModelDecision::Text { text: "ok".into() })],
2219 Some(0),
2220 4,
2221 false,
2222 );
2223 fact.user_text("hello").await.unwrap();
2224 assert_eq!(calls.load(Ordering::SeqCst), 1);
2225 assert_eq!(model.tool_counts.lock().unwrap().as_slice(), &[0]);
2226 assert!(model
2227 .messages
2228 .lock()
2229 .unwrap()
2230 .iter()
2231 .flatten()
2232 .all(|line| !line.contains("TOOL_BUDGET_FINALIZATION")));
2233 }
2234
2235 #[tokio::test]
2236 async fn fact_log_steer_is_a_user_message_fact() {
2237 let (_dir, fact, calls, _, _) = run_with(
2238 vec![
2239 Ok(ModelDecision::Text {
2240 text: "first".into(),
2241 }),
2242 Ok(ModelDecision::Text {
2243 text: "steered".into(),
2244 }),
2245 ],
2246 None,
2247 4,
2248 false,
2249 );
2250 fact.user_text("hello").await.unwrap();
2251 let steered = fact.steer("turn left").await.unwrap();
2252 assert_eq!(steered.view.assistant.as_deref(), Some("steered"));
2253 assert_eq!(calls.load(Ordering::SeqCst), 2);
2254 let facts = fact.read_facts().unwrap();
2255 assert_eq!(
2256 facts
2257 .iter()
2258 .filter(|fact| fact.kind == "user.message")
2259 .count(),
2260 2
2261 );
2262 }
2263
2264 struct HangClient;
2265
2266 #[async_trait::async_trait]
2267 impl LlmClient for HangClient {
2268 async fn complete(
2269 &self,
2270 _messages: &[Message],
2271 _system: Option<&str>,
2272 _tools: &[ToolDefinition],
2273 ) -> Result<crate::llm::LlmResponse> {
2274 std::future::pending().await
2275 }
2276
2277 async fn complete_streaming(
2278 &self,
2279 _messages: &[Message],
2280 _system: Option<&str>,
2281 _tools: &[ToolDefinition],
2282 _cancel_token: tokio_util::sync::CancellationToken,
2283 ) -> Result<tokio::sync::mpsc::Receiver<crate::llm::StreamEvent>> {
2284 let (sender, receiver) = tokio::sync::mpsc::channel(1);
2285 tokio::spawn(async move {
2286 std::future::pending::<()>().await;
2287 drop(sender);
2288 });
2289 Ok(receiver)
2290 }
2291 }
2292
2293 #[tokio::test]
2294 async fn fact_log_model_stream_timeout_stops_at_the_provider_deadline() {
2295 let workspace = tempfile::tempdir().expect("workspace");
2296 let executor = Arc::new(crate::tools::ToolExecutor::new(
2297 workspace.path().to_string_lossy().into_owned(),
2298 ));
2299 let mut config = crate::agent::AgentConfig::default();
2300 config.llm_api_timeout_ms = Some(80);
2301 config.planning_mode = crate::prompts::PlanningMode::Disabled;
2302 let agent = crate::agent::AgentLoop::new(
2303 Arc::new(HangClient),
2304 Arc::clone(&executor),
2305 crate::tools::ToolContext::new(workspace.path().to_path_buf()),
2306 config,
2307 );
2308 let surface = SessionSurface {
2309 agent,
2310 session_id: "timeout-session".into(),
2311 checkpoint: None,
2312 events: None,
2313 cancel: tokio_util::sync::CancellationToken::new(),
2314 transcript: Arc::new(Mutex::new(Vec::new())),
2315 usage: Arc::new(Mutex::new(crate::llm::TokenUsage::default())),
2316 confirmation: None,
2317 run_store: None,
2318 run_id: None,
2319 ledger: Arc::new(Mutex::new(crate::harness_loop::MutationLedger::default())),
2320 reports: Arc::new(Mutex::new(Vec::new())),
2321 run_control: None,
2322 harness: None,
2323 host_harness_registry: None,
2324 host_harness_assembler: None,
2325 };
2326 let run = FactRun::open_projected(
2327 workspace.path(),
2328 "timeout-session",
2329 Arc::new(HangClient),
2330 executor,
2331 PermissionPolicy::new().allow("*"),
2332 &[],
2333 4,
2334 surface,
2335 )
2336 .expect("fact run");
2337 let result =
2338 tokio::time::timeout(std::time::Duration::from_secs(3), run.user_text("hello"))
2339 .await
2340 .expect("provider deadline did not stop the model call");
2341 let Err(error) = result else {
2342 panic!("a hanging model must not settle");
2343 };
2344 let rendered = format!("{error:#}");
2345 assert!(
2346 rendered.contains("timed out"),
2347 "expected the provider deadline, got {rendered}"
2348 );
2349 }
2350
2351 #[test]
2352 fn fact_log_plan_guardrail_denies_a_write_the_default_policy_would_ask() {
2353 let policy = PermissionPolicy::new();
2354 let checker = crate::permissions::InteractiveToolGuardrail::for_mode("plan");
2355 let write_args = serde_json::json!({
2356 "file_path": "hello.txt",
2357 "content": "hello"
2358 });
2359 assert_eq!(
2360 execution_permission(&policy, Some(&checker), "write", &write_args),
2361 PermissionDecision::Deny
2362 );
2363 assert_eq!(
2364 execution_permission(
2365 &policy,
2366 Some(&checker),
2367 "read",
2368 &serde_json::json!({ "file_path": "hello.txt" })
2369 ),
2370 PermissionDecision::Allow
2371 );
2372 let yolo = PermissionPolicy::new().allow_yolo_lanes([SessionLane::Execute]);
2373 let asking = PermissionPolicy::new();
2374 assert_eq!(
2375 execution_permission(&yolo, Some(&asking), "bash", &serde_json::json!({})),
2376 PermissionDecision::Allow
2377 );
2378 let denied = PermissionPolicy::new().deny("write(*)");
2379 let allowing = PermissionPolicy::new().allow("write(*)");
2380 assert_eq!(
2381 execution_permission(&denied, Some(&allowing), "write", &write_args),
2382 PermissionDecision::Deny
2383 );
2384 }
2385
2386 #[test]
2387 fn fact_log_yolo_lane_is_allow_not_a_second_policy() {
2388 let policy = PermissionPolicy::new().allow_yolo_lanes([SessionLane::Execute]);
2389 assert!(!project_needs_confirmation(&policy, "bash"));
2390 assert!(!project_needs_confirmation(&policy, "write"));
2391 assert!(!project_needs_confirmation(&policy, "unknown_tool"));
2392 assert!(project_needs_confirmation(&policy, "read"));
2393 assert!(project_needs_confirmation(&policy, "search"));
2394 let only_read = PermissionPolicy::new().allow("read(*)");
2395 assert!(project_needs_confirmation(&only_read, "bash"));
2396 let denied = PermissionPolicy::new()
2397 .deny("bash(*)")
2398 .allow_yolo_lanes([SessionLane::Execute]);
2399 assert_eq!(
2400 denied.check("bash", &serde_json::json!({})),
2401 PermissionDecision::Deny
2402 );
2403 assert_eq!(
2404 denied.check("write", &serde_json::json!({})),
2405 PermissionDecision::Allow
2406 );
2407 assert!(!crate::hitl::ConfirmationPolicy::enabled()
2408 .with_yolo_lanes([SessionLane::Execute])
2409 .is_yolo("bash"));
2410 }
2411
2412 #[tokio::test]
2413 async fn fact_log_question_cap_rejects_before_a_fourth_fact() {
2414 let questions = (1..=4)
2415 .map(|index| {
2416 Ok(ModelDecision::Question {
2417 question_id: format!("q{index}"),
2418 question: format!("q{index}?"),
2419 allow_free_text: false,
2420 options: vec!["a".into()],
2421 })
2422 })
2423 .collect();
2424 let (_dir, fact, _, _, _) = run_with(questions, None, 4, false);
2425 fact.user_text("one").await.unwrap();
2426 fact.answer("a").await.unwrap();
2427 fact.answer("b").await.unwrap();
2428 let rejected = match fact.answer("c").await {
2429 Ok(_) => panic!("fourth question must fail before a model turn"),
2430 Err(error) => error,
2431 };
2432 let rendered = format!("{rejected:#}");
2433 assert!(
2434 rendered.contains("ask_user question cap exceeded"),
2435 "{rendered}"
2436 );
2437 let facts = fact.read_facts().unwrap();
2438 assert_eq!(
2439 facts
2440 .iter()
2441 .filter(|fact| fact.kind == "model.turn")
2442 .count(),
2443 3
2444 );
2445 }
2446
2447 #[tokio::test]
2448 async fn fact_log_confirm_is_true_only_for_the_pending_tool() {
2449 let (_dir, fact, _, tools, _) = run_with(
2450 vec![Ok(ModelDecision::Tool {
2451 call: ToolCall {
2452 id: "t1".into(),
2453 name: "read".into(),
2454 args: serde_json::json!({}),
2455 needs_confirmation: true,
2456 text: None,
2457 reasoning: None,
2458 },
2459 })],
2460 None,
2461 4,
2462 false,
2463 );
2464 fact.user_text("read it").await.unwrap();
2465 assert!(!fact.confirm_if_pending("other", true).await.unwrap());
2466 assert!(fact
2467 .read_facts()
2468 .unwrap()
2469 .iter()
2470 .all(|fact| fact.kind != "confirmation.answered"));
2471 assert!(fact.confirm_if_pending("t1", false).await.unwrap());
2472 assert_eq!(tools.load(Ordering::SeqCst), 0);
2473 assert_eq!(
2474 fact.read_facts()
2475 .unwrap()
2476 .iter()
2477 .filter(|fact| fact.kind == "confirmation.answered")
2478 .count(),
2479 1
2480 );
2481 }
2482
2483 #[tokio::test]
2484 async fn fact_log_history_seeds_a_caused_model_turn_and_one_new_call() {
2485 let (_dir, fact, calls, _, _) = run_with(
2486 vec![Ok(ModelDecision::Text {
2487 text: "next".into(),
2488 })],
2489 None,
2490 4,
2491 false,
2492 );
2493 fact.seed_history(&[Message::user("old"), Message::assistant("stored")])
2494 .unwrap();
2495 let settled = fact.user_text("new").await.unwrap();
2496 assert_eq!(settled.view.assistant.as_deref(), Some("next"));
2497 assert_eq!(calls.load(Ordering::SeqCst), 1);
2498 let facts = fact.read_facts().unwrap();
2499 assert_eq!(
2500 facts
2501 .iter()
2502 .filter(|fact| fact.kind == "user.message")
2503 .count(),
2504 2
2505 );
2506 let stored = facts
2507 .iter()
2508 .find(|fact| fact.cause.as_deref() == Some("infer:1:0"))
2509 .expect("caused history turn");
2510 assert_eq!(stored.payload["text"], "stored");
2511 assert_eq!(
2512 facts
2513 .iter()
2514 .filter(|fact| fact.kind == "model.turn")
2515 .count(),
2516 2
2517 );
2518 }
2519
2520 fn tool_use(name: &str, id: &str, args: serde_json::Value) -> Message {
2521 Message {
2522 role: "assistant".into(),
2523 content: vec![crate::llm::ContentBlock::ToolUse {
2524 id: id.into(),
2525 name: name.into(),
2526 input: args,
2527 }],
2528 reasoning_content: None,
2529 transcript_text: None,
2530 transcript_visibility: crate::llm::TranscriptVisibility::Product,
2531 }
2532 }
2533
2534 struct SeqClient {
2535 calls: Arc<AtomicUsize>,
2536 tools: Mutex<Vec<Vec<String>>>,
2537 seen: Mutex<Vec<Vec<Message>>>,
2538 responses: Mutex<Vec<Message>>,
2539 }
2540
2541 #[async_trait::async_trait]
2542 impl LlmClient for SeqClient {
2543 async fn complete(
2544 &self,
2545 messages: &[Message],
2546 _system: Option<&str>,
2547 tools: &[crate::llm::ToolDefinition],
2548 ) -> Result<crate::llm::LlmResponse> {
2549 self.calls.fetch_add(1, Ordering::SeqCst);
2550 self.seen.lock().unwrap().push(messages.to_vec());
2551 self.tools
2552 .lock()
2553 .unwrap()
2554 .push(tools.iter().map(|tool| tool.name.clone()).collect());
2555 let message = self
2556 .responses
2557 .lock()
2558 .unwrap()
2559 .pop()
2560 .unwrap_or_else(|| Message::assistant("done"));
2561 Ok(crate::llm::LlmResponse {
2562 message,
2563 usage: Default::default(),
2564 stop_reason: None,
2565 token_logprobs: Vec::new(),
2566 meta: None,
2567 })
2568 }
2569
2570 async fn complete_streaming(
2571 &self,
2572 _messages: &[Message],
2573 _system: Option<&str>,
2574 _tools: &[crate::llm::ToolDefinition],
2575 _cancel_token: tokio_util::sync::CancellationToken,
2576 ) -> Result<tokio::sync::mpsc::Receiver<crate::llm::StreamEvent>> {
2577 anyhow::bail!("unused")
2578 }
2579 }
2580
2581 fn open_live(
2582 responses: Vec<Message>,
2583 policy: PermissionPolicy,
2584 tool_name: &str,
2585 ) -> (tempfile::TempDir, FactRun, Arc<AtomicUsize>, Arc<SeqClient>) {
2586 let dir = tempfile::tempdir().unwrap();
2587 let calls = Arc::new(AtomicUsize::new(0));
2588 let client = Arc::new(SeqClient {
2589 calls: Arc::clone(&calls),
2590 tools: Mutex::new(Vec::new()),
2591 seen: Mutex::new(Vec::new()),
2592 responses: Mutex::new(responses.into_iter().rev().collect()),
2593 });
2594 let (completion, _) = LiveCompletion::new(client.clone(), policy, Vec::new());
2595 let config = HarnessConfig::new(
2596 4,
2597 1_000_000,
2598 16,
2599 1,
2600 Vec::new(),
2601 vec![ToolSpec {
2602 name: tool_name.into(),
2603 description: "tool".into(),
2604 }],
2605 )
2606 .unwrap();
2607 let fact = FactRun::open(
2608 dir.path(),
2609 Arc::new(completion),
2610 Arc::new(ScriptTools {
2611 calls: Arc::new(AtomicUsize::new(0)),
2612 fail_first: AtomicUsize::new(0),
2613 }),
2614 config,
2615 format!("live-{}", calls.as_ptr() as usize),
2616 )
2617 .unwrap();
2618 (dir, fact, calls, client)
2619 }
2620
2621 #[tokio::test]
2622 async fn fact_log_live_completion_stores_a_tool_call_and_passes_tools() {
2623 let (_dir, fact, calls, client) = open_live(
2624 vec![tool_use(
2625 "bash",
2626 "call-1",
2627 serde_json::json!({"command": "ls"}),
2628 )],
2629 PermissionPolicy::new(),
2630 "bash",
2631 );
2632 let settled = fact.user_text("list").await.unwrap();
2633 assert_eq!(settled.view.phase, CodingPhase::Confirm);
2634 assert_eq!(calls.load(Ordering::SeqCst), 1);
2635 assert_eq!(client.tools.lock().unwrap().as_slice(), &[vec!["bash"]]);
2636 let turn = fact
2637 .read_facts()
2638 .unwrap()
2639 .into_iter()
2640 .find(|fact| fact.kind == "model.turn")
2641 .expect("model turn");
2642 assert_eq!(turn.payload["kind"], "tool");
2643 assert_eq!(turn.payload["call"]["needs_confirmation"], true);
2644 }
2645
2646 #[tokio::test]
2647 async fn fact_log_live_completion_yolo_lane_runs_the_tool() {
2648 let (_dir, fact, _, client) = open_live(
2649 vec![
2650 tool_use("bash", "call-1", serde_json::json!({})),
2651 Message::assistant("done"),
2652 ],
2653 PermissionPolicy::new().allow_yolo_lanes([SessionLane::Execute]),
2654 "bash",
2655 );
2656 let settled = fact.user_text("list").await.unwrap();
2657 assert_eq!(settled.view.phase, CodingPhase::Done);
2658 assert!(settled.log.iter().any(|fact| fact.kind == "tool.result"));
2659 assert!(!settled
2660 .log
2661 .iter()
2662 .any(|fact| fact.kind == "confirmation.answered"));
2663 assert!(client
2664 .tools
2665 .lock()
2666 .unwrap()
2667 .iter()
2668 .any(|tools| !tools.is_empty()));
2669 }
2670
2671 #[tokio::test]
2672 async fn fact_log_live_completion_ask_user_parks_with_options() {
2673 let (_dir, fact, _, _) = open_live(
2674 vec![tool_use(
2675 "ask_user",
2676 "q-live",
2677 serde_json::json!({
2678 "question": "Which?",
2679 "options": ["left", "right"],
2680 "allow_free_text": true
2681 }),
2682 )],
2683 PermissionPolicy::new(),
2684 "ask_user",
2685 );
2686 let settled = fact.user_text("ask").await.unwrap();
2687 assert_eq!(settled.view.phase, CodingPhase::Question);
2688 let question = settled.view.pending_question.expect("question");
2689 assert!(question.allow_free_text);
2690 assert_eq!(
2691 question.options,
2692 vec!["left".to_string(), "right".to_string()]
2693 );
2694 }
2695
2696 #[tokio::test]
2697 async fn fact_log_sessions_do_not_share_a_thread() {
2698 let workspace = tempfile::tempdir().unwrap();
2699 let calls = Arc::new(AtomicUsize::new(0));
2700 let client = Arc::new(SeqClient {
2701 calls: Arc::clone(&calls),
2702 tools: Mutex::new(Vec::new()),
2703 seen: Mutex::new(Vec::new()),
2704 responses: Mutex::new(vec![Message::assistant("bee"), Message::assistant("aye")]),
2705 });
2706 let executor = Arc::new(ToolExecutor::new(
2707 workspace.path().to_string_lossy().to_string(),
2708 ));
2709 let first = FactRun::open_session(
2710 workspace.path(),
2711 "session-a",
2712 client.clone(),
2713 Arc::clone(&executor),
2714 PermissionPolicy::new(),
2715 &[],
2716 4,
2717 )
2718 .unwrap();
2719 let second = FactRun::open_session(
2720 workspace.path(),
2721 "session-b",
2722 client,
2723 executor,
2724 PermissionPolicy::new(),
2725 &[],
2726 4,
2727 )
2728 .unwrap();
2729 first.user_text("from-a").await.unwrap();
2730 second.user_text("from-b").await.unwrap();
2731 let facts_a = read_workspace_facts(workspace.path(), "session-a").unwrap();
2732 let facts_b = read_workspace_facts(workspace.path(), "session-b").unwrap();
2733 assert!(facts_a.iter().any(|fact| fact.payload["text"] == "from-a"));
2734 assert!(facts_b.iter().any(|fact| fact.payload["text"] == "from-b"));
2735 assert!(facts_a.iter().all(|fact| fact.payload["text"] != "from-b"));
2736 assert!(facts_b.iter().all(|fact| fact.payload["text"] != "from-a"));
2737 }
2738
2739 fn dsml_bash(command: &str) -> String {
2740 let ns = "\u{FF5C}DSML\u{FF5C}";
2741 format!(
2742 "<{ns}invoke name=\"bash\"><{ns}parameter name=\"command\" string=\"true\">{command}</{ns}parameter></{ns}invoke>"
2743 )
2744 }
2745
2746 #[tokio::test]
2747 async fn fact_log_leaked_dsml_text_is_a_tool_decision() {
2748 let (_dir, fact, _, _) = open_live(
2749 vec![Message::assistant(&dsml_bash("pwd"))],
2750 PermissionPolicy::new(),
2751 "bash",
2752 );
2753 let settled = fact.user_text("list").await.unwrap();
2754 assert_eq!(settled.view.phase, CodingPhase::Confirm);
2755 let turn = settled
2756 .log
2757 .iter()
2758 .find(|fact| fact.kind == "model.turn")
2759 .expect("model turn");
2760 assert_eq!(turn.payload["kind"], "tool");
2761 assert_eq!(turn.payload["call"]["name"], "bash");
2762 assert_eq!(turn.payload["call"]["args"]["command"], "pwd");
2763 }
2764
2765 #[tokio::test]
2766 async fn fact_log_leaked_dsml_reasoning_is_a_tool_decision() {
2767 let markup = dsml_bash("pwd");
2768 let (_dir, fact, _, _) = open_live(
2769 vec![Message {
2770 role: "assistant".into(),
2771 content: Vec::new(),
2772 reasoning_content: Some(markup),
2773 transcript_text: None,
2774 transcript_visibility: crate::llm::TranscriptVisibility::Wire,
2775 }],
2776 PermissionPolicy::new(),
2777 "bash",
2778 );
2779 let settled = fact.user_text("list").await.unwrap();
2780 assert_eq!(settled.view.phase, CodingPhase::Confirm);
2781 let turn = settled
2782 .log
2783 .iter()
2784 .find(|fact| fact.kind == "model.turn")
2785 .expect("model turn");
2786 assert_eq!(turn.payload["call"]["name"], "bash");
2787 assert_eq!(turn.payload["call"]["args"]["command"], "pwd");
2788 }
2789
2790 #[tokio::test]
2791 async fn fact_log_follow_up_sends_the_recorded_tool_call() {
2792 let workspace = tempfile::tempdir().unwrap();
2793 let log = log_dir(workspace.path());
2794 std::fs::create_dir_all(&log).unwrap();
2795 let calls = Arc::new(AtomicUsize::new(0));
2796 let client = Arc::new(SeqClient {
2797 calls: Arc::clone(&calls),
2798 tools: Mutex::new(Vec::new()),
2799 seen: Mutex::new(Vec::new()),
2800 responses: Mutex::new(vec![Message::assistant("done"), {
2801 let mut planned = tool_use(
2802 "read",
2803 "call-read",
2804 serde_json::json!({"file_path": "note.txt"}),
2805 );
2806 planned.content.insert(
2807 0,
2808 crate::llm::ContentBlock::Text {
2809 text: "Next, rerun the test.".into(),
2810 },
2811 );
2812 planned.reasoning_content = Some("the stats file is still wrong".into());
2813 planned
2814 }]),
2815 });
2816 let executor = Arc::new(crate::tools::ToolExecutor::new(
2817 workspace.path().to_string_lossy().into_owned(),
2818 ));
2819 let mut config = crate::agent::AgentConfig::default();
2820 config.planning_mode = crate::prompts::PlanningMode::Disabled;
2821 let agent = crate::agent::AgentLoop::new(
2822 Arc::new(SeqClient {
2823 calls: Arc::new(AtomicUsize::new(0)),
2824 tools: Mutex::new(Vec::new()),
2825 seen: Mutex::new(Vec::new()),
2826 responses: Mutex::new(Vec::new()),
2827 }),
2828 executor,
2829 crate::tools::ToolContext::new(workspace.path().to_path_buf()),
2830 config,
2831 );
2832 let surface = SessionSurface {
2833 agent,
2834 session_id: THREAD.into(),
2835 checkpoint: None,
2836 events: None,
2837 cancel: tokio_util::sync::CancellationToken::new(),
2838 transcript: Arc::new(Mutex::new(Vec::new())),
2839 usage: Arc::new(Mutex::new(crate::llm::TokenUsage::default())),
2840 confirmation: None,
2841 run_store: None,
2842 run_id: None,
2843 ledger: Arc::new(Mutex::new(crate::harness_loop::MutationLedger::default())),
2844 reports: Arc::new(Mutex::new(Vec::new())),
2845 run_control: None,
2846 harness: None,
2847 host_harness_registry: None,
2848 host_harness_assembler: None,
2849 };
2850 let (completion, _) = LiveCompletion::new(
2851 client.clone(),
2852 PermissionPolicy::new().allow("*"),
2853 Vec::new(),
2854 );
2855 let completion = completion.with_surface(surface);
2856 let fact = FactRun::open(
2857 log,
2858 Arc::new(completion),
2859 Arc::new(ScriptTools {
2860 calls: Arc::new(AtomicUsize::new(0)),
2861 fail_first: AtomicUsize::new(0),
2862 }),
2863 HarnessConfig::new(
2864 4,
2865 1_000_000,
2866 16,
2867 1,
2868 Vec::new(),
2869 vec![ToolSpec {
2870 name: "read".into(),
2871 description: "Read".into(),
2872 }],
2873 )
2874 .unwrap(),
2875 "recorded-tool",
2876 )
2877 .unwrap();
2878 let settled = fact.user_text("read the note").await.unwrap();
2879 assert_eq!(settled.view.phase, CodingPhase::Done);
2880 assert_eq!(calls.load(Ordering::SeqCst), 2);
2881 let seen = client.seen.lock().unwrap();
2882 let follow = seen.last().expect("follow-up completion");
2883 let recorded = follow
2884 .iter()
2885 .find_map(|message| {
2886 let calls = message.tool_calls();
2887 if calls.is_empty() {
2888 None
2889 } else {
2890 Some(calls)
2891 }
2892 })
2893 .expect("recorded tool call");
2894 assert_eq!(recorded[0].id, "call-read");
2895 assert_eq!(recorded[0].name, "read");
2896 assert_eq!(recorded[0].args["file_path"], "note.txt");
2897 let replay = follow
2898 .iter()
2899 .find(|message| !message.tool_calls().is_empty())
2900 .expect("assistant tool message");
2901 assert!(replay.text().contains("Next, rerun the test."));
2902 assert_eq!(
2903 replay.reasoning_content.as_deref(),
2904 Some("the stats file is still wrong")
2905 );
2906 assert!(follow.iter().any(|message| {
2907 message.content.iter().any(|block| {
2908 matches!(
2909 block,
2910 crate::llm::ContentBlock::ToolResult { tool_use_id, .. }
2911 if tool_use_id == "call-read"
2912 )
2913 })
2914 }));
2915 }
2916
2917 #[test]
2918 fn fact_log_recorded_tool_calls_start_after_the_latest_compaction() {
2919 use a3s_effect::LogStore;
2920 let workspace = tempfile::tempdir().unwrap();
2921 let dir = log_dir(workspace.path());
2922 std::fs::create_dir_all(&dir).unwrap();
2923 let log = a3s_effect::FileLog::open(&dir).unwrap();
2924 let tool = |id: &str, name: &str| {
2925 serde_json::json!({
2926 "kind": "tool",
2927 "call": {
2928 "id": id,
2929 "name": name,
2930 "args": {"file_path": name},
2931 "needs_confirmation": false
2932 }
2933 })
2934 };
2935 let append = |key: &str, kind: &str, payload: serde_json::Value| {
2936 log.append(
2937 THREAD,
2938 &[a3s_effect::NewFact {
2939 kind: kind.into(),
2940 key: key.into(),
2941 payload,
2942 }],
2943 None,
2944 )
2945 .unwrap();
2946 };
2947 append("m1", "model.turn", tool("old", "edit"));
2948 append(
2949 "t1",
2950 "tool.result",
2951 serde_json::json!({"toolCallId": "old", "ok": true, "output": "old"}),
2952 );
2953 append(
2954 "c1",
2955 "compaction.done",
2956 serde_json::json!({"summary": "earlier edit"}),
2957 );
2958 append("m2", "model.turn", tool("new", "write"));
2959 append(
2960 "t2",
2961 "tool.result",
2962 serde_json::json!({"toolCallId": "new", "ok": true, "output": "new"}),
2963 );
2964 let calls = recorded_tool_calls(workspace.path(), THREAD);
2965 assert_eq!(calls.len(), 1);
2966 assert_eq!(calls[0].id, "new");
2967 assert_eq!(calls[0].name, "write");
2968 }
2969}