- 1
//! Heartbeat behaviors (docs/design/29-personal-os.md P7): a "nothing" - 2
//! reply records nothing anywhere, findings park exactly one inbox entry - 3
//! with zero deliveries, an URGENT marker adds the delivery leg through the - 4
//! gateway chokepoint, a breached day-cap denies the dispatch entirely, and - 5
//! a disabled flag never schedules. - 6
- 7
#![allow(clippy::unwrap_used, clippy::expect_used, clippy::panic)] - 8
- 9
use std::io::BufRead; - 10
use std::path::{Path, PathBuf}; - 11
use std::sync::{ - 12
Arc, Mutex, - 13
atomic::{AtomicUsize, Ordering}, - 14
}; - 15
- 16
use chrono::Timelike; - 17
use tokio_util::sync::CancellationToken; - 18
- 19
use vak_core::Core; - 20
use vak_llm::stream; - 21
use vak_llm::types::{AssistantMessage, ChatRequest, ContentBlock, Usage}; - 22
use vak_llm::{EventStream, LlmError, Provider}; - 23
- 24
/// Provider whose final text is settable per-phase and which counts every - 25
/// dispatch, mirroring the counting mock the scheduler tests use. - 26
struct Scripted { - 27
dispatches: Arc<AtomicUsize>, - 28
reply: Arc<Mutex<String>>, - 29
} - 30
- 31
#[async_trait::async_trait] - 32
impl Provider for Scripted { - 33
fn name(&self) -> &str { - 34
"scripted" - 35
} - 36
- 37
async fn stream( - 38
&self, - 39
_request: ChatRequest, - 40
_cancel: CancellationToken, - 41
) -> Result<EventStream, LlmError> { - 42
self.dispatches.fetch_add(1, Ordering::SeqCst); - 43
let text = self - 44
.reply - 45
.lock() - 46
.unwrap_or_else(std::sync::PoisonError::into_inner) - 47
.clone(); - 48
let (mut sink, rx) = stream::channel(8); - 49
let done = AssistantMessage { - 50
content: vec![ContentBlock::text(text)], - 51
stop_reason: vak_llm::types::StopReason::EndTurn, - 52
usage: Usage { - 53
input_tokens: 9, - 54
output_tokens: 1, - 55
..Default::default() - 56
}, - 57
model: "scripted-model".into(), - 58
response_id: None, - 59
}; - 60
sink.push(stream::StreamEvent::Start { - 61
partial: done.clone(), - 62
}); - 63
sink.close_message(done).await; - 64
Ok(rx) - 65
} - 66
} - 67
- 68
struct Fixture { - 69
home: PathBuf, - 70
ws: PathBuf, - 71
dispatches: Arc<AtomicUsize>, - 72
_dir: Arc<tempfile::TempDir>, - 73
} - 74
- 75
/// Hermetic workspace + home with the FULL stack started (bearer auth + - 76
/// background scheduler, heartbeat loop included when enabled). `seed` - 77
/// runs against the home dir before any Core exists — the injection point - 78
/// for cost-ledger rows that must predate the first beat. - 79
async fn spawn_full_seeded( - 80
config_toml: &str, - 81
reply_text: &str, - 82
seed: Option<&dyn Fn(&Path)>, - 83
) -> Fixture { - 84
let dir = Arc::new(tempfile::tempdir().unwrap()); - 85
let ws = dir.path().join("ws"); - 86
std::fs::create_dir_all(ws.join(".vak")).unwrap(); - 87
std::fs::write( - 88
ws.join(".vak/config.toml"), - 89
format!("[memory]\nreflection = false\n{config_toml}"), - 90
) - 91
.unwrap(); - 92
- 93
let home = dir.path().join("home"); - 94
std::fs::create_dir_all(&home).unwrap(); - 95
if let Some(seed) = seed { - 96
seed(&home); - 97
} - 98
let dispatches = Arc::new(AtomicUsize::new(0)); - 99
let reply = Arc::new(Mutex::new(reply_text.to_string())); - 100
vak_config::paths::isolate_home_for_tests(); - 101
let core = Core::new_with_trust(ws.clone(), true).unwrap(); - 102
core.set_sessions_home(home.clone()); - 103
core.set_permission_mode(vak_config::PermissionMode::FullAccess); - 104
core.set_tool_worker_exe(PathBuf::from(env!("CARGO_BIN_EXE_vak-tool-worker"))); - 105
core.set_provider_instance(Arc::new(Scripted { - 106
dispatches: dispatches.clone(), - 107
reply: reply.clone(), - 108
})); - 109
- 110
let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); - 111
let addr = listener.local_addr().unwrap(); - 112
let (app, _token) = vak_server::secured_router_with(core, false); - 113
tokio::spawn(async move { - 114
axum::serve(listener, app).await.unwrap(); - 115
}); - 116
let _ = addr; - 117
- 118
Fixture { - 119
home, - 120
ws, - 121
dispatches, - 122
_dir: dir, - 123
} - 124
} - 125
- 126
async fn spawn_full(config_toml: &str, reply_text: &str) -> Fixture { - 127
spawn_full_seeded(config_toml, reply_text, None).await - 128
} - 129
- 130
fn agent_home(home: &Path) -> PathBuf { - 131
let agent = home.join("agents").join("vak"); - 132
if agent.exists() { - 133
agent - 134
} else { - 135
home.to_path_buf() - 136
} - 137
} - 138
- 139
fn heartbeat_entries(home: &Path) -> Vec<vak_core::inbox::Entry> { - 140
vak_core::inbox::list_scanned(home, 100) - 141
.entries - 142
.into_iter() - 143
.filter(|e| e.kind == vak_core::inbox::Kind::Heartbeat) - 144
.collect() - 145
} - 146
- 147
fn delivery_lines(home: &Path) -> Vec<(String, String)> { - 148
let Ok(f) = std::fs::File::open(home.join("gateway").join("deliveries.jsonl")) else { - 149
return Vec::new(); - 150
}; - 151
let mut out = Vec::new(); - 152
for line in std::io::BufReader::new(f).lines().map_while(Result::ok) { - 153
if let Ok(v) = serde_json::from_str::<serde_json::Value>(&line) { - 154
out.push(( - 155
v["target"].as_str().unwrap_or_default().to_string(), - 156
v["text"].as_str().unwrap_or_default().to_string(), - 157
)); - 158
} - 159
} - 160
out - 161
} - 162
- 163
fn heartbeat_ledger(fx: &Fixture) -> PathBuf { - 164
let resolved = agent_home(&fx.home); - 165
vak_session::SessionPath::new_session_file(&resolved, &fx.ws, "heartbeat") - 166
} - 167
- 168
async fn wait_for_dispatch(fx: &Fixture, secs: u64) -> bool { - 169
let deadline = std::time::Instant::now() + std::time::Duration::from_secs(secs); - 170
while std::time::Instant::now() < deadline { - 171
if fx.dispatches.load(Ordering::SeqCst) > 0 { - 172
return true; - 173
} - 174
tokio::time::sleep(std::time::Duration::from_millis(150)).await; - 175
} - 176
fx.dispatches.load(Ordering::SeqCst) > 0 - 177
} - 178
- 179
async fn settle() { - 180
tokio::time::sleep(std::time::Duration::from_millis(1200)).await; - 181
} - 182
- 183
// ---- Nothing-to-report costs tokens, zero noise --------------------------------- - 184
- 185
#[tokio::test(flavor = "multi_thread", worker_threads = 4)] - 186
async fn nothing_reply_records_nothing_anywhere() { - 187
let fx = spawn_full("[heartbeat]\nenabled = true\n", "nothing").await; - 188
- 189
assert!( - 190
wait_for_dispatch(&fx, 15).await, - 191
"enabled heartbeat must dispatch its first beat" - 192
); - 193
settle().await; - 194
- 195
assert!( - 196
heartbeat_entries(&fx.home).is_empty(), - 197
"a nothing-reply must never enter the inbox" - 198
); - 199
assert!( - 200
delivery_lines(&fx.home).is_empty(), - 201
"a nothing-reply must never be delivered" - 202
); - 203
assert!( - 204
heartbeat_ledger(&fx).is_file(), - 205
"the dedicated persistent session must exist on disk" - 206
); - 207
} - 208
- 209
// ---- Findings park once, push never ---------------------------------------------- - 210
- 211
#[tokio::test(flavor = "multi_thread", worker_threads = 4)] - 212
async fn findings_record_exactly_one_entry_with_zero_deliveries() { - 213
let fx = spawn_full( - 214
"[heartbeat]\nenabled = true\n", - 215
"check the flaky cron\nclose stale PR #4", - 216
) - 217
.await; - 218
- 219
assert!(wait_for_dispatch(&fx, 15).await); - 220
settle().await; - 221
- 222
let entries = heartbeat_entries(&fx.home); - 223
assert_eq!(entries.len(), 1, "findings park exactly one entry"); - 224
assert_eq!(entries[0].title, "heartbeat: 2 findings"); - 225
assert_eq!(entries[0].body, "check the flaky cron\nclose stale PR #4"); - 226
assert_eq!(entries[0].session_id.as_deref(), Some("heartbeat")); - 227
assert!( - 228
delivery_lines(&fx.home).is_empty(), - 229
"non-urgent findings must never ride a transport" - 230
); - 231
} - 232
- 233
#[tokio::test(flavor = "multi_thread", worker_threads = 4)] - 234
async fn urgent_finding_adds_exactly_one_delivery_and_strips_markers() { - 235
let fx = spawn_full( - 236
"[heartbeat]\nenabled = true\n", - 237
"URGENT: disk almost full\nminor lint debt", - 238
) - 239
.await; - 240
- 241
// No tasks configured, so the urgent beat falls back to log:vak. - 242
let deadline = std::time::Instant::now() + std::time::Duration::from_secs(15); - 243
loop { - 244
assert!( - 245
std::time::Instant::now() < deadline, - 246
"urgent finding was never delivered" - 247
); - 248
if delivery_lines(&fx.home) - 249
.iter() - 250
.any(|(_, x)| x.contains("disk almost full")) - 251
{ - 252
break; - 253
} - 254
tokio::time::sleep(std::time::Duration::from_millis(150)).await; - 255
} - 256
settle().await; - 257
- 258
let deliveries = delivery_lines(&fx.home) - 259
.iter() - 260
.filter(|(_, x)| x.contains("disk almost full")) - 261
.count(); - 262
assert_eq!(deliveries, 1, "urgent beat delivers once"); - 263
- 264
let entries = heartbeat_entries(&fx.home); - 265
assert!(!entries.is_empty(), "findings still park an inbox entry"); - 266
assert!( - 267
entries.iter().all(|e| !e.body.contains("URGENT:")), - 268
"markers are stripped from recorded bodies: {:?}", - 269
entries.iter().map(|e| &e.body).collect::<Vec<_>>() - 270
); - 271
assert!( - 272
entries.iter().any(|e| e.body.contains("disk almost full")), - 273
"the finding text survives the strip" - 274
); - 275
} - 276
- 277
// ---- Budget gate ------------------------------------------------------------------ - 278
- 279
#[tokio::test(flavor = "multi_thread", worker_threads = 4)] - 280
async fn breached_day_cap_denies_the_dispatch_silently() { - 281
// Seed the cost ledger BEFORE the server exists: the first beat fires - 282
// immediately, so the cap must already be breached by then. - 283
let seed = |home: &Path| { - 284
vak_core::finops::FinOpsLedger::new(home) - 285
.append(&vak_core::finops::CostRow { - 286
ts: chrono::Utc::now(), - 287
model: "claude-sonnet".into(), - 288
provider: "anthropic".into(), - 289
input_tokens: 1000, - 290
output_tokens: 500, - 291
cache_read_input_tokens: None, - 292
usd: Some(5.0), - 293
source: "estimated".into(), - 294
session_id: "seed".into(), - 295
}) - 296
.unwrap(); - 297
}; - 298
let fx = spawn_full_seeded( - 299
"[heartbeat]\nenabled = true\n\n[finops]\nmax_day_usd = 1.0\n", - 300
"would-be findings", - 301
Some(&seed), - 302
) - 303
.await; - 304
- 305
settle().await; - 306
settle().await; - 307
assert_eq!( - 308
fx.dispatches.load(Ordering::SeqCst), - 309
0, - 310
"budget-denied cycle must dispatch zero" - 311
); - 312
assert!( - 313
!heartbeat_ledger(&fx).exists(), - 314
"denied cycle must not even open the heartbeat session" - 315
); - 316
} - 317
- 318
// ---- Quiet hours -------------------------------------------------------------------- - 319
- 320
#[tokio::test(flavor = "multi_thread", worker_threads = 4)] - 321
async fn quiet_hours_window_skips_the_cycle() { - 322
// Window covering the current local minute plus the next two, so the - 323
// test's whole lifetime sits inside quiet hours regardless of where - 324
// within the minute the process starts. - 325
let now = chrono::Local::now(); - 326
let m = now.hour() * 60 + now.minute(); - 327
let end = (m + 3) % (24 * 60); - 328
let quiet = format!( - 329
"[heartbeat]\nenabled = true\nquiet_hours = \"{:02}:{:02}-{:02}:{:02}\"\n", - 330
m / 60, - 331
m % 60, - 332
end / 60, - 333
end % 60 - 334
); - 335
let fx = spawn_full(&quiet, "would-be findings").await; - 336
- 337
settle().await; - 338
settle().await; - 339
assert_eq!( - 340
fx.dispatches.load(Ordering::SeqCst), - 341
0, - 342
"quiet-hours cycle must not dispatch" - 343
); - 344
assert!(!heartbeat_ledger(&fx).exists()); - 345
} - 346
- 347
// ---- Disabled ------------------------------------------------------------------------ - 348
- 349
#[tokio::test(flavor = "multi_thread", worker_threads = 4)] - 350
async fn disabled_flag_never_schedules() { - 351
let fx = spawn_full("", "findings that must never happen").await; - 352
- 353
settle().await; - 354
settle().await; - 355
assert_eq!(fx.dispatches.load(Ordering::SeqCst), 0); - 356
assert!(!heartbeat_ledger(&fx).exists()); - 357
assert!(heartbeat_entries(&fx.home).is_empty()); - 358
} - 359
Indexing the workspace…
Vakyartha documentation is discovering safe artifacts, anchors, and source references.