- 1
#![allow(clippy::unwrap_used, clippy::expect_used, clippy::panic)] - 2
- 3
//! Worker tail (docs/design/68-context-engine.md §6/§10): a spawned child - 4
//! agent must carry the same turn context block (clock instant + epistemic - 5
//! stance) as its parent turn, not an empty one — see `TaskDeps::tail` and - 6
//! `Core`'s `TaskDeps` construction, which now thread the parent's - 7
//! `AgentConfig::tail` through instead of leaving the child's default. - 8
//! Assembly and single-turn stability of the tail itself are covered by - 9
//! `context_tail.rs`; this test only checks that a WORKER's own request - 10
//! carries it too. - 11
- 12
use std::collections::VecDeque; - 13
use std::sync::{Arc, Mutex}; - 14
- 15
use tempfile::tempdir; - 16
use tokio::sync::mpsc; - 17
use tokio_util::sync::CancellationToken; - 18
- 19
use vak_agent::{Agent, AgentConfig, AutoApprove, TailInput, TaskDeps, TaskTool, TurnOutcome}; - 20
use vak_llm::stream; - 21
use vak_llm::types::{AssistantMessage, ChatRequest, ContentBlock, Role, StopReason, Usage}; - 22
use vak_llm::{EventStream, LlmError, Provider}; - 23
use vak_permission::{Mode, PermissionEngine}; - 24
use vak_session::types::{FrozenContract, SessionHeader}; - 25
use vak_session::{SessionLog, SessionPath}; - 26
use vak_tools::read::ReadTool; - 27
- 28
/// Records every request it sees and replays a scripted queue of - 29
/// responses — shared between the parent agent and its spawned child, so - 30
/// both turns' requests land in one inspectable list, in dispatch order. - 31
struct Scripted { - 32
responses: Mutex<VecDeque<AssistantMessage>>, - 33
requests: Arc<Mutex<Vec<ChatRequest>>>, - 34
} - 35
- 36
#[async_trait::async_trait] - 37
impl Provider for Scripted { - 38
fn name(&self) -> &str { - 39
"scripted" - 40
} - 41
- 42
async fn stream( - 43
&self, - 44
request: ChatRequest, - 45
_cancel: CancellationToken, - 46
) -> Result<EventStream, LlmError> { - 47
self.requests.lock().unwrap().push(request); - 48
let next = self.responses.lock().unwrap().pop_front(); - 49
let (mut sink, rx) = stream::channel(64); - 50
match next { - 51
Some(m) => { - 52
sink.push(stream::StreamEvent::Start { partial: m.clone() }); - 53
sink.close_message(m).await; - 54
} - 55
None => { - 56
sink.close_error(LlmError::Parse("script exhausted".into())) - 57
.await - 58
} - 59
} - 60
Ok(rx) - 61
} - 62
} - 63
- 64
fn text_msg(t: &str) -> AssistantMessage { - 65
AssistantMessage { - 66
content: vec![ContentBlock::text(t)], - 67
stop_reason: StopReason::EndTurn, - 68
usage: Usage { - 69
input_tokens: 1, - 70
output_tokens: 1, - 71
..Default::default() - 72
}, - 73
model: "test-model".into(), - 74
response_id: None, - 75
} - 76
} - 77
- 78
fn task_call(id: &str, prompt: &str) -> AssistantMessage { - 79
AssistantMessage { - 80
content: vec![ContentBlock::ToolUse { - 81
id: id.into(), - 82
name: "task".into(), - 83
input: serde_json::json!({"prompt": prompt, "label": "explore"}), - 84
}], - 85
stop_reason: StopReason::ToolUse, - 86
usage: Usage::default(), - 87
model: "test-model".into(), - 88
response_id: None, - 89
} - 90
} - 91
- 92
#[tokio::test(flavor = "multi_thread", worker_threads = 4)] - 93
async fn spawned_worker_request_carries_the_parent_turns_tail() { - 94
let dir = tempdir().unwrap(); - 95
let home = dir.path().join("home"); - 96
std::fs::create_dir_all(&home).unwrap(); - 97
let parent_id = "parent-tail".to_string(); - 98
- 99
let header = SessionHeader { - 100
agent: None, - 101
session_id: parent_id.clone(), - 102
created_at: chrono::Utc::now(), - 103
cwd: dir.path().to_path_buf(), - 104
parent_session_id: None, - 105
contract_id: None, - 106
work_item_id: None, - 107
conversation: None, - 108
contract: FrozenContract { - 109
app_version: "0".into(), - 110
provider: "scripted".into(), - 111
model: "test-model".into(), - 112
route_ladder: Vec::new(), - 113
route_objective: String::new(), - 114
route_annotations: Vec::new(), - 115
system_prompt: "sys".into(), - 116
permission_mode: "workspace-write".into(), - 117
capabilities: Vec::new(), - 118
prompt_layers: Vec::new(), - 119
}, - 120
}; - 121
let log = SessionLog::create( - 122
SessionPath::new_session_file(&home, dir.path(), &parent_id), - 123
header, - 124
) - 125
.unwrap(); - 126
- 127
let requests = Arc::new(Mutex::new(Vec::new())); - 128
let scripted = Arc::new(Scripted { - 129
responses: Mutex::new(VecDeque::from(vec![ - 130
task_call("t1", "find the answer"), - 131
text_msg("child final answer"), - 132
text_msg("parent done"), - 133
])), - 134
requests: requests.clone(), - 135
}); - 136
- 137
let parent_tail = TailInput { - 138
temporal: "current UTC instant 2026-09-19T00:00:00Z".into(), - 139
stance: "Provide a clear, direct answer.".into(), - 140
}; - 141
- 142
let mut cfg = AgentConfig::new("sys"); - 143
cfg.model = "test-model".into(); - 144
cfg.tail = parent_tail.clone(); - 145
cfg.tools = vec![Arc::new(TaskTool::new(TaskDeps { - 146
parent_agent_identity: None, - 147
role_prompts: Default::default(), - 148
provider: scripted.clone(), - 149
system_prompt: "child-sys".into(), - 150
// The field under test: `Core` threads `cfg.tail.clone()` through - 151
// here (crates/vak-core/src/lib.rs, TaskDeps construction) rather - 152
// than leaving the child with an empty `TailInput::default()`. - 153
tail: parent_tail.clone(), - 154
model: "test-model".into(), - 155
tools: vec![Arc::new(ReadTool)], - 156
capabilities: Vec::new(), - 157
hooks: None, - 158
revocation_check: None, - 159
presentation_rebuild: None, - 160
mcp_tool_index: None, - 161
input_normalizer: None, - 162
read_only_tools: vec![Arc::new(ReadTool)], - 163
max_turns: 5, - 164
outcome_objective: None, - 165
outcome: None, - 166
max_retries: 0, - 167
retry_base_backoff_ms: 0, - 168
request_timeout: None, - 169
circuit_breaker: None, - 170
run_retry_attempts: 0, - 171
run_retry_base_backoff_ms: 0, - 172
dispatch_ceiling: 1, - 173
spend_gate: None, - 174
permission: Some(Arc::new(PermissionEngine::default())), - 175
mode: Mode::WorkspaceWrite, - 176
approval_mode: vak_agent::ApprovalMode::Ask, - 177
approver: None, - 178
sandbox: None, - 179
cwd: dir.path().to_path_buf(), - 180
sessions_home: home.clone(), - 181
parent_session_id: parent_id.clone(), - 182
contract_id: None, - 183
work_item_id: None, - 184
work_item_ids: vec![], - 185
events: None, - 186
registry: Some(Arc::new(vak_agent::WorkerRegistry::new())), - 187
}))]; - 188
cfg.permission = Some(Arc::new( - 189
PermissionEngine::from_rule_strings(&["+task".to_string()]).unwrap(), - 190
)); - 191
cfg.approver = Some(Arc::new(AutoApprove)); - 192
- 193
let mut agent = Agent::new(scripted, log, cfg); - 194
let outcome = agent - 195
.run( - 196
"delegate please", - 197
&Default::default(), - 198
CancellationToken::new(), - 199
mpsc::channel(64).0, - 200
) - 201
.await; - 202
assert!(matches!(outcome, TurnOutcome::Completed { .. })); - 203
- 204
let reqs = requests.lock().unwrap(); - 205
assert!( - 206
reqs.len() >= 2, - 207
"parent call + child call expected, got {}", - 208
reqs.len() - 209
); - 210
let child_request = &reqs[1]; - 211
assert_eq!( - 212
child_request.system.as_deref(), - 213
Some("child-sys"), - 214
"reqs[1] must be the child's own request" - 215
); - 216
- 217
let last = child_request - 218
.messages - 219
.last() - 220
.expect("child request has at least one message"); - 221
assert_eq!( - 222
last.role, - 223
Role::User, - 224
"the tail rides the child's last USER message" - 225
); - 226
// The tail precedes the directive text in the last user message - 227
// (docs/design/68-context-engine.md §6). - 228
let tail_text = last - 229
.content - 230
.iter() - 231
.filter_map(|b| match b { - 232
ContentBlock::Text { text } => Some(text.clone()), - 233
_ => None, - 234
}) - 235
.find(|text| text.contains("<turn_context>")) - 236
.expect("the tail block on the child's last user message"); - 237
assert!( - 238
tail_text.contains("<turn_context>"), - 239
"worker request tail: {tail_text}" - 240
); - 241
assert!( - 242
tail_text.contains("current UTC instant 2026-09-19T00:00:00Z"), - 243
"worker must carry the PARENT turn's temporal context, not an empty one: {tail_text}" - 244
); - 245
assert!(tail_text.contains("<stance>")); - 246
assert!( - 247
tail_text.contains("Provide a clear, direct answer."), - 248
"worker must carry the parent turn's epistemic stance: {tail_text}" - 249
); - 250
} - 251
Indexing the workspace…
Vakyartha documentation is discovering safe artifacts, anchors, and source references.