- 1
//! Scheduled runs that cannot start say so, and runs that do are findable - 2
//! (docs/plans/data-architecture-plan.md, M0): a refused routine leaves a - 3
//! `routine_failed` inbox entry with the reason and remedy instead of a log - 4
//! line; a folder that is not a git repository is refused loudly; two tasks - 5
//! due in one tick both run; and a run recorded on a task opens after a - 6
//! restart, because its handle is its ledger id. - 7
- 8
#![allow(clippy::unwrap_used, clippy::expect_used, clippy::panic)] - 9
- 10
use std::path::{Path, PathBuf}; - 11
use std::sync::{ - 12
Arc, - 13
atomic::{AtomicUsize, Ordering}, - 14
}; - 15
- 16
use tokio_util::sync::CancellationToken; - 17
- 18
use vak_core::Core; - 19
use vak_llm::stream; - 20
use vak_llm::types::{AssistantMessage, ChatRequest, ContentBlock, Usage}; - 21
use vak_llm::{EventStream, LlmError, Provider}; - 22
- 23
struct Counting(Arc<AtomicUsize>); - 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.0.fetch_add(1, Ordering::SeqCst); - 37
let (mut sink, rx) = stream::channel(8); - 38
let done = AssistantMessage { - 39
content: vec![ContentBlock::text("routine done")], - 40
stop_reason: vak_llm::types::StopReason::EndTurn, - 41
usage: Usage::default(), - 42
model: "counted-model".into(), - 43
response_id: None, - 44
}; - 45
sink.push(stream::StreamEvent::Start { - 46
partial: done.clone(), - 47
}); - 48
sink.close_message(done).await; - 49
Ok(rx) - 50
} - 51
} - 52
- 53
struct Server { - 54
base: String, - 55
token: String, - 56
dispatches: Arc<AtomicUsize>, - 57
} - 58
- 59
impl Server { - 60
fn client(&self) -> reqwest::Client { - 61
let mut headers = reqwest::header::HeaderMap::new(); - 62
headers.insert( - 63
reqwest::header::AUTHORIZATION, - 64
format!("Bearer {}", self.token).parse().unwrap(), - 65
); - 66
reqwest::ClientBuilder::new() - 67
.default_headers(headers) - 68
.build() - 69
.unwrap() - 70
} - 71
- 72
async fn get(&self, path: &str) -> (reqwest::StatusCode, serde_json::Value) { - 73
let response = self - 74
.client() - 75
.get(format!("{}{path}", self.base)) - 76
.send() - 77
.await - 78
.unwrap(); - 79
let status = response.status(); - 80
(status, response.json().await.unwrap_or_default()) - 81
} - 82
- 83
async fn task(&self, id: &str) -> serde_json::Value { - 84
let (_, body) = self.get("/tasks").await; - 85
body["tasks"] - 86
.as_array() - 87
.unwrap() - 88
.iter() - 89
.find(|task| task["id"] == id) - 90
.cloned() - 91
.unwrap_or_default() - 92
} - 93
- 94
async fn run_now(&self, id: &str) -> reqwest::StatusCode { - 95
self.client() - 96
.post(format!("{}/tasks/{id}/run-now", self.base)) - 97
.send() - 98
.await - 99
.unwrap() - 100
.status() - 101
} - 102
- 103
async fn inbox(&self) -> Vec<serde_json::Value> { - 104
let (_, body) = self.get("/inbox").await; - 105
body["entries"].as_array().cloned().unwrap_or_default() - 106
} - 107
} - 108
- 109
fn git_seed(cwd: &Path) { - 110
let run = |args: &[&str]| { - 111
let out = std::process::Command::new("git") - 112
.args(args) - 113
.current_dir(cwd) - 114
.env("GIT_AUTHOR_NAME", "t") - 115
.env("GIT_AUTHOR_EMAIL", "t@t") - 116
.env("GIT_COMMITTER_NAME", "t") - 117
.env("GIT_COMMITTER_EMAIL", "t@t") - 118
.output() - 119
.unwrap(); - 120
assert!(out.status.success(), "git {args:?} failed"); - 121
}; - 122
run(&["init", "-q"]); - 123
std::fs::write(cwd.join("README.md"), "seed\n").unwrap(); - 124
run(&["add", "."]); - 125
run(&["commit", "-q", "-m", "seed"]); - 126
} - 127
- 128
/// A never-run prompt task in `ws`, due at the first tick. - 129
fn prompt_task(id: &str, ws: &Path, agent_id: Option<&str>) -> serde_json::Value { - 130
serde_json::json!({ - 131
"id": id, - 132
"name": id, - 133
"prompt": "summarise the day", - 134
"interval_secs": 3600, - 135
"enabled": true, - 136
"cwd": ws.display().to_string(), - 137
"created_at": chrono::Utc::now().to_rfc3339(), - 138
"last_run_at": null, - 139
"last_session_id": null, - 140
"last_summary": null, - 141
"last_wt": null, - 142
"agent_id": agent_id, - 143
}) - 144
} - 145
- 146
/// A workspace (a git repository when `git`) and a home seeded with - 147
/// `tasks`, returned before any server starts. - 148
fn space( - 149
git: bool, - 150
tasks: impl FnOnce(&Path) -> serde_json::Value, - 151
) -> (tempfile::TempDir, PathBuf, PathBuf) { - 152
let dir = tempfile::tempdir().unwrap(); - 153
let ws = dir.path().join("ws"); - 154
std::fs::create_dir_all(ws.join(".vak")).unwrap(); - 155
std::fs::write( - 156
ws.join(".vak/config.toml"), - 157
"[memory]\nreflection = false\n", - 158
) - 159
.unwrap(); - 160
if git { - 161
git_seed(&ws); - 162
} - 163
let home = dir.path().join("home"); - 164
std::fs::create_dir_all(&home).unwrap(); - 165
std::fs::write( - 166
vak_core::tasks::tasks_file(&home), - 167
serde_json::to_string_pretty(&tasks(&ws)).unwrap(), - 168
) - 169
.unwrap(); - 170
(dir, ws, home) - 171
} - 172
- 173
async fn serve(ws: &Path, home: &Path) -> Server { - 174
vak_config::paths::isolate_home_for_tests(); - 175
let core = Core::new_with_trust(ws.to_path_buf(), true).unwrap(); - 176
core.set_sessions_home(home.to_path_buf()); - 177
core.set_permission_mode(vak_config::PermissionMode::FullAccess); - 178
core.set_tool_worker_exe(PathBuf::from(env!("CARGO_BIN_EXE_vak-tool-worker"))); - 179
let dispatches = Arc::new(AtomicUsize::new(0)); - 180
core.set_provider_instance(Arc::new(Counting(dispatches.clone()))); - 181
let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); - 182
let addr = listener.local_addr().unwrap(); - 183
let (app, token) = vak_server::secured_router_with(core, false); - 184
tokio::spawn(async move { axum::serve(listener, app).await.unwrap() }); - 185
Server { - 186
base: format!("http://{addr}"), - 187
token, - 188
dispatches, - 189
} - 190
} - 191
- 192
async fn eventually<F, Fut>(secs: u64, mut check: F) -> bool - 193
where - 194
F: FnMut() -> Fut, - 195
Fut: std::future::Future<Output = bool>, - 196
{ - 197
let deadline = std::time::Instant::now() + std::time::Duration::from_secs(secs); - 198
while std::time::Instant::now() < deadline { - 199
if check().await { - 200
return true; - 201
} - 202
tokio::time::sleep(std::time::Duration::from_millis(150)).await; - 203
} - 204
check().await - 205
} - 206
- 207
fn routine_failures(inbox: &[serde_json::Value], task: &str) -> Vec<serde_json::Value> { - 208
inbox - 209
.iter() - 210
.filter(|entry| entry["kind"] == "routine_failed" && entry["task_id"] == task) - 211
.cloned() - 212
.collect() - 213
} - 214
- 215
#[tokio::test(flavor = "multi_thread", worker_threads = 4)] - 216
async fn fire_task_records_refusal() { - 217
let (_dir, ws, home) = space(true, |ws| { - 218
serde_json::json!([prompt_task("for-nobody", ws, Some("no-such-agent"))]) - 219
}); - 220
let server = serve(&ws, &home).await; - 221
assert_eq!( - 222
server.run_now("for-nobody").await, - 223
reqwest::StatusCode::UNPROCESSABLE_ENTITY - 224
); - 225
let failures = routine_failures(&server.inbox().await, "for-nobody"); - 226
assert_eq!( - 227
failures.len(), - 228
1, - 229
"one entry for the missed slot, however often it is retried" - 230
); - 231
assert!( - 232
failures[0]["body"] - 233
.as_str() - 234
.unwrap() - 235
.contains("no-such-agent"), - 236
"{failures:?}" - 237
); - 238
assert_eq!( - 239
server.task("for-nobody").await["last_run_status"], - 240
"refused" - 241
); - 242
assert_eq!(server.dispatches.load(Ordering::SeqCst), 0); - 243
} - 244
- 245
#[tokio::test(flavor = "multi_thread", worker_threads = 4)] - 246
async fn non_git_space_routine_is_refused_loudly() { - 247
let (_dir, ws, home) = space(false, |ws| { - 248
serde_json::json!([prompt_task("plain-folder", ws, None)]) - 249
}); - 250
let server = serve(&ws, &home).await; - 251
assert!( - 252
eventually(10, || async { - 253
!routine_failures(&server.inbox().await, "plain-folder").is_empty() - 254
}) - 255
.await, - 256
"the scheduler's refusal reached the inbox" - 257
); - 258
let failures = routine_failures(&server.inbox().await, "plain-folder"); - 259
let body = failures[0]["body"].as_str().unwrap(); - 260
assert!(body.contains("is not a git repository"), "{body}"); - 261
assert!(body.contains("git init"), "the remedy is named: {body}"); - 262
assert_eq!( - 263
server.run_now("plain-folder").await, - 264
reqwest::StatusCode::UNPROCESSABLE_ENTITY - 265
); - 266
assert_eq!( - 267
routine_failures(&server.inbox().await, "plain-folder").len(), - 268
1, - 269
"retries of the same slot do not repeat the entry" - 270
); - 271
assert_eq!(server.dispatches.load(Ordering::SeqCst), 0); - 272
} - 273
- 274
#[tokio::test(flavor = "multi_thread", worker_threads = 4)] - 275
async fn two_tasks_due_same_tick_both_fire() { - 276
let (_dir, ws, home) = space(true, |ws| { - 277
serde_json::json!([ - 278
prompt_task("first", ws, None), - 279
prompt_task("second", ws, None) - 280
]) - 281
}); - 282
let server = serve(&ws, &home).await; - 283
assert!( - 284
eventually(20, || async { - 285
server.task("first").await["last_run_status"] == "complete" - 286
&& server.task("second").await["last_run_status"] == "complete" - 287
}) - 288
.await, - 289
"both routines ran: {} / {}", - 290
server.task("first").await, - 291
server.task("second").await - 292
); - 293
let first = server.task("first").await["last_session_id"].clone(); - 294
let second = server.task("second").await["last_session_id"].clone(); - 295
assert_ne!(first, second); - 296
assert!(server.dispatches.load(Ordering::SeqCst) >= 2); - 297
} - 298
- 299
#[tokio::test(flavor = "multi_thread", worker_threads = 4)] - 300
async fn scheduled_run_resolves_after_restart() { - 301
let (_dir, ws, home) = space(true, |ws| { - 302
serde_json::json!([prompt_task("nightly", ws, None)]) - 303
}); - 304
let server = serve(&ws, &home).await; - 305
assert!( - 306
eventually(20, || async { - 307
server.task("nightly").await["last_run_status"] == "complete" - 308
}) - 309
.await, - 310
"the routine ran" - 311
); - 312
let session = server.task("nightly").await["last_session_id"] - 313
.as_str() - 314
.unwrap() - 315
.to_string(); - 316
- 317
let restarted = serve(&ws, &home).await; - 318
let (status, transcript) = restarted - 319
.get(&format!("/sessions/{session}/transcript")) - 320
.await; - 321
assert_eq!(status, reqwest::StatusCode::OK, "{transcript}"); - 322
assert!( - 323
transcript.to_string().contains("routine done"), - 324
"{transcript}" - 325
); - 326
} - 327
Indexing the workspace…
Vakyartha documentation is discovering safe artifacts, anchors, and source references.