#![allow(clippy::unwrap_used, clippy::expect_used, clippy::panic)] use std::collections::VecDeque; use std::sync::{Arc, Mutex}; use tokio::sync::mpsc; use tokio_util::sync::CancellationToken; use tempfile::tempdir; use vak_agent::{Agent, AgentConfig, TurnOutcome}; use vak_hooks::HookDef; use vak_llm::stream; use vak_llm::types::{AssistantMessage, ChatRequest, ContentBlock, StopReason, Usage}; use vak_llm::{EventStream, LlmError, Provider}; use vak_session::SessionLog; use vak_session::types::{FrozenContract, SessionHeader}; use vak_tools::bash::BashTool; struct Scripted { responses: Mutex>, } #[async_trait::async_trait] impl Provider for Scripted { fn name(&self) -> &str { "scripted" } async fn stream( &self, _request: ChatRequest, _cancel: CancellationToken, ) -> Result { let next = self.responses.lock().unwrap().pop_front(); let (mut sink, rx) = stream::channel(64); match next { Some(m) => { sink.push(stream::StreamEvent::Start { partial: m.clone() }); sink.close_message(m).await; } None => sink.close_error(LlmError::Parse("exhausted".into())).await, } Ok(rx) } } fn text_msg(t: &str) -> AssistantMessage { AssistantMessage { content: vec![ContentBlock::text(t)], stop_reason: StopReason::EndTurn, usage: Usage { input_tokens: 1, output_tokens: 1, ..Default::default() }, model: "test-model".into(), response_id: None, } } fn bash_call(id: &str, cmd: &str) -> AssistantMessage { AssistantMessage { content: vec![ContentBlock::ToolUse { id: id.into(), name: "bash".into(), input: serde_json::json!({"command": cmd}), }], stop_reason: StopReason::ToolUse, usage: Usage::default(), model: "test-model".into(), response_id: None, } } fn build(responses: Vec, hooks: Option>) -> Agent { let dir = tempdir().unwrap(); let header = SessionHeader { agent: None, session_id: "s-hooks".into(), created_at: chrono::Utc::now(), cwd: dir.path().to_path_buf(), parent_session_id: None, contract_id: None, work_item_id: None, conversation: None, contract: FrozenContract { app_version: "0".into(), provider: "scripted".into(), model: "test-model".into(), route_ladder: Vec::new(), route_objective: String::new(), route_annotations: Vec::new(), system_prompt: "sys".into(), permission_mode: "full-access".into(), capabilities: Vec::new(), prompt_layers: Vec::new(), }, }; let log = SessionLog::create(dir.path().join("s.jsonl"), header).unwrap(); let mut cfg = AgentConfig::new("sys"); cfg.tools = vec![Arc::new(BashTool)]; cfg.mode = vak_permission::Mode::FullAccess; cfg.hooks = hooks.map(Arc::new); cfg.stop_policy = None; std::mem::forget(dir); Agent::new( Arc::new(Scripted { responses: Mutex::new(responses.into_iter().collect()), }), log, cfg, ) } #[tokio::test] async fn pre_tool_use_hook_blocks_execution() { let marker = tempdir().unwrap(); let marker_path = marker.path().join("ran"); let hooks = vec![HookDef { event: vak_hooks::HookEvent::PreToolUse, matcher: Some(vak_permission::Rule::parse("Bash(touch *)").unwrap()), command: r#"echo '{"decision":"block","reason":"no touching"}'"#.to_string(), timeout_ms: 5000, failure_mode: vak_hooks::HookFailureMode::Open, refusal: None, }]; let marker_path = marker_path.display().to_string(); let mut agent = build( vec![ bash_call("t1", &format!("touch {marker_path}")), text_msg("adapted"), ], Some(hooks), ); let outcome = agent .run( "go", &Default::default(), CancellationToken::new(), mpsc::channel(64).0, ) .await; assert!(matches!(outcome, TurnOutcome::Completed { .. })); assert!( !std::path::Path::new(&marker_path).exists(), "blocked command must never have executed" ); let session = agent.session.lock().await; // Raw ledger: the closed turn's result is a trace line in the // projection now (docs/design/68-context-engine.md §10). let result = session .message_chain() .iter() .flat_map(|(_, m)| m.content.iter()) .find_map(|b| match b { ContentBlock::ToolResult { content, is_error, .. } => Some((content.clone(), *is_error)), _ => None, }) .unwrap(); assert!(result.1); assert!(result.0.contains("no touching")); } #[tokio::test] async fn stop_hook_forces_continuation_once() { let dir = tempdir().unwrap(); let flag = dir.path().join("block-once"); let flag_display = flag.display().to_string(); // Block the first Stop; after the flag exists, allow it. let cmd = format!( "if [ -f {flag_display} ]; then exit 0; else touch {flag_display}; echo '{{\"decision\":\"block\",\"reason\":\"say goodbye first\"}}'; fi" ); let hooks = vec![HookDef { event: vak_hooks::HookEvent::Stop, matcher: None, command: cmd, timeout_ms: 5000, failure_mode: vak_hooks::HookFailureMode::Open, refusal: None, }]; let mut agent = build( vec![text_msg("first attempt"), text_msg("goodbye")], Some(hooks), ); let outcome = agent .run( "go", &Default::default(), CancellationToken::new(), mpsc::channel(64).0, ) .await; match outcome { TurnOutcome::Completed { response } => { assert_eq!(response.text_content(), "goodbye"); } other => panic!("expected completed after continuation, got {other:?}"), } let session = agent.session.lock().await; // Raw ledger: a control nudge is scaffolding for the turn still in // progress and is dropped once the turn closes // (docs/design/68-context-engine.md §10). let texts: Vec = session .message_chain() .iter() .filter(|(_, m)| m.text_content().contains("[stop-hook]")) .map(|(_, m)| m.text_content()) .collect(); assert_eq!(texts.len(), 1, "continuation message must be logged once"); }