use std::collections::HashMap;
use std::future::Future;
use std::path::PathBuf;
use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
use std::sync::{Arc, Mutex as StdMutex};
use std::time::Duration;
use car_engine::Runtime;
use car_inference::tasks::generate::{ContentBlock, Message, Provenance, ToolCall};
use car_ir::{ActionProposal, ActionStatus};
use serde_json::{json, Value};
use tokio::sync::{mpsc, oneshot, Mutex as AsyncMutex};
use super::agent_loop::{
run_assistant_goal_loop_in_session_durable, run_assistant_loop_cancellable_in_session_durable,
ApprovalDecision, ApprovalGate, AssistantEvent,
};
use super::governance::AssistantDurability;
use super::AssistantConfig;
use crate::coder::native_loop::TurnGenerator;
const APPROVAL_TIMEOUT: Duration = Duration::from_secs(300);
#[derive(Clone, Debug)]
pub struct ChatGoal {
pub check: String,
pub max_iterations: u32,
}
fn shell_check_proposal(command: &str) -> car_ir::ActionProposal {
serde_json::from_value(json!({
"source": "chat-goal-check",
"actions": [{
"id": "goal_check",
"type": "tool_call",
"tool": "shell",
"parameters": { "command": command },
}],
}))
.expect("static shell-check proposal shape")
}
async fn run_shell_check_with_approval(
runtime: &Runtime,
cfg: &AssistantConfig,
approval: Option<&dyn ApprovalGate>,
command: &str,
) -> i32 {
if cfg.gated_tools.iter().any(|tool| tool == "shell") {
let params = json!({ "command": command, "purpose": "goal_check" });
match approval {
Some(gate) => match gate.request("shell", ¶ms).await {
ApprovalDecision::Approved => {}
ApprovalDecision::Denied(_) => return 1,
},
None => return 1,
}
}
let exec = runtime.execute(&shell_check_proposal(command)).await;
exec.results
.first()
.and_then(|r| r.output.as_ref())
.and_then(|o| o.get("exit_code"))
.and_then(|v| v.as_i64())
.unwrap_or(1) as i32
}
pub struct AssistantService {
generator: Arc<dyn TurnGenerator>,
runtime: Arc<Runtime>,
cfg: AssistantConfig,
system: String,
threads: AsyncMutex<HashMap<String, Vec<Message>>>,
cancels: StdMutex<HashMap<String, Arc<AtomicBool>>>,
runtime_sessions: AsyncMutex<HashMap<String, String>>,
approvals: Arc<StdMutex<HashMap<String, oneshot::Sender<bool>>>>,
durability: Option<Arc<dyn AssistantDurability>>,
repository_root: Option<PathBuf>,
}
impl AssistantService {
pub fn new(
generator: Arc<dyn TurnGenerator>,
runtime: Arc<Runtime>,
cfg: AssistantConfig,
system: String,
) -> Self {
Self {
generator,
runtime,
cfg,
system,
threads: AsyncMutex::new(HashMap::new()),
cancels: StdMutex::new(HashMap::new()),
runtime_sessions: AsyncMutex::new(HashMap::new()),
approvals: Arc::new(StdMutex::new(HashMap::new())),
durability: None,
repository_root: None,
}
}
pub fn new_durable(
generator: Arc<dyn TurnGenerator>,
runtime: Arc<Runtime>,
cfg: AssistantConfig,
system: String,
durability: Arc<dyn AssistantDurability>,
repository_root: PathBuf,
) -> Self {
let mut service = Self::new(generator, runtime, cfg, system);
service.durability = Some(durability);
service.repository_root = Some(repository_root);
service
}
fn config_for_model(&self, model: Option<&str>) -> AssistantConfig {
let mut cfg = self.cfg.clone();
if let Some(model) = model.map(str::trim).filter(|model| !model.is_empty()) {
cfg.model = Some(model.to_string());
cfg.strict_model = true;
}
cfg
}
async fn runtime_session_for(&self, session_id: &str) -> String {
let mut sessions = self.runtime_sessions.lock().await;
if let Some(runtime_session) = sessions.get(session_id) {
return runtime_session.clone();
}
let runtime_session = self.runtime.open_session().await;
sessions.insert(session_id.to_string(), runtime_session.clone());
runtime_session
}
async fn reconcile_dangling_actions(
&self,
session_id: &str,
runtime_session: &str,
messages: &mut Vec<Message>,
) -> Result<(), String> {
let Some(store) = &self.durability else {
return Ok(());
};
let answered: std::collections::HashSet<String> = messages
.iter()
.filter_map(|message| match message {
Message::ToolResult { tool_use_id, .. } => Some(tool_use_id.clone()),
_ => None,
})
.collect();
let calls: Vec<ToolCall> = messages
.iter()
.flat_map(|message| match message {
Message::Assistant { tool_calls, .. } => tool_calls.clone(),
_ => Vec::new(),
})
.filter(|call| call.id.as_ref().is_some_and(|id| !answered.contains(id)))
.collect();
for call in calls {
let call_id = call.id.clone().expect("filtered to calls with ids");
let params = serde_json::to_value(&call.arguments).unwrap_or(Value::Null);
let Some(scope) = action_scope(self.repository_root.as_ref(), &call.name, ¶ms)
else {
messages.push(Message::ToolResult {
tool_use_id: call_id,
content: json!({"error": "tool call was interrupted before a durable action scope existed; not replayed"}).to_string(),
provenance: Provenance::Internal,
});
continue;
};
let candidate =
super::governance::SupervisedActionRecord::propose(session_id, &call_id, scope);
let Some(mut record) = store.load_action(&candidate.id).await? else {
messages.push(Message::ToolResult {
tool_use_id: call_id,
content:
json!({"error": "tool call was interrupted before approval; not replayed"})
.to_string(),
provenance: Provenance::Internal,
});
continue;
};
let content = match record.state {
super::governance::ActionState::Approved => {
record.transition(super::governance::ActionState::Dispatched, None)?;
store.record_action(&record).await?;
let proposal: ActionProposal = serde_json::from_value(json!({
"source": "durable-resume",
"actions": [{
"id": call_id,
"type": "tool_call",
"tool": call.name,
"parameters": params,
}],
}))
.map_err(|e| format!("cannot rebuild approved action on resume: {e}"))?;
let exec = self
.runtime
.execute_with_session(&proposal, runtime_session)
.await;
let result = exec.results.first();
let ok = result.is_some_and(|result| {
matches!(result.status, ActionStatus::Succeeded)
&& (call.name != "shell"
|| result
.output
.as_ref()
.and_then(|output| output.get("exit_code"))
.and_then(Value::as_i64)
== Some(0))
});
let receipt = json!({
"ok": ok,
"action_id": result.map(|result| result.action_id.clone()),
"output": result.and_then(|result| result.output.clone()),
});
record.transition(
if ok {
super::governance::ActionState::Completed
} else {
super::governance::ActionState::Failed
},
Some(receipt.clone()),
)?;
store.record_action(&record).await?;
receipt.to_string()
}
super::governance::ActionState::Dispatched => {
record.transition(
super::governance::ActionState::Indeterminate,
Some(json!({"reason": "process restarted after dispatch without a terminal receipt"})),
)?;
store.record_action(&record).await?;
json!({"error": "action outcome is indeterminate after restart; it was not replayed"}).to_string()
}
super::governance::ActionState::Completed
| super::governance::ActionState::Failed => record
.receipt
.clone()
.unwrap_or_else(|| json!({"status": format!("{:?}", record.state)}))
.to_string(),
super::governance::ActionState::Proposed => {
json!({"error": "approval was interrupted; action was not dispatched"})
.to_string()
}
super::governance::ActionState::Denied
| super::governance::ActionState::Indeterminate => {
json!({"error": format!("durable action is {:?}; not replayed", record.state)})
.to_string()
}
};
messages.push(Message::ToolResult {
tool_use_id: call_id,
content,
provenance: Provenance::Internal,
});
}
Ok(())
}
pub fn cancel(&self, session_id: &str) {
if let Ok(g) = self.cancels.lock() {
if let Some(flag) = g.get(session_id) {
flag.store(true, Ordering::Relaxed);
}
}
}
pub fn resolve_approval(&self, approval_id: &str, approved: bool) -> bool {
let tx = self
.approvals
.lock()
.ok()
.and_then(|mut g| g.remove(approval_id));
match tx {
Some(tx) => tx.send(approved).is_ok(),
None => false,
}
}
pub async fn handle_turn<E, Fut>(
&self,
session_id: &str,
prompt: &str,
attachments: Option<Vec<Value>>,
emit: E,
) where
E: Fn(Value) -> Fut + Send + Sync + 'static,
Fut: Future<Output = ()> + Send + 'static,
{
self.handle_turn_with_model(session_id, prompt, attachments, None, emit)
.await;
}
pub async fn handle_turn_with_model<E, Fut>(
&self,
session_id: &str,
prompt: &str,
attachments: Option<Vec<Value>>,
model: Option<&str>,
emit: E,
) where
E: Fn(Value) -> Fut + Send + Sync + 'static,
Fut: Future<Output = ()> + Send + 'static,
{
let cfg = self.config_for_model(model);
let runtime_session = self.runtime_session_for(session_id).await;
let images: Vec<ContentBlock> = attachments
.unwrap_or_default()
.into_iter()
.filter_map(|a| serde_json::from_value::<ContentBlock>(a).ok())
.filter(|c| {
matches!(
c,
ContentBlock::ImageBase64 { .. } | ContentBlock::ImageUrl { .. }
)
})
.collect();
let cancel = Arc::new(AtomicBool::new(false));
if let Ok(mut g) = self.cancels.lock() {
g.insert(session_id.to_string(), cancel.clone());
}
let cached = self.threads.lock().await.get(session_id).cloned();
let mut messages = match cached {
Some(existing) => existing,
None => {
let restored = match &self.durability {
Some(store) => match store.load_checkpoint(session_id).await {
Ok(checkpoint) => checkpoint.map(|checkpoint| checkpoint.messages),
Err(e) => {
emit(json!({
"kind": "error",
"error": format!("durable transcript resume failed: {e}"),
"session_id": session_id,
}))
.await;
if let Ok(mut g) = self.cancels.lock() {
g.remove(session_id);
}
return;
}
},
None => None,
};
let seeded = restored.unwrap_or_else(|| {
vec![Message::System {
content: self.system.clone(),
}]
});
self.threads
.lock()
.await
.insert(session_id.to_string(), seeded.clone());
seeded
}
};
if let Err(e) = self
.reconcile_dangling_actions(session_id, &runtime_session, &mut messages)
.await
{
emit(json!({
"kind": "error",
"error": format!("durable action reconciliation failed: {e}"),
"session_id": session_id,
}))
.await;
return;
}
messages.push(Message::User {
content: prompt.to_string(),
});
if let Some(store) = &self.durability {
if let Err(e) = store
.checkpoint(session_id, &messages, "user_turn", None)
.await
{
emit(json!({
"kind": "error",
"error": format!("durable checkpoint failed before inference: {e}"),
"session_id": session_id,
}))
.await;
return;
}
}
let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel::<Value>();
let drain = tokio::spawn(async move {
while let Some(v) = rx.recv().await {
emit(v).await;
}
});
let sid = session_id.to_string();
let tx_term = tx.clone();
let gate = ChatApprovalGate {
session_id: sid.clone(),
tx: tx.clone(),
approvals: self.approvals.clone(),
counter: Arc::new(AtomicU64::new(0)),
durability: self.durability.clone(),
repository_root: self.repository_root.clone(),
};
let outcome = run_assistant_loop_cancellable_in_session_durable(
&*self.generator,
&self.runtime,
&cfg,
&mut messages,
&cancel,
Some(&gate),
if images.is_empty() {
None
} else {
Some(images.as_slice())
},
Some(&runtime_session),
Some(session_id),
self.durability.as_deref(),
true,
{
let sid = sid.clone();
move |ev| {
if matches!(ev, AssistantEvent::Done { .. } | AssistantEvent::Error(_)) {
return;
}
if let Some(mut payload) = event_to_wire(ev) {
payload["session_id"] = json!(sid);
let _ = tx.send(payload);
}
}
},
)
.await;
drop(gate);
let completion = super::governance::completion_matrix_from_messages(&messages);
let ungrounded_claims =
super::agent_loop::ungrounded_summary_claims(&outcome.summary, &outcome.tool_receipts);
let _ = tx_term.send(json!({
"kind": "receipt_report",
"completion": completion,
"ungrounded_claims": ungrounded_claims,
"session_id": sid,
}));
let terminal_summary = super::agent_loop::annotate_summary_with_claim_note(
&outcome.summary,
&ungrounded_claims,
);
let terminal = match outcome.status {
"cancelled" | "error" => {
json!({ "kind": "error", "error": terminal_summary, "session_id": sid })
}
_ => json!({ "kind": "done", "text": terminal_summary, "session_id": sid }),
};
let _ = tx_term.send(terminal);
drop(tx_term);
let _ = drain.await;
if outcome.status != "cancelled" {
let mut g = self.threads.lock().await;
g.insert(session_id.to_string(), messages);
}
if let Ok(mut g) = self.cancels.lock() {
g.remove(session_id);
}
}
pub async fn handle_goal_turn<E, Fut>(
&self,
session_id: &str,
prompt: &str,
_attachments: Option<Vec<Value>>,
goal: ChatGoal,
emit: E,
) where
E: Fn(Value) -> Fut + Send + Sync + 'static,
Fut: Future<Output = ()> + Send + 'static,
{
self.handle_goal_turn_with_model(session_id, prompt, _attachments, goal, None, emit)
.await;
}
pub async fn handle_goal_turn_with_model<E, Fut>(
&self,
session_id: &str,
prompt: &str,
_attachments: Option<Vec<Value>>,
goal: ChatGoal,
model: Option<&str>,
emit: E,
) where
E: Fn(Value) -> Fut + Send + Sync + 'static,
Fut: Future<Output = ()> + Send + 'static,
{
let cfg = self.config_for_model(model);
use car_verify::goal::{GoalCondition, GoalGovernor, GoalSpec, GoalStatus};
let runtime_session = self.runtime_session_for(session_id).await;
let cancel = Arc::new(AtomicBool::new(false));
if let Ok(mut g) = self.cancels.lock() {
g.insert(session_id.to_string(), cancel.clone());
}
let mut messages = match &self.durability {
Some(store) => match store.load_checkpoint(session_id).await {
Ok(Some(checkpoint)) => checkpoint.messages,
Ok(None) => vec![Message::System {
content: format!(
"{}\n\nYou are working toward a goal. Completion is verified \
deterministically by running this shell command:\n {}\nIt is \
done only when that command exits 0. Keep working until it does.",
self.system, goal.check
),
}],
Err(e) => {
emit(json!({
"kind": "error",
"error": format!("durable goal resume failed: {e}"),
"session_id": session_id,
}))
.await;
return;
}
},
None => vec![Message::System {
content: format!(
"{}\n\nYou are working toward a goal. Completion is verified \
deterministically by running this shell command:\n {}\nIt is \
done only when that command exits 0. Keep working until it does.",
self.system, goal.check
),
}],
};
if let Err(e) = self
.reconcile_dangling_actions(session_id, &runtime_session, &mut messages)
.await
{
emit(json!({
"kind": "error",
"error": format!("durable goal action reconciliation failed: {e}"),
"session_id": session_id,
}))
.await;
return;
}
let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel::<Value>();
let drain = tokio::spawn(async move {
while let Some(v) = rx.recv().await {
emit(v).await;
}
});
let sid = session_id.to_string();
let tx_term = tx.clone();
let gate = ChatApprovalGate {
session_id: sid.clone(),
tx: tx.clone(),
approvals: self.approvals.clone(),
counter: Arc::new(AtomicU64::new(0)),
durability: self.durability.clone(),
repository_root: self.repository_root.clone(),
};
let check_gate = gate.clone();
let spec = GoalSpec {
goal: prompt.to_string(),
condition: GoalCondition::Command {
id: "goal_check".into(),
expect_exit: 0,
},
governor: GoalGovernor {
max_turns: Some(goal.max_iterations.max(1)),
..Default::default()
},
};
let check = goal.check.clone();
let check_cfg = cfg.clone();
let result = run_assistant_goal_loop_in_session_durable(
&*self.generator,
&self.runtime,
&cfg,
&mut messages,
&cancel,
Some(&gate),
&spec,
Some(&runtime_session),
Some(session_id),
self.durability.as_deref(),
move |_outcome| {
let cmd = check.clone();
let check_gate = check_gate.clone();
let check_cfg = check_cfg.clone();
async move {
let exit = run_shell_check_with_approval(
&self.runtime,
&check_cfg,
Some(&check_gate),
&cmd,
)
.await;
let mut g = car_engine::GoalGather::default();
g.command_exits.insert("goal_check".into(), exit);
g
}
},
{
let sid = sid.clone();
move |ev| {
if matches!(ev, AssistantEvent::Done { .. } | AssistantEvent::Error(_)) {
return;
}
if let Some(mut payload) = event_to_wire(ev) {
payload["session_id"] = json!(sid);
let _ = tx.send(payload);
}
}
},
)
.await;
drop(gate);
let completion = super::governance::completion_matrix_from_messages(&messages);
let _ = tx_term.send(json!({
"kind": "receipt_report",
"completion": completion,
"session_id": sid,
}));
let terminal = match result.run.status {
GoalStatus::Achieved => {
json!({ "kind": "done", "text": result.outcome.summary, "session_id": sid })
}
GoalStatus::Halted { halt } => json!({
"kind": "error",
"error": format!(
"goal not reached: {} after {} iteration(s); last check: {}",
halt.as_str(),
result.run.iterations,
result.run.last_reason
),
"session_id": sid,
}),
};
let _ = tx_term.send(terminal);
drop(tx_term);
let _ = drain.await;
if !matches!(result.run.status, GoalStatus::Halted { .. }) {
let mut g = self.threads.lock().await;
g.insert(session_id.to_string(), messages);
}
if let Ok(mut g) = self.cancels.lock() {
g.remove(session_id);
}
}
}
fn action_scope(
repository_root: Option<&PathBuf>,
tool: &str,
params: &Value,
) -> Option<super::governance::ActionScope> {
let repository_root = repository_root?.clone();
let command = params
.get("command")
.and_then(Value::as_str)
.unwrap_or_default();
let target = params
.get("target")
.or_else(|| params.get("url"))
.or_else(|| params.get("path"))
.and_then(Value::as_str)
.unwrap_or(command)
.to_string();
let environment = params
.get("environment")
.and_then(Value::as_str)
.unwrap_or("unspecified")
.to_string();
let lower = format!("{tool} {command}").to_ascii_lowercase();
let mut capabilities = Vec::new();
if lower.contains("git push") {
capabilities.push(super::governance::CredentialCapability(
"git:configured-remote".into(),
));
}
if lower.contains("az ") || lower.contains("azure") {
capabilities.push(super::governance::CredentialCapability(
"azure:active-account".into(),
));
}
if lower.contains("sql") || lower.contains("database") || lower.contains("migration") {
capabilities.push(super::governance::CredentialCapability(
"database:project-configured".into(),
));
}
Some(super::governance::ActionScope {
tool: tool.to_string(),
parameters: params.clone(),
repository_root,
target,
environment,
credential_capabilities: capabilities,
})
}
#[derive(Clone)]
struct ChatApprovalGate {
session_id: String,
tx: mpsc::UnboundedSender<Value>,
approvals: Arc<StdMutex<HashMap<String, oneshot::Sender<bool>>>>,
counter: Arc<AtomicU64>,
durability: Option<Arc<dyn AssistantDurability>>,
repository_root: Option<PathBuf>,
}
impl ChatApprovalGate {
fn action_scope(&self, tool: &str, params: &Value) -> Option<super::governance::ActionScope> {
action_scope(self.repository_root.as_ref(), tool, params)
}
async fn durable_action(
&self,
call_id: &str,
tool: &str,
params: &Value,
) -> Result<Option<super::governance::SupervisedActionRecord>, String> {
let (Some(store), Some(scope)) = (&self.durability, self.action_scope(tool, params)) else {
return Ok(None);
};
let action =
super::governance::SupervisedActionRecord::propose(&self.session_id, call_id, scope);
Ok(store.load_action(&action.id).await?.or(Some(action)))
}
}
#[async_trait::async_trait]
impl ApprovalGate for ChatApprovalGate {
async fn request(&self, tool: &str, params: &Value) -> ApprovalDecision {
let n = self.counter.fetch_add(1, Ordering::Relaxed);
self.request_action(&format!("unbound-{n}"), tool, params)
.await
}
async fn request_action(&self, call_id: &str, tool: &str, params: &Value) -> ApprovalDecision {
let mut action = match self.durable_action(call_id, tool, params).await {
Ok(action) => action,
Err(e) => {
return ApprovalDecision::Denied(format!("cannot persist approval scope: {e}"))
}
};
if let Some(existing) = &action {
match existing.state {
super::governance::ActionState::Approved => return ApprovalDecision::Approved,
super::governance::ActionState::Dispatched => {
let mut indeterminate = existing.clone();
let _ = indeterminate.transition(
super::governance::ActionState::Indeterminate,
Some(json!({"reason": "resumed after dispatch without terminal receipt"})),
);
if let Some(store) = &self.durability {
let _ = store.record_action(&indeterminate).await;
}
return ApprovalDecision::Denied(
"action was dispatched before restart and is indeterminate; reconcile it before retrying".into(),
);
}
super::governance::ActionState::Completed
| super::governance::ActionState::Failed
| super::governance::ActionState::Denied
| super::governance::ActionState::Indeterminate => {
return ApprovalDecision::Denied(
"this durable action identity is terminal and cannot be replayed".into(),
);
}
super::governance::ActionState::Proposed => {}
}
}
if let (Some(store), Some(proposed)) = (&self.durability, &action) {
if store
.load_action(&proposed.id)
.await
.ok()
.flatten()
.is_none()
{
if let Err(e) = store.record_action(proposed).await {
return ApprovalDecision::Denied(format!("cannot record action proposal: {e}"));
}
}
}
let n = self.counter.fetch_add(1, Ordering::Relaxed);
let approval_id = format!("{}-appr-{n}", self.session_id);
let (otx, orx) = oneshot::channel();
if let Ok(mut g) = self.approvals.lock() {
g.insert(approval_id.clone(), otx);
}
let _ = self.tx.send(json!({
"kind": "approval_pending",
"approval_id": approval_id,
"tool": tool,
"params": params,
"session_id": self.session_id,
"action_id": action.as_ref().map(|record| record.id.clone()),
"scope": action.as_ref().map(|record| record.scope.clone()),
}));
let decision = match tokio::time::timeout(APPROVAL_TIMEOUT, orx).await {
Ok(Ok(true)) => ApprovalDecision::Approved,
Ok(Ok(false)) => ApprovalDecision::Denied("declined by user".into()),
_ => {
if let Ok(mut g) = self.approvals.lock() {
g.remove(&approval_id);
}
ApprovalDecision::Denied("approval timed out".into())
}
};
if let (Some(store), Some(record)) = (&self.durability, action.as_mut()) {
let next = match &decision {
ApprovalDecision::Approved => super::governance::ActionState::Approved,
ApprovalDecision::Denied(_) => super::governance::ActionState::Denied,
};
if let Err(e) = record.transition(next, Some(json!({ "approval_id": approval_id }))) {
return ApprovalDecision::Denied(format!("cannot durably record approval: {e}"));
}
if let Err(e) = store.record_action(record).await {
return ApprovalDecision::Denied(format!("cannot durably record approval: {e}"));
}
}
decision
}
async fn before_dispatch(
&self,
call_id: &str,
tool: &str,
params: &Value,
) -> Result<(), String> {
let Some(mut record) = self.durable_action(call_id, tool, params).await? else {
return Ok(());
};
if record.state != super::governance::ActionState::Approved {
return Err(format!(
"action {} is {:?}, not approved",
record.id, record.state
));
}
record.transition(super::governance::ActionState::Dispatched, None)?;
self.durability
.as_ref()
.expect("durable action has a store")
.record_action(&record)
.await
}
async fn after_dispatch(
&self,
call_id: &str,
tool: &str,
params: &Value,
ok: bool,
receipt: &Value,
) -> Result<(), String> {
let Some(mut record) = self.durable_action(call_id, tool, params).await? else {
return Ok(());
};
record.transition(
if ok {
super::governance::ActionState::Completed
} else {
super::governance::ActionState::Failed
},
Some(receipt.clone()),
)?;
self.durability
.as_ref()
.expect("durable action has a store")
.record_action(&record)
.await
}
}
fn event_to_wire(ev: AssistantEvent) -> Option<Value> {
match ev {
AssistantEvent::Text(t) => Some(json!({ "kind": "token", "delta": t })),
AssistantEvent::ToolCall { name, params } => {
Some(json!({ "kind": "tool_call", "tool": name, "params": params }))
}
AssistantEvent::ToolResult { .. } => None,
AssistantEvent::Done { text } => Some(json!({ "kind": "done", "text": text })),
AssistantEvent::Error(e) => Some(json!({ "kind": "error", "error": e })),
AssistantEvent::GoalEvaluated {
iteration,
met,
grounded,
reason,
} => Some(json!({
"kind": "goal_evaluated",
"iteration": iteration,
"met": met,
"grounded": grounded,
"reason": reason,
})),
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::assistant::executor::GeneralExecutor;
use async_trait::async_trait;
use car_engine::{LocalSubstrate, Substrate, ToolExecutor};
use car_inference::{GenerateRequest, InferenceEngine, InferenceResult};
use std::sync::atomic::AtomicUsize;
fn turn(text: &str, tool_calls: Value) -> InferenceResult {
serde_json::from_value(json!({
"text": text, "tool_calls": tool_calls,
"trace_id": "t", "model_used": "scripted", "latency_ms": 0,
}))
.unwrap()
}
struct Script {
turns: Vec<InferenceResult>,
cursor: AtomicUsize,
}
#[async_trait]
impl TurnGenerator for Script {
async fn generate(&self, _r: GenerateRequest) -> Result<InferenceResult, String> {
let i = self.cursor.fetch_add(1, Ordering::SeqCst);
self.turns.get(i).cloned().ok_or("exhausted".into())
}
}
struct RecordingScript {
requests: Arc<StdMutex<Vec<GenerateRequest>>>,
}
#[async_trait]
impl TurnGenerator for RecordingScript {
async fn generate(&self, request: GenerateRequest) -> Result<InferenceResult, String> {
self.requests.lock().unwrap().push(request);
Ok(turn("done", json!([])))
}
}
async fn runtime(dir: &std::path::Path) -> Arc<Runtime> {
let substrate: Arc<dyn Substrate> = Arc::new(LocalSubstrate::new());
let exec: Arc<dyn ToolExecutor> =
Arc::new(GeneralExecutor::new(substrate.clone(), dir, true));
let rt = Runtime::new()
.with_inference(Arc::new(InferenceEngine::new(Default::default())))
.with_executor(exec)
.with_substrate(substrate);
rt.register_agent_basics().await;
rt.register_tool_entry(
car_engine::ToolEntry::new(car_ir::builtins::shell()).with_side_effects(true),
)
.await;
Arc::new(rt)
}
struct FixedApproval(bool);
#[async_trait]
impl ApprovalGate for FixedApproval {
async fn request(&self, _tool: &str, _params: &Value) -> ApprovalDecision {
if self.0 {
ApprovalDecision::Approved
} else {
ApprovalDecision::Denied("denied".into())
}
}
}
#[derive(Default)]
struct MemoryDurability {
actions: AsyncMutex<HashMap<String, super::super::governance::SupervisedActionRecord>>,
checkpoints: AsyncMutex<HashMap<String, super::super::governance::AssistantCheckpoint>>,
}
#[async_trait]
impl AssistantDurability for MemoryDurability {
async fn load_checkpoint(
&self,
session_id: &str,
) -> Result<Option<super::super::governance::AssistantCheckpoint>, String> {
Ok(self.checkpoints.lock().await.get(session_id).cloned())
}
async fn checkpoint(
&self,
session_id: &str,
messages: &[Message],
reason: &str,
goal: Option<Value>,
) -> Result<(), String> {
let mut checkpoints = self.checkpoints.lock().await;
let revision = checkpoints
.get(session_id)
.map(|checkpoint| checkpoint.revision + 1)
.unwrap_or(1);
checkpoints.insert(
session_id.to_string(),
super::super::governance::AssistantCheckpoint {
id: session_id.to_string(),
session_id: session_id.to_string(),
revision,
repository_root: PathBuf::from("/fixture/repo"),
messages: messages.to_vec(),
goal,
compaction: Some(json!({ "reason": reason })),
completion: super::super::governance::completion_matrix_from_messages(messages),
},
);
Ok(())
}
async fn load_action(
&self,
action_id: &str,
) -> Result<Option<super::super::governance::SupervisedActionRecord>, String> {
Ok(self.actions.lock().await.get(action_id).cloned())
}
async fn record_action(
&self,
record: &super::super::governance::SupervisedActionRecord,
) -> Result<(), String> {
self.actions
.lock()
.await
.insert(record.id.clone(), record.clone());
Ok(())
}
}
#[tokio::test]
async fn restarted_service_resumes_by_stable_host_session_id() {
let dir = tempfile::tempdir().unwrap();
let durability = Arc::new(MemoryDurability::default());
let cfg = AssistantConfig {
model: Some("scripted".into()),
strict_model: false,
max_turns: 2,
tools: GeneralExecutor::tool_defs(),
gated_tools: Vec::new(),
approval_policy: None,
proactive_memory: None,
tool_labels: None,
todos: None,
value_store_previews: false,
};
let first = AssistantService::new_durable(
Arc::new(Script {
turns: vec![turn("phase-one-evidence", json!([]))],
cursor: AtomicUsize::new(0),
}),
runtime(dir.path()).await,
cfg.clone(),
"sys".into(),
durability.clone(),
dir.path().to_path_buf(),
);
first
.handle_turn("stable-host-session", "investigate", None, |_| async {})
.await;
drop(first);
let second = AssistantService::new_durable(
Arc::new(Script {
turns: vec![turn("continuity-confirmed", json!([]))],
cursor: AtomicUsize::new(0),
}),
runtime(dir.path()).await,
cfg,
"sys".into(),
durability.clone(),
dir.path().to_path_buf(),
);
second
.handle_turn(
"stable-host-session",
"continue without repeating",
None,
|_| async {},
)
.await;
let checkpoints = durability.checkpoints.lock().await;
assert_eq!(
checkpoints.len(),
1,
"runtime UUIDs must not become checkpoint keys"
);
let resumed = checkpoints
.get("stable-host-session")
.expect("stable session checkpoint");
let transcript = serde_json::to_string(&resumed.messages).unwrap();
assert!(transcript.contains("phase-one-evidence"));
assert!(transcript.contains("continue without repeating"));
assert!(transcript.contains("continuity-confirmed"));
}
#[tokio::test]
async fn unsupported_final_claim_is_redriven_before_done() {
let dir = tempfile::tempdir().unwrap();
let script = Arc::new(Script {
turns: vec![
turn("The repository is clean.", json!([])),
turn(
"No git status receipt is available, so repository state remains unknown.",
json!([]),
),
],
cursor: AtomicUsize::new(0),
});
let cfg = AssistantConfig {
model: Some("scripted".into()),
strict_model: false,
max_turns: 3,
tools: GeneralExecutor::tool_defs(),
gated_tools: Vec::new(),
approval_policy: None,
proactive_memory: None,
tool_labels: None,
todos: None,
value_store_previews: false,
};
let service =
AssistantService::new(script.clone(), runtime(dir.path()).await, cfg, "sys".into());
let events = Arc::new(StdMutex::new(Vec::<Value>::new()));
let captured = events.clone();
service
.handle_turn("claims", "inspect", None, move |event| {
let captured = captured.clone();
async move { captured.lock().unwrap().push(event) }
})
.await;
assert_eq!(script.cursor.load(Ordering::SeqCst), 2);
let events = events.lock().unwrap();
let done = events.iter().find(|event| event["kind"] == "done").unwrap();
assert!(done["text"].as_str().unwrap().contains("remains unknown"));
assert!(!done["text"].as_str().unwrap().contains("[claim check]"));
}
#[tokio::test]
async fn scoped_approval_is_durable_exact_and_auditable() {
let repo = tempfile::tempdir().unwrap();
std::fs::create_dir(repo.path().join(".git")).unwrap();
let durability = Arc::new(MemoryDurability::default());
let approvals = Arc::new(StdMutex::new(HashMap::new()));
let (tx, mut rx) = mpsc::unbounded_channel();
let gate = ChatApprovalGate {
session_id: "s1".into(),
tx,
approvals: approvals.clone(),
counter: Arc::new(AtomicU64::new(0)),
durability: Some(durability.clone()),
repository_root: Some(repo.path().to_path_buf()),
};
let params = json!({
"command": "git push origin HEAD:main",
"target": "origin/main",
"environment": "fixture"
});
let pending_gate = gate.clone();
let pending_params = params.clone();
let pending = tokio::spawn(async move {
pending_gate
.request_action("call-1", "shell", &pending_params)
.await
});
let event = rx.recv().await.expect("approval event");
assert_eq!(event["kind"], "approval_pending");
assert_eq!(event["scope"]["target"], "origin/main");
assert_eq!(event["scope"]["environment"], "fixture");
assert_eq!(
event["scope"]["credential_capabilities"][0],
"git:configured-remote"
);
let approval_id = event["approval_id"].as_str().unwrap();
approvals
.lock()
.unwrap()
.remove(approval_id)
.unwrap()
.send(true)
.unwrap();
assert!(matches!(pending.await.unwrap(), ApprovalDecision::Approved));
gate.before_dispatch("call-1", "shell", ¶ms)
.await
.unwrap();
gate.after_dispatch(
"call-1",
"shell",
¶ms,
true,
&json!({"remote_sha": "abc"}),
)
.await
.unwrap();
let action_id = event["action_id"].as_str().unwrap();
let action = durability.load_action(action_id).await.unwrap().unwrap();
assert_eq!(
action.state,
super::super::governance::ActionState::Completed
);
let changed = json!({
"command": "git push origin HEAD:other",
"target": "origin/other",
"environment": "fixture"
});
assert!(gate
.before_dispatch("call-1", "shell", &changed)
.await
.is_err());
let denied_gate = gate.clone();
let denied_params = changed.clone();
let denied = tokio::spawn(async move {
denied_gate
.request_action("call-2", "shell", &denied_params)
.await
});
let denied_event = rx.recv().await.expect("denial approval event");
let denied_id = denied_event["approval_id"].as_str().unwrap();
approvals
.lock()
.unwrap()
.remove(denied_id)
.unwrap()
.send(false)
.unwrap();
assert!(matches!(denied.await.unwrap(), ApprovalDecision::Denied(_)));
let denied_action = durability
.load_action(denied_event["action_id"].as_str().unwrap())
.await
.unwrap()
.unwrap();
assert_eq!(
denied_action.state,
super::super::governance::ActionState::Denied
);
assert!(denied_action.receipt.is_some(), "denial must be auditable");
assert!(gate
.before_dispatch("call-2", "shell", &changed)
.await
.is_err());
}
fn dangling_shell(call_id: &str, command: &str) -> Vec<Message> {
vec![
Message::System {
content: "sys".into(),
},
Message::User {
content: "do it".into(),
},
Message::Assistant {
content: String::new(),
tool_calls: vec![serde_json::from_value(json!({
"id": call_id,
"name": "shell",
"arguments": {"command": command},
}))
.unwrap()],
thinking: vec![],
},
]
}
#[tokio::test]
async fn restart_before_dispatch_runs_approved_action_once() {
let repo = tempfile::tempdir().unwrap();
std::fs::create_dir(repo.path().join(".git")).unwrap();
let rt = runtime(repo.path()).await;
let durability = Arc::new(MemoryDurability::default());
let generator: Arc<dyn TurnGenerator> = Arc::new(Script {
turns: vec![],
cursor: AtomicUsize::new(0),
});
let service = AssistantService::new_durable(
generator,
rt,
test_cfg_with_gated_shell(),
"sys".into(),
durability.clone(),
repo.path().to_path_buf(),
);
let command = "printf x >> effect.txt";
let params = json!({"command": command});
let scope = action_scope(Some(&repo.path().to_path_buf()), "shell", ¶ms).unwrap();
let mut action = super::super::governance::SupervisedActionRecord::propose(
"restart-before",
"call-1",
scope,
);
action
.transition(super::super::governance::ActionState::Approved, None)
.unwrap();
durability.record_action(&action).await.unwrap();
let mut messages = dangling_shell("call-1", command);
let runtime_session = service.runtime_session_for("restart-before").await;
service
.reconcile_dangling_actions("restart-before", &runtime_session, &mut messages)
.await
.unwrap();
service
.reconcile_dangling_actions("restart-before", &runtime_session, &mut messages)
.await
.unwrap();
assert_eq!(
std::fs::read_to_string(repo.path().join("effect.txt")).unwrap(),
"x",
"the approved effect must execute exactly once"
);
let recovered = durability.load_action(&action.id).await.unwrap().unwrap();
assert_eq!(
recovered.state,
super::super::governance::ActionState::Completed
);
}
#[tokio::test]
async fn restart_after_dispatch_marks_indeterminate_without_replay() {
let repo = tempfile::tempdir().unwrap();
std::fs::create_dir(repo.path().join(".git")).unwrap();
let rt = runtime(repo.path()).await;
let durability = Arc::new(MemoryDurability::default());
let generator: Arc<dyn TurnGenerator> = Arc::new(Script {
turns: vec![],
cursor: AtomicUsize::new(0),
});
let service = AssistantService::new_durable(
generator,
rt,
test_cfg_with_gated_shell(),
"sys".into(),
durability.clone(),
repo.path().to_path_buf(),
);
let command = "printf x >> must-not-exist.txt";
let params = json!({"command": command});
let scope = action_scope(Some(&repo.path().to_path_buf()), "shell", ¶ms).unwrap();
let mut action = super::super::governance::SupervisedActionRecord::propose(
"restart-after",
"call-2",
scope,
);
action
.transition(super::super::governance::ActionState::Approved, None)
.unwrap();
action
.transition(super::super::governance::ActionState::Dispatched, None)
.unwrap();
durability.record_action(&action).await.unwrap();
let mut messages = dangling_shell("call-2", command);
let runtime_session = service.runtime_session_for("restart-after").await;
service
.reconcile_dangling_actions("restart-after", &runtime_session, &mut messages)
.await
.unwrap();
assert!(!repo.path().join("must-not-exist.txt").exists());
let recovered = durability.load_action(&action.id).await.unwrap().unwrap();
assert_eq!(
recovered.state,
super::super::governance::ActionState::Indeterminate
);
}
fn test_cfg_with_gated_shell() -> AssistantConfig {
AssistantConfig {
model: Some("scripted".into()),
strict_model: false,
max_turns: 4,
tools: GeneralExecutor::tool_defs(),
gated_tools: vec!["shell".into()],
approval_policy: None,
proactive_memory: None,
tool_labels: None,
todos: None,
value_store_previews: false,
}
}
#[tokio::test]
async fn goal_shell_check_does_not_run_without_required_approval() {
let dir = tempfile::tempdir().unwrap();
let rt = runtime(dir.path()).await;
let cfg = test_cfg_with_gated_shell();
let target = dir.path().join("should-not-exist");
let exit = run_shell_check_with_approval(
&rt,
&cfg,
None,
&crate::coder::test_cmds::touch("should-not-exist"),
)
.await;
assert_eq!(exit, 1);
assert!(
!target.exists(),
"gated goal verifier command must not run without approval"
);
}
#[tokio::test]
async fn goal_shell_check_runs_after_required_approval() {
let dir = tempfile::tempdir().unwrap();
let rt = runtime(dir.path()).await;
let cfg = test_cfg_with_gated_shell();
let target = dir.path().join("approved-check");
let exit = run_shell_check_with_approval(
&rt,
&cfg,
Some(&FixedApproval(true)),
&crate::coder::test_cmds::touch("approved-check"),
)
.await;
assert_eq!(exit, 0);
assert!(target.exists(), "approved verifier command should run");
}
#[tokio::test]
async fn chat_turn_streams_tokens_and_done() {
let dir = tempfile::tempdir().unwrap();
let rt = runtime(dir.path()).await;
let generator: Arc<dyn TurnGenerator> = Arc::new(Script {
turns: vec![
turn(
"let me compute",
json!([{ "id": "c1", "name": "calculate", "arguments": { "expression": "2+2" } }]),
),
turn("It's 4.", json!([])),
],
cursor: AtomicUsize::new(0),
});
let cfg = AssistantConfig {
model: Some("scripted".into()),
strict_model: false,
max_turns: 4,
tools: GeneralExecutor::tool_defs(),
gated_tools: Vec::new(),
approval_policy: None,
proactive_memory: None,
tool_labels: None,
todos: None,
value_store_previews: false,
};
let svc = AssistantService::new(generator, rt, cfg, "sys".into());
let events = Arc::new(StdMutex::new(Vec::<Value>::new()));
let ev2 = events.clone();
svc.handle_turn("s1", "what is 2+2?", None, move |v| {
let ev = ev2.clone();
async move {
ev.lock().unwrap().push(v);
}
})
.await;
let got = events.lock().unwrap().clone();
assert!(got.iter().all(|e| e["session_id"] == "s1"));
assert!(got
.iter()
.any(|e| e["kind"] == "tool_call" && e["tool"] == "calculate"));
let last = got.last().unwrap();
assert_eq!(last["kind"], "done");
assert_eq!(last["text"], "It's 4.");
let thread_len = svc.threads.lock().await.get("s1").map(|m| m.len()).unwrap();
assert!(thread_len >= 3, "thread should persist across the turn");
}
#[tokio::test]
async fn every_chat_event_is_forwardable_by_the_daemon() {
let dir = tempfile::tempdir().unwrap();
let rt = runtime(dir.path()).await;
let generator: Arc<dyn TurnGenerator> = Arc::new(Script {
turns: vec![
turn(
"let me compute",
json!([{ "id": "c1", "name": "calculate", "arguments": { "expression": "1+1" } }]),
),
turn("It's 2.", json!([])),
],
cursor: AtomicUsize::new(0),
});
let cfg = AssistantConfig {
model: Some("scripted".into()),
strict_model: false,
max_turns: 4,
tools: GeneralExecutor::tool_defs(),
gated_tools: Vec::new(),
approval_policy: None,
proactive_memory: None,
tool_labels: None,
todos: None,
value_store_previews: false,
};
let svc = AssistantService::new(generator, rt, cfg, "sys".into());
let events = Arc::new(StdMutex::new(Vec::<Value>::new()));
let ev2 = events.clone();
svc.handle_turn("sess-42", "1+1?", None, move |v| {
let ev = ev2.clone();
async move {
ev.lock().unwrap().push(v);
}
})
.await;
const KNOWN_KINDS: [&str; 7] = [
"token",
"tool_call",
"approval_pending",
"goal_evaluated",
"receipt_report",
"done",
"error",
];
let got = events.lock().unwrap().clone();
assert!(!got.is_empty());
for e in &got {
assert_eq!(
e.get("session_id").and_then(Value::as_str),
Some("sess-42"),
"every event must carry its session_id: {e}"
);
let kind = e.get("kind").and_then(Value::as_str).unwrap_or("");
assert!(KNOWN_KINDS.contains(&kind), "unknown event kind: {e}");
}
assert_eq!(got.last().unwrap()["kind"], "done");
}
#[tokio::test]
async fn explicit_chat_model_reaches_inference_and_unset_preserves_agent_default() {
let dir = tempfile::tempdir().unwrap();
let rt = runtime(dir.path()).await;
let requests = Arc::new(StdMutex::new(Vec::new()));
let generator: Arc<dyn TurnGenerator> = Arc::new(RecordingScript {
requests: requests.clone(),
});
let cfg = AssistantConfig {
model: Some("agent/default".into()),
strict_model: false,
max_turns: 2,
tools: Vec::new(),
gated_tools: Vec::new(),
approval_policy: None,
proactive_memory: None,
tool_labels: None,
todos: None,
value_store_previews: false,
};
let svc = AssistantService::new(generator, rt, cfg, "sys".into());
svc.handle_turn_with_model(
"selected",
"hello",
None,
Some("openrouter/deepseek/deepseek-v3.2"),
|_| async {},
)
.await;
svc.handle_turn("adaptive", "hello", None, |_| async {})
.await;
let got = requests.lock().unwrap();
assert_eq!(got.len(), 2);
assert_eq!(
got[0].model.as_deref(),
Some("openrouter/deepseek/deepseek-v3.2")
);
assert!(
got[0].params.strict_model,
"a selected native model must not silently fall back"
);
assert_eq!(got[1].model.as_deref(), Some("agent/default"));
assert!(
!got[1].params.strict_model,
"an unset native preference preserves the agent's routing policy"
);
}
#[test]
fn goal_evaluated_event_has_host_wire_shape() {
let wire = event_to_wire(AssistantEvent::GoalEvaluated {
iteration: 2,
met: false,
grounded: true,
reason: "command goal_check exited 1".into(),
})
.expect("goal verifier events should be surfaced to hosts");
assert_eq!(wire["kind"], "goal_evaluated");
assert_eq!(wire["iteration"], 2);
assert_eq!(wire["met"], false);
assert_eq!(wire["grounded"], true);
assert_eq!(wire["reason"], "command goal_check exited 1");
}
#[tokio::test]
async fn goal_turn_streams_verifier_events_and_one_terminal() {
let dir = tempfile::tempdir().unwrap();
let rt = runtime(dir.path()).await;
let create = crate::coder::test_cmds::touch("goal.done");
let generator: Arc<dyn TurnGenerator> = Arc::new(Script {
turns: vec![
turn("starting", json!([])),
turn(
"creating sentinel",
json!([{ "id": "s1", "name": "shell", "arguments": { "command": create } }]),
),
turn("done", json!([])),
],
cursor: AtomicUsize::new(0),
});
let cfg = AssistantConfig {
model: Some("scripted".into()),
strict_model: false,
max_turns: 4,
tools: GeneralExecutor::tool_defs(),
gated_tools: Vec::new(),
approval_policy: None,
proactive_memory: None,
tool_labels: None,
todos: None,
value_store_previews: false,
};
let svc = AssistantService::new(generator, rt, cfg, "sys".into());
let events = Arc::new(StdMutex::new(Vec::<Value>::new()));
let ev2 = events.clone();
svc.handle_goal_turn(
"goal-s1",
"create goal.done",
None,
ChatGoal {
check: crate::coder::test_cmds::file_exists("goal.done"),
max_iterations: 4,
},
move |v| {
let ev = ev2.clone();
async move {
ev.lock().unwrap().push(v);
}
},
)
.await;
let got = events.lock().unwrap().clone();
let verifier: Vec<_> = got
.iter()
.filter(|e| e["kind"] == "goal_evaluated")
.collect();
assert_eq!(verifier.len(), 2, "one verifier event per goal iteration");
assert_eq!(verifier[0]["met"], false);
assert_eq!(verifier[1]["met"], true);
assert_eq!(verifier[1]["grounded"], true);
assert_eq!(
got.iter().filter(|e| e["kind"] == "done").count(),
1,
"iteration-local done events must not leak as terminal chat events"
);
assert_eq!(got.last().unwrap()["kind"], "done");
assert_eq!(
std::fs::read_to_string(dir.path().join("goal.done")).unwrap_or_default(),
""
);
}
#[tokio::test]
async fn chat_gated_write_emits_approval_and_resumes_on_approve() {
let dir = tempfile::tempdir().unwrap();
let rt = runtime(dir.path()).await;
let generator: Arc<dyn TurnGenerator> = Arc::new(Script {
turns: vec![
turn(
"",
json!([{ "id": "w1", "name": "write_file", "arguments": { "path": "z.txt", "content": "zephyr" } }]),
),
turn("done", json!([])),
],
cursor: AtomicUsize::new(0),
});
let cfg = AssistantConfig {
model: Some("scripted".into()),
strict_model: false,
max_turns: 4,
tools: GeneralExecutor::tool_defs(),
gated_tools: vec!["write_file".into()],
approval_policy: None,
proactive_memory: None,
tool_labels: None,
todos: None,
value_store_previews: false,
};
let svc = Arc::new(AssistantService::new(generator, rt, cfg, "sys".into()));
let events = Arc::new(StdMutex::new(Vec::<Value>::new()));
let ev2 = events.clone();
let svc_run = svc.clone();
let turn_task = tokio::spawn(async move {
svc_run
.handle_turn("s1", "write z.txt", None, move |v| {
let ev = ev2.clone();
async move {
ev.lock().unwrap().push(v);
}
})
.await;
});
let approved = {
let mut ok = false;
for _ in 0..200 {
let id = events
.lock()
.unwrap()
.iter()
.find(|e| e["kind"] == "approval_pending")
.and_then(|e| e["approval_id"].as_str().map(String::from));
if let Some(id) = id {
assert!(svc.resolve_approval(&id, true));
ok = true;
break;
}
tokio::time::sleep(std::time::Duration::from_millis(10)).await;
}
ok
};
assert!(
approved,
"an approval_pending event should have been emitted"
);
turn_task.await.unwrap();
assert_eq!(
std::fs::read_to_string(dir.path().join("z.txt")).unwrap(),
"zephyr"
);
}
}