- 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::{ - 12
Agent, AgentConfig, ApprovalMode, AutoApprove, SteeringQueues, TurnOutcome, WorkMode, - 13
}; - 14
use vak_llm::stream; - 15
use vak_llm::types::{AssistantMessage, ChatRequest, ContentBlock, StopReason, Usage}; - 16
use vak_llm::{EventStream, LlmError, Provider, WorkPurpose}; - 17
use vak_session::SessionLog; - 18
use vak_session::types::{ - 19
CapabilityDescriptor, CapabilityInvocation, CapabilityKind, FrozenContract, SessionHeader, - 20
}; - 21
use vak_tools::Tool; - 22
use vak_tools::bash::BashTool; - 23
use vak_tools::write::WriteTool; - 24
- 25
enum ScriptedResponse { - 26
Message(AssistantMessage), - 27
Error(LlmError), - 28
} - 29
- 30
struct Scripted { - 31
responses: Mutex<VecDeque<ScriptedResponse>>, - 32
requests: Arc<Mutex<Vec<ChatRequest>>>, - 33
} - 34
- 35
#[async_trait::async_trait] - 36
impl Provider for Scripted { - 37
fn name(&self) -> &str { - 38
"scripted" - 39
} - 40
- 41
async fn stream( - 42
&self, - 43
request: ChatRequest, - 44
_cancel: CancellationToken, - 45
) -> Result<EventStream, LlmError> { - 46
self.requests.lock().unwrap().push(request); - 47
let next = self.responses.lock().unwrap().pop_front(); - 48
let (mut sink, rx) = stream::channel(64); - 49
match next { - 50
Some(ScriptedResponse::Message(m)) => { - 51
sink.push(stream::StreamEvent::Start { partial: m.clone() }); - 52
sink.close_message(m).await; - 53
} - 54
Some(ScriptedResponse::Error(e)) => sink.close_error(e).await, - 55
None => { - 56
sink.close_error(LlmError::Parse("script exhausted".into())) - 57
.await - 58
} - 59
} - 60
Ok(rx) - 61
} - 62
} - 63
fn assistant_text(text: &str) -> AssistantMessage { - 64
AssistantMessage { - 65
content: vec![ContentBlock::text(text)], - 66
stop_reason: StopReason::EndTurn, - 67
usage: Usage { - 68
input_tokens: 10, - 69
output_tokens: 5, - 70
..Default::default() - 71
}, - 72
model: "test-model".into(), - 73
response_id: None, - 74
} - 75
} - 76
- 77
fn tool_call_msg(id: &str, name: &str, input: serde_json::Value) -> AssistantMessage { - 78
AssistantMessage { - 79
content: vec![ContentBlock::ToolUse { - 80
id: id.into(), - 81
name: name.into(), - 82
input, - 83
}], - 84
stop_reason: StopReason::ToolUse, - 85
usage: Usage::default(), - 86
model: "test-model".into(), - 87
response_id: None, - 88
} - 89
} - 90
- 91
struct Harness { - 92
agent: Agent, - 93
requests: Arc<Mutex<Vec<ChatRequest>>>, - 94
events_tx: mpsc::Sender<vak_agent::AgentEvent>, - 95
_events_rx: mpsc::Receiver<vak_agent::AgentEvent>, - 96
} - 97
- 98
fn harness(responses: Vec<ScriptedResponse>, tools: Vec<Arc<dyn Tool>>) -> Harness { - 99
let dir = tempdir().unwrap(); - 100
let header = SessionHeader { - 101
agent: None, - 102
session_id: "s-test".into(), - 103
created_at: chrono::Utc::now(), - 104
cwd: dir.path().to_path_buf(), - 105
parent_session_id: None, - 106
contract_id: None, - 107
work_item_id: None, - 108
conversation: None, - 109
contract: FrozenContract { - 110
app_version: "0.1.0".into(), - 111
provider: "scripted".into(), - 112
model: "test-model".into(), - 113
route_ladder: Vec::new(), - 114
route_objective: String::new(), - 115
route_annotations: Vec::new(), - 116
system_prompt: "sys".into(), - 117
permission_mode: "full-access".into(), - 118
capabilities: vec![CapabilityDescriptor { - 119
name: "code-task".into(), - 120
kind: CapabilityKind::Skill, - 121
invocation: CapabilityInvocation::SkillLoader, - 122
description: "test skill".into(), - 123
source: None, - 124
digest: None, - 125
provenance: None, - 126
configuration: serde_json::Value::Null, - 127
}], - 128
prompt_layers: Vec::new(), - 129
}, - 130
}; - 131
let log = SessionLog::create(dir.path().join("s.jsonl"), header).unwrap(); - 132
let mut cfg = AgentConfig::new("sys"); - 133
cfg.tools = tools; - 134
let (tx, rx) = mpsc::channel(4096); - 135
let requests: Arc<Mutex<Vec<ChatRequest>>> = Arc::new(Mutex::new(Vec::new())); - 136
let provider = Arc::new(Scripted { - 137
responses: Mutex::new(responses.into_iter().collect()), - 138
requests: requests.clone(), - 139
}); - 140
std::mem::forget(dir); - 141
Harness { - 142
agent: Agent::new(provider, log, cfg), - 143
requests, - 144
events_tx: tx, - 145
_events_rx: rx, - 146
} - 147
} - 148
- 149
fn next_tool_id() -> String { - 150
use std::sync::atomic::{AtomicU32, Ordering}; - 151
static C: AtomicU32 = AtomicU32::new(0); - 152
format!("t{}", C.fetch_add(1, Ordering::Relaxed)) - 153
} - 154
- 155
#[tokio::test] - 156
async fn single_turn_no_tools_completes() { - 157
let mut h = harness( - 158
vec![ScriptedResponse::Message(assistant_text("done"))], - 159
vec![], - 160
); - 161
let outcome = h - 162
.agent - 163
.run( - 164
"hello", - 165
&Default::default(), - 166
CancellationToken::new(), - 167
h.events_tx.clone(), - 168
) - 169
.await; - 170
match outcome { - 171
TurnOutcome::Completed { response } => assert_eq!(response.text_content(), "done"), - 172
other => panic!("expected completed, got {other:?}"), - 173
} - 174
let session = h.agent.session.lock().await; - 175
assert_eq!(session.derive_messages().len(), 2); - 176
} - 177
- 178
#[tokio::test] - 179
async fn outcome_turn_cap_counts_model_calls_not_tool_round_trips() { - 180
let mut h = harness( - 181
vec![ - 182
ScriptedResponse::Message(tool_call_msg( - 183
"tool-1", - 184
"bash", - 185
serde_json::json!({"command": "printf tool"}), - 186
)), - 187
ScriptedResponse::Message(assistant_text("done")), - 188
], - 189
vec![Arc::new(BashTool)], - 190
); - 191
h.agent.config.mode = vak_permission::Mode::FullAccess; - 192
h.agent.config.approval_mode = ApprovalMode::AutoApprove; - 193
h.agent.config.approver = Some(Arc::new(AutoApprove)); - 194
h.agent.config.outcome = Some(vak_intent::OutcomeSpec { - 195
schema_version: 1, - 196
revision: 0, - 197
objective: "complete one tool-assisted answer".into(), - 198
assumptions: Vec::new(), - 199
requirements: Vec::new(), - 200
resolver_version: 1, - 201
evidence_max_age_secs: None, - 202
acts: Default::default(), - 203
stop: Default::default(), - 204
max_turns: Some(2), - 205
}); - 206
- 207
let outcome = h - 208
.agent - 209
.run( - 210
"run the tool", - 211
&Default::default(), - 212
CancellationToken::new(), - 213
h.events_tx.clone(), - 214
) - 215
.await; - 216
- 217
assert!( - 218
matches!(outcome, TurnOutcome::Completed { .. }), - 219
"{outcome:?}" - 220
); - 221
assert_eq!(h.requests.lock().unwrap().len(), 2); - 222
} - 223
- 224
#[tokio::test] - 225
async fn queued_outcome_revision_is_applied_and_logged_before_dispatch() { - 226
let mut h = harness( - 227
vec![ScriptedResponse::Message(assistant_text("done"))], - 228
vec![], - 229
); - 230
let steering = SteeringQueues::new(); - 231
let intent = vak_intent::Intent::general(1); - 232
steering.push_outcome_update(vak_session::types::IntentRecord { - 233
reading: intent.reading, - 234
engagement: intent.engagement, - 235
provenance: intent.provenance, - 236
outcome: Some(vak_intent::OutcomeSpec { - 237
schema_version: 1, - 238
revision: 7, - 239
objective: "updated objective".into(), - 240
assumptions: Vec::new(), - 241
requirements: Vec::new(), - 242
resolver_version: 1, - 243
evidence_max_age_secs: None, - 244
acts: Default::default(), - 245
stop: Default::default(), - 246
max_turns: Some(1), - 247
}), - 248
model_visible: None, - 249
commitment_id: None, - 250
strands: Vec::new(), - 251
strand_commitments: Default::default(), - 252
}); - 253
- 254
let outcome = h - 255
.agent - 256
.run( - 257
"continue", - 258
&steering, - 259
CancellationToken::new(), - 260
h.events_tx.clone(), - 261
) - 262
.await; - 263
assert!( - 264
matches!(outcome, TurnOutcome::Completed { .. }), - 265
"{outcome:?}" - 266
); - 267
let session = h.agent.session.lock().await; - 268
let revision = session.chain_to_root().iter().rev().find_map(|entry| { - 269
if let vak_session::EntryPayload::Intent(record) = &entry.payload { - 270
record.outcome.as_ref().map(|outcome| outcome.revision) - 271
} else { - 272
None - 273
} - 274
}); - 275
assert_eq!(revision, Some(7)); - 276
} - 277
- 278
#[tokio::test] - 279
async fn managed_turn_authors_and_persists_validated_contract() { - 280
let authored = serde_json::json!({ - 281
"objective": "make the change", - 282
"constraints": [], - 283
"assumptions": [], - 284
"criteria": [{ - 285
"criterion_id": "change", - 286
"statement": "the change tool completed", - 287
"kind": {"kind": "tool_succeeded", "tool": "bash"}, - 288
"required": true - 289
}], - 290
"items": [{ - 291
"item_id": "change", - 292
"title": "Make the change", - 293
"instructions": "make it", - 294
"dependencies": [], - 295
"owner": "parent_agent", - 296
"required": true, - 297
"readonly": false, - 298
"path_claims": [], - 299
"criterion_ids": ["change"] - 300
}] - 301
}); - 302
let mut h = harness( - 303
vec![ - 304
ScriptedResponse::Message(assistant_text(&authored.to_string())), - 305
ScriptedResponse::Message(tool_call_msg( - 306
"b1", - 307
"bash", - 308
serde_json::json!({"command": "true"}), - 309
)), - 310
ScriptedResponse::Message(assistant_text("completed")), - 311
], - 312
vec![Arc::new(BashTool)], - 313
); - 314
h.agent.config.work_mode = WorkMode::Managed; - 315
let outcome = h - 316
.agent - 317
.run( - 318
"please make the change", - 319
&SteeringQueues::new(), - 320
CancellationToken::new(), - 321
h.events_tx, - 322
) - 323
.await; - 324
assert!(matches!(outcome, TurnOutcome::Completed { .. })); - 325
let session = h.agent.into_session().await; - 326
let projection = session.work_projection().unwrap().unwrap(); - 327
assert_eq!(projection.contract.objective, "make the change"); - 328
assert_eq!( - 329
projection.status, - 330
vak_session::types::WorkContractStatus::Completed - 331
); - 332
assert_eq!( - 333
projection.items["change"].status, - 334
vak_session::types::WorkItemStatus::Succeeded - 335
); - 336
assert_eq!(session.receipts()[0].purpose, WorkPurpose::Plan); - 337
} - 338
- 339
#[tokio::test] - 340
async fn managed_semantic_criterion_uses_the_independent_judge() { - 341
let authored = serde_json::json!({ - 342
"objective": "make the change", - 343
"constraints": [], - 344
"assumptions": [], - 345
"criteria": [{ - 346
"criterion_id": "review", - 347
"statement": "the requested change is complete", - 348
"kind": {"kind": "semantic"}, - 349
"required": true - 350
}], - 351
"items": [{ - 352
"item_id": "change", - 353
"title": "Make the change", - 354
"instructions": "make it", - 355
"dependencies": [], - 356
"owner": "parent_agent", - 357
"required": true, - 358
"readonly": false, - 359
"path_claims": [], - 360
"criterion_ids": ["review"] - 361
}] - 362
}); - 363
let judge = serde_json::json!({ - 364
"results": [{ - 365
"criterion": "[review] the requested change is complete", - 366
"verdict": "pass", - 367
"evidence": "the tool result and final response agree" - 368
}] - 369
}); - 370
let mut h = harness( - 371
vec![ - 372
ScriptedResponse::Message(assistant_text(&authored.to_string())), - 373
ScriptedResponse::Message(tool_call_msg( - 374
"b1", - 375
"bash", - 376
serde_json::json!({"command": "true"}), - 377
)), - 378
ScriptedResponse::Message(assistant_text("done")), - 379
ScriptedResponse::Message(assistant_text(&judge.to_string())), - 380
], - 381
vec![Arc::new(BashTool)], - 382
); - 383
h.agent.config.work_mode = WorkMode::Managed; - 384
let outcome = h - 385
.agent - 386
.run( - 387
"please make the change", - 388
&SteeringQueues::new(), - 389
CancellationToken::new(), - 390
h.events_tx, - 391
) - 392
.await; - 393
assert!(matches!(outcome, TurnOutcome::Completed { .. })); - 394
let session = h.agent.into_session().await; - 395
let projection = session.work_projection().unwrap().unwrap(); - 396
assert_eq!( - 397
projection.criteria["review"], - 398
vak_session::types::CriterionResult::Passed { - 399
evidence: "the tool result and final response agree".into() - 400
} - 401
); - 402
assert_eq!( - 403
projection.status, - 404
vak_session::types::WorkContractStatus::Completed - 405
); - 406
} - 407
- 408
#[tokio::test] - 409
async fn normalizes_steering_before_it_reaches_the_ledger_or_provider() { - 410
let mut h = harness( - 411
vec![ScriptedResponse::Message(assistant_text("done"))], - 412
vec![], - 413
); - 414
h.agent.config.input_normalizer = Some(Arc::new(|mut message| { - 415
if let Some(ContentBlock::Text { text }) = message - 416
.content - 417
.iter_mut() - 418
.find(|block| matches!(block, ContentBlock::Text { .. })) - 419
{ - 420
*text = format!("normalized: {text}"); - 421
} - 422
Ok(message) - 423
})); - 424
let steering = vak_agent::SteeringQueues::new(); - 425
steering.push_steering("queued command"); - 426
let outcome = h - 427
.agent - 428
.run( - 429
"initial", - 430
&steering, - 431
CancellationToken::new(), - 432
h.events_tx.clone(), - 433
) - 434
.await; - 435
assert!(matches!(outcome, TurnOutcome::Completed { .. })); - 436
let request = h.requests.lock().unwrap().first().cloned().unwrap(); - 437
assert!( - 438
request - 439
.messages - 440
.iter() - 441
.any(|message| message.text_content() == "normalized: queued command") - 442
); - 443
let session = h.agent.session.lock().await; - 444
assert!( - 445
session - 446
.derive_messages() - 447
.iter() - 448
.any(|message| message.text_content() == "normalized: queued command") - 449
); - 450
} - 451
- 452
#[tokio::test] - 453
async fn tool_roundtrip_executes_and_feeds_result_back() { - 454
let h = harness( - 455
vec![ - 456
ScriptedResponse::Message(tool_call_msg( - 457
&next_tool_id(), - 458
"bash", - 459
serde_json::json!({"command": "echo ran"}), - 460
)), - 461
ScriptedResponse::Message(assistant_text("finished")), - 462
], - 463
vec![Arc::new(BashTool)], - 464
); - 465
let mut agent = h.agent; - 466
let outcome = agent - 467
.run( - 468
"run it", - 469
&Default::default(), - 470
CancellationToken::new(), - 471
h.events_tx.clone(), - 472
) - 473
.await; - 474
assert!(matches!(outcome, TurnOutcome::Completed { .. })); - 475
- 476
let (req_count, has_tool_result) = { - 477
let reqs = h.requests.lock().unwrap(); - 478
let hit = reqs.get(1).map(|r| { - 479
r.messages.iter().any(|m| { - 480
m.content - 481
.iter() - 482
.any(|b| matches!(b, ContentBlock::ToolResult { content, is_error: false, .. } if content.contains("ran"))) - 483
}) - 484
}).unwrap_or(false); - 485
(reqs.len(), hit) - 486
}; - 487
assert_eq!(req_count, 2); - 488
assert!( - 489
has_tool_result, - 490
"second request must contain the tool result" - 491
); - 492
- 493
let session = agent.session.lock().await; - 494
// Four raw entries on the ledger (directive, tool_use, tool_result, - 495
// final answer); the turn is now closed, so the model-visible - 496
// projection is its full record: the same four messages with the - 497
// result digested and thinking dropped - 498
// (docs/design/68-context-engine.md §10). - 499
assert_eq!(session.message_chain().len(), 4); - 500
let projected = session.derive_messages(); - 501
assert_eq!(projected.len(), 4); - 502
assert!(projected.iter().any(|m| m.content.iter().any(|b| matches!( - 503
b, - 504
ContentBlock::ToolResult { content, .. } if content.contains("[evidence:") - 505
)))); - 506
} - 507
- 508
#[tokio::test] - 509
async fn malformed_tool_input_is_rejected_by_the_admitted_schema() { - 510
let h = harness( - 511
vec![ - 512
ScriptedResponse::Message(tool_call_msg("t-schema", "write", serde_json::json!({}))), - 513
ScriptedResponse::Message(assistant_text("recovered")), - 514
], - 515
vec![Arc::new(WriteTool)], - 516
); - 517
let mut agent = h.agent; - 518
let outcome = agent - 519
.run( - 520
"write it", - 521
&Default::default(), - 522
CancellationToken::new(), - 523
h.events_tx.clone(), - 524
) - 525
.await; - 526
assert!(matches!(outcome, TurnOutcome::Completed { .. })); - 527
let session = agent.session.lock().await; - 528
// Raw ledger, not the model-visible projection: once the turn closes, - 529
// a closed turn's tool results collapse into trace lines - 530
// (docs/design/68-context-engine.md §10) — the mechanics under test - 531
// here are what got recorded, not how a later turn would see it. - 532
let result = session - 533
.message_chain() - 534
.iter() - 535
.flat_map(|(_, message)| message.content.iter()) - 536
.find_map(|block| match block { - 537
ContentBlock::ToolResult { content, .. } => Some(content.clone()), - 538
_ => None, - 539
}) - 540
.expect("schema rejection must be recorded as a tool result"); - 541
assert!(result.contains("required parameter `path` was omitted")); - 542
} - 543
- 544
#[tokio::test] - 545
async fn unknown_tool_becomes_error_value_not_crash() { - 546
let h = harness( - 547
vec![ - 548
ScriptedResponse::Message(tool_call_msg("t1", "nonexistent", serde_json::json!({}))), - 549
ScriptedResponse::Message(assistant_text("recovered")), - 550
], - 551
vec![Arc::new(BashTool)], - 552
); - 553
let mut agent = h.agent; - 554
let outcome = agent - 555
.run( - 556
"go", - 557
&Default::default(), - 558
CancellationToken::new(), - 559
h.events_tx.clone(), - 560
) - 561
.await; - 562
assert!(matches!(outcome, TurnOutcome::Completed { .. })); - 563
let session = agent.session.lock().await; - 564
// Raw ledger: the closed turn's tool result is a trace line in the - 565
// projection now, not a raw block (docs/design/68-context-engine.md - 566
// §10) — this test is about what the tool actually recorded. - 567
let msgs = session.message_chain(); - 568
let results_block = msgs - 569
.iter() - 570
.rev() - 571
.flat_map(|(_, m)| m.content.iter()) - 572
.find_map(|b| match b { - 573
ContentBlock::ToolResult { content, .. } => Some(content.clone()), - 574
_ => None, - 575
}) - 576
.expect("a tool result must exist"); - 577
assert!(results_block.contains("unknown_capability")); - 578
assert!(results_block.contains("available_tools")); - 579
} - 580
- 581
#[tokio::test] - 582
async fn skill_name_tool_call_returns_typed_loader_recovery() { - 583
let mut h = harness( - 584
vec![ - 585
ScriptedResponse::Message(tool_call_msg("t1", "code-task", serde_json::json!({}))), - 586
ScriptedResponse::Message(assistant_text("recovered")), - 587
], - 588
vec![Arc::new(BashTool)], - 589
); - 590
let outcome = h - 591
.agent - 592
.run( - 593
"go", - 594
&Default::default(), - 595
CancellationToken::new(), - 596
h.events_tx.clone(), - 597
) - 598
.await; - 599
assert!(matches!(outcome, TurnOutcome::Completed { .. })); - 600
let session = h.agent.session.lock().await; - 601
// Raw ledger: see the comment on `unknown_tool_becomes_error_value_not_crash`. - 602
let result = session - 603
.message_chain() - 604
.iter() - 605
.flat_map(|(_, m)| m.content.iter()) - 606
.find_map(|b| match b { - 607
ContentBlock::ToolResult { content, .. } => Some(content.clone()), - 608
_ => None, - 609
}) - 610
.expect("a tool result must exist"); - 611
assert!(result.contains("capability_kind_mismatch")); - 612
assert!(result.contains("\"tool\":\"skill\"")); - 613
assert!(result.contains("\"name\":\"code-task\"")); - 614
} - 615
- 616
#[tokio::test] - 617
async fn parallel_tools_preserve_source_order() { - 618
let msg = AssistantMessage { - 619
content: vec![ - 620
ContentBlock::ToolUse { - 621
id: "a".into(), - 622
name: "write".into(), - 623
input: serde_json::json!({"path": "a.txt", "content": "A"}), - 624
}, - 625
ContentBlock::ToolUse { - 626
id: "b".into(), - 627
name: "write".into(), - 628
input: serde_json::json!({"path": "b.txt", "content": "B"}), - 629
}, - 630
], - 631
stop_reason: StopReason::ToolUse, - 632
usage: Usage::default(), - 633
model: "test-model".into(), - 634
response_id: None, - 635
}; - 636
let h = harness( - 637
vec![ - 638
ScriptedResponse::Message(msg), - 639
ScriptedResponse::Message(assistant_text("both done")), - 640
], - 641
vec![Arc::new(WriteTool), Arc::new(BashTool)], - 642
); - 643
let mut h = h; - 644
h.agent.config.parallel_tools = true; - 645
let mut agent = h.agent; - 646
let outcome = agent - 647
.run( - 648
"fan out", - 649
&Default::default(), - 650
CancellationToken::new(), - 651
h.events_tx.clone(), - 652
) - 653
.await; - 654
assert!(matches!(outcome, TurnOutcome::Completed { .. })); - 655
- 656
let session = agent.session.lock().await; - 657
// Raw ledger: the closed turn's results are trace lines in the - 658
// projection now (docs/design/68-context-engine.md §10); source order - 659
// is a property of what was recorded, checked here against the raw - 660
// entries. - 661
let msgs: Vec<vak_llm::types::Message> = session - 662
.message_chain() - 663
.into_iter() - 664
.map(|(_, m)| m) - 665
.collect(); - 666
let results_msg = msgs - 667
.iter() - 668
.rev() - 669
.find(|m| matches!(m.content.first(), Some(ContentBlock::ToolResult { .. }))) - 670
.expect("tool results message must exist"); - 671
assert_eq!(results_msg.content.len(), 2); - 672
if let ContentBlock::ToolResult { tool_use_id, .. } = &results_msg.content[0] { - 673
assert_eq!( - 674
tool_use_id, "a", - 675
"results must follow assistant source order" - 676
); - 677
} else { - 678
panic!("expected tool result block"); - 679
} - 680
} - 681
- 682
#[tokio::test] - 683
async fn abort_mid_stream_returns_partial_and_persists_it() { - 684
let h = harness( - 685
vec![ScriptedResponse::Error(LlmError::Aborted { - 686
partial: Some(Box::new(assistant_text("partial work"))), - 687
})], - 688
vec![], - 689
); - 690
let mut agent = h.agent; - 691
let outcome = agent - 692
.run( - 693
"start", - 694
&Default::default(), - 695
CancellationToken::new(), - 696
h.events_tx.clone(), - 697
) - 698
.await; - 699
match outcome { - 700
TurnOutcome::Aborted { partial } => { - 701
assert_eq!(partial.unwrap().text_content(), "partial work"); - 702
} - 703
other => panic!("expected aborted, got {other:?}"), - 704
} - 705
let session = agent.session.lock().await; - 706
assert_eq!( - 707
session.derive_messages().last().unwrap().text_content(), - 708
"partial work" - 709
); - 710
} - 711
- 712
#[tokio::test] - 713
async fn max_turns_guard_stops_runaway_loops() { - 714
let responses: Vec<_> = (0..6) - 715
.map(|_| { - 716
ScriptedResponse::Message(tool_call_msg( - 717
&next_tool_id(), - 718
"bash", - 719
serde_json::json!({"command": "true"}), - 720
)) - 721
}) - 722
.collect(); - 723
let mut h = harness(responses, vec![Arc::new(BashTool)]); - 724
h.agent.config.max_turns = 3; - 725
let mut agent = h.agent; - 726
let outcome = agent - 727
.run( - 728
"loop forever", - 729
&Default::default(), - 730
CancellationToken::new(), - 731
h.events_tx.clone(), - 732
) - 733
.await; - 734
assert!(matches!(outcome, TurnOutcome::MaxTurnsReached)); - 735
} - 736
- 737
#[tokio::test] - 738
async fn projection_invariant_every_request_message_is_logged() { - 739
let h = harness( - 740
vec![ - 741
ScriptedResponse::Message(tool_call_msg( - 742
&next_tool_id(), - 743
"bash", - 744
serde_json::json!({"command": "echo x"}), - 745
)), - 746
ScriptedResponse::Message(assistant_text("ok")), - 747
], - 748
vec![Arc::new(BashTool)], - 749
); - 750
let mut agent = h.agent; - 751
let _ = agent - 752
.run( - 753
"invariant check", - 754
&Default::default(), - 755
CancellationToken::new(), - 756
h.events_tx.clone(), - 757
) - 758
.await; - 759
- 760
// Raw ledger, not the model-visible projection: once this turn closes, - 761
// `derive_messages()` collapses it to its two-message full record - 762
// (docs/design/68-context-engine.md §10) and no longer contains the - 763
// raw tool_use/tool_result blocks a mid-turn request carried — those - 764
// are still reconstructable (a trace line + `SessionLog::evidence`), - 765
// just not via byte-identical containment in the CURRENT projection. - 766
// Invariant 1 ("model-visible means logged") is checked here against - 767
// what was actually appended to the ledger. - 768
let raw_msgs: Vec<vak_llm::types::Message> = agent - 769
.session - 770
.lock() - 771
.await - 772
.message_chain() - 773
.into_iter() - 774
.map(|(_, m)| m) - 775
.collect(); - 776
for req in h.requests.lock().unwrap().iter() { - 777
for msg in &req.messages { - 778
assert!( - 779
raw_msgs.contains(msg), - 780
"model-visible means logged: request message missing from the ledger" - 781
); - 782
} - 783
} - 784
} - 785
Indexing the workspace…
Vakyartha documentation is discovering safe artifacts, anchors, and source references.