- 1
//! The engagement governs the turn: what `vak-intent` decides reaches the - 2
//! runtime knobs it was designed for (docs/design/47-commitment-kernel.md, - 3
//! *What the review changed*). Each test here pins one seam that was - 4
//! computed and never consumed before resolver version 2. - 5
- 6
#![allow(clippy::unwrap_used, clippy::expect_used, clippy::panic)] - 7
- 8
use std::collections::VecDeque; - 9
use std::sync::{Arc, Mutex}; - 10
- 11
use tokio_util::sync::CancellationToken; - 12
- 13
use vak_core::Core; - 14
use vak_llm::stream; - 15
use vak_llm::types::{AssistantMessage, ChatRequest, ContentBlock, Usage}; - 16
use vak_llm::{EventStream, LlmError, Provider}; - 17
- 18
struct Scripted { - 19
responses: Mutex<VecDeque<AssistantMessage>>, - 20
/// Every request this provider was sent, so a test can read what the - 21
/// model actually saw. - 22
requests: Mutex<Vec<ChatRequest>>, - 23
} - 24
- 25
impl Scripted { - 26
fn new(responses: Vec<AssistantMessage>) -> Self { - 27
Scripted { - 28
responses: Mutex::new(VecDeque::from(responses)), - 29
requests: Mutex::new(Vec::new()), - 30
} - 31
} - 32
} - 33
- 34
#[async_trait::async_trait] - 35
impl Provider for Scripted { - 36
fn name(&self) -> &str { - 37
"scripted" - 38
} - 39
- 40
async fn stream( - 41
&self, - 42
request: ChatRequest, - 43
_cancel: CancellationToken, - 44
) -> Result<EventStream, LlmError> { - 45
self.requests.lock().unwrap().push(request); - 46
let next = self.responses.lock().unwrap().pop_front(); - 47
let (mut sink, rx) = stream::channel(64); - 48
match next { - 49
Some(m) => { - 50
sink.push(stream::StreamEvent::Start { partial: m.clone() }); - 51
sink.close_message(m).await; - 52
} - 53
None => sink.close_error(LlmError::Parse("exhausted".into())).await, - 54
} - 55
Ok(rx) - 56
} - 57
} - 58
- 59
fn text(t: &str) -> AssistantMessage { - 60
AssistantMessage { - 61
content: vec![ContentBlock::text(t)], - 62
stop_reason: vak_llm::types::StopReason::EndTurn, - 63
usage: Usage::default(), - 64
model: "test-model".into(), - 65
response_id: None, - 66
} - 67
} - 68
- 69
fn core_with(config: &str, responses: Vec<AssistantMessage>) -> (Core, std::path::PathBuf) { - 70
let (core, cwd, _) = core_with_provider(config, responses); - 71
(core, cwd) - 72
} - 73
- 74
fn core_with_provider( - 75
config: &str, - 76
responses: Vec<AssistantMessage>, - 77
) -> (Core, std::path::PathBuf, Arc<Scripted>) { - 78
let dir = tempfile::tempdir().unwrap(); - 79
let cwd = dir.path().to_path_buf(); - 80
let _ = std::fs::create_dir_all(cwd.join(".vak")); - 81
std::fs::write(cwd.join(".vak/config.toml"), config).unwrap(); - 82
vak_config::paths::isolate_home_for_tests(); - 83
let core = Core::new_with_trust(cwd.clone(), true).unwrap(); - 84
core.set_sessions_home(dir.path().join("home")); - 85
core.set_permission_mode(vak_config::PermissionMode::WorkspaceWrite); - 86
let provider = Arc::new(Scripted::new(responses)); - 87
core.set_provider_instance(provider.clone()); - 88
std::mem::forget(dir); - 89
(core, cwd, provider) - 90
} - 91
- 92
/// `ContextProfile::Working`: a verifying (or modifying) turn sees what - 93
/// changed in the workspace since the session began, from the ledger, in - 94
/// its tail; a question does not pay for it. - 95
#[tokio::test(flavor = "multi_thread", worker_threads = 4)] - 96
async fn a_working_turn_sees_the_workspace_delta_and_a_question_does_not() { - 97
// The stop gate is off: this test is about the tail, and a scripted - 98
// provider cannot produce the execution receipt a verify turn owes. - 99
let (core, cwd, provider) = core_with_provider( - 100
"[memory]\nreflection = false\n\n[stop_policy]\nenabled = false\n", - 101
vec![text("hello"), text("it is a note"), text("checked")], - 102
); - 103
// The model steps of one turn: every request whose last message is the - 104
// user's (the judge's transcript request is not one). - 105
let step_texts = |from: usize| -> Vec<String> { - 106
provider - 107
.requests - 108
.lock() - 109
.unwrap() - 110
.iter() - 111
.skip(from) - 112
.filter(|request| { - 113
request - 114
.messages - 115
.last() - 116
.is_some_and(|m| m.role == vak_llm::Role::User) - 117
}) - 118
.map(|request| { - 119
request - 120
.messages - 121
.iter() - 122
.map(|m| m.text_content()) - 123
.collect::<Vec<_>>() - 124
.join("\n") - 125
}) - 126
.collect() - 127
}; - 128
std::fs::write(cwd.join("notes.md"), "seed\n").unwrap(); - 129
let session = core.start_session().await.unwrap(); - 130
let (tx, _rx) = tokio::sync::mpsc::channel(64); - 131
let (_, session) = core - 132
.run_turn_with( - 133
session, - 134
"hi", - 135
CancellationToken::new(), - 136
None, - 137
None, - 138
None, - 139
tx, - 140
) - 141
.await - 142
.unwrap(); - 143
- 144
// The workspace changes between turns (the session's own tools would do - 145
// this; the change is what matters, not who made it). - 146
std::fs::write(cwd.join("notes.md"), "seed\nedited\n").unwrap(); - 147
- 148
let before = provider.requests.lock().unwrap().len(); - 149
let (tx, _rx) = tokio::sync::mpsc::channel(64); - 150
let (_, session) = core - 151
.run_turn_with( - 152
session, - 153
"what is in the notes file?", - 154
CancellationToken::new(), - 155
None, - 156
None, - 157
None, - 158
tx, - 159
) - 160
.await - 161
.unwrap(); - 162
let question_steps = step_texts(before); - 163
assert!(!question_steps.is_empty()); - 164
assert!( - 165
question_steps - 166
.iter() - 167
.all(|text| !text.contains("<workspace_delta>")), - 168
"a question does not carry the delta: {question_steps:?}" - 169
); - 170
- 171
let before = provider.requests.lock().unwrap().len(); - 172
let (tx, _rx) = tokio::sync::mpsc::channel(64); - 173
let (_, session) = core - 174
.run_turn_with( - 175
session, - 176
"check the notes file", - 177
CancellationToken::new(), - 178
None, - 179
None, - 180
None, - 181
tx, - 182
) - 183
.await - 184
.unwrap(); - 185
let working_steps = step_texts(before); - 186
assert!( - 187
working_steps - 188
.iter() - 189
.any(|text| text.contains("<workspace_delta>") && text.contains("notes.md")), - 190
"a working turn sees what the session changed: {working_steps:?}" - 191
); - 192
// And the bytes the model saw are in the ledger. - 193
let logged = session.chain_to_root().iter().any(|entry| { - 194
matches!(&entry.payload, vak_session::EntryPayload::Activity(activity) - 195
if activity.data.get("section").map(String::as_str) == Some("workspace_delta") - 196
&& activity.detail.as_deref().is_some_and(|d| d.contains("notes.md"))) - 197
}); - 198
assert!(logged, "the delta is a ledger entry"); - 199
} - 200
- 201
/// Invariant 10: an image turn is served only by a leg declared able to see - 202
/// it. With hints declared and no capable leg, the turn fails typed rather - 203
/// than quietly sending the image to a text-only model. - 204
#[tokio::test(flavor = "multi_thread", worker_threads = 4)] - 205
async fn an_unsupported_modality_fails_typed_instead_of_dropping_the_image() { - 206
let (core, _cwd) = core_with( - 207
"[memory]\nreflection = false\n\n[route]\nmodality_hints = [\"vision\"]\n", - 208
vec![text("I see nothing")], - 209
); - 210
let session = core.start_session().await.unwrap(); - 211
let (tx, _rx) = tokio::sync::mpsc::channel(64); - 212
let prompt = vak_llm::Message { - 213
role: vak_llm::Role::User, - 214
content: vec![ - 215
ContentBlock::text("what is in this picture"), - 216
ContentBlock::image_base64("image/png", "iVBORw0KGgo="), - 217
], - 218
}; - 219
let result = core - 220
.run_turn_with_message( - 221
session, - 222
vak_session::MessageRecord { - 223
message: prompt, - 224
meta: None, - 225
}, - 226
CancellationToken::new(), - 227
None, - 228
None, - 229
None, - 230
tx, - 231
) - 232
.await; - 233
match result { - 234
Err(vak_core::CoreError::UnsupportedModality { modalities, .. }) => { - 235
assert!(modalities.contains("image"), "{modalities}"); - 236
} - 237
other => panic!( - 238
"expected a typed modality error, got {:?}", - 239
other.map(|_| ()) - 240
), - 241
} - 242
- 243
// With no hints declared, every leg is assumed capable and the turn runs. - 244
let (core, _cwd) = core_with("[memory]\nreflection = false\n", vec![text("I see a cat")]); - 245
let session = core.start_session().await.unwrap(); - 246
let (tx, _rx) = tokio::sync::mpsc::channel(64); - 247
let prompt = vak_llm::Message { - 248
role: vak_llm::Role::User, - 249
content: vec![ - 250
ContentBlock::text("what is in this picture"), - 251
ContentBlock::image_base64("image/png", "iVBORw0KGgo="), - 252
], - 253
}; - 254
let (outcome, _) = core - 255
.run_turn_with_message( - 256
session, - 257
vak_session::MessageRecord { - 258
message: prompt, - 259
meta: None, - 260
}, - 261
CancellationToken::new(), - 262
None, - 263
None, - 264
None, - 265
tx, - 266
) - 267
.await - 268
.unwrap(); - 269
assert!(matches!(outcome, vak_agent::TurnOutcome::Completed { .. })); - 270
} - 271
- 272
/// A turn's intent entry carries its strands, and the next turn's lineage - 273
/// sees them as open threads. - 274
#[tokio::test(flavor = "multi_thread", worker_threads = 4)] - 275
async fn strands_are_recorded_and_become_open_threads() { - 276
let (core, _cwd) = core_with( - 277
"[memory]\nreflection = false\n", - 278
vec![text("done"), text("done again")], - 279
); - 280
let session = core.start_session().await.unwrap(); - 281
let (tx, _rx) = tokio::sync::mpsc::channel(64); - 282
let (_, session) = core - 283
.run_turn_with( - 284
session, - 285
"explain the parser, then summarise the changelog", - 286
CancellationToken::new(), - 287
None, - 288
None, - 289
None, - 290
tx, - 291
) - 292
.await - 293
.unwrap(); - 294
let record = session - 295
.chain_to_root() - 296
.into_iter() - 297
.rev() - 298
.find_map(|entry| match &entry.payload { - 299
vak_session::EntryPayload::Intent(record) => Some(record.as_ref().clone()), - 300
_ => None, - 301
}) - 302
.expect("an intent entry"); - 303
assert_eq!(record.strands.len(), 2, "{:?}", record.strands); - 304
assert!( - 305
record - 306
.model_visible - 307
.as_deref() - 308
.is_some_and(|note| note.contains("2 parts")), - 309
"{:?}", - 310
record.model_visible - 311
); - 312
- 313
let threads = vak_core::intent::open_threads(&session); - 314
assert_eq!(threads.len(), 2); - 315
assert!(threads.iter().any(|t| t.keywords.contains("parser"))); - 316
- 317
// "now explain it in more detail" continues a thread rather than opening - 318
// a twin. - 319
let (tx, _rx) = tokio::sync::mpsc::channel(64); - 320
let (_, session) = core - 321
.run_turn_with( - 322
session, - 323
"now explain it in more detail", - 324
CancellationToken::new(), - 325
None, - 326
None, - 327
None, - 328
tx, - 329
) - 330
.await - 331
.unwrap(); - 332
let latest = session - 333
.chain_to_root() - 334
.into_iter() - 335
.rev() - 336
.find_map(|entry| match &entry.payload { - 337
vak_session::EntryPayload::Intent(record) => Some(record.as_ref().clone()), - 338
_ => None, - 339
}) - 340
.unwrap(); - 341
assert!( - 342
matches!( - 343
latest.strands[0].lineage, - 344
vak_intent::Lineage::Continues { .. } - 345
), - 346
"{:?}", - 347
latest.strands[0].lineage - 348
); - 349
} - 350
- 351
/// `Defer`: a gate nobody can answer is parked in the inbox and suspends the - 352
/// commitment, and the turn still fails closed. - 353
#[tokio::test] - 354
async fn a_deferred_gate_reaches_the_inbox_and_suspends_the_commitment() { - 355
use vak_agent::Approver; - 356
let dir = tempfile::tempdir().unwrap(); - 357
let sessions_home = dir.path().join("sessions"); - 358
let shared_home = dir.path().join("shared"); - 359
let ledger = vak_commit::CommitmentLedger::new(&sessions_home); - 360
let reading = vak_intent::Reading { - 361
horizon: vak_intent::Horizon::Durable, - 362
stakes: vak_intent::Stakes::Costly, - 363
..vak_intent::Reading::general() - 364
}; - 365
let commitment_id = ledger - 366
.open_commitment(vak_commit::spec_from_reading( - 367
"rotate the production keys", - 368
reading, - 369
Vec::new(), - 370
dir.path().to_path_buf(), - 371
vak_commit::Economics::default(), - 372
)) - 373
.unwrap(); - 374
- 375
struct Nobody; - 376
#[async_trait::async_trait] - 377
impl Approver for Nobody { - 378
async fn approve(&self, _: &str, _: &str, _: &str) -> bool { - 379
panic!("an unanswerable approver must never be asked") - 380
} - 381
fn answerable(&self) -> bool { - 382
false - 383
} - 384
} - 385
let approver = vak_core::intent::DeferringApprover::new( - 386
Some(Arc::new(Nobody)), - 387
shared_home.clone(), - 388
sessions_home.clone(), - 389
"s1".into(), - 390
commitment_id.clone(), - 391
vak_intent::Escalation::WaitIndefinitely, - 392
); - 393
assert!(!approver.answerable()); - 394
let allowed = approver - 395
.approve("bash", "{\"command\":\"rotate-keys\"}", "irreversible") - 396
.await; - 397
assert!(!allowed, "the turn fails closed"); - 398
- 399
let commitment = ledger.get(&commitment_id).unwrap().unwrap(); - 400
assert!( - 401
matches!( - 402
commitment.suspension, - 403
Some(vak_commit::Suspension::Human { .. }) - 404
), - 405
"{:?}", - 406
commitment.suspension - 407
); - 408
let inbox = vak_core::inbox::unread(&shared_home); - 409
assert_eq!(inbox.len(), 1); - 410
assert_eq!(inbox[0].kind, vak_core::inbox::Kind::ApprovalPending); - 411
assert!(inbox[0].body.contains("rotate-keys")); - 412
} - 413
Indexing the workspace…
Vakyartha documentation is discovering safe artifacts, anchors, and source references.