- 1
#![allow(clippy::unwrap_used, clippy::expect_used, clippy::panic)] - 2
- 3
use std::collections::VecDeque; - 4
use std::sync::{Arc, Mutex}; - 5
- 6
use tokio::sync::mpsc; - 7
use tokio_util::sync::CancellationToken; - 8
- 9
use tempfile::tempdir; - 10
- 11
use vak_agent::{Agent, AgentConfig, TaskDeps, TaskTool, TurnOutcome, WorkerRegistry}; - 12
use vak_llm::stream; - 13
use vak_llm::types::{AssistantMessage, ChatRequest, ContentBlock, StopReason, Usage}; - 14
use vak_llm::{EventStream, LlmError, Provider}; - 15
use vak_session::types::{FrozenContract, SessionHeader}; - 16
use vak_session::{SessionLog, SessionPath}; - 17
use vak_tools::read::ReadTool; - 18
- 19
struct Scripted { - 20
responses: Mutex<VecDeque<AssistantMessage>>, - 21
requests: Arc<Mutex<Vec<ChatRequest>>>, - 22
} - 23
- 24
#[async_trait::async_trait] - 25
impl Provider for Scripted { - 26
fn name(&self) -> &str { - 27
"scripted" - 28
} - 29
- 30
async fn stream( - 31
&self, - 32
request: ChatRequest, - 33
_cancel: CancellationToken, - 34
) -> Result<EventStream, LlmError> { - 35
self.requests.lock().unwrap().push(request); - 36
let next = self.responses.lock().unwrap().pop_front(); - 37
let (mut sink, rx) = stream::channel(64); - 38
match next { - 39
Some(m) => { - 40
sink.push(stream::StreamEvent::Start { partial: m.clone() }); - 41
sink.close_message(m).await; - 42
} - 43
None => { - 44
sink.close_error(LlmError::Parse("script exhausted".into())) - 45
.await - 46
} - 47
} - 48
Ok(rx) - 49
} - 50
} - 51
- 52
fn text_msg(t: &str) -> AssistantMessage { - 53
AssistantMessage { - 54
content: vec![ContentBlock::text(t)], - 55
stop_reason: StopReason::EndTurn, - 56
usage: Usage { - 57
input_tokens: 1, - 58
output_tokens: 1, - 59
..Default::default() - 60
}, - 61
model: "test-model".into(), - 62
response_id: None, - 63
} - 64
} - 65
- 66
fn task_call(id: &str, prompt: &str) -> AssistantMessage { - 67
AssistantMessage { - 68
content: vec![ContentBlock::ToolUse { - 69
id: id.into(), - 70
name: "task".into(), - 71
input: serde_json::json!({"prompt": prompt, "label": "explore"}), - 72
}], - 73
stop_reason: StopReason::ToolUse, - 74
usage: Usage::default(), - 75
model: "test-model".into(), - 76
response_id: None, - 77
} - 78
} - 79
- 80
use vak_permission::PermissionEngine; - 81
- 82
#[tokio::test] - 83
async fn worker_roundtrip_with_shared_scripted_provider() { - 84
let dir = tempdir().unwrap(); - 85
let home = dir.path().join("home"); - 86
std::fs::create_dir_all(&home).unwrap(); - 87
let parent_id = "parent-x".to_string(); - 88
- 89
let header = SessionHeader { - 90
agent: None, - 91
session_id: parent_id.clone(), - 92
created_at: chrono::Utc::now(), - 93
cwd: dir.path().to_path_buf(), - 94
parent_session_id: None, - 95
contract_id: None, - 96
work_item_id: None, - 97
conversation: None, - 98
contract: FrozenContract { - 99
app_version: "0".into(), - 100
provider: "scripted".into(), - 101
model: "test-model".into(), - 102
route_ladder: Vec::new(), - 103
route_objective: String::new(), - 104
route_annotations: Vec::new(), - 105
system_prompt: "sys".into(), - 106
permission_mode: "workspace-write".into(), - 107
capabilities: Vec::new(), - 108
prompt_layers: Vec::new(), - 109
}, - 110
}; - 111
let log = SessionLog::create( - 112
SessionPath::new_session_file(&home, dir.path(), &parent_id), - 113
header, - 114
) - 115
.unwrap(); - 116
- 117
let requests = Arc::new(Mutex::new(Vec::new())); - 118
let scripted = Arc::new(Scripted { - 119
responses: Mutex::new(VecDeque::from(vec![ - 120
task_call("t1", "find the answer"), - 121
text_msg("child final answer"), - 122
text_msg("parent done"), - 123
])), - 124
requests: requests.clone(), - 125
}); - 126
- 127
let mut cfg = AgentConfig::new("sys"); - 128
cfg.model = "test-model".into(); - 129
cfg.tools = vec![Arc::new(TaskTool::new(TaskDeps { - 130
parent_agent_identity: None, - 131
role_prompts: Default::default(), - 132
provider: scripted.clone(), - 133
system_prompt: "child-sys".into(), - 134
tail: Default::default(), - 135
model: "test-model".into(), - 136
tools: vec![Arc::new(ReadTool)], - 137
capabilities: Vec::new(), - 138
hooks: None, - 139
revocation_check: None, - 140
presentation_rebuild: None, - 141
mcp_tool_index: None, - 142
input_normalizer: None, - 143
read_only_tools: vec![Arc::new(ReadTool)], - 144
max_turns: 5, - 145
outcome_objective: None, - 146
outcome: None, - 147
max_retries: 0, - 148
retry_base_backoff_ms: 0, - 149
request_timeout: None, - 150
circuit_breaker: None, - 151
run_retry_attempts: 0, - 152
run_retry_base_backoff_ms: 0, - 153
dispatch_ceiling: 1, - 154
spend_gate: None, - 155
permission: Some(Arc::new(PermissionEngine::default())), - 156
mode: vak_permission::Mode::WorkspaceWrite, - 157
approval_mode: vak_agent::ApprovalMode::Ask, - 158
approver: None, - 159
sandbox: None, - 160
cwd: dir.path().to_path_buf(), - 161
sessions_home: home.clone(), - 162
parent_session_id: parent_id.clone(), - 163
contract_id: None, - 164
work_item_id: None, - 165
work_item_ids: vec![], - 166
events: None, - 167
registry: Some(Arc::new(WorkerRegistry::new())), - 168
}))]; - 169
cfg.permission = Some(Arc::new( - 170
PermissionEngine::from_rule_strings(&["+task".to_string()]).unwrap(), - 171
)); - 172
cfg.approver = Some(Arc::new(vak_agent::AutoApprove)); - 173
- 174
let mut agent = Agent::new(scripted, log, cfg); - 175
let outcome = agent - 176
.run( - 177
"delegate please", - 178
&Default::default(), - 179
CancellationToken::new(), - 180
mpsc::channel(64).0, - 181
) - 182
.await; - 183
- 184
match outcome { - 185
TurnOutcome::Completed { response } => { - 186
assert_eq!(response.text_content(), "parent done"); - 187
} - 188
other => panic!("expected completed, got {other:?}"), - 189
} - 190
- 191
// Raw ledger: the closed turn's result is a trace line in the - 192
// projection now (docs/design/68-context-engine.md §10). - 193
let session = agent.session.lock().await; - 194
let tool_result = session - 195
.message_chain() - 196
.iter() - 197
.flat_map(|(_, m)| m.content.iter()) - 198
.find_map(|b| match b { - 199
ContentBlock::ToolResult { - 200
content, is_error, .. - 201
} => Some((content.clone(), *is_error)), - 202
_ => None, - 203
}) - 204
.expect("tool result from task must exist"); - 205
assert!(!tool_result.1); - 206
assert_eq!(tool_result.0, "child final answer"); - 207
- 208
let reqs = requests.lock().unwrap(); - 209
assert!( - 210
reqs.len() >= 3, - 211
"parent call + child call(s) + parent continuation expected, got {}", - 212
reqs.len() - 213
); - 214
let child_req = &reqs[1]; - 215
assert_eq!( - 216
child_req.system.as_deref(), - 217
Some("child-sys"), - 218
"child must run with its own narrowed system prompt" - 219
); - 220
assert!( - 221
child_req.tools.iter().all(|t| t.name != "task"), - 222
"children must not be able to spawn further workers" - 223
); - 224
- 225
let mut found_child = false; - 226
for entry in std::fs::read_dir(SessionPath::sessions_dir(&home, dir.path())) - 227
.unwrap() - 228
.flatten() - 229
{ - 230
let name = entry.file_name().to_string_lossy().into_owned(); - 231
if name.starts_with("child-") { - 232
found_child = true; - 233
let content = std::fs::read_to_string(entry.path()).unwrap(); - 234
let first: serde_json::Value = - 235
serde_json::from_str(content.lines().next().unwrap()).unwrap(); - 236
assert_eq!( - 237
first["parent_session_id"].as_str(), - 238
Some(parent_id.as_str()), - 239
"child session must link to its parent" - 240
); - 241
} - 242
} - 243
assert!(found_child, "a child session file must exist"); - 244
} - 245
Indexing the workspace…
Vakyartha documentation is discovering safe artifacts, anchors, and source references.