- 1
//! Inbox chokepoint + endpoints (docs/design/29-personal-os.md P6): every - 2
//! delivery has a durable pull-side twin, watchdog output survives with - 3
//! zero transports configured, and ack/read-state is idempotent over HTTP. - 4
- 5
#![allow(clippy::unwrap_used, clippy::expect_used, clippy::panic)] - 6
- 7
use std::io::BufRead; - 8
use std::path::{Path, PathBuf}; - 9
use std::sync::{ - 10
Arc, - 11
atomic::{AtomicUsize, Ordering}, - 12
}; - 13
- 14
use tokio_util::sync::CancellationToken; - 15
- 16
use vak_core::Core; - 17
use vak_llm::stream; - 18
use vak_llm::types::{AssistantMessage, ChatRequest, ContentBlock, Usage}; - 19
use vak_llm::{EventStream, LlmError, Provider}; - 20
- 21
struct Counting { - 22
dispatches: Arc<AtomicUsize>, - 23
} - 24
- 25
#[async_trait::async_trait] - 26
impl Provider for Counting { - 27
fn name(&self) -> &str { - 28
"counting" - 29
} - 30
- 31
async fn stream( - 32
&self, - 33
_request: ChatRequest, - 34
_cancel: CancellationToken, - 35
) -> Result<EventStream, LlmError> { - 36
self.dispatches.fetch_add(1, Ordering::SeqCst); - 37
let (mut sink, rx) = stream::channel(8); - 38
let done = AssistantMessage { - 39
content: vec![ContentBlock::text("done")], - 40
stop_reason: vak_llm::types::StopReason::EndTurn, - 41
usage: Usage { - 42
input_tokens: 5, - 43
output_tokens: 1, - 44
..Default::default() - 45
}, - 46
model: "counted-model".into(), - 47
response_id: None, - 48
}; - 49
sink.push(stream::StreamEvent::Start { - 50
partial: done.clone(), - 51
}); - 52
sink.close_message(done).await; - 53
Ok(rx) - 54
} - 55
} - 56
- 57
struct Server { - 58
base: String, - 59
home: PathBuf, - 60
_dir: Arc<tempfile::TempDir>, - 61
client: reqwest::Client, - 62
} - 63
- 64
/// Spawn the FULL secured stack (bearer auth + scheduler) over a hermetic - 65
/// workspace, exactly like the desktop shell does. - 66
async fn spawn_server(config_toml: &str) -> Server { - 67
let dir = Arc::new(tempfile::tempdir().unwrap()); - 68
let cwd = dir.path().join("ws"); - 69
std::fs::create_dir_all(cwd.join(".vak")).unwrap(); - 70
std::fs::write( - 71
cwd.join(".vak/config.toml"), - 72
format!("[memory]\nreflection = false\n{config_toml}"), - 73
) - 74
.unwrap(); - 75
- 76
vak_config::paths::isolate_home_for_tests(); - 77
let core = Core::new_with_trust(cwd, true).unwrap(); - 78
let home = dir.path().join("home"); - 79
core.set_sessions_home(home.clone()); - 80
core.set_permission_mode(vak_config::PermissionMode::FullAccess); - 81
// A REAL worker executable: the harness cannot speak the broker - 82
// protocol, and watchdog scripts run through it. - 83
core.set_tool_worker_exe(PathBuf::from(env!("CARGO_BIN_EXE_vak-tool-worker"))); - 84
core.set_provider_instance(Arc::new(Counting { - 85
dispatches: Arc::new(AtomicUsize::new(0)), - 86
})); - 87
- 88
let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); - 89
let addr = listener.local_addr().unwrap(); - 90
let (app, token) = vak_server::secured_router(core); - 91
tokio::spawn(async move { - 92
axum::serve(listener, app).await.unwrap(); - 93
}); - 94
let client = reqwest::ClientBuilder::new() - 95
.default_headers({ - 96
let mut h = reqwest::header::HeaderMap::new(); - 97
h.insert( - 98
reqwest::header::AUTHORIZATION, - 99
format!("Bearer {token}").parse().unwrap(), - 100
); - 101
h - 102
}) - 103
.build() - 104
.unwrap(); - 105
Server { - 106
base: format!("http://{addr}"), - 107
home, - 108
_dir: dir, - 109
client, - 110
} - 111
} - 112
- 113
impl Server { - 114
async fn get_json(&self, path: &str) -> serde_json::Value { - 115
self.client - 116
.get(format!("{}{path}", self.base)) - 117
.send() - 118
.await - 119
.unwrap() - 120
.json() - 121
.await - 122
.unwrap() - 123
} - 124
- 125
async fn post_status(&self, path: &str) -> reqwest::StatusCode { - 126
self.client - 127
.post(format!("{}{path}", self.base)) - 128
.send() - 129
.await - 130
.unwrap() - 131
.status() - 132
} - 133
- 134
async fn post_json( - 135
&self, - 136
path: &str, - 137
body: serde_json::Value, - 138
) -> (reqwest::StatusCode, serde_json::Value) { - 139
let res = self - 140
.client - 141
.post(format!("{}{path}", self.base)) - 142
.json(&body) - 143
.send() - 144
.await - 145
.unwrap(); - 146
let status = res.status(); - 147
(status, res.json::<serde_json::Value>().await.unwrap()) - 148
} - 149
- 150
async fn create_task(&self, body: serde_json::Value) -> String { - 151
let (status, _) = self.post_json("/tasks", body).await; - 152
assert_eq!(status, 200, "create failed"); - 153
let list = self.get_json("/tasks").await; - 154
list["tasks"][0]["id"].as_str().unwrap().to_string() - 155
} - 156
- 157
async fn run_now(&self, id: &str) -> reqwest::StatusCode { - 158
self.post_status(&format!("/tasks/{id}/run-now")).await - 159
} - 160
- 161
async fn inbox(&self) -> serde_json::Value { - 162
self.get_json("/inbox").await - 163
} - 164
} - 165
- 166
fn delivery_lines(home: &Path) -> Vec<(String, String)> { - 167
let Ok(f) = std::fs::File::open(home.join("gateway").join("deliveries.jsonl")) else { - 168
return Vec::new(); - 169
}; - 170
let mut out = Vec::new(); - 171
for line in std::io::BufReader::new(f).lines().map_while(Result::ok) { - 172
if let Ok(v) = serde_json::from_str::<serde_json::Value>(&line) { - 173
out.push(( - 174
v["target"].as_str().unwrap_or_default().to_string(), - 175
v["text"].as_str().unwrap_or_default().to_string(), - 176
)); - 177
} - 178
} - 179
out - 180
} - 181
- 182
async fn wait_until(secs: u64, mut pred: impl FnMut() -> bool) -> bool { - 183
let deadline = std::time::Instant::now() + std::time::Duration::from_secs(secs); - 184
while std::time::Instant::now() < deadline { - 185
if pred() { - 186
return true; - 187
} - 188
tokio::time::sleep(std::time::Duration::from_millis(150)).await; - 189
} - 190
pred() - 191
} - 192
- 193
/// Synchronous ledger read for poll loops (async endpoints cannot be - 194
/// awaited inside a plain closure). - 195
fn ledger_entries(home: &Path) -> Vec<vak_core::inbox::Entry> { - 196
vak_core::inbox::list(home, vak_core::inbox::MAX_SCAN) - 197
} - 198
- 199
// ---- P6 exit criterion ------------------------------------------------------ - 200
- 201
#[tokio::test(flavor = "multi_thread", worker_threads = 4)] - 202
async fn watchdog_summary_lands_in_inbox_with_zero_transports() { - 203
let srv = spawn_server("").await; - 204
- 205
// No deliver_to anywhere: no transport can carry this output. - 206
let tid = srv - 207
.create_task(serde_json::json!({ - 208
"name": "quiet-watch", "script": "echo inbox-signal-42", - 209
"interval_secs": 3600 - 210
})) - 211
.await; - 212
assert_eq!(srv.run_now(&tid).await, 202); - 213
- 214
assert!( - 215
wait_until(15, || { - 216
ledger_entries(&srv.home).iter().any(|e| { - 217
e.kind == vak_core::inbox::Kind::TaskSummary - 218
&& e.title == "watchdog 'quiet-watch'" - 219
&& e.body.contains("inbox-signal-42") - 220
&& e.task_id.as_deref() == Some(tid.as_str()) - 221
}) - 222
}) - 223
.await, - 224
"watchdog summary never landed in the inbox" - 225
); - 226
// The durable ledger itself exists on disk — not just an endpoint view. - 227
assert!(srv.home.join("inbox.jsonl").is_file()); - 228
// Zero transports configured: nothing was delivered anywhere else. - 229
assert!( - 230
!wait_until(2, || !delivery_lines(&srv.home).is_empty()).await, - 231
"no delivery may happen without a configured transport" - 232
); - 233
} - 234
- 235
// ---- Ack idempotence + unread math ------------------------------------------ - 236
- 237
async fn ack_entry(srv: &Server, id: &str) -> (reqwest::StatusCode, serde_json::Value) { - 238
srv.post_json(&format!("/inbox/{id}/ack"), serde_json::json!({})) - 239
.await - 240
} - 241
- 242
#[tokio::test(flavor = "multi_thread", worker_threads = 2)] - 243
async fn ack_is_idempotent_over_http_and_404s_unknown_ids() { - 244
let srv = spawn_server("").await; - 245
let a = vak_core::inbox::record( - 246
&srv.home, - 247
vak_core::inbox::Kind::Digest, - 248
"daily", - 249
"", - 250
None, - 251
None, - 252
) - 253
.unwrap(); - 254
let b = vak_core::inbox::record( - 255
&srv.home, - 256
vak_core::inbox::Kind::TaskSummary, - 257
"run", - 258
"body", - 259
None, - 260
None, - 261
) - 262
.unwrap(); - 263
- 264
let view = srv.inbox().await; - 265
assert_eq!(view["entries"].as_array().unwrap().len(), 2); - 266
assert_eq!(view["unread_count"], 2); - 267
- 268
let (status, body) = ack_entry(&srv, &a.id).await; - 269
assert_eq!(status, 200); - 270
assert_eq!(body["acked"], true, "{body}"); - 271
let (status, body) = ack_entry(&srv, &a.id).await; - 272
assert_eq!(status, 200); - 273
assert_eq!(body["acked"], false, "re-ack is a no-op: {body}"); - 274
- 275
let (status, _) = ack_entry(&srv, "deadbeef0000").await; - 276
assert_eq!(status, 404, "unknown id must 404, not report acked:false"); - 277
- 278
assert_eq!(srv.get_json("/inbox/unread_count").await["count"], 1); - 279
let unread = srv.get_json("/inbox?unread=true").await; - 280
let left = unread["entries"].as_array().unwrap(); - 281
assert_eq!(left.len(), 1); - 282
assert_eq!(left[0]["id"], b.id.as_str()); - 283
} - 284
- 285
#[tokio::test(flavor = "multi_thread", worker_threads = 2)] - 286
async fn unread_count_matches_entries_and_limit_bounds_only_the_list() { - 287
let srv = spawn_server("").await; - 288
for i in 0..3 { - 289
vak_core::inbox::record( - 290
&srv.home, - 291
vak_core::inbox::Kind::Heartbeat, - 292
&format!("beat-{i}"), - 293
"", - 294
None, - 295
None, - 296
) - 297
.unwrap(); - 298
} - 299
- 300
let view = srv.get_json("/inbox").await; - 301
assert_eq!(view["entries"].as_array().unwrap().len(), 3); - 302
assert_eq!(view["unread_count"], 3); - 303
- 304
// limit truncates the window without touching the count. - 305
let capped = srv.get_json("/inbox?limit=2").await; - 306
assert_eq!(capped["entries"].as_array().unwrap().len(), 2); - 307
assert_eq!(capped["unread_count"], 3); - 308
- 309
// Ack everything through the endpoint; unread drains to zero while the - 310
// plain list keeps every entry (tombstones never delete). - 311
for e in view["entries"].as_array().unwrap() { - 312
let id = e["id"].as_str().unwrap(); - 313
assert_eq!(ack_entry(&srv, id).await.1["acked"], true); - 314
} - 315
assert_eq!(srv.get_json("/inbox/unread_count").await["count"], 0); - 316
assert!( - 317
srv.get_json("/inbox?unread=true").await["entries"] - 318
.as_array() - 319
.unwrap() - 320
.is_empty() - 321
); - 322
assert_eq!(srv.inbox().await["entries"].as_array().unwrap().len(), 3); - 323
} - 324
- 325
// ---- Budget alert twin -------------------------------------------------------- - 326
- 327
#[tokio::test(flavor = "multi_thread", worker_threads = 4)] - 328
async fn budget_alert_recorded_once_per_window_alongside_delivery() { - 329
let srv = spawn_server("[finops]\nmax_day_usd = 10.0\n").await; - 330
- 331
// Day spend already past the 80% threshold of the $10 cap. - 332
let ledger = vak_core::finops::FinOpsLedger::new(&srv.home); - 333
ledger - 334
.append(&vak_core::finops::CostRow { - 335
ts: chrono::Utc::now(), - 336
model: "claude-sonnet".into(), - 337
provider: "anthropic".into(), - 338
input_tokens: 1000, - 339
output_tokens: 500, - 340
cache_read_input_tokens: None, - 341
usd: Some(9.0), - 342
source: "estimated".into(), - 343
session_id: "seed".into(), - 344
}) - 345
.unwrap(); - 346
- 347
let tid = srv - 348
.create_task(serde_json::json!({ - 349
"name": "budget-probe", "script": "true", - 350
"interval_secs": 3600, "deliver_to": "log:budget" - 351
})) - 352
.await; - 353
- 354
// First fire crosses the threshold: one delivery AND one inbox entry. - 355
assert_eq!(srv.run_now(&tid).await, 202); - 356
assert!( - 357
wait_until(15, || { - 358
delivery_lines(&srv.home) - 359
.iter() - 360
.any(|(t, x)| t == "log:budget" && x.contains("budget alert [eighty]")) - 361
}) - 362
.await, - 363
"budget alert never delivered" - 364
); - 365
assert!( - 366
wait_until(5, || { - 367
ledger_entries(&srv.home).iter().any(|e| { - 368
e.kind == vak_core::inbox::Kind::BudgetAlert - 369
&& e.session_id.as_deref() == Some(tid.as_str()) - 370
&& e.body.contains("budget alert [eighty]") - 371
}) - 372
}) - 373
.await, - 374
"budget alert never recorded in the inbox" - 375
); - 376
let budget_entries = |view: &serde_json::Value| { - 377
view["entries"] - 378
.as_array() - 379
.unwrap() - 380
.iter() - 381
.filter(|e| e["kind"] == "budget_alert") - 382
.count() - 383
}; - 384
assert_eq!(budget_entries(&srv.inbox().await), 1); - 385
- 386
// Second fire inside the same day window: no redelivery, no new entry. - 387
assert_eq!(srv.run_now(&tid).await, 202); - 388
tokio::time::sleep(std::time::Duration::from_millis(1500)).await; - 389
let deliveries = delivery_lines(&srv.home) - 390
.iter() - 391
.filter(|(_, x)| x.contains("budget alert [eighty]")) - 392
.count(); - 393
assert_eq!(deliveries, 1, "same-level alert must not redeliver"); - 394
assert_eq!( - 395
budget_entries(&srv.inbox().await), - 396
1, - 397
"same-level alert must not re-record" - 398
); - 399
} - 400
Indexing the workspace…
Vakyartha documentation is discovering safe artifacts, anchors, and source references.