- 8164
for project in entries.flatten().filter(|entry| entry.path().is_dir()) { - 8165
let path = project.path().join(format!("{id}.jsonl")); - 8166
if let Some(header) = read(&path) { - 8167
return Some(header); - 8168
} - 8169
} - 8170
} - 8171
if let Ok(agents) = std::fs::read_dir(shared.join("agents")) { - 8172
for agent in agents.flatten().filter(|entry| entry.path().is_dir()) { - 8173
if let Ok(projects) = std::fs::read_dir(agent.path().join("sessions")) { - 8174
for project in projects.flatten().filter(|entry| entry.path().is_dir()) { - 8175
let path = project.path().join(format!("{id}.jsonl")); - 8176
if let Some(header) = read(&path) { - 8177
return Some(header); - 8178
} - 8179
} - 8180
} - 8181
} - 8182
} - 8183
None - 8184
} - 8185
- 8186
/// Reopen a session whose in-memory handle was consumed by a turn that - 8187
/// then failed: `run_turn_with` returns `Err(CoreError)` without the log, - 8188
/// but the append-only ledger file is durable — restore from it so the - 8189
/// session does not stay wedged as "run in progress" forever. - 8190
fn reopen_ledger(core: &vak_core::Core, id: &str) -> Option<vak_session::SessionLog> { - 8191
let path = core - 8192
.sessions_home() - 8193
.join("sessions") - 8194
.join(vak_core::memory::hash_cwd(core.cwd())) - 8195
.join(format!("{id}.jsonl")); - 8196
vak_session::SessionLog::open(path).ok() - 8197
} - 8198
- 8199
// ---- Personal-OS surfaces (docs/design/29-personal-os.md P1–P4) ------------- - 8200
- 8201
#[derive(serde::Deserialize)] - 8202
struct DoctorQuery { - 8203
#[serde(default)] - 8204
session: Option<String>, - 8205
} - 8206
- 8207
fn health_report_json(report: vak_core::health::HealthReport) -> serde_json::Value { - 8208
let ok = report.failures == 0; - 8209
let report_text = report - 8210
.checks - 8211
.iter() - 8212
.map(|check| match &check.detail { - 8213
Ok(detail) => format!("✓ {} — {detail}", check.label), - 8214
Err(detail) => format!("✗ {} — {detail}", check.label), - 8215
}) - 8216
.chain(report.facts.iter().map(|fact| format!("· {fact}"))) - 8217
.collect::<Vec<_>>() - 8218
.join("\n"); - 8219
serde_json::json!({ - 8220
"ok": ok, - 8221
"report": report_text, - 8222
"failures": report.failures, - 8223
"checks": report.checks.iter().map(|c| serde_json::json!({ - 8224
"label": c.label, - 8225
"ok": c.detail.is_ok(), - 8226
"detail": match &c.detail { Ok(d) => d, Err(e) => e }, - 8227
})).collect::<Vec<_>>(), - 8228
"facts": report.facts, - 8229
"ladder": report.ladder.map(|l| serde_json::json!({ - 8230
"legs": l.legs, - 8231
"rendered": l.rendered, - 8232
"objective": l.objective, - 8233
"fallback_legs": l.fallback_legs, - 8234
"annotations": l.annotations, - 8235
})), - 8236
}) - 8237
} - 8238
- 8239
/// `GET /onboarding` — the derived setup projection - 8240
/// (`docs/design/46-stabilization-install-and-onboarding.md` Part III). - 8241
/// - 8242
/// The same `vak_core::onboarding::derive` the CLI renders, so web, - 8243
/// desktop, and terminal cannot disagree about what is configured. The - 8244
/// install manifest is probed by the CLI (which owns install layout) and - 8245
/// is therefore reported here as not-probed; services are probed from the - 8246
/// same service manager the ops endpoints already use. - 8247
async fn onboarding_state(State(state): State<AppState>) -> axum::response::Response { - 8248
use axum::response::IntoResponse; - 8249
// The service manager probe shells out, so it belongs on a blocking - 8250
// worker rather than inside an async handler (invariant 26). - 8251
let services = tokio::task::spawn_blocking(|| { - 8252
let cfg = vak_ops::OpsConfig::detect(); - 8253
let probed: Vec<(String, bool)> = [vak_ops::Service::Gateway, vak_ops::Service::Bridges] - 8254
.into_iter() - 8255
.filter_map(|service| { - 8256
let st = vak_ops::status(service, &cfg); - 8257
(st != vak_ops::State::NotInstalled) - 8258
.then(|| (service.label().to_string(), st == vak_ops::State::Running)) - 8259
}) - 8260
.collect(); - 8261
(!probed.is_empty()).then_some(probed) - 8262
}) - 8263
.await - 8264
.unwrap_or(None); - 8265
- 8266
let awaiting_activation = service_control::activation_drift(&state.core) - 8267
.await - 8268
.map(|drift| drift.awaiting_activation) - 8269
.unwrap_or_default(); - 8270
let projection = vak_core::onboarding::derive( - 8271
&state.core, - 8272
&vak_core::onboarding::ProbedFacts { - 8273
services, - 8274
install: None, - 8275
awaiting_activation, - 8276
}, - 8277
); - 8278
Json(projection).into_response() - 8279
} - 8280
- 8281
/// `POST /onboarding/seed` — install the Shared starter capabilities. - 8282
/// - 8283
/// Idempotent: new standard skills/plugins are added, untouched shipped - 8284
/// content may advance on update, and edited or independently installed - 8285
/// content is preserved. Hooks and network defaults are seeded only when - 8286
/// their configuration layer is empty. Explicit because seeding is a setup - 8287
/// action, never an install side effect (doc 46 D6). - 8288
async fn onboarding_seed(State(state): State<AppState>) -> axum::response::Response { - 8289
use axum::response::IntoResponse; - 8290
// Touches the filesystem and the plugin store; not an async handler's - 8291
// work (invariant 26). - 8292
let outcome = tokio::task::spawn_blocking(vak_core::seed::seed_shared_capabilities).await; - 8293
match outcome { - 8294
Ok(Ok(())) => { - 8295
state.hub.emit_config_changed("capabilities_seeded", ""); - 8296
Json(serde_json::json!({ "ok": true })).into_response() - 8297
} - 8298
Ok(Err(e)) => ( - 8299
StatusCode::INTERNAL_SERVER_ERROR, - 8300
Json(serde_json::json!({ "error": e })), - 8301
) - 8302
.into_response(), - 8303
Err(e) => ( - 8304
StatusCode::INTERNAL_SERVER_ERROR, - 8305
Json(serde_json::json!({ "error": format!("seeding did not complete: {e}") })), - 8306
) - 8307
.into_response(), - 8308
} - 8309
} - 8310
- 8311
#[derive(serde::Deserialize)] - 8312
struct WorkspacePath { - 8313
path: String, - 8314
} - 8315
- 8316
/// `POST /onboarding/workspace-review` — what a folder would ask for. - 8317
/// - 8318
/// Reports privileged sections **without loading them** (doc 46, Step 2). - 8319
/// Describing a project's config by parsing it through the normal loader - 8320
/// would activate the very thing the operator is being asked about. - 8321
async fn onboarding_workspace_review(Json(body): Json<WorkspacePath>) -> axum::response::Response { - 8322
use axum::response::IntoResponse; - 8323
let path = std::path::PathBuf::from(body.path.trim()); - 8324
if std::fs::read_dir(&path).is_err() { - 8325
return ( - 8326
StatusCode::BAD_REQUEST, - 8327
Json(serde_json::json!({ "error": format!("cannot read {}", path.display()) })), - 8328
) - 8329
.into_response(); - 8330
} - 8331
Json(serde_json::json!({ - 8332
"path": path, - 8333
"git": path.join(".git").exists(), - 8334
"requests_privilege": vak_core::trust::requests_privilege(&path), - 8335
"privileges": vak_core::trust::requested_privileges(&path), - 8336
"trusted": vak_core::trust::is_trusted(&path), - 8337
})) - 8338
.into_response() - 8339
} - 8340
- 8341
/// `POST /onboarding/trust` — record an explicit trust decision. - 8342
/// - 8343
/// Only ever *grants*: opening safely is the absence of a decision, and is - 8344
/// already the default, so there is nothing to write for it. Selecting a - 8345
/// folder is never itself consent (doc 46 security invariant 2). - 8346
async fn onboarding_trust( - 8347
State(state): State<AppState>, - 8348
Json(body): Json<WorkspacePath>, - 8349
) -> axum::response::Response { - 8350
use axum::response::IntoResponse; - 8351
let path = std::path::PathBuf::from(body.path.trim()); - 8352
match vak_core::trust::record(&path) { - 8353
Ok(()) => { - 8354
vak_core::security_events::record( - 8355
&state.core.sessions_home(), - 8356
vak_core::security_events::EventKind::ConfigChange, - 8357
"workspace_trusted", - 8358
&path.display().to_string(), - 8359
None, - 8360
); - 8361
Json(serde_json::json!({ "trusted": true, "path": path })).into_response() - 8362
} - 8363
Err(e) => ( - 8364
StatusCode::INTERNAL_SERVER_ERROR, - 8365
Json(serde_json::json!({ "error": e.to_string() })), - 8366
) - 8367
.into_response(), - 8368
} - 8369
} - 8370
- 8371
/// The starter task. Read-only by construction, and deliberately not - 8372
/// something the caller supplies: a prompt this endpoint accepted would be - 8373
/// a way to run arbitrary work under the onboarding path. - 8374
const FIRST_TASK_PROMPT: &str = "Map this codebase and explain its architecture, key flows, \ - 8375
and highest-risk areas. Do not modify files or run any destructive command."; - 8376
- 8377
/// `POST /onboarding/first-task` — create the guided starter session. - 8378
/// - 8379
/// Capped to read-only **regardless of the workspace's configured mode** - 8380
/// (doc 46 security invariant 5). The cap is not advisory and not the - 8381
/// caller's to choose: the session is created against a `Core` resolved - 8382
/// through `CorePool` with a read-only override, which `capped_by` folds - 8383
/// against the workspace ceiling as a `min` — so the result is provably - 8384
/// never more permissive than the workspace, and never less strict than - 8385
/// read-only. The handle carries that `Core`, and every run path uses the - 8386
/// handle's `Core`, so the cap holds for the actual dispatch rather than - 8387
/// only at creation. - 8388
/// - 8389
/// Returns the session id; the caller drives it through the normal run and - 8390
/// SSE endpoints, which is what makes its receipt an ordinary receipt. - 8391
async fn onboarding_first_task(State(state): State<AppState>) -> axum::response::Response { - 8392
use axum::response::IntoResponse; - 8393
let workspace = state.core.cwd().clone(); - 8394
let capped = match state.gateway.core_pool.resolve_at( - 8395
&workspace, - 8396
Some(vak_config::PermissionMode::ReadOnly), - 8397
std::time::Instant::now(), - 8398
) { - 8399
Ok(core) => core, - 8400
Err(e) => { - 8401
return ( - 8402
StatusCode::INTERNAL_SERVER_ERROR, - 8403
Json(serde_json::json!({ "error": e })), - 8404
) - 8405
.into_response(); - 8406
} - 8407
}; - 8408
- 8409
let session = match capped.start_session().await { - 8410
Ok(s) => s, - 8411
Err(e) => { - 8412
return ( - 8413
StatusCode::INTERNAL_SERVER_ERROR, - 8414
Json(serde_json::json!({ "error": e.to_string() })), - 8415
) - 8416
.into_response(); - 8417
} - 8418
}; - 8419
let id = session - 8420
.header() - 8421
.map(|h| h.session_id.clone()) - 8422
.unwrap_or_default(); - 8423
register_handle(&state, id.clone(), session, workspace, capped.clone()); - 8424
state.hub.emit_session_created(&id, ""); - 8425
index_session_later(state.store.clone(), state.core.sessions_home(), id.clone()); - 8426
- 8427
Json(serde_json::json!({ - 8428
"session_id": id, - 8429
"prompt": FIRST_TASK_PROMPT, - 8430
"permission_mode": format!("{:?}", capped.effective_permission_mode()), - 8431
})) - 8432
.into_response() - 8433
} - 8434
- 8435
async fn doctor_report( - 8436
State(state): State<AppState>, - 8437
axum::extract::Query(q): axum::extract::Query<DoctorQuery>, - 8438
) -> axum::response::Response { - 8439
use axum::response::IntoResponse; - 8440
// The optional session adds its frozen-ladder section; a live run owns - 8441
// the ledger, in which case doctor reports without that section rather - 8442
// than failing. - 8443
let session_handle = q.session.and_then(|sid| state.get(&sid)); - 8444
let session_guard = session_handle.as_deref().map(|h| { - 8445
h.session - 8446
.lock() - 8447
.unwrap_or_else(std::sync::PoisonError::into_inner) - 8448
}); - 8449
let report = - 8450
vak_core::health::collect(&state.core, session_guard.as_ref().and_then(|g| g.as_ref())); - 8451
(StatusCode::OK, Json(health_report_json(report))).into_response() - 8452
} - 8453
- 8454
#[derive(serde::Deserialize)] - 8455
struct BackupExportBody { - 8456
dest_dir: String, - 8457
#[serde(default)] - 8458
include_secrets: bool, - 8459
} - 8460
- 8461
#[derive(serde::Deserialize)] - 8462
struct BackupImportBody { - 8463
src_dir: String, - 8464
#[serde(default)] - 8465
conflict: Option<String>, - 8466
} - 8467
- 8468
/// Equality under canonicalization when both sides resolve; raw compare as - 8469
/// a fallback for paths that do not exist yet. - 8470
fn same_path(a: &std::path::Path, b: &std::path::Path) -> bool { - 8471
match (a.canonicalize(), b.canonicalize()) { - 8472
(Ok(ca), Ok(cb)) => ca == cb, - 8473
_ => a == b, - 8474
} - 8475
} - 8476
- 8477
async fn backup_export( - 8478
State(state): State<AppState>, - 8479
Json(body): Json<BackupExportBody>, - 8480
) -> axum::response::Response { - 8481
use axum::response::IntoResponse; - 8482
let home = state.core.sessions_home(); - 8483
let dest = std::path::PathBuf::from(body.dest_dir.trim()); - 8484
if dest.as_os_str().is_empty() || same_path(&dest, &home) { - 8485
return ( - 8486
StatusCode::BAD_REQUEST, - 8487
Json(serde_json::json!({ - 8488
"error": "backup destination must differ from the vak home itself" - 8489
})), - 8490
) - 8491
.into_response(); - 8492
} - 8493
match tokio::task::spawn_blocking(move || { - 8494
vak_core::backup::export_to(&home, &dest, body.include_secrets) - 8495
}) - 8496
.await - 8497
{ - 8498
Ok(Ok(manifest)) => ( - 8499
StatusCode::OK, - 8500
Json(serde_json::json!({ - 8501
"manifest": manifest, - 8502
"included_secrets": body.include_secrets, - 8503
})), - 8504
) - 8505
.into_response(), - 8506
Ok(Err(e)) => ( - 8507
StatusCode::BAD_REQUEST, - 8508
Json(serde_json::json!({ "error": e.to_string() })), - 8509
) - 8510
.into_response(), - 8511
Err(e) => ( - 8512
StatusCode::INTERNAL_SERVER_ERROR, - 8513
Json(serde_json::json!({ "error": e.to_string() })), - 8514
) - 8515
.into_response(), - 8516
} - 8517
} - 8518
- 8519
async fn backup_import( - 8520
State(state): State<AppState>, - 8521
Json(body): Json<BackupImportBody>, - 8522
) -> axum::response::Response { - 8523
use axum::response::IntoResponse; - 8524
let home = state.core.sessions_home(); - 8525
let src = std::path::PathBuf::from(body.src_dir.trim()); - 8526
if src.as_os_str().is_empty() || same_path(&src, &home) { - 8527
return ( - 8528
StatusCode::BAD_REQUEST, - 8529
Json(serde_json::json!({ - 8530
"error": "backup source must differ from the vak home itself" - 8531
})), - 8532
) - 8533
.into_response(); - 8534
} - 8535
let conflict = match body.conflict.as_deref() { - 8536
None | Some("skip") => vak_core::backup::Conflict::Skip, - 8537
Some("rename") => vak_core::backup::Conflict::Rename, - 8538
Some(other) => { - 8539
return ( - 8540
StatusCode::BAD_REQUEST, - 8541
Json(serde_json::json!({ - 8542
"error": format!("unknown conflict policy '{other}': expected \"skip\" or \"rename\"") - 8543
})), - 8544
) - 8545
.into_response(); - 8546
} - 8547
}; - 8548
match tokio::task::spawn_blocking(move || vak_core::backup::import_from(&src, &home, conflict)) - 8549
.await - 8550
{ - 8551
Ok(Ok(report)) => ( - 8552
StatusCode::OK, - 8553
Json(serde_json::json!({ - 8554
"copied": report.copied, - 8555
"renamed": report.renamed, - 8556
"skipped": report.skipped, - 8557
})), - 8558
) - 8559
.into_response(), - 8560
Ok(Err(e)) => ( - 8561
StatusCode::BAD_REQUEST, - 8562
Json(serde_json::json!({ "error": e.to_string() })), - 8563
) - 8564
.into_response(), - 8565
Err(e) => ( - 8566
StatusCode::INTERNAL_SERVER_ERROR, - 8567
Json(serde_json::json!({ "error": e.to_string() })), - 8568
) - 8569
.into_response(), - 8570
} - 8571
} - 8572
- 8573
#[derive(serde::Deserialize)] - 8574
struct DigestQuery { - 8575
#[serde(default)] - 8576
days: Option<u32>, - 8577
} - 8578
- 8579
async fn digest_report( - 8580
State(state): State<AppState>, - 8581
axum::extract::Query(q): axum::extract::Query<DigestQuery>, - 8582
) -> Json<vak_core::digest::DigestReport> { - 8583
let days = q.days.unwrap_or(7).clamp(1, 90); - 8584
Json(vak_core::digest::digest( - 8585
&state.core.sessions_home(), - 8586
&state.core.shared_data_home(), - 8587
days, - 8588
)) - 8589
} - 8590
- 8591
// ---- Intent kernel + commitments (docs/design/47-commitment-kernel.md) ----- - 8592
- 8593
#[derive(serde::Deserialize)] - 8594
struct IntentExplainQuery { - 8595
prompt: String, - 8596
#[serde(default)] - 8597
session_id: Option<String>, - 8598
#[serde(default)] - 8599
surface: Option<String>, - 8600
#[serde(default)] - 8601
act: Option<String>, - 8602
#[serde(default)] - 8603
horizon: Option<String>, - 8604
#[serde(default)] - 8605
stakes: Option<String>, - 8606
#[serde(default)] - 8607
evidence: Option<String>, - 8608
} - 8609
- 8610
/// Resolve a prompt without running it. - 8611
/// - 8612
/// The same free tiers the runtime uses, so what this returns is what that - 8613
/// prompt would actually get. Costs nothing and dispatches nothing, which is - 8614
/// what makes it safe to call from a composer as the user types. - 8615
async fn intent_explain( - 8616
State(state): State<AppState>, - 8617
axum::extract::Query(q): axum::extract::Query<IntentExplainQuery>, - 8618
) -> axum::response::Response { - 8619
// A registered session already carries the exact Core it was opened - 8620
// under (agent identity, isolated workspace) — resolving through it - 8621
// keeps this preview consistent with what that session's own turns - 8622
// would actually resolve, instead of always describing the default - 8623
// "vak" workspace regardless of which Agent's composer called this. - 8624
let core = q - 8625
.session_id - 8626
.as_deref() - 8627
.and_then(|sid| state.get(sid)) - 8628
.map(|handle| handle.core.clone()) - 8629
.unwrap_or_else(|| state.core.clone()); - 8630
let mut declared = vak_intent::Declared::default(); - 8631
let mut bad = Vec::new(); - 8632
if let Some(raw) = &q.act { - 8633
match vak_intent::Act::parse(raw) { - 8634
Some(value) => declared.act = Some(value), - 8635
None => bad.push(format!("act '{raw}'")), - 8636
} - 8637
} - 8638
if let Some(raw) = &q.horizon { - 8639
match vak_intent::Horizon::parse(raw) { - 8640
Some(value) => declared.horizon = Some(value), - 8641
None => bad.push(format!("horizon '{raw}'")), - 8642
} - 8643
} - 8644
if let Some(raw) = &q.stakes { - 8645
match vak_intent::Stakes::parse(raw) { - 8646
Some(value) => declared.stakes = Some(value), - 8647
None => bad.push(format!("stakes '{raw}'")), - 8648
} - 8649
} - 8650
if let Some(raw) = &q.evidence { - 8651
match vak_intent::Evidence::parse(raw) { - 8652
Some(value) => declared.evidence = Some(value), - 8653
None => bad.push(format!("evidence '{raw}'")), - 8654
} - 8655
} - 8656
if !bad.is_empty() { - 8657
return ( - 8658
StatusCode::BAD_REQUEST, - 8659
Json(serde_json::json!({ "error": format!("unknown {}", bad.join(", ")) })), - 8660
) - 8661
.into_response(); - 8662
} - 8663
- 8664
let surface = match q.surface.as_deref() { - 8665
None => core.surface().clone(), - 8666
Some(raw) => match vak_intent::Surface::parse(raw) { - 8667
Some(vak_intent::Surface::Cli) => vak_core::Surface::Cli, - 8668
Some(vak_intent::Surface::Desktop) => vak_core::Surface::Desktop, - 8669
Some(vak_intent::Surface::Server) => vak_core::Surface::Server, - 8670
Some(vak_intent::Surface::Chat) => vak_core::Surface::Chat { - 8671
channel: "chat".into(), - 8672
}, - 8673
Some(vak_intent::Surface::Cron | vak_intent::Surface::Heartbeat) => { - 8674
vak_core::Surface::Background - 8675
} - 8676
Some(vak_intent::Surface::Worker) => vak_core::Surface::Worker, - 8677
None => { - 8678
return ( - 8679
StatusCode::BAD_REQUEST, - 8680
Json(serde_json::json!({ "error": format!("unknown surface '{raw}'") })), - 8681
) - 8682
.into_response(); - 8683
} - 8684
}, - 8685
}; - 8686
- 8687
let history = match q.session_id.as_deref() { - 8688
Some(sid) => match core.open_session(sid).await { - 8689
Ok(session) => vak_core::intent::history_facts(&session), - 8690
Err(_) => vak_intent::HistoryFacts::default(), - 8691
}, - 8692
None => vak_intent::HistoryFacts::default(), - 8693
}; - 8694
- 8695
let resolution = vak_core::intent::resolve_turn( - 8696
&q.prompt, - 8697
"", - 8698
&surface, - 8699
&[], - 8700
vak_core::intent::workspace_facts(core.cwd()), - 8701
history.clone(), - 8702
&declared, - 8703
&core.turn_authority_for(&surface), - 8704
&vak_core::intent::resolver_config(core.config()), - 8705
); - 8706
let escalation = match &resolution { - 8707
vak_intent::Resolution::Escalate { reason, .. } => Some(reason.clone()), - 8708
vak_intent::Resolution::Settled(_) => None, - 8709
}; - 8710
let intent = resolution.intent(); - 8711
Json(serde_json::json!({ - 8712
"reading": intent.reading, - 8713
"strands": intent.strands, - 8714
"engagement": intent.engagement, - 8715
"provenance": intent.provenance, - 8716
"narrows": intent - 8717
.engagement - 8718
.limits - 8719
.diff_from(&vak_intent::Limits::unrestricted()), - 8720
"escalation_recommended": escalation, - 8721
"model_visible": intent.model_visible(), - 8722
"history": { - 8723
"turn_index": history.turn_index, - 8724
"previous_act": history.previous_act.map(|a| a.as_str()), - 8725
"open_threads": history.open_threads.len(), - 8726
}, - 8727
})) - 8728
.into_response() - 8729
} - 8730
- 8731
/// The resolved intent and commitment policy for this workspace. - 8732
async fn intent_policy(State(state): State<AppState>) -> Json<serde_json::Value> { - 8733
let config = state.core.config(); - 8734
Json(serde_json::json!({ - 8735
"intent": { - 8736
"enabled": config.intent.enabled, - 8737
"accept_confidence": config.intent.accept_confidence, - 8738
"provisional_confidence": config.intent.provisional_confidence, - 8739
"slice_capabilities": config.intent.slice_capabilities, - 8740
"posture": config.intent.posture, - 8741
"escalate": config.intent.escalate, - 8742
"max_classify_usd": config.intent.max_classify_usd, - 8743
"classify_timeout_secs": config.intent.classify_timeout_secs, - 8744
"autonomy": config.intent.autonomy, - 8745
"evidence_max_age_secs": config.intent.evidence_max_age_secs, - 8746
}, - 8747
"commitment": { - 8748
"enabled": config.commitment.enabled, - 8749
"lifetime_budget_usd": config.commitment.lifetime_budget_usd, - 8750
"stall_limit": config.commitment.stall_limit, - 8751
"review_every_hours": config.commitment.review_every_hours, - 8752
"default_ttl_days": config.commitment.default_ttl_days, - 8753
}, - 8754
})) - 8755
} - 8756
- 8757
#[derive(serde::Deserialize)] - 8758
struct CommitmentQuery { - 8759
/// Include closed commitments. - 8760
#[serde(default)] - 8761
all: bool, - 8762
#[serde(default)] - 8763
agent: Option<String>, - 8764
} - 8765
- 8766
/// The portfolio, in the order the scheduler would work it. - 8767
async fn list_commitments( - 8768
State(state): State<AppState>, - 8769
axum::extract::Query(q): axum::extract::Query<CommitmentQuery>, - 8770
) -> axum::response::Response { - 8771
use axum::response::IntoResponse; - 8772
let core = scoped_core!(&state, None, q.agent.as_deref()); - 8773
let ledger = vak_commit::CommitmentLedger::new(&core.sessions_home()); - 8774
let commitments = if q.all { ledger.all() } else { ledger.open() }; - 8775
let ranked = vak_commit::rank(&commitments, &vak_commit::SchedulerContext::default()); - 8776
Json(serde_json::json!({ - 8777
"commitments": commitments, - 8778
// Priorities ride alongside rather than being baked into the rows: - 8779
// the ordering is a scheduling opinion, and a UI should be able to - 8780
// show why as well as what. - 8781
"priorities": ranked, - 8782
})) - 8783
.into_response() - 8784
} - 8785
- 8786
async fn get_commitment( - 8787
State(state): State<AppState>, - 8788
axum::extract::Path(id): axum::extract::Path<String>, - 8789
axum::extract::Query(q): axum::extract::Query<AgentScopeQuery>, - 8790
) -> axum::response::Response { - 8791
let core = scoped_core!(&state, None, q.agent.as_deref()); - 8792
let ledger = vak_commit::CommitmentLedger::new(&core.sessions_home()); - 8793
match ledger.get(&id) { - 8794
Ok(Some(commitment)) => Json(serde_json::json!({ - 8795
"commitment": commitment, - 8796
"events": ledger.events_for(&id), - 8797
})) - 8798
.into_response(), - 8799
Ok(None) => ( - 8800
StatusCode::NOT_FOUND, - 8801
Json(serde_json::json!({ "error": "no such commitment" })), - 8802
) - 8803
.into_response(), - 8804
Err(error) => ( - 8805
StatusCode::INTERNAL_SERVER_ERROR, - 8806
Json(serde_json::json!({ "error": error.to_string() })), - 8807
) - 8808
.into_response(), - 8809
} - 8810
} - 8811
- 8812
#[derive(serde::Deserialize)] - 8813
struct CloseCommitmentBody { - 8814
verdict: String, - 8815
#[serde(default)] - 8816
note: String, - 8817
#[serde(default)] - 8818
agent: Option<String>, - 8819
} - 8820
- 8821
async fn close_commitment( - 8822
State(state): State<AppState>, - 8823
axum::extract::Path(id): axum::extract::Path<String>, - 8824
Json(body): Json<CloseCommitmentBody>, - 8825
) -> axum::response::Response { - 8826
let core = scoped_core!(&state, None, body.agent.as_deref()); - 8827
let ledger = vak_commit::CommitmentLedger::new(&core.sessions_home()); - 8828
let Ok(Some(commitment)) = ledger.get(&id) else { - 8829
return ( - 8830
StatusCode::NOT_FOUND, - 8831
Json(serde_json::json!({ "error": "no such commitment" })), - 8832
) - 8833
.into_response(); - 8834
}; - 8835
let verdict = match body.verdict.as_str() { - 8836
"fulfilled" => vak_commit::Verdict::Fulfilled, - 8837
"partial" => vak_commit::Verdict::Partial, - 8838
"failed" => vak_commit::Verdict::Failed, - 8839
"abandoned" => vak_commit::Verdict::Abandoned, - 8840
"expired" => vak_commit::Verdict::Expired, - 8841
"unknown" => vak_commit::Verdict::Unknown, - 8842
other => { - 8843
return ( - 8844
StatusCode::BAD_REQUEST, - 8845
Json(serde_json::json!({ "error": format!("unknown verdict '{other}'") })), - 8846
) - 8847
.into_response(); - 8848
} - 8849
}; - 8850
let strength = commitment.achieved_strength(); - 8851
match ledger.append(&vak_commit::Event::new( - 8852
&commitment.commitment_id, - 8853
vak_commit::EventKind::Closed { - 8854
verdict, - 8855
strength, - 8856
evidence: Vec::new(), - 8857
note: body.note, - 8858
}, - 8859
)) { - 8860
Ok(()) => { - 8861
Json(serde_json::json!({ "ok": true, "verdict": verdict.as_str() })).into_response() - 8862
} - 8863
// A refused closure is a 409, not a 500: the request was well-formed - 8864
// and the server is fine — the evidence simply does not support the - 8865
// claim. The message says which evidence was missing. - 8866
Err(error) => ( - 8867
StatusCode::CONFLICT, - 8868
Json(serde_json::json!({ "error": error.to_string() })), - 8869
) - 8870
.into_response(), - 8871
} - 8872
} - 8873
- 8874
// ---- Inbox (durable attention layer, docs/design/29-personal-os.md P6) ------ - 8875
- 8876
const DEFAULT_INBOX_LIMIT: usize = 200; - 8877
- 8878
#[derive(serde::Deserialize)] - 8879
struct InboxQuery { - 8880
#[serde(default)] - 8881
limit: Option<usize>, - 8882
/// Only entries without an ack tombstone. - 8883
#[serde(default)] - 8884
unread: bool, - 8885
} - 8886
- 8887
/// Newest-first inbox entries plus the live unread total. The count always - 8888
/// reflects the full unfiltered set; `limit` bounds the returned window only. - 8889
async fn inbox_list( - 8890
State(state): State<AppState>, - 8891
axum::extract::Query(q): axum::extract::Query<InboxQuery>, - 8892
) -> Json<serde_json::Value> { - 8893
let home = state.core.shared_data_home(); - 8894
let unread_count = vak_core::inbox::unread_count(&home); - 8895
let limit = q - 8896
.limit - 8897
.unwrap_or(DEFAULT_INBOX_LIMIT) - 8898
.clamp(1, vak_core::inbox::MAX_SCAN); - 8899
let entries = if q.unread { - 8900
vak_core::inbox::unread(&home) - 8901
} else { - 8902
vak_core::inbox::list(&home, limit) - 8903
} - 8904
.into_iter() - 8905
.take(limit) - 8906
.collect::<Vec<_>>(); - 8907
let entries = entries - 8908
.into_iter() - 8909
.map(|entry| { - 8910
let mut value = serde_json::to_value(&entry).unwrap_or_else(|_| serde_json::json!({})); - 8911
if let Some(session_id) = entry.session_id.as_deref() { - 8912
let available = state.get(session_id).is_some() - 8913
|| open_historical_session(&state, session_id).is_some(); - 8914
value["origin_state"] = serde_json::json!(if available { - 8915
"available" - 8916
} else { - 8917
"unavailable" - 8918
}); - 8919
} - 8920
value - 8921
}) - 8922
.collect::<Vec<_>>(); - 8923
Json(serde_json::json!({ "entries": entries, "unread_count": unread_count })) - 8924
} - 8925
- 8926
async fn inbox_unread_count(State(state): State<AppState>) -> Json<serde_json::Value> { - 8927
Json(serde_json::json!({ - 8928
"count": vak_core::inbox::unread_count(&state.core.shared_data_home()) - 8929
})) - 8930
} - 8931
- 8932
/// Idempotent read-state: a tombstone append via `inbox::ack`. An unknown id - 8933
/// is a 404; re-acking reports `{acked:false}` instead of writing twice. - 8934
async fn inbox_ack( - 8935
State(state): State<AppState>, - 8936
Path(id): Path<String>, - 8937
) -> axum::response::Response { - 8938
use axum::response::IntoResponse; - 8939
let home = state.core.shared_data_home(); - 8940
if !vak_core::inbox::list(&home, vak_core::inbox::MAX_SCAN) - 8941
.iter() - 8942
.any(|e| e.id == id) - 8943
{ - 8944
return ( - 8945
StatusCode::NOT_FOUND, - 8946
Json(serde_json::json!({ "error": format!("unknown inbox entry '{id}'") })), - 8947
) - 8948
.into_response(); - 8949
} - 8950
match vak_core::inbox::ack(&home, &id) { - 8951
Ok(acked) => Json(serde_json::json!({ "acked": acked })).into_response(), - 8952
Err(e) => ( - 8953
StatusCode::INTERNAL_SERVER_ERROR, - 8954
Json(serde_json::json!({ "error": e.to_string() })), - 8955
) - 8956
.into_response(), - 8957
} - 8958
} - 8959
- 8960
async fn git_output(cwd: &std::path::Path, args: &[&str]) -> Option<String> { - 8961
let out = tokio::process::Command::new("git") - 8962
.args(args) - 8963
.current_dir(cwd) - 8964
.output() - 8965
.await - 8966
.ok()?; - 8967
if !out.status.success() { - 8968
return None; - 8969
} - 8970
Some(String::from_utf8_lossy(&out.stdout).into_owned()) - 8971
} - 8972
- 8973
/// Workspace diff for the review pane. Untracked files appear in `status` - 8974
/// as `??` lines; patches are split per-file client-side. - 8975
async fn session_diff( - 8976
State(state): State<AppState>, - 8977
Path(id): Path<String>, - 8978
) -> axum::response::Response { - 8979
use axum::response::IntoResponse; - 8980
let Some(handle) = state.get(&id) else { - 8981
return ( - 8982
StatusCode::NOT_FOUND, - 8983
Json(serde_json::json!({ "error": "unknown session" })), - 8984
) - 8985
.into_response(); - 8986
}; - 8987
let cwd = handle.cwd.clone(); - 8988
let Some(status) = git_output(&cwd, &["status", "--porcelain"]).await else { - 8989
return ( - 8990
StatusCode::UNPROCESSABLE_ENTITY, - 8991
Json(serde_json::json!({ "error": "not a git repository" })), - 8992
) - 8993
.into_response(); - 8994
}; - 8995
// `--no-color` is a `diff` option, not a global git flag — placed - 8996
// before the subcommand (as this read for a long time) git rejects it - 8997
// outright ("unknown option: --no-color", exit 129), and `git_output` - 8998
// turns that failure into a silently empty string via - 8999
// `unwrap_or_default()`. Every consumer of this endpoint — this - 9000
// console's Worktree Diff tab, the desktop DiffPane, `openFileSmart`'s - 9001
// diff-vs-editor routing — has been reading an empty diff regardless - 9002
// of what actually changed. - 9003
let diff = git_output(&cwd, &["diff", "--no-color", "--unified=3"]) - 9004
.await - 9005
.unwrap_or_default(); - 9006
let staged = git_output(&cwd, &["diff", "--no-color", "--cached", "--unified=3"]) - 9007
.await - 9008
.unwrap_or_default(); - 9009
Json(serde_json::json!({ - 9010
"root": cwd, - 9011
"diff": diff, - 9012
"staged_diff": staged, - 9013
"status": status, - 9014
})) - 9015
.into_response() - 9016
} - 9017
- 9018
// ---- checkpoints (time travel) ---------------------------------------------- - 9019
- 9020
async fn list_checkpoints( - 9021
State(state): State<AppState>, - 9022
Path(id): Path<String>, - 9023
axum::extract::Query(q): axum::extract::Query<AgentScopeQuery>, - 9024
) -> axum::response::Response { - 9025
// A registered session's own Core is reused when it's still open. Once - 9026
// it's closed, `session_id` alone can no longer identify its owning - 9027
// Agent, so an explicit `?agent=` is the only way to keep finding its - 9028
// checkpoints instead of silently falling back to the default Agent. - 9029
let core = scoped_core!(&state, Some(&id), q.agent.as_deref()); - 9030
// A session with no snapshots yet has no directory; that's an empty - 9031
// list, not an error. - 9032
let list = match vak_core::checkpoints::list(&core.sessions_home(), &id) { - 9033
Ok(list) if !list.is_empty() => list, - 9034
_ => match vak_core::checkpoints::list(&core.shared_data_home(), &id) { - 9035
Ok(list) => list, - 9036
Err(e) if e.kind() == std::io::ErrorKind::NotFound => Vec::new(), - 9037
Err(e) => { - 9038
return ( - 9039
StatusCode::INTERNAL_SERVER_ERROR, - 9040
Json(serde_json::json!({ "error": e.to_string() })), - 9041
) - 9042
.into_response(); - 9043
} - 9044
}, - 9045
}; - 9046
let checkpoints: Vec<serde_json::Value> = list - 9047
.iter() - 9048
.map(|cp| { - 9049
serde_json::json!({ - 9050
"seq": cp.seq, - 9051
"label": cp.label, - 9052
"created_at": cp.created_at.to_rfc3339(), - 9053
"files": cp.files.len(), - 9054
}) - 9055
}) - 9056
.collect(); - 9057
Json(serde_json::json!({ "checkpoints": checkpoints })).into_response() - 9058
} - 9059
- 9060
async fn restore_checkpoint( - 9061
State(state): State<AppState>, - 9062
Path((id, seq)): Path<(String, u32)>, - 9063
axum::extract::Query(q): axum::extract::Query<AgentScopeQuery>, - 9064
) -> axum::response::Response { - 9065
// A live run must never have its workspace mutated underneath it. Reuse - 9066
// this lookup for the cwd fallback below instead of re-fetching it. - 9067
let handle = state.get(&id); - 9068
if let Some(handle) = &handle - 9069
&& handle - 9070
.session - 9071
.lock() - 9072
.unwrap_or_else(std::sync::PoisonError::into_inner) - 9073
.is_none() - 9074
{ - 9075
return ( - 9076
StatusCode::CONFLICT, - 9077
Json(serde_json::json!({ "error": "a run is active on this session" })), - 9078
) - 9079
.into_response(); - 9080
} - 9081
// See `list_checkpoints`: a closed session needs an explicit `?agent=` - 9082
// to keep resolving its real owning Agent rather than the default one. - 9083
let core = scoped_core!(&state, Some(&id), q.agent.as_deref()); - 9084
// Best-of-N children captured inside their worktrees; attached handles - 9085
// know that cwd. Everything else restores into the workspace root. - 9086
let cwd = handle - 9087
.map(|h| h.cwd.clone()) - 9088
.unwrap_or_else(|| core.cwd().clone()); - 9089
// The blob store the manifest's hashes resolve against lives under - 9090
// whichever home the manifest itself was found in. - 9091
let (cp, checkpoints_home) = match vak_core::checkpoints::load(&core.sessions_home(), &id, seq) - 9092
{ - 9093
Ok(cp) => (cp, core.sessions_home()), - 9094
Err(_) => match vak_core::checkpoints::load(&core.shared_data_home(), &id, seq) { - 9095
Ok(cp) => (cp, core.shared_data_home()), - 9096
Err(_) => { - 9097
return ( - 9098
StatusCode::NOT_FOUND, - 9099
Json(serde_json::json!({ "error": format!("checkpoint {seq} not found") })), - 9100
) - 9101
.into_response(); - 9102
} - 9103
}, - 9104
}; - 9105
match tokio::task::spawn_blocking(move || { - 9106
vak_core::checkpoints::restore(&cwd, &checkpoints_home, &cp) - 9107
}) - 9108
.await - 9109
{ - 9110
Ok(Ok((restored, deleted))) => Json(serde_json::json!({ - 9111
"restored": restored, - 9112
"deleted": deleted, - 9113
"seq": seq, - 9114
})) - 9115
.into_response(), - 9116
Ok(Err(e)) => ( - 9117
StatusCode::INTERNAL_SERVER_ERROR, - 9118
Json(serde_json::json!({ "error": e.to_string() })), - 9119
) - 9120
.into_response(), - 9121
Err(e) => ( - 9122
StatusCode::INTERNAL_SERVER_ERROR, - 9123
Json(serde_json::json!({ "error": e.to_string() })), - 9124
) - 9125
.into_response(), - 9126
} - 9127
} - 9128
- 9129
// ---- archive (sidebar visibility; ledgers stay untouched) -------------------- - 9130
- 9131
fn archive_path(core: &Core) -> PathBuf { - 9132
core.shared_data_home().join("archive.json") - 9133
} - 9134
- 9135
fn read_archive(core: &Core) -> HashMap<String, bool> { - 9136
std::fs::read_to_string(archive_path(core)) - 9137
.ok() - 9138
.and_then(|text| serde_json::from_str(&text).ok()) - 9139
.unwrap_or_default() - 9140
} - 9141
- 9142
fn write_archive(core: &Core, map: &HashMap<String, bool>) { - 9143
if let Some(parent) = archive_path(core).parent() { - 9144
let _ = std::fs::create_dir_all(parent); - 9145
} - 9146
let tmp = archive_path(core).with_extension("json.tmp"); - 9147
if std::fs::write(&tmp, serde_json::to_string(map).unwrap_or_default()).is_ok() { - 9148
let _ = std::fs::rename(&tmp, archive_path(core)); - 9149
} - 9150
} - 9151
- 9152
fn find_session_in_cwd(core: &Core, id: &str) -> bool { - 9153
let direct = vak_session::SessionPath::sessions_dir(&core.sessions_home(), core.cwd()) - 9154
.join(format!("{id}.jsonl")); - 9155
if direct.is_file() { - 9156
return true; - 9157
} - 9158
let shared = core.shared_data_home(); - 9159
if let Ok(agents) = std::fs::read_dir(shared.join("agents")) { - 9160
for agent in agents.flatten() { - 9161
let candidate = vak_session::SessionPath::sessions_dir(&agent.path(), core.cwd()) - 9162
.join(format!("{id}.jsonl")); - 9163
if candidate.is_file() {
Indexing the workspace…
Vakyartha documentation is discovering safe artifacts, anchors, and source references.