- 1
#![allow(clippy::unwrap_used, clippy::expect_used, clippy::panic)] - 2
- 3
use std::collections::VecDeque; - 4
use std::sync::{Arc, Mutex}; - 5
use std::time::Duration; - 6
- 7
use tokio_util::sync::CancellationToken; - 8
- 9
use vak_core::Core; - 10
use vak_llm::stream; - 11
use vak_llm::types::{AssistantMessage, ChatRequest, Usage}; - 12
use vak_llm::{EventStream, LlmError, Provider}; - 13
- 14
struct Scripted { - 15
responses: Mutex<VecDeque<AssistantMessage>>, - 16
} - 17
- 18
#[async_trait::async_trait] - 19
impl Provider for Scripted { - 20
fn name(&self) -> &str { - 21
"scripted" - 22
} - 23
- 24
async fn stream( - 25
&self, - 26
_request: ChatRequest, - 27
_cancel: CancellationToken, - 28
) -> Result<EventStream, LlmError> { - 29
let next = self.responses.lock().unwrap().pop_front(); - 30
let (mut sink, rx) = stream::channel(64); - 31
match next { - 32
Some(m) => { - 33
sink.push(stream::StreamEvent::Start { partial: m.clone() }); - 34
sink.close_message(m).await; - 35
} - 36
None => sink.close_error(LlmError::Parse("exhausted".into())).await, - 37
} - 38
Ok(rx) - 39
} - 40
} - 41
- 42
/// Spawn `secured_router` exactly like the desktop shell will: bind an - 43
/// ephemeral loopback port and keep the token for authenticated calls. - 44
async fn spawn_secured( - 45
provider: Arc<dyn Provider>, - 46
) -> ( - 47
String, - 48
String, - 49
std::path::PathBuf, - 50
tokio::task::JoinHandle<()>, - 51
) { - 52
let dir = tempfile::tempdir().unwrap(); - 53
let cwd = dir.path().to_path_buf(); - 54
// Isolate from the developer's real global config: layered config - 55
// loads $HOME/.config/vak/config.toml, and a personal - 56
// `[memory] reflection = true` would make the post-turn reflection - 57
// seam consume scripted provider responses mid-test. - 58
let project = cwd.join(".vak"); - 59
std::fs::create_dir_all(&project).unwrap(); - 60
std::fs::write( - 61
project.join("config.toml"), - 62
"[memory]\nreflection = false\n", - 63
) - 64
.unwrap(); - 65
vak_config::paths::isolate_home_for_tests(); - 66
let core = Core::new(cwd.clone()).unwrap(); - 67
core.set_sessions_home(dir.path().join("home")); - 68
// The real brokered-tool worker: under `cargo test`, `current_exe()` is - 69
// the test harness, which answers a preview launch with "0 tests" and - 70
// exits (as in gateway.rs and scheduler_personal_os.rs). - 71
core.set_tool_worker_exe(std::path::PathBuf::from(env!( - 72
"CARGO_BIN_EXE_vak-tool-worker" - 73
))); - 74
core.set_provider_instance(provider); - 75
// keep tempdir alive for the process lifetime of the test - 76
std::mem::forget(dir); - 77
- 78
let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); - 79
let addr = listener.local_addr().unwrap(); - 80
let (app, token) = vak_server::secured_router(core); - 81
let handle = tokio::spawn(async move { - 82
axum::serve(listener, app).await.unwrap(); - 83
}); - 84
(format!("http://{addr}"), token, cwd, handle) - 85
} - 86
- 87
fn client_with(token: &str) -> reqwest::Client { - 88
reqwest::ClientBuilder::new() - 89
.default_headers({ - 90
let mut h = reqwest::header::HeaderMap::new(); - 91
h.insert( - 92
reqwest::header::AUTHORIZATION, - 93
format!("Bearer {token}").parse().unwrap(), - 94
); - 95
h - 96
}) - 97
.build() - 98
.unwrap() - 99
} - 100
- 101
fn text(t: &str) -> AssistantMessage { - 102
AssistantMessage { - 103
content: vec![vak_llm::types::ContentBlock::text(t)], - 104
stop_reason: vak_llm::types::StopReason::EndTurn, - 105
usage: Usage { - 106
input_tokens: 7, - 107
output_tokens: 3, - 108
..Default::default() - 109
}, - 110
model: "test-model".into(), - 111
response_id: None, - 112
} - 113
} - 114
- 115
async fn wait_transcript(client: &reqwest::Client, base: &str, id: &str) -> serde_json::Value { - 116
let deadline = std::time::Instant::now() + std::time::Duration::from_secs(15); - 117
loop { - 118
assert!(std::time::Instant::now() < deadline, "transcript timeout"); - 119
if let Ok(res) = client - 120
.get(format!("{base}/sessions/{id}/transcript")) - 121
.send() - 122
.await - 123
&& res.status() == reqwest::StatusCode::OK - 124
{ - 125
let body: serde_json::Value = res.json().await.unwrap(); - 126
// While a run is live the endpoint answers 200 with an error - 127
// payload; keep polling until the real transcript lands. - 128
if body.get("error").is_none() { - 129
return body; - 130
} - 131
} - 132
tokio::time::sleep(std::time::Duration::from_millis(100)).await; - 133
} - 134
} - 135
- 136
#[tokio::test(flavor = "multi_thread", worker_threads = 4)] - 137
async fn secured_router_enforces_token_and_cors() { - 138
let (base, token, _cwd, _server) = spawn_secured(Arc::new(Scripted { - 139
responses: Mutex::new(VecDeque::new()), - 140
})) - 141
.await; - 142
- 143
// /health is open by design. - 144
let health = reqwest::get(format!("{base}/health")).await.unwrap(); - 145
assert_eq!(health.status(), 200); - 146
- 147
// Missing or wrong bearer tokens are rejected. - 148
let anon = reqwest::get(format!("{base}/sessions")).await.unwrap(); - 149
assert_eq!(anon.status(), 401); - 150
let wrong = reqwest::Client::new() - 151
.get(format!("{base}/sessions")) - 152
.bearer_auth("vk_wrong") - 153
.send() - 154
.await - 155
.unwrap(); - 156
assert_eq!(wrong.status(), 401); - 157
- 158
// Both Tauri webview origins must be allowed through CORS preflight. - 159
for origin in ["tauri://localhost", "http://tauri.localhost"] { - 160
let preflight = reqwest::Client::new() - 161
.request(reqwest::Method::OPTIONS, format!("{base}/sessions")) - 162
.header(reqwest::header::ORIGIN, origin) - 163
.header(reqwest::header::ACCESS_CONTROL_REQUEST_METHOD, "POST") - 164
.send() - 165
.await - 166
.unwrap(); - 167
assert_eq!( - 168
preflight - 169
.headers() - 170
.get(reqwest::header::ACCESS_CONTROL_ALLOW_ORIGIN), - 171
Some(&reqwest::header::HeaderValue::from_bytes(origin.as_bytes()).unwrap()), - 172
); - 173
} - 174
- 175
// The real token grants access. - 176
let ok = client_with(&token) - 177
.get(format!("{base}/sessions")) - 178
.send() - 179
.await - 180
.unwrap(); - 181
assert_eq!(ok.status(), 200); - 182
} - 183
- 184
#[tokio::test(flavor = "multi_thread", worker_threads = 4)] - 185
async fn sessions_list_attach_and_title_roundtrip() { - 186
// Build the persisted ledger directly, then drop it so the lock is - 187
// released — mirroring "TUI wrote this session earlier and exited". - 188
let dir = tempfile::tempdir().unwrap(); - 189
let cwd = dir.path().to_path_buf(); - 190
vak_config::paths::isolate_home_for_tests(); - 191
let core_a = Core::new(cwd.clone()).unwrap(); - 192
core_a.set_sessions_home(cwd.join("home")); - 193
std::mem::forget(dir); - 194
- 195
let mut log = core_a.start_session().await.unwrap(); - 196
let session_id = log.header().unwrap().session_id.clone(); - 197
use vak_llm::types::{ContentBlock as CB, Role}; - 198
log.append_message(vak_session::MessageRecord { - 199
message: vak_llm::Message { - 200
role: Role::User, - 201
content: vec![CB::text("title probe here")], - 202
}, - 203
meta: None, - 204
}) - 205
.unwrap(); - 206
log.append_message(vak_session::MessageRecord { - 207
message: vak_llm::Message { - 208
role: Role::Assistant, - 209
content: vec![CB::text("all done")], - 210
}, - 211
meta: None, - 212
}) - 213
.unwrap(); - 214
drop(log); - 215
- 216
// A fresh server over the same store: lists the session, resumes it. - 217
vak_config::paths::isolate_home_for_tests(); - 218
let core = Core::new(cwd.clone()).unwrap(); - 219
core.set_sessions_home(cwd.join("home")); - 220
let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); - 221
let addr = listener.local_addr().unwrap(); - 222
let (app, token) = vak_server::secured_router(core); - 223
tokio::spawn(async move { - 224
axum::serve(listener, app).await.unwrap(); - 225
}); - 226
let base = format!("http://{addr}"); - 227
let client = client_with(&token); - 228
- 229
// The sidebar projection must see the persisted file with a usable title. - 230
let listed = client.get(format!("{base}/sessions")).send().await.unwrap(); - 231
let body: serde_json::Value = listed.json().await.unwrap(); - 232
let entry = body["sessions"] - 233
.as_array() - 234
.unwrap() - 235
.iter() - 236
.find(|s| s["session_id"] == session_id.as_str()) - 237
.expect("persisted session must be listed"); - 238
assert_eq!(entry["title"], "title probe here"); - 239
assert!(entry["updated_at"].is_string()); - 240
assert_eq!(entry["running"], false); - 241
- 242
let attached = client - 243
.post(format!("{base}/sessions/{session_id}/attach")) - 244
.json(&serde_json::json!({"session_id": session_id})) - 245
.send() - 246
.await - 247
.unwrap(); - 248
assert_eq!(attached.status(), 200); - 249
- 250
let resumed: serde_json::Value = client - 251
.get(format!("{base}/sessions/{session_id}/transcript")) - 252
.send() - 253
.await - 254
.unwrap() - 255
.json() - 256
.await - 257
.unwrap(); - 258
assert_eq!(resumed["count"].as_u64(), Some(2)); - 259
} - 260
- 261
#[tokio::test(flavor = "multi_thread", worker_threads = 4)] - 262
async fn persisted_conversation_accepts_followup_and_streams_without_explicit_attach() { - 263
let dir = tempfile::tempdir().unwrap(); - 264
let cwd = dir.path().to_path_buf(); - 265
vak_config::paths::isolate_home_for_tests(); - 266
let core = Core::new(cwd.clone()).unwrap(); - 267
core.set_sessions_home(cwd.join("home")); - 268
let mut session = core.start_session().await.unwrap(); - 269
let id = session.header().unwrap().session_id.clone(); - 270
session - 271
.append_message(vak_session::MessageRecord { - 272
message: vak_llm::Message { - 273
role: vak_llm::types::Role::User, - 274
content: vec![vak_llm::types::ContentBlock::text("Earlier question")], - 275
}, - 276
meta: None, - 277
}) - 278
.unwrap(); - 279
session - 280
.append_message(vak_session::MessageRecord { - 281
message: vak_llm::Message { - 282
role: vak_llm::types::Role::Assistant, - 283
content: vec![vak_llm::types::ContentBlock::text("Earlier answer")], - 284
}, - 285
meta: None, - 286
}) - 287
.unwrap(); - 288
drop(session); - 289
- 290
// A new process has no live handle. EventSources reconnect before the - 291
// browser sends a follow-up, so both streams and /run must recover it. - 292
let resumed = Core::new(cwd.clone()).unwrap(); - 293
resumed.set_sessions_home(cwd.join("home")); - 294
resumed.set_provider_instance(Arc::new(Scripted { - 295
responses: Mutex::new(VecDeque::from(vec![text("Follow-up answer")])), - 296
})); - 297
let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); - 298
let addr = listener.local_addr().unwrap(); - 299
let (app, token) = vak_server::secured_router(resumed); - 300
tokio::spawn(async move { - 301
axum::serve(listener, app).await.unwrap(); - 302
}); - 303
let base = format!("http://{addr}"); - 304
let client = client_with(&token); - 305
- 306
let events = client - 307
.get(format!("{base}/sessions/{id}/events")) - 308
.send() - 309
.await - 310
.unwrap(); - 311
assert_eq!(events.status(), 200); - 312
let presentation = client - 313
.get(format!("{base}/sessions/{id}/presentation/events")) - 314
.send() - 315
.await - 316
.unwrap(); - 317
assert_eq!(presentation.status(), 200); - 318
drop(events); - 319
drop(presentation); - 320
- 321
let started = client - 322
.post(format!("{base}/sessions/{id}/run")) - 323
.json(&serde_json::json!({"prompt": "Follow-up question"})) - 324
.send() - 325
.await - 326
.unwrap(); - 327
assert_eq!(started.status(), 202); - 328
let deadline = std::time::Instant::now() + Duration::from_secs(15); - 329
loop { - 330
assert!(std::time::Instant::now() < deadline, "follow-up timeout"); - 331
let transcript = wait_transcript(&client, &base, &id).await; - 332
if transcript["count"].as_u64() == Some(4) { - 333
break; - 334
} - 335
tokio::time::sleep(Duration::from_millis(50)).await; - 336
} - 337
} - 338
- 339
#[tokio::test(flavor = "multi_thread", worker_threads = 4)] - 340
async fn sessions_list_hides_abandoned_header_only_drafts() { - 341
let dir = tempfile::tempdir().unwrap(); - 342
let cwd = dir.path().to_path_buf(); - 343
vak_config::paths::isolate_home_for_tests(); - 344
let core = Core::new(cwd.clone()).unwrap(); - 345
core.set_sessions_home(cwd.join("home")); - 346
let draft = core.start_session().await.unwrap(); - 347
let draft_id = draft.header().unwrap().session_id.clone(); - 348
drop(draft); - 349
- 350
let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); - 351
let addr = listener.local_addr().unwrap(); - 352
let (app, token) = vak_server::secured_router(core); - 353
tokio::spawn(async move { - 354
axum::serve(listener, app).await.unwrap(); - 355
}); - 356
- 357
let body: serde_json::Value = client_with(&token) - 358
.get(format!("http://{addr}/sessions")) - 359
.send() - 360
.await - 361
.unwrap() - 362
.json() - 363
.await - 364
.unwrap(); - 365
assert!( - 366
body["sessions"] - 367
.as_array() - 368
.unwrap() - 369
.iter() - 370
.all(|session| session["session_id"] != draft_id) - 371
); - 372
} - 373
- 374
#[tokio::test(flavor = "multi_thread", worker_threads = 4)] - 375
async fn fs_endpoints_are_confined_to_workspace() { - 376
let (base, token, cwd, _server) = spawn_secured(Arc::new(Scripted { - 377
responses: Mutex::new(VecDeque::new()), - 378
})) - 379
.await; - 380
let client = client_with(&token); - 381
- 382
// Write inside the workspace. - 383
let put = client - 384
.put(format!("{base}/fs/file")) - 385
.json(&serde_json::json!({"path": "notes/hello.txt", "content": "hi"})) - 386
.send() - 387
.await - 388
.unwrap(); - 389
assert_eq!(put.status(), 200); - 390
assert_eq!( - 391
tokio::fs::read_to_string(cwd.join("notes/hello.txt")) - 392
.await - 393
.unwrap(), - 394
"hi" - 395
); - 396
- 397
let got: serde_json::Value = client - 398
.get(format!("{base}/fs/file?path=notes/hello.txt")) - 399
.send() - 400
.await - 401
.unwrap() - 402
.json() - 403
.await - 404
.unwrap(); - 405
assert_eq!(got["content"], "hi"); - 406
- 407
// Traversal escape is rejected. - 408
let escape_put = client - 409
.put(format!("{base}/fs/file")) - 410
.json(&serde_json::json!({"path": "../evil.txt", "content": "nope"})) - 411
.send() - 412
.await - 413
.unwrap(); - 414
assert_eq!(escape_put.status(), 403); - 415
- 416
let escape_get = client - 417
.get(format!("{base}/fs/file?path=/etc/hostname")) - 418
.send() - 419
.await - 420
.unwrap(); - 421
assert_eq!(escape_get.status(), 403); - 422
- 423
let missing = client - 424
.get(format!("{base}/fs/file?path=no/such.txt")) - 425
.send() - 426
.await - 427
.unwrap(); - 428
assert_eq!(missing.status(), 404); - 429
} - 430
- 431
#[tokio::test(flavor = "multi_thread", worker_threads = 4)] - 432
async fn mode_switch_and_diff_endpoint() { - 433
let (base, token, cwd, _server) = spawn_secured(Arc::new(Scripted { - 434
responses: Mutex::new(VecDeque::new()), - 435
})) - 436
.await; - 437
let client = client_with(&token); - 438
- 439
let session_id: String = client - 440
.post(format!("{base}/sessions")) - 441
.send() - 442
.await - 443
.unwrap() - 444
.json::<serde_json::Value>() - 445
.await - 446
.unwrap()["session_id"] - 447
.as_str() - 448
.unwrap() - 449
.to_string(); - 450
- 451
// Unknown mode is a bad request. - 452
let bad = client - 453
.post(format!("{base}/config/mode")) - 454
.json(&serde_json::json!({"mode": "yolo"})) - 455
.send() - 456
.await - 457
.unwrap(); - 458
assert_eq!(bad.status(), 400); - 459
- 460
let switched = client - 461
.post(format!("{base}/config/mode")) - 462
.json(&serde_json::json!({"mode": "read-only"})) - 463
.send() - 464
.await - 465
.unwrap(); - 466
assert_eq!(switched.status(), 200); - 467
let health: serde_json::Value = reqwest::get(format!("{base}/health")) - 468
.await - 469
.unwrap() - 470
.json() - 471
.await - 472
.unwrap(); - 473
assert_eq!(health["permission_mode"], "ReadOnly"); - 474
- 475
let config: serde_json::Value = client - 476
.get(format!("{base}/config")) - 477
.send() - 478
.await - 479
.unwrap() - 480
.json() - 481
.await - 482
.unwrap(); - 483
assert_eq!(config["permission_mode"], "ReadOnly"); - 484
assert!(config["paths"]["project_config"].is_string()); - 485
assert!(config["integrations"]["skills"].is_array()); - 486
- 487
let patched = client - 488
.patch(format!("{base}/config")) - 489
.json(&serde_json::json!({ - 490
"provider": "google", - 491
"model": "gemini-test", - 492
"max_turns": 17, - 493
"permission_mode": "workspace-write" - 494
})) - 495
.send() - 496
.await - 497
.unwrap(); - 498
assert_eq!(patched.status(), 200); - 499
let updated: serde_json::Value = client - 500
.get(format!("{base}/config")) - 501
.send() - 502
.await - 503
.unwrap() - 504
.json() - 505
.await - 506
.unwrap(); - 507
assert_eq!(updated["provider"], "google"); - 508
assert_eq!(updated["model"], "gemini-test"); - 509
assert_eq!(updated["max_turns"], 17); - 510
assert_eq!(updated["permission_mode"], "WorkspaceWrite"); - 511
assert_eq!(updated["provider_source"], "project_config"); - 512
assert_eq!(updated["model_source"], "project_config"); - 513
assert!( - 514
updated["route_revision"] - 515
.as_str() - 516
.is_some_and(|r| r.starts_with('r')) - 517
); - 518
let persisted = tokio::fs::read_to_string(cwd.join(".vak/config.toml")) - 519
.await - 520
.unwrap(); - 521
assert!(persisted.contains("provider = \"google\"")); - 522
assert!(persisted.contains("model = \"gemini-test\"")); - 523
vak_config::paths::isolate_home_for_tests(); - 524
let restarted = Core::new_with_trust(cwd, true).unwrap(); - 525
assert_eq!(restarted.effective_provider(), "google"); - 526
assert_eq!(restarted.effective_model(), "gemini-test"); - 527
- 528
// Tempdir is not a git repo: diff endpoint reports that as a value. - 529
let diff: serde_json::Value = client - 530
.get(format!("{base}/sessions/{session_id}/diff")) - 531
.send() - 532
.await - 533
.unwrap() - 534
.json() - 535
.await - 536
.unwrap(); - 537
assert_eq!(diff["error"], "not a git repository"); - 538
} - 539
- 540
#[tokio::test(flavor = "multi_thread", worker_threads = 4)] - 541
async fn fs_tree_lists_and_skips_vendored_dirs() { - 542
let (base, token, cwd, _server) = spawn_secured(Arc::new(Scripted { - 543
responses: Mutex::new(VecDeque::new()), - 544
})) - 545
.await; - 546
let client = client_with(&token); - 547
- 548
tokio::fs::create_dir_all(cwd.join("src/deep")) - 549
.await - 550
.unwrap(); - 551
tokio::fs::write(cwd.join("src/main.rs"), "fn main() {}") - 552
.await - 553
.unwrap(); - 554
tokio::fs::write(cwd.join("src/deep/util.rs"), "pub fn u() {}") - 555
.await - 556
.unwrap(); - 557
tokio::fs::create_dir_all(cwd.join("node_modules/left-pad")) - 558
.await - 559
.unwrap(); - 560
tokio::fs::write(cwd.join("node_modules/left-pad/index.js"), "x") - 561
.await - 562
.unwrap(); - 563
- 564
let body: serde_json::Value = client - 565
.get(format!("{base}/fs/tree")) - 566
.send() - 567
.await - 568
.unwrap() - 569
.json() - 570
.await - 571
.unwrap(); - 572
- 573
let files: Vec<String> = body["files"] - 574
.as_array() - 575
.unwrap() - 576
.iter() - 577
.map(|f| f.as_str().unwrap().to_string()) - 578
.collect(); - 579
assert!( - 580
files.contains(&"src/main.rs".to_string()), - 581
"files: {files:?}" - 582
); - 583
assert!(files.contains(&"src/deep/util.rs".to_string())); - 584
assert!( - 585
!files.iter().any(|f| f.starts_with("node_modules")), - 586
"vendored dirs must be skipped" - 587
); - 588
} - 589
- 590
#[tokio::test(flavor = "multi_thread", worker_threads = 4)] - 591
async fn failed_side_turn_does_not_wedge_the_session() { - 592
// Regression (v0.6 deployment gate): whatever way a side turn ends - 593
// without a restorable log — provider exhausted into endless endurance - 594
// retries, or an explicit cancel — the session handle must become - 595
// usable again instead of answering "run in progress" forever. - 596
let provider = Arc::new(Scripted { - 597
responses: Mutex::new(VecDeque::from(vec![ - 598
text("main answer"), - 599
text("side answer"), - 600
])), - 601
}); - 602
let (base, token, _cwd, _server) = spawn_secured(provider).await; - 603
let client = client_with(&token); - 604
- 605
let session_id: String = client - 606
.post(format!("{base}/sessions")) - 607
.send() - 608
.await - 609
.unwrap() - 610
.json::<serde_json::Value>() - 611
.await - 612
.unwrap()["session_id"] - 613
.as_str() - 614
.unwrap() - 615
.to_string(); - 616
- 617
client - 618
.post(format!("{base}/sessions/{session_id}/run")) - 619
.json(&serde_json::json!({"prompt": "say main answer"})) - 620
.send() - 621
.await - 622
.unwrap(); - 623
let before = wait_transcript(&client, &base, &session_id).await; - 624
assert_eq!(before["count"].as_u64(), Some(2)); - 625
- 626
// Start a side turn and cancel it mid-flight. - 627
let started = client - 628
.post(format!("{base}/sessions/{session_id}/side")) - 629
.json(&serde_json::json!({"question": "long aside?"})) - 630
.send() - 631
.await - 632
.unwrap(); - 633
assert_eq!(started.status(), 202); - 634
let cancelled = client - 635
.post(format!("{base}/sessions/{session_id}/side/cancel")) - 636
.send() - 637
.await - 638
.unwrap(); - 639
assert_eq!(cancelled.status(), 202); - 640
- 641
// The handle must come back: transcript answers with real content and - 642
// a follow-up side turn is accepted rather than 409-conflicted. - 643
let after = wait_transcript(&client, &base, &session_id).await; - 644
assert_eq!( - 645
after["count"].as_u64(), - 646
Some(2), - 647
"cancelled side turn must leave the main chain intact" - 648
); - 649
let again = client - 650
.post(format!("{base}/sessions/{session_id}/side")) - 651
.json(&serde_json::json!({"question": "still here?"})) - 652
.send() - 653
.await - 654
.unwrap(); - 655
assert_eq!(again.status(), 202, "session handle must be restorable"); - 656
} - 657
- 658
#[tokio::test(flavor = "multi_thread", worker_threads = 4)] - 659
async fn side_chat_branches_off_and_restores_main_chain() { - 660
let provider = Arc::new(Scripted { - 661
responses: Mutex::new(VecDeque::from(vec![ - 662
text("main answer"), - 663
text("side answer"), - 664
])), - 665
}); - 666
let (base, token, cwd, _server) = spawn_secured(provider).await; - 667
let client = client_with(&token); - 668
- 669
let session_id: String = client - 670
.post(format!("{base}/sessions")) - 671
.send() - 672
.await - 673
.unwrap() - 674
.json::<serde_json::Value>() - 675
.await - 676
.unwrap()["session_id"] - 677
.as_str() - 678
.unwrap() - 679
.to_string(); - 680
- 681
client - 682
.post(format!("{base}/sessions/{session_id}/run")) - 683
.json(&serde_json::json!({"prompt": "say main answer"})) - 684
.send() - 685
.await - 686
.unwrap(); - 687
let before = wait_transcript(&client, &base, &session_id).await; - 688
assert_eq!(before["count"].as_u64(), Some(2)); - 689
- 690
let started = client - 691
.post(format!("{base}/sessions/{session_id}/side")) - 692
.json(&serde_json::json!({"question": "quick aside?"})) - 693
.send() - 694
.await - 695
.unwrap(); - 696
assert_eq!(started.status(), 202); - 697
- 698
// Once the side turn completes, the main view is exactly what it was. - 699
let after = wait_transcript(&client, &base, &session_id).await; - 700
assert_eq!(after["count"].as_u64(), Some(2), "main chain untouched"); - 701
assert!( - 702
!serde_json::to_string(&after) - 703
.unwrap() - 704
.contains("side answer"), - 705
"side answer must not leak into derive_messages()" - 706
); - 707
assert!( - 708
serde_json::to_string(&after) - 709
.unwrap() - 710
.contains("main answer"), - 711
"main context must remain intact" - 712
); - 713
- 714
// ...but the sibling branch stays reconstructable in the JSONL. - 715
tokio::time::sleep(Duration::from_millis(100)).await; - 716
let agent_home = cwd.join("home").join("agents").join("vak"); - 717
let home = if agent_home.exists() { - 718
agent_home - 719
} else { - 720
cwd.join("home") - 721
}; - 722
let mut dir = std::fs::read_dir(vak_session::SessionPath::sessions_dir(&home, &cwd)).unwrap(); - 723
let path = dir.next().unwrap().unwrap().path(); - 724
let raw = std::fs::read_to_string(path).unwrap(); - 725
assert!(raw.contains("quick aside?"), "question persisted"); - 726
assert!( - 727
raw.contains("side answer"), - 728
"answer persisted as branch entry" - 729
); - 730
} - 731
- 732
fn git_seed(cwd: &std::path::Path) { - 733
let run = |args: &[&str]| { - 734
let out = std::process::Command::new("git") - 735
.args(args) - 736
.current_dir(cwd) - 737
.env("GIT_AUTHOR_NAME", "t") - 738
.env("GIT_AUTHOR_EMAIL", "t@t") - 739
.env("GIT_COMMITTER_NAME", "t") - 740
.env("GIT_COMMITTER_EMAIL", "t@t") - 741
.output() - 742
.unwrap(); - 743
assert!( - 744
out.status.success(), - 745
"git {args:?} failed: {}", - 746
String::from_utf8_lossy(&out.stderr) - 747
); - 748
}; - 749
run(&["init", "-q"]); - 750
std::fs::write(cwd.join("README.md"), "seed\n").unwrap(); - 751
run(&["add", "."]); - 752
run(&["commit", "-q", "-m", "seed"]); - 753
} - 754
- 755
#[tokio::test(flavor = "multi_thread", worker_threads = 4)] - 756
async fn bestofn_fans_out_keep_and_discard() { - 757
let provider = Arc::new(Scripted { - 758
responses: Mutex::new(VecDeque::from(vec![ - 759
text("candidate A"), - 760
text("candidate B"), - 761
])), - 762
}); - 763
let (base, token, cwd, _server) = spawn_secured(provider).await; - 764
let client = client_with(&token); - 765
git_seed(&cwd); - 766
- 767
let anchor: String = client - 768
.post(format!("{base}/sessions")) - 769
.send() - 770
.await - 771
.unwrap() - 772
.json::<serde_json::Value>() - 773
.await - 774
.unwrap()["session_id"] - 775
.as_str() - 776
.unwrap() - 777
.to_string(); - 778
- 779
// Not a repo guard would have fired before seed; now it must succeed. - 780
let res = client - 781
.post(format!("{base}/sessions/{anchor}/bestofn")) - 782
.json(&serde_json::json!({"prompt": "solve it", "n": 2})) - 783
.send() - 784
.await - 785
.unwrap(); - 786
let status = res.status(); - 787
let body: serde_json::Value = res.json().await.unwrap(); - 788
assert_eq!(status, 200, "start failed: {body}"); - 789
let runs = body["runs"].as_array().unwrap().clone(); - 790
assert_eq!(runs.len(), 2); - 791
- 792
let wt_root = cwd.join(".vak/worktrees"); - 793
let wait_child = |cid: String| { - 794
let client = client.clone(); - 795
let base = base.clone(); - 796
async move { wait_transcript(&client, &base, &cid).await } - 797
}; - 798
- 799
let mut child_ids = Vec::new(); - 800
for r in &runs { - 801
let cid = r["session_id"].as_str().unwrap().to_string(); - 802
let t = wait_child(cid.clone()).await; - 803
assert_eq!(t["count"].as_u64(), Some(2)); - 804
child_ids.push(cid); - 805
} - 806
- 807
// Both worktrees exist on disk. - 808
assert_eq!(std::fs::read_dir(&wt_root).unwrap().count(), 2); - 809
- 810
// Discard one: worktree + branch gone. - 811
let disc = client - 812
.post(format!("{base}/sessions/{}/discard", child_ids[1])) - 813
.send() - 814
.await - 815
.unwrap(); - 816
assert_eq!(disc.status(), 200); - 817
assert_eq!(std::fs::read_dir(&wt_root).unwrap().count(), 1); - 818
- 819
// Keep the other: endpoint merges (up-to-date is fine for text-only - 820
// candidates) and always cleans up the worktree + branch. - 821
let keep = client - 822
.post(format!("{base}/sessions/{}/keep", child_ids[0])) - 823
.send() - 824
.await - 825
.unwrap(); - 826
assert_eq!( - 827
keep.status(), - 828
200, - 829
"keep failed: {}", - 830
keep.text().await.unwrap_or_default() - 831
); - 832
- 833
assert_eq!( - 834
std::fs::read_dir(&wt_root).unwrap().count(), - 835
0, - 836
"worktrees all cleaned" - 837
); - 838
let branches = { - 839
let o = std::process::Command::new("git") - 840
.args(["branch", "--list", "vak/*"]) - 841
.current_dir(&cwd) - 842
.output() - 843
.unwrap(); - 844
String::from_utf8_lossy(&o.stdout).trim().to_string() - 845
}; - 846
assert!( - 847
branches.is_empty(), - 848
"candidate branches must be gone: {branches}" - 849
); - 850
} - 851
- 852
#[tokio::test(flavor = "multi_thread", worker_threads = 4)] - 853
async fn pr_endpoints_surface_structured_results() { - 854
let (base, token, cwd, _server) = spawn_secured(Arc::new(Scripted { - 855
responses: Mutex::new(VecDeque::new()), - 856
})) - 857
.await; - 858
let client = client_with(&token); - 859
git_seed_main(&cwd); - 860
- 861
let anchor: String = client - 862
.post(format!("{base}/sessions")) - 863
.send() - 864
.await - 865
.unwrap() - 866
.json::<serde_json::Value>() - 867
.await - 868
.unwrap()["session_id"] - 869
.as_str() - 870
.unwrap() - 871
.to_string(); - 872
- 873
// Status is always structured JSON: either a PR view, a no_pr/gh reason, - 874
// or an error — never a hang. - 875
let body: serde_json::Value = client - 876
.get(format!("{base}/sessions/{anchor}/pr")) - 877
.send() - 878
.await - 879
.unwrap() - 880
.json() - 881
.await - 882
.unwrap(); - 883
assert_eq!(body["branch"], "main"); - 884
assert!(body.get("pr").is_some(), "must carry pr slot: {body}"); - 885
if body["pr"].is_null() { - 886
assert!(body.get("reason").is_some() || body.get("error").is_some()); - 887
} - 888
- 889
// Merge endpoint surfaces gh failures as values (409), not hangs. - 890
let merged = client - 891
.post(format!("{base}/sessions/{anchor}/pr/merge")) - 892
.json(&serde_json::json!({"number": 999})) - 893
.send() - 894
.await - 895
.unwrap(); - 896
assert!( - 897
merged.status() == 200 || merged.status() == 409, - 898
"merge must resolve deterministically" - 899
); - 900
- 901
// Unknown session → 404 on merge. - 902
let missing = client - 903
.post(format!("{base}/sessions/nope/pr/merge")) - 904
.json(&serde_json::json!({"number": 1})) - 905
.send() - 906
.await - 907
.unwrap(); - 908
assert_eq!(missing.status(), 404); - 909
} - 910
- 911
fn git_seed_main(cwd: &std::path::Path) { - 912
let run = |args: &[&str]| { - 913
let out = std::process::Command::new("git") - 914
.args(args) - 915
.current_dir(cwd) - 916
.env("GIT_AUTHOR_NAME", "t") - 917
.env("GIT_AUTHOR_EMAIL", "t@t") - 918
.env("GIT_COMMITTER_NAME", "t") - 919
.env("GIT_COMMITTER_EMAIL", "t@t") - 920
.output() - 921
.unwrap(); - 922
assert!(out.status.success(), "git {args:?} failed"); - 923
}; - 924
run(&["init", "-q", "-b", "main"]); - 925
std::fs::write(cwd.join("README.md"), "seed\n").unwrap(); - 926
run(&["add", "."]); - 927
run(&["commit", "-q", "-m", "seed"]); - 928
} - 929
- 930
#[tokio::test(flavor = "multi_thread", worker_threads = 4)] - 931
async fn scheduled_tasks_crud_runnow_and_worktree_churn() { - 932
let provider = Arc::new(Scripted { - 933
responses: Mutex::new(VecDeque::from(vec![ - 934
text("task ran"), - 935
text("task ran again"), - 936
])), - 937
}); - 938
let (base, token, cwd, _server) = spawn_secured(provider).await; - 939
let client = client_with(&token); - 940
git_seed_main(&cwd); - 941
- 942
// Interval validation is a value, not a panic. - 943
let bad = client - 944
.post(format!("{base}/tasks")) - 945
.json(&serde_json::json!({"name":"too hot","prompt":"x","interval_secs":10})) - 946
.send() - 947
.await - 948
.unwrap(); - 949
assert_eq!(bad.status(), 400); - 950
- 951
let created = client - 952
.post(format!("{base}/tasks")) - 953
.json(&serde_json::json!({"name":"nightly sweep","prompt":"sweep the repo","interval_secs":3600})) - 954
.send() - 955
.await - 956
.unwrap(); - 957
assert_eq!(created.status(), 200); - 958
- 959
let listed: serde_json::Value = client - 960
.get(format!("{base}/tasks")) - 961
.send() - 962
.await - 963
.unwrap() - 964
.json() - 965
.await - 966
.unwrap(); - 967
let tasks = listed["tasks"].as_array().unwrap(); - 968
assert_eq!(tasks.len(), 1); - 969
assert_eq!(tasks[0]["name"], "nightly sweep"); - 970
assert_eq!(tasks[0]["enabled"], true); - 971
let tid = tasks[0]["id"].as_str().unwrap().to_string(); - 972
- 973
// Run now: child runs in a worktree; transcript completes; task records it. - 974
let fired = client - 975
.post(format!("{base}/tasks/{tid}/run-now")) - 976
.send() - 977
.await - 978
.unwrap(); - 979
assert_eq!(fired.status(), 202); - 980
- 981
#[allow(unused_assignments)] - 982
let mut child_id = String::new(); - 983
let deadline = std::time::Instant::now() + std::time::Duration::from_secs(15); - 984
loop { - 985
assert!(std::time::Instant::now() < deadline, "task run timeout"); - 986
let t: serde_json::Value = client - 987
.get(format!("{base}/tasks")) - 988
.send() - 989
.await - 990
.unwrap() - 991
.json() - 992
.await - 993
.unwrap(); - 994
let t0 = &t["tasks"][0]; - 995
if let Some(cid) = t0["last_session_id"].as_str() { - 996
child_id = cid.to_string(); - 997
let _ = &child_id; - 998
// The watcher records the run's real final answer, not a - 999
// status word (docs/design/22-gateway.md). - 1000
if t0["last_summary"] == "task ran" {
Indexing the workspace…
Vakyartha documentation is discovering safe artifacts, anchors, and source references.