- 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, ContentBlock, Usage}; - 12
use vak_llm::{EventStream, LlmError, Provider}; - 13
use vak_session::{Entry, EntryPayload, SessionPath}; - 14
- 15
struct Scripted { - 16
responses: Mutex<VecDeque<AssistantMessage>>, - 17
} - 18
- 19
#[async_trait::async_trait] - 20
impl Provider for Scripted { - 21
fn name(&self) -> &str { - 22
"scripted" - 23
} - 24
- 25
async fn stream( - 26
&self, - 27
_request: ChatRequest, - 28
_cancel: CancellationToken, - 29
) -> Result<EventStream, LlmError> { - 30
let next = self.responses.lock().unwrap().pop_front(); - 31
let (mut sink, rx) = stream::channel(64); - 32
match next { - 33
Some(m) => { - 34
sink.push(stream::StreamEvent::Start { partial: m.clone() }); - 35
sink.close_message(m).await; - 36
} - 37
None => sink.close_error(LlmError::Parse("exhausted".into())).await, - 38
} - 39
Ok(rx) - 40
} - 41
} - 42
- 43
fn text(t: &str) -> AssistantMessage { - 44
AssistantMessage { - 45
content: vec![ContentBlock::text(t)], - 46
stop_reason: vak_llm::types::StopReason::EndTurn, - 47
usage: Usage { - 48
input_tokens: 7, - 49
output_tokens: 3, - 50
..Default::default() - 51
}, - 52
model: "test-model".into(), - 53
response_id: None, - 54
} - 55
} - 56
- 57
fn tool_call(id: &str, name: &str, input: serde_json::Value) -> AssistantMessage { - 58
AssistantMessage { - 59
content: vec![ContentBlock::ToolUse { - 60
id: id.into(), - 61
name: name.into(), - 62
input, - 63
}], - 64
stop_reason: vak_llm::types::StopReason::ToolUse, - 65
usage: Usage::default(), - 66
model: "test-model".into(), - 67
response_id: None, - 68
} - 69
} - 70
- 71
/// Spawn the secured stack exactly like `serve --gateway`: bearer token + - 72
/// forced gateway enable, independent of project config. - 73
async fn spawn_gateway( - 74
provider: Arc<dyn Provider>, - 75
) -> ( - 76
String, - 77
String, - 78
std::path::PathBuf, - 79
tokio::task::JoinHandle<()>, - 80
) { - 81
let dir = tempfile::tempdir().unwrap(); - 82
let cwd = dir.path().to_path_buf(); - 83
// Hermetic against the developer's global config (e.g. reflection=true): - 84
// pin learning flags off for deterministic scripted flows. - 85
let _ = std::fs::create_dir_all(cwd.join(".vak")); - 86
let _ = std::fs::write( - 87
cwd.join(".vak/config.toml"), - 88
"permission_mode = \"full-access\"\n[memory]\nreflection = false\n[gateway]\nchat_allowlist_open = true\n", - 89
); - 90
vak_config::paths::isolate_home_for_tests(); - 91
let core = Core::new_with_trust(cwd.clone(), true).unwrap(); - 92
core.set_sessions_home(dir.path().join("home")); - 93
// Gateway turns run unattended with AutoDeny; give the fixture bash - 94
// execution so scripted tool flows behave like an interactive session. - 95
core.set_permission_mode(vak_config::PermissionMode::FullAccess); - 96
// Pin the REAL brokered-tool worker instead of `current_exe()`: under - 97
// `cargo test` the latter is the test harness itself, which cannot speak - 98
// the broker protocol and makes brokered bash flake (or hard-fail) under - 99
// CPU contention. Mirrors the established fixture pattern in - 100
// scheduler_personal_os.rs / inbox_endpoints.rs / vak-tool-worker.rs. - 101
core.set_tool_worker_exe(std::path::PathBuf::from(env!( - 102
"CARGO_BIN_EXE_vak-tool-worker" - 103
))); - 104
core.set_provider_instance(provider); - 105
std::mem::forget(dir); - 106
- 107
let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); - 108
let addr = listener.local_addr().unwrap(); - 109
let (app, token) = vak_server::secured_router_with(core, true); - 110
let handle = tokio::spawn(async move { - 111
axum::serve(listener, app).await.unwrap(); - 112
}); - 113
(format!("http://{addr}"), token, cwd, handle) - 114
} - 115
- 116
/// Scheduler-free variant (gateway_router): dropping it releases session - 117
/// locks immediately, like a real process exit. - 118
async fn spawn_gateway_bare( - 119
provider: Arc<dyn Provider>, - 120
) -> (String, std::path::PathBuf, tokio::task::JoinHandle<()>) { - 121
let dir = tempfile::tempdir().unwrap(); - 122
let cwd = dir.path().to_path_buf(); - 123
// Hermetic against the developer's global config (e.g. reflection=true): - 124
// pin learning flags off for deterministic scripted flows. - 125
let _ = std::fs::create_dir_all(cwd.join(".vak")); - 126
let _ = std::fs::write( - 127
cwd.join(".vak/config.toml"), - 128
"permission_mode = \"full-access\"\n[memory]\nreflection = false\n[gateway]\nchat_allowlist_open = true\n", - 129
); - 130
vak_config::paths::isolate_home_for_tests(); - 131
let core = Core::new_with_trust(cwd.clone(), true).unwrap(); - 132
core.set_sessions_home(dir.path().join("home")); - 133
core.set_permission_mode(vak_config::PermissionMode::FullAccess); - 134
core.set_provider_instance(provider); - 135
std::mem::forget(dir); - 136
- 137
let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); - 138
let addr = listener.local_addr().unwrap(); - 139
let app = vak_server::gateway_router(core); - 140
let handle = tokio::spawn(async move { - 141
axum::serve(listener, app).await.unwrap(); - 142
}); - 143
(format!("http://{addr}"), cwd, handle) - 144
} - 145
- 146
async fn spawn_plain(provider: Arc<dyn Provider>) -> (String, tokio::task::JoinHandle<()>) { - 147
let dir = tempfile::tempdir().unwrap(); - 148
let cwd = dir.path().to_path_buf(); - 149
vak_config::paths::isolate_home_for_tests(); - 150
let core = Core::new(cwd.clone()).unwrap(); - 151
core.set_sessions_home(dir.path().join("home")); - 152
core.set_provider_instance(provider); - 153
std::mem::forget(dir); - 154
- 155
let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); - 156
let addr = listener.local_addr().unwrap(); - 157
let app = vak_server::router(core); - 158
let handle = tokio::spawn(async move { - 159
axum::serve(listener, app).await.unwrap(); - 160
}); - 161
(format!("http://{addr}"), handle) - 162
} - 163
- 164
fn client_with(token: &str) -> reqwest::Client { - 165
reqwest::ClientBuilder::new() - 166
.default_headers({ - 167
let mut h = reqwest::header::HeaderMap::new(); - 168
h.insert( - 169
reqwest::header::AUTHORIZATION, - 170
format!("Bearer {token}").parse().unwrap(), - 171
); - 172
h - 173
}) - 174
.build() - 175
.unwrap() - 176
} - 177
- 178
async fn inbound( - 179
client: &reqwest::Client, - 180
base: &str, - 181
body: serde_json::Value, - 182
) -> reqwest::Response { - 183
client - 184
.post(format!("{base}/gateway/inbound")) - 185
.json(&body) - 186
.send() - 187
.await - 188
.unwrap() - 189
} - 190
- 191
#[tokio::test(flavor = "multi_thread", worker_threads = 4)] - 192
async fn inbound_wait_roundtrip_reuses_binding() { - 193
let provider = Arc::new(Scripted { - 194
responses: Mutex::new(VecDeque::from(vec![ - 195
text("gateway hello"), - 196
text("second reply"), - 197
])), - 198
}); - 199
let (base, token, home, _server) = spawn_gateway(provider).await; - 200
let client = client_with(&token); - 201
let msg = |text: &str| { - 202
serde_json::json!({ - 203
"surface": "webhook", - 204
"chat": "ci", - 205
"sender": "bot", - 206
"text": text, - 207
"wait": true - 208
}) - 209
}; - 210
- 211
// Auth middleware covers gateway routes too. - 212
let anon = reqwest::Client::new(); - 213
let res = anon - 214
.post(format!("{base}/gateway/inbound")) - 215
.json(&msg("hi")) - 216
.send() - 217
.await - 218
.unwrap(); - 219
assert_eq!(res.status(), 401); - 220
- 221
let mut first = msg("hello from chat"); - 222
first["request_id"] = serde_json::json!("gateway-req-1"); - 223
let res = inbound(&client, &base, first.clone()).await; - 224
assert_eq!(res.status(), 200, "wait roundtrip must complete"); - 225
let body: serde_json::Value = res.json().await.unwrap(); - 226
assert_eq!(body["state"], "completed"); - 227
assert_eq!(body["text"], "gateway hello"); - 228
let sid = body["session_id"].as_str().unwrap().to_string(); - 229
assert!(!sid.is_empty()); - 230
- 231
let agent_home = vak_config::paths::agent_home_at(&home.join("home"), "vak"); - 232
let ledger = SessionPath::new_session_file(&agent_home, &home, &sid); - 233
let raw_ledger = std::fs::read_to_string(ledger).unwrap(); - 234
assert!(raw_ledger.lines().any(|line| { - 235
serde_json::from_str::<Entry>(line) - 236
.ok() - 237
.and_then(|entry| match entry.payload { - 238
EntryPayload::Activity(activity) => Some(activity.data), - 239
_ => None, - 240
}) - 241
.is_some_and(|data| { - 242
data.get("request_id").map(String::as_str) == Some("gateway-req-1") - 243
&& data.get("agent_id").map(String::as_str) == Some("vak") - 244
&& data.get("audience_id").map(String::as_str) == Some("webhook:ci") - 245
}) - 246
})); - 247
- 248
// Retries with the same id are acknowledged as duplicates and do not - 249
// consume another provider response or append a second user turn. - 250
let duplicate = inbound(&client, &base, first).await; - 251
assert_eq!(duplicate.status(), 202); - 252
let duplicate_body: serde_json::Value = duplicate.json().await.unwrap(); - 253
assert_eq!(duplicate_body["decision"], "duplicate"); - 254
assert_eq!(duplicate_body["request_id"], "gateway-req-1"); - 255
- 256
// Binding table knows the route; ledger holds the exchange. - 257
let status: serde_json::Value = client - 258
.get(format!("{base}/gateway/status")) - 259
.send() - 260
.await - 261
.unwrap() - 262
.json() - 263
.await - 264
.unwrap(); - 265
assert_eq!(status["enabled"], true); - 266
let bindings = status["bindings"].as_array().unwrap(); - 267
assert_eq!(bindings.len(), 1); - 268
assert_eq!(bindings[0]["target"], "webhook:ci"); - 269
assert_eq!(bindings[0]["session_id"], sid.as_str()); - 270
- 271
let raw = std::fs::read_to_string(home.join("home/gateway/bindings.json")).unwrap(); - 272
assert!(raw.contains("webhook:ci") && raw.contains(&sid)); - 273
- 274
// Second message resumes the SAME session. - 275
let mut second = msg("again"); - 276
second["request_id"] = serde_json::json!("gateway-req-2"); - 277
let res = inbound(&client, &base, second).await; - 278
assert_eq!(res.status(), 200); - 279
let body: serde_json::Value = res.json().await.unwrap(); - 280
assert_eq!(body["text"], "second reply"); - 281
assert_eq!(body["session_id"], sid.as_str()); - 282
- 283
let t: serde_json::Value = client - 284
.get(format!("{base}/sessions/{sid}/transcript")) - 285
.send() - 286
.await - 287
.unwrap() - 288
.json() - 289
.await - 290
.unwrap(); - 291
// Two exchanges project as four messages; the `<conversation_thread>` - 292
// block now rides in the per-turn tail, not the projection - 293
// (docs/design/68-context-engine.md §6). - 294
assert_eq!(t["count"].as_u64(), Some(4)); - 295
} - 296
- 297
#[tokio::test(flavor = "multi_thread", worker_threads = 4)] - 298
async fn telegram_reply_is_an_ordered_multi_message_packet() { - 299
let answer = format!("# Long result\n\n{}", "🧪 result line\n".repeat(700)); - 300
let provider = Arc::new(Scripted { - 301
responses: Mutex::new(VecDeque::from(vec![text(&answer)])), - 302
}); - 303
let (base, token, _home, _server) = spawn_gateway(provider).await; - 304
let client = client_with(&token); - 305
let response = inbound( - 306
&client, - 307
&base, - 308
serde_json::json!({ - 309
"surface": "telegram", - 310
"chat": "42", - 311
"sender": "tester", - 312
"text": "give me the full result", - 313
"wait": true - 314
}), - 315
) - 316
.await; - 317
assert_eq!(response.status(), 200); - 318
let body: serde_json::Value = response.json().await.unwrap(); - 319
assert_eq!(body["delivery"]["fallback_markdown"], answer); - 320
let chunks = body["delivery"]["chunks"].as_array().unwrap(); - 321
assert!( - 322
chunks.len() >= 3, - 323
"long answer must become multiple messages" - 324
); - 325
assert!(chunks.iter().all(|chunk| { - 326
chunk - 327
.as_str() - 328
.is_some_and(|text| text.chars().count() <= 4000) - 329
})); - 330
assert!(chunks.iter().all(|chunk| { - 331
let text = chunk.as_str().unwrap_or_default(); - 332
text.matches('<').count() == text.matches('>').count() - 333
})); - 334
} - 335
- 336
#[tokio::test(flavor = "multi_thread", worker_threads = 4)] - 337
async fn unbind_removes_route_and_404s_after() { - 338
let provider = Arc::new(Scripted { - 339
responses: Mutex::new(VecDeque::from(vec![text("only")])), - 340
}); - 341
let (base, token, _home, _server) = spawn_gateway(provider).await; - 342
let client = client_with(&token); - 343
- 344
let res = inbound( - 345
&client, - 346
&base, - 347
serde_json::json!({"surface":"log","chat":"ops","text":"hi","wait":true}), - 348
) - 349
.await; - 350
assert_eq!(res.status(), 200); - 351
- 352
let res = client - 353
.delete(format!( - 354
"{base}/gateway/bindings/{}", - 355
urlencoding_escape("log:ops") - 356
)) - 357
.send() - 358
.await - 359
.unwrap(); - 360
assert_eq!(res.status(), 200); - 361
- 362
let status: serde_json::Value = client - 363
.get(format!("{base}/gateway/status")) - 364
.send() - 365
.await - 366
.unwrap() - 367
.json() - 368
.await - 369
.unwrap(); - 370
assert_eq!(status["bindings"].as_array().unwrap().len(), 0); - 371
- 372
let res = client - 373
.delete(format!( - 374
"{base}/gateway/bindings/{}", - 375
urlencoding_escape("log:ops") - 376
)) - 377
.send() - 378
.await - 379
.unwrap(); - 380
assert_eq!(res.status(), 404); - 381
} - 382
- 383
#[tokio::test(flavor = "multi_thread", worker_threads = 4)] - 384
async fn bindings_survive_process_restart() { - 385
let provider = Arc::new(Scripted { - 386
responses: Mutex::new(VecDeque::from(vec![text("first life")])), - 387
}); - 388
let (base, cwd, server) = spawn_gateway_bare(provider).await; - 389
let client = reqwest::Client::new(); - 390
- 391
let res = inbound( - 392
&client, - 393
&base, - 394
serde_json::json!({"surface":"telegram","chat":"48211","text":"hi","wait":true}), - 395
) - 396
.await; - 397
assert_eq!(res.status(), 200); - 398
let sid: String = res.json::<serde_json::Value>().await.unwrap()["session_id"] - 399
.as_str() - 400
.unwrap() - 401
.to_string(); - 402
- 403
// Simulate a process restart: kill the first runtime, boot a fresh - 404
// router over the SAME sessions home. The persisted binding must route - 405
// to the same ledger, reattached from disk. - 406
server.abort(); - 407
let _ = server.await; - 408
- 409
vak_config::paths::isolate_home_for_tests(); - 410
let core2 = Core::new(cwd.clone()).unwrap(); - 411
core2.set_sessions_home(cwd.join("home")); - 412
core2.set_provider_instance(Arc::new(Scripted { - 413
responses: Mutex::new(VecDeque::from(vec![text("reborn")])), - 414
})); - 415
let listener2 = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); - 416
let addr2 = listener2.local_addr().unwrap(); - 417
let (app2, token2) = vak_server::secured_router_with(core2, true); - 418
tokio::spawn(async move { - 419
axum::serve(listener2, app2).await.unwrap(); - 420
}); - 421
let base2 = format!("http://{addr2}"); - 422
let client2 = client_with(&token2); - 423
- 424
// No handles live in process #2: the binding must reopen the ledger. - 425
let status: serde_json::Value = client2 - 426
.get(format!("{base2}/gateway/status")) - 427
.send() - 428
.await - 429
.unwrap() - 430
.json() - 431
.await - 432
.unwrap(); - 433
assert_eq!(status["bindings"][0]["session_id"], sid.as_str()); - 434
- 435
let res = inbound( - 436
&client2, - 437
&base2, - 438
serde_json::json!({"surface":"telegram","chat":"48211","text":"still there?","wait":true}), - 439
) - 440
.await; - 441
assert_eq!(res.status(), 200); - 442
let body: serde_json::Value = res.json().await.unwrap(); - 443
assert_eq!( - 444
body["session_id"], - 445
sid.as_str(), - 446
"same session after restart" - 447
); - 448
assert_eq!(body["text"], "reborn"); - 449
- 450
let t: serde_json::Value = client2 - 451
.get(format!("{base2}/sessions/{sid}/transcript")) - 452
.send() - 453
.await - 454
.unwrap() - 455
.json() - 456
.await - 457
.unwrap(); - 458
assert_eq!(t["count"].as_u64(), Some(4), "history continued"); - 459
} - 460
- 461
/// Whether the session's turn is running: `GET /sessions` reports it. The - 462
/// transcript used to answer "run in progress" while a turn held the ledger - 463
/// and now answers normally, which left these tests waiting for a busy - 464
/// signal that never came. - 465
async fn session_running(client: &reqwest::Client, base: &str, sid: &str) -> bool { - 466
let Ok(res) = client.get(format!("{base}/sessions")).send().await else { - 467
return false; - 468
}; - 469
let Ok(body) = res.json::<serde_json::Value>().await else { - 470
return false; - 471
}; - 472
body["sessions"].as_array().is_some_and(|sessions| { - 473
sessions - 474
.iter() - 475
.any(|s| s["session_id"] == sid && s["running"] == true) - 476
}) - 477
} - 478
- 479
fn urlencoding_escape(s: &str) -> String { - 480
s.replace(':', "%3A") - 481
} - 482
- 483
#[tokio::test(flavor = "multi_thread", worker_threads = 4)] - 484
async fn busy_message_is_steered_not_dropped() { - 485
let provider = Arc::new(Scripted { - 486
responses: Mutex::new(VecDeque::from(vec![ - 487
tool_call("t1", "bash", serde_json::json!({"command": "sleep 8"})), - 488
text("done two"), - 489
])), - 490
}); - 491
let (base, token, _home, _server) = spawn_gateway(provider).await; - 492
let client = client_with(&token); - 493
- 494
// Start a run that occupies the session for a while. - 495
let res = inbound( - 496
&client, - 497
&base, - 498
serde_json::json!({"surface":"webhook","chat":"ci","text":"first msg"}), - 499
) - 500
.await; - 501
assert_eq!(res.status(), 202); - 502
let body: serde_json::Value = res.json().await.unwrap(); - 503
assert_eq!(body["state"], "started"); - 504
let sid = { - 505
let status: serde_json::Value = client - 506
.get(format!("{base}/gateway/status")) - 507
.send() - 508
.await - 509
.unwrap() - 510
.json() - 511
.await - 512
.unwrap(); - 513
status["bindings"][0]["session_id"] - 514
.as_str() - 515
.unwrap() - 516
.to_string() - 517
}; - 518
- 519
// Wait until the turn actually owns the ledger (transcript answers - 520
// "run in progress" while busy), then send a follow-up: it must be - 521
// queued as logged steering, never rejected. - 522
let deadline_busy = std::time::Instant::now() + Duration::from_secs(5); - 523
loop { - 524
assert!( - 525
std::time::Instant::now() < deadline_busy, - 526
"never became busy" - 527
); - 528
if session_running(&client, &base, &sid).await { - 529
break; - 530
} - 531
tokio::time::sleep(Duration::from_millis(50)).await; - 532
} - 533
let res = inbound( - 534
&client, - 535
&base, - 536
serde_json::json!({ - 537
"surface": "webhook", - 538
"chat": "ci", - 539
"sender": "@alice", - 540
"text": "second msg", - 541
"attachments": [{"mime": "image/png", "data": "TUVPT1c="}], - 542
}), - 543
) - 544
.await; - 545
assert_eq!(res.status(), 202); - 546
let body: serde_json::Value = res.json().await.unwrap(); - 547
assert_eq!(body["state"], "steering_queued"); - 548
- 549
// Eventually every message is visible on the chain and the run ends: - 550
// user, assistant(tool), user(result), user(steered), assistant(final). - 551
// The queued message keeps its image blocks AND its sender attribution - 552
// — busy queueing must never degrade the payload to bare text. - 553
let deadline = std::time::Instant::now() + Duration::from_secs(20); - 554
let mut raw = String::new(); - 555
while std::time::Instant::now() < deadline { - 556
if let Ok(res) = client - 557
.get(format!("{base}/sessions/{sid}/transcript")) - 558
.send() - 559
.await - 560
&& res.status() == reqwest::StatusCode::OK - 561
{ - 562
let t: serde_json::Value = res.json().await.unwrap(); - 563
if t.get("error").is_none() { - 564
raw = serde_json::to_string(&t).unwrap(); - 565
if raw.contains("second msg") && raw.contains("done two") { - 566
break; - 567
} - 568
} - 569
} - 570
tokio::time::sleep(Duration::from_millis(150)).await; - 571
} - 572
assert!(raw.contains("first msg"), "prompt missing: {raw}"); - 573
assert!( - 574
raw.contains("[from @alice] second msg"), - 575
"sender attributed on queued turn: {raw}" - 576
); - 577
assert!( - 578
raw.contains("TUVPT1c="), - 579
"image attachment survives busy queueing: {raw}" - 580
); - 581
assert!(raw.contains("done two"), "run must complete after steering"); - 582
} - 583
- 584
/// Control plane (docs/design/47-commitment-kernel.md): a human's free text - 585
/// while the run is busy is steering, whatever word it starts with; only an - 586
/// explicit command is control. - 587
#[tokio::test(flavor = "multi_thread", worker_threads = 4)] - 588
async fn busy_free_text_steers_and_only_an_explicit_command_cancels() { - 589
let provider = Arc::new(Scripted { - 590
responses: Mutex::new(VecDeque::from(vec![ - 591
tool_call("t1", "bash", serde_json::json!({"command": "sleep 8"})), - 592
text("never reached"), - 593
])), - 594
}); - 595
let (base, token, _home, _server) = spawn_gateway(provider).await; - 596
let client = client_with(&token); - 597
let res = inbound( - 598
&client, - 599
&base, - 600
serde_json::json!({"surface":"webhook","chat":"ci","text":"first msg"}), - 601
) - 602
.await; - 603
assert_eq!(res.status(), 202); - 604
let sid = { - 605
let status: serde_json::Value = client - 606
.get(format!("{base}/gateway/status")) - 607
.send() - 608
.await - 609
.unwrap() - 610
.json() - 611
.await - 612
.unwrap(); - 613
status["bindings"][0]["session_id"] - 614
.as_str() - 615
.unwrap() - 616
.to_string() - 617
}; - 618
let deadline_busy = std::time::Instant::now() + Duration::from_secs(5); - 619
loop { - 620
assert!( - 621
std::time::Instant::now() < deadline_busy, - 622
"never became busy" - 623
); - 624
if session_running(&client, &base, &sid).await { - 625
break; - 626
} - 627
tokio::time::sleep(Duration::from_millis(50)).await; - 628
} - 629
- 630
// Free text that begins with "stop" is a steer, not a cancellation. - 631
let res = inbound( - 632
&client, - 633
&base, - 634
serde_json::json!({ - 635
"surface": "webhook", - 636
"chat": "ci", - 637
"text": "stop using semicolons in the output", - 638
}), - 639
) - 640
.await; - 641
assert_eq!(res.status(), 202); - 642
let body: serde_json::Value = res.json().await.unwrap(); - 643
assert_eq!(body["state"], "steering_queued", "{body}"); - 644
- 645
// A bare "stop" is the explicit command. - 646
let res = inbound( - 647
&client, - 648
&base, - 649
serde_json::json!({"surface": "webhook", "chat": "ci", "text": "stop"}), - 650
) - 651
.await; - 652
let body: serde_json::Value = res.json().await.unwrap(); - 653
assert_eq!(body["state"], "cancelled", "{body}"); - 654
} - 655
- 656
#[tokio::test(flavor = "multi_thread", worker_threads = 4)] - 657
async fn gateway_disabled_by_default_returns_conflict() { - 658
let provider = Arc::new(Scripted { - 659
responses: Mutex::new(VecDeque::from(vec![text("nope")])), - 660
}); - 661
let (base, _server) = spawn_plain(provider).await; - 662
let client = reqwest::Client::new(); // plain router has no auth gate - 663
- 664
let res = inbound( - 665
&client, - 666
&base, - 667
serde_json::json!({"surface":"webhook","chat":"x","text":"hi"}), - 668
) - 669
.await; - 670
assert_eq!(res.status(), 409); - 671
let body: serde_json::Value = res.json().await.unwrap(); - 672
assert!( - 673
body["error"] - 674
.as_str() - 675
.unwrap_or_default() - 676
.contains("gateway disabled") - 677
); - 678
} - 679
- 680
#[tokio::test(flavor = "multi_thread", worker_threads = 4)] - 681
async fn empty_chat_allowlist_denies_by_default() { - 682
// 0c-02: an empty chat_allowlist must fail closed, not open — no - 683
// `chat_allowlist_open = true` here, unlike the shared test helpers. - 684
let provider = Arc::new(Scripted { - 685
responses: Mutex::new(VecDeque::from(vec![text("should never run")])), - 686
}); - 687
let dir = tempfile::tempdir().unwrap(); - 688
let cwd = dir.path().to_path_buf(); - 689
let _ = std::fs::create_dir_all(cwd.join(".vak")); - 690
let _ = std::fs::write( - 691
cwd.join(".vak/config.toml"), - 692
"[memory]\nreflection = false\n", - 693
); - 694
vak_config::paths::isolate_home_for_tests(); - 695
let core = Core::new_with_trust(cwd.clone(), true).unwrap(); - 696
core.set_sessions_home(dir.path().join("home")); - 697
core.set_permission_mode(vak_config::PermissionMode::FullAccess); - 698
core.set_provider_instance(provider); - 699
std::mem::forget(dir); - 700
- 701
let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); - 702
let addr = listener.local_addr().unwrap(); - 703
let app = vak_server::gateway_router(core); - 704
let _server = tokio::spawn(async move { - 705
axum::serve(listener, app).await.unwrap(); - 706
}); - 707
let base = format!("http://{addr}"); - 708
let client = reqwest::Client::new(); - 709
- 710
let res = inbound( - 711
&client, - 712
&base, - 713
serde_json::json!({"surface":"telegram","chat":"999","sender":"999","text":"hi"}), - 714
) - 715
.await; - 716
assert_eq!(res.status(), 403); - 717
let body: serde_json::Value = res.json().await.unwrap(); - 718
// docs/design/34: an unknown chat becomes a reviewable pending entry - 719
// instead of a flat rejection with a config-editing hint. - 720
assert_eq!(body["state"].as_str(), Some("pending")); - 721
} - 722
- 723
#[tokio::test(flavor = "multi_thread", worker_threads = 4)] - 724
async fn cron_task_delivers_summary_to_log_surface() { - 725
let provider = Arc::new(Scripted { - 726
responses: Mutex::new(VecDeque::from(vec![text("task finished cleanly")])), - 727
}); - 728
let (base, token, cwd, _server) = spawn_gateway(provider).await; - 729
let client = client_with(&token); - 730
git_seed(&cwd); - 731
- 732
// Create a routine with a delivery target and fire it immediately. - 733
let res = client - 734
.post(format!("{base}/tasks")) - 735
.json(&serde_json::json!({ - 736
"name": "nightly", - 737
"prompt": "what is the nightly status?", - 738
"interval_secs": 3600, - 739
"deliver_to": "log:ops" - 740
})) - 741
.send() - 742
.await - 743
.unwrap(); - 744
assert_eq!(res.status(), 200, "task create failed"); - 745
- 746
// Malformed targets are rejected at creation time. - 747
let res = client - 748
.post(format!("{base}/tasks")) - 749
.json(&serde_json::json!({ - 750
"name": "bad", - 751
"prompt": "x", - 752
"interval_secs": 3600, - 753
"deliver_to": "nologseparator" - 754
})) - 755
.send() - 756
.await - 757
.unwrap(); - 758
assert_eq!(res.status(), 400); - 759
- 760
let tasks: serde_json::Value = client - 761
.get(format!("{base}/tasks")) - 762
.send() - 763
.await - 764
.unwrap() - 765
.json() - 766
.await - 767
.unwrap(); - 768
let tid = tasks["tasks"][0]["id"].as_str().unwrap().to_string(); - 769
- 770
let res = client - 771
.post(format!("{base}/tasks/{tid}/run-now")) - 772
.send() - 773
.await - 774
.unwrap(); - 775
assert_eq!(res.status(), 202); - 776
- 777
// The summary lands in the log surface's delivery journal. - 778
let path = cwd.join("home/gateway/deliveries.jsonl"); - 779
let deadline = std::time::Instant::now() + Duration::from_secs(20); - 780
let mut found = false; - 781
while std::time::Instant::now() < deadline { - 782
if let Ok(raw) = std::fs::read_to_string(&path) - 783
&& raw.contains("routine 'nightly' finished") - 784
&& raw.contains("task finished cleanly") - 785
{ - 786
found = true; - 787
break; - 788
} - 789
tokio::time::sleep(Duration::from_millis(200)).await; - 790
} - 791
assert!(found, "delivery journal never received the routine summary"); - 792
} - 793
- 794
fn git_seed(cwd: &std::path::Path) { - 795
let run = |args: &[&str]| { - 796
let out = std::process::Command::new("git") - 797
.args(args) - 798
.current_dir(cwd) - 799
.env("GIT_AUTHOR_NAME", "t") - 800
.env("GIT_AUTHOR_EMAIL", "t@t") - 801
.env("GIT_COMMITTER_NAME", "t") - 802
.env("GIT_COMMITTER_EMAIL", "t@t") - 803
.output() - 804
.unwrap(); - 805
assert!( - 806
out.status.success(), - 807
"git {args:?} failed: {}", - 808
String::from_utf8_lossy(&out.stderr) - 809
); - 810
}; - 811
run(&["init", "-q"]); - 812
std::fs::write(cwd.join("README.md"), "seed\n").unwrap(); - 813
run(&["add", "."]); - 814
run(&["commit", "-q", "-m", "seed"]); - 815
} - 816
Indexing the workspace…
Vakyartha documentation is discovering safe artifacts, anchors, and source references.