- 3164
effective - 3165
.select_preferred("timeline", "user", owner) - 3166
.is_none() - 3167
); - 3168
assert!(effective.select_preferred("table", "user", owner).is_some()); - 3169
let dir = tempfile::tempdir().expect("tempdir"); - 3170
let store = - 3171
vak_store::presentation::PresentationStore::new(dir.path().join("presentations.json")); - 3172
store.save(&library).expect("save suppression"); - 3173
library = store.load().expect("reload suppression"); - 3174
assert!( - 3175
super::effective_presentation_library(&library, owner) - 3176
.select_preferred("timeline", "user", owner) - 3177
.is_none() - 3178
); - 3179
let seed = vak_presentation::seeds::built_in_seed_pack() - 3180
.into_iter() - 3181
.find(|record| record.spec.id == "seed.timeline") - 3182
.expect("timeline seed"); - 3183
library.register(seed).expect("register seed"); - 3184
library - 3185
.activate( - 3186
"seed.timeline", - 3187
6, - 3188
vak_presentation::LibraryScope::Workspace, - 3189
owner, - 3190
) - 3191
.expect("activate seed"); - 3192
let effective = super::effective_presentation_library(&library, owner); - 3193
assert_eq!( - 3194
effective - 3195
.select_preferred("timeline", "user", owner) - 3196
.map(|record| record.spec.id.as_str()), - 3197
Some("seed.timeline") - 3198
); - 3199
} - 3200
} - 3201
- 3202
/// Paths that must be reachable without a token: health probe, the SPA - 3203
/// shell (static assets carry no data), and the login endpoint itself. - 3204
fn auth_exempt_path(path: &str) -> bool { - 3205
path == "/health" - 3206
|| path == "/admin" - 3207
|| path == "/admin/" - 3208
|| path == "/admin/favicon.svg" - 3209
|| path == "/admin/vak-icon.png" - 3210
|| path.starts_with("/admin/assets/") - 3211
// The workspace client's shell and its hashed assets carry no data - 3212
// and must load before a session exists — the login form is part of - 3213
// the bundle. Every route it then calls is authenticated. - 3214
|| path == "/app" - 3215
|| path == "/app/" - 3216
|| path == "/app/vak-icon.png" - 3217
|| path == "/app/manifest.webmanifest" - 3218
|| path == "/app/sw.js" - 3219
|| path.starts_with("/app/assets/") - 3220
|| path.starts_with("/app/characters/") - 3221
// The login exchange itself, and the probe that decides whether to - 3222
// show it. `/auth/session` answers `{authenticated:false}` rather - 3223
// than 401 so an unauthenticated client can tell "no session" from - 3224
// "server unreachable". - 3225
|| path == "/auth/login" - 3226
|| path == "/auth/session" - 3227
// The public site and its build stamp. Every page answers an - 3228
// unauthenticated stranger by design — a blank 401 at `/` told a - 3229
// visitor nothing at all, not even that anything was listening — - 3230
// and `site.rs` owns what those pages may say. - 3231
|| site::ROUTES.iter().any(|(uri, _)| { - 3232
*uri == path || (*uri != "/" && path.len() == uri.len() + 1 && path.starts_with(uri) && path.ends_with('/')) - 3233
}) - 3234
|| path.starts_with("/site/") - 3235
|| path == "/version" - 3236
|| path == "/favicon.ico" - 3237
|| path == "/favicon.svg" - 3238
} - 3239
- 3240
#[cfg(test)] - 3241
mod auth_exempt_path_tests { - 3242
use super::{auth_exempt_path, host_is_loopback}; - 3243
- 3244
/// A page that resolves its own domain to 127.0.0.1 becomes same-origin - 3245
/// with this server, and `SameSite=Strict` does not help — after - 3246
/// rebinding the request is not cross-site. The Host name it carries is - 3247
/// still the attacker's, which is what this rejects. - 3248
#[test] - 3249
fn only_loopback_hostnames_are_answered() { - 3250
for host in [ - 3251
"localhost", - 3252
"localhost:8901", - 3253
"127.0.0.1", - 3254
"127.0.0.1:8901", - 3255
"[::1]:8901", - 3256
"::1", - 3257
] { - 3258
assert!(host_is_loopback(Some(host)), "{host} is loopback"); - 3259
} - 3260
for host in [ - 3261
"rebind.attacker.example", - 3262
"rebind.attacker.example:8901", - 3263
"vak.internal", - 3264
"127.0.0.1.attacker.example", - 3265
"evil.com:80", - 3266
] { - 3267
assert!(!host_is_loopback(Some(host)), "{host} must be rejected"); - 3268
} - 3269
} - 3270
- 3271
/// Rebinding is a browser attack and browsers always send Host, so an - 3272
/// absent value is a non-browser client rather than something to defend - 3273
/// against. - 3274
#[test] - 3275
fn an_absent_host_is_not_treated_as_an_attack() { - 3276
assert!(host_is_loopback(None)); - 3277
} - 3278
- 3279
/// The login screen's own logo must be reachable before a cookie - 3280
/// exists to authenticate the request that would fetch it — the same - 3281
/// reasoning that already exempts favicon.svg. This regressed once - 3282
/// already: the admin console shipped a login-screen `<img>` pointing - 3283
/// at a root-level dist file, and only /admin/assets/* (the hashed - 3284
/// JS/CSS bundle) was exempt, so the logo 401'd on every fresh login. - 3285
#[test] - 3286
fn login_screen_logo_is_exempt() { - 3287
assert!(auth_exempt_path("/admin/vak-icon.png")); - 3288
} - 3289
- 3290
#[test] - 3291
fn admin_api_routes_still_require_auth() { - 3292
assert!(!auth_exempt_path("/admin/api/config")); - 3293
assert!(!auth_exempt_path("/admin/api/gateway/status")); - 3294
} - 3295
} - 3296
- 3297
/// Whether a `Host` header names this server legitimately. - 3298
/// - 3299
/// Loopback names always pass. Anything else must be listed verbatim in - 3300
/// `[server] trusted_hosts`, which is a privileged setting an untrusted - 3301
/// project cannot write (docs/design/48-web-client.md §4.2). - 3302
/// - 3303
/// This is the DNS-rebinding defence: binding 127.0.0.1 does not stop a - 3304
/// page the user visits from resolving its *own* domain to 127.0.0.1 and - 3305
/// becoming same-origin with this server, and `SameSite=Strict` does not - 3306
/// help once that has happened, because the request is then not cross-site. - 3307
/// Pinning the `Host` name is the check that does — and the reason - 3308
/// `trusted_hosts` takes exact names rather than patterns: a wildcard here - 3309
/// re-opens exactly the hole the list closes. - 3310
fn host_is_trusted(host: Option<&str>, trusted: &[String]) -> bool { - 3311
if host_is_loopback(host) { - 3312
return true; - 3313
} - 3314
let Some(host) = host else { return true }; - 3315
let name = strip_port(host).to_ascii_lowercase(); - 3316
trusted.iter().any(|allowed| allowed == &name) - 3317
} - 3318
- 3319
/// Host name without its port, leaving a bare IPv6 literal intact. - 3320
fn strip_port(host: &str) -> &str { - 3321
if let Some(rest) = host.strip_prefix('[') { - 3322
rest.split_once(']').map_or(rest, |(name, _)| name) - 3323
} else if host.matches(':').count() > 1 { - 3324
host - 3325
} else { - 3326
host.split_once(':').map_or(host, |(name, _)| name) - 3327
} - 3328
} - 3329
- 3330
/// Whether `Host` names this machine's own loopback interface. - 3331
fn host_is_loopback(host: Option<&str>) -> bool { - 3332
let Some(host) = host else { - 3333
// Absent entirely: not a browser. Rebinding is a browser attack and - 3334
// every browser sends Host (or HTTP/2 `:authority`, which axum - 3335
// surfaces as the URI authority), so an absent value carries no - 3336
// attacker-chosen name to defend against. Rejecting here would only - 3337
// break well-behaved non-browser clients that omit it. - 3338
return true; - 3339
}; - 3340
// Strip the port without mangling a bare IPv6 literal, which is all - 3341
// colons: `"::1".rsplit_once(':')` yields `"::"`. - 3342
matches!( - 3343
strip_port(host), - 3344
"localhost" | "127.0.0.1" | "::1" | "0:0:0:0:0:0:0:1" - 3345
) - 3346
} - 3347
- 3348
/// Whether a state-changing request's `Origin` is one of ours. - 3349
/// - 3350
/// Cookies alone are not enough to authorize a mutation: a cookie is - 3351
/// attached by the browser to whoever asks, and `SameSite=Strict` covers - 3352
/// the common cases but not a same-site subdomain or a rebound name. So a - 3353
/// mutation carrying a cookie must ALSO carry an `Origin` we recognise. - 3354
/// - 3355
/// A request with no `Origin` at all is not a browser form post — browsers - 3356
/// always send one on cross-origin mutations — so it is allowed through - 3357
/// here and still has to satisfy `require_bearer` with a real header - 3358
/// token. That is what keeps curl, the CLI, and the bridges working - 3359
/// without giving a page any new power. - 3360
fn origin_is_trusted(origin: Option<&str>, trusted: &[String]) -> bool { - 3361
let Some(origin) = origin else { return true }; - 3362
// The Tauri webview's own origins: the desktop is a first-party client - 3363
// and its scheme is not something an attacker can mint. - 3364
if matches!( - 3365
origin, - 3366
"tauri://localhost" | "http://tauri.localhost" | "https://tauri.localhost" - 3367
) { - 3368
return true; - 3369
} - 3370
let authority = origin - 3371
.split_once("://") - 3372
.map(|(_, rest)| rest) - 3373
.unwrap_or(origin); - 3374
host_is_trusted(Some(authority), trusted) - 3375
} - 3376
- 3377
/// What the auth layer needs to know about this deployment's exposure. - 3378
#[derive(Clone)] - 3379
pub(crate) struct AuthPolicy { - 3380
pub(crate) token: String, - 3381
pub(crate) home: std::path::PathBuf, - 3382
/// Non-loopback `Host` names this server answers to (`[server] - 3383
/// trusted_hosts`). Empty on a default install. - 3384
pub(crate) trusted_hosts: Vec<String>, - 3385
} - 3386
- 3387
/// Identity established by the HTTP boundary. Participant grants are kept - 3388
/// distinct from the operator credential so downstream handlers can add - 3389
/// attribution and narrower decisions without ever treating a collaborator - 3390
/// as the workspace owner. - 3391
#[derive(Clone, Debug)] - 3392
pub(crate) enum AuthenticatedPrincipal { - 3393
Operator, - 3394
Participant(coworking::VerifiedPrincipal), - 3395
} - 3396
- 3397
fn participant_read_route_allowed( - 3398
method: &axum::http::Method, - 3399
path: &str, - 3400
principal: &coworking::VerifiedPrincipal, - 3401
) -> bool { - 3402
let segments: Vec<&str> = path.trim_matches('/').split('/').collect(); - 3403
let ["sessions", conversation_id, rest @ ..] = segments.as_slice() else { - 3404
return false; - 3405
}; - 3406
if *conversation_id != principal.conversation_id { - 3407
return false; - 3408
} - 3409
if method == axum::http::Method::POST { - 3410
let can_read = principal - 3411
.capabilities - 3412
.iter() - 3413
.any(|capability| capability == "read"); - 3414
if !can_read { - 3415
return false; - 3416
} - 3417
if matches!(*rest, ["coworking", "messages"]) { - 3418
return principal - 3419
.capabilities - 3420
.iter() - 3421
.any(|capability| capability == "message"); - 3422
} - 3423
if matches!(*rest, ["coworking", "approvals", _]) { - 3424
return true; - 3425
} - 3426
if matches!(*rest, ["office-workspaces", _, "presence"]) { - 3427
return true; - 3428
} - 3429
if matches!(*rest, ["office-workspaces", _]) { - 3430
return principal - 3431
.capabilities - 3432
.iter() - 3433
.any(|capability| capability == "edit"); - 3434
} - 3435
return principal - 3436
.capabilities - 3437
.iter() - 3438
.any(|capability| capability == "comment") - 3439
&& matches!(*rest, ["sandbox", "candidates", _, "comments"]); - 3440
} - 3441
if method != axum::http::Method::GET - 3442
|| !principal - 3443
.capabilities - 3444
.iter() - 3445
.any(|capability| capability == "read") - 3446
{ - 3447
return false; - 3448
} - 3449
matches!( - 3450
*rest, - 3451
["transcript"] - 3452
| ["transcript.md"] - 3453
| ["presentation"] - 3454
| ["results", _] - 3455
| ["sandbox", "records"] - 3456
| ["sandbox", "candidates", _, "files"] - 3457
| ["sandbox", "candidates", _, "files", "raw"] - 3458
| ["sandbox", "candidates", _, "office-review"] - 3459
| ["sandbox", "candidates", _, "office"] - 3460
| ["sandbox", "candidates", _, "comments"] - 3461
| ["coworking", "me"] - 3462
| ["coworking", "presence"] - 3463
| ["coworking", "approvals"] - 3464
| ["coworking", "updates"] - 3465
| ["office-workspaces"] - 3466
) - 3467
} - 3468
- 3469
pub(crate) async fn require_bearer( - 3470
State(policy): State<AuthPolicy>, - 3471
mut req: axum::extract::Request, - 3472
next: axum::middleware::Next, - 3473
) -> axum::response::Response { - 3474
let AuthPolicy { - 3475
token, - 3476
home, - 3477
trusted_hosts, - 3478
} = policy; - 3479
let host = req - 3480
.headers() - 3481
.get(axum::http::header::HOST) - 3482
.and_then(|value| value.to_str().ok()) - 3483
.or_else(|| req.uri().host()); - 3484
let loopback = host_is_loopback(host); - 3485
if !host_is_trusted(host, &trusted_hosts) { - 3486
return ( - 3487
StatusCode::MISDIRECTED_REQUEST, - 3488
Json(serde_json::json!({ - 3489
"error": "this server does not answer to that hostname; \ - 3490
add it to [server] trusted_hosts to allow it", - 3491
})), - 3492
) - 3493
.into_response(); - 3494
} - 3495
// Cross-origin mutations are refused before routing, whatever - 3496
// credential they carry. See `origin_is_trusted` for why a *missing* - 3497
// Origin is not treated as a failure. - 3498
let mutating = !matches!( - 3499
*req.method(), - 3500
axum::http::Method::GET | axum::http::Method::HEAD | axum::http::Method::OPTIONS - 3501
); - 3502
let origin = req - 3503
.headers() - 3504
.get(axum::http::header::ORIGIN) - 3505
.and_then(|v| v.to_str().ok()); - 3506
if mutating && !origin_is_trusted(origin, &trusted_hosts) { - 3507
vak_core::security_events::record( - 3508
&home, - 3509
vak_core::security_events::EventKind::AuthFailure, - 3510
"cross_origin_rejected", - 3511
&format!( - 3512
"origin={} path={}", - 3513
origin.unwrap_or("<none>"), - 3514
req.uri().path() - 3515
), - 3516
None, - 3517
); - 3518
return ( - 3519
StatusCode::FORBIDDEN, - 3520
Json(serde_json::json!({ "error": "cross-origin request refused" })), - 3521
) - 3522
.into_response(); - 3523
} - 3524
if auth_exempt_path(req.uri().path()) { - 3525
return next.run(req).await; - 3526
} - 3527
use subtle::ConstantTimeEq; - 3528
let header_token = req - 3529
.headers() - 3530
.get(axum::http::header::AUTHORIZATION) - 3531
.and_then(|v| v.to_str().ok()) - 3532
.and_then(|v| v.strip_prefix("Bearer ")) - 3533
.map(String::from); - 3534
// Browser surfaces authenticate once via /auth/login which sets an - 3535
// HttpOnly cookie; EventSource cannot send Authorization headers, so - 3536
// the cookie is the only workable channel for SSE. - 3537
let cookie_token = req - 3538
.headers() - 3539
.get(axum::http::header::COOKIE) - 3540
.and_then(|v| v.to_str().ok()) - 3541
.and_then(|cookies| { - 3542
cookies.split(';').find_map(|pair| { - 3543
let pair = pair.trim(); - 3544
pair.strip_prefix("vak_session=") - 3545
.map(|v| v.trim().to_string()) - 3546
}) - 3547
}); - 3548
// `EventSource` cannot set request headers, and the desktop app never - 3549
// performs the `/auth/login` cookie exchange -- that is the browser - 3550
// surfaces' flow, not the desktop's. The query parameter is - 3551
// therefore the ONLY channel the desktop's SSE streams can - 3552
// authenticate on, and `openEventStream`/`openSideStream` have always - 3553
// used it. It was never accepted here, so every desktop event stream - 3554
// was rejected 401: the agent completed turns and durably logged them - 3555
// while the UI received not one event -- no reply, "Working" forever, - 3556
// usage stuck at 0 in / 0 out, and nothing in the console, because a - 3557
// 401 on an EventSource surfaces only as a bare `onerror`. - 3558
// - 3559
// The startup banner has advertised `?token=` since before this - 3560
// middleware existed; this makes the implementation match the - 3561
// contract rather than narrowing the contract to the implementation. - 3562
// A token in a query string is a real (if bounded) exposure -- it can - 3563
// reach access logs and `Referer` headers -- but this server is - 3564
// loopback-only with a token that is either ephemeral per boot or - 3565
// pinned into the credential store, and no other channel exists for - 3566
// the one client that needs it. - 3567
// - 3568
// Loopback ONLY. A token in a query string can reach access logs, - 3569
// `Referer` headers, and browser history; on a loopback server with an - 3570
// in-process client and no proxy between them, none of those exist. - 3571
// On any deployment reachable by a real hostname they all do, and the - 3572
// web client does not need this channel anyway — it is same-origin, so - 3573
// its cookie covers `EventSource` (docs/design/48-web-client.md §4.3). - 3574
let query_token = loopback - 3575
.then(|| { - 3576
req.uri().query().and_then(|q| { - 3577
q.split('&').find_map(|pair| { - 3578
let (key, value) = pair.split_once('=')?; - 3579
if key != "token" { - 3580
return None; - 3581
} - 3582
Some( - 3583
percent_encoding::percent_decode_str(value) - 3584
.decode_utf8_lossy() - 3585
.into_owned(), - 3586
) - 3587
}) - 3588
}) - 3589
}) - 3590
.flatten(); - 3591
let participant_token = header_token.as_deref(); - 3592
let provided = header_token.clone().or(cookie_token).or(query_token); - 3593
let ok = provided - 3594
.as_deref() - 3595
.map(|p| p.as_bytes().ct_eq(token.as_bytes()).into()) - 3596
.unwrap_or(false); - 3597
if ok { - 3598
req.extensions_mut() - 3599
.insert(AuthenticatedPrincipal::Operator); - 3600
next.run(req).await - 3601
} else if let Some(participant_token) = participant_token { - 3602
match coworking::verify( - 3603
&coworking::store_path(&home), - 3604
participant_token, - 3605
chrono::Utc::now(), - 3606
) { - 3607
Ok(Some(principal)) => { - 3608
if !participant_read_route_allowed(req.method(), req.uri().path(), &principal) { - 3609
return StatusCode::FORBIDDEN.into_response(); - 3610
} - 3611
req.extensions_mut() - 3612
.insert(AuthenticatedPrincipal::Participant(principal)); - 3613
next.run(req).await - 3614
} - 3615
Ok(None) => unauthorized_response(&home, &req, provided.as_deref()), - 3616
Err(error) => { - 3617
vak_core::security_events::record( - 3618
&home, - 3619
vak_core::security_events::EventKind::AuthFailure, - 3620
"coworking_grant_store_unavailable", - 3621
&format!("path={} error={error}", req.uri().path()), - 3622
None, - 3623
); - 3624
StatusCode::SERVICE_UNAVAILABLE.into_response() - 3625
} - 3626
} - 3627
} else { - 3628
unauthorized_response(&home, &req, provided.as_deref()) - 3629
} - 3630
} - 3631
- 3632
fn unauthorized_response( - 3633
home: &std::path::Path, - 3634
req: &axum::extract::Request, - 3635
provided: Option<&str>, - 3636
) -> axum::response::Response { - 3637
let ip = req - 3638
.headers() - 3639
.get("x-forwarded-for") - 3640
.and_then(|v| v.to_str().ok()) - 3641
.and_then(|v| v.split(',').next()) - 3642
.map(str::trim) - 3643
.filter(|s| !s.is_empty()); - 3644
let detail = format!( - 3645
"path={} provided={}", - 3646
req.uri().path(), - 3647
provided - 3648
.map(|p| format!( - 3649
"{}...{}", - 3650
&p[..4.min(p.len())], - 3651
&p[p.len().saturating_sub(4)..] - 3652
)) - 3653
.unwrap_or_else(|| "<none>".into()) - 3654
); - 3655
vak_core::security_events::record( - 3656
home, - 3657
vak_core::security_events::EventKind::AuthFailure, - 3658
"auth_failure", - 3659
&detail, - 3660
ip, - 3661
); - 3662
if let Some(hub) = events::global() { - 3663
hub.emit_security("AuthFailure", req.uri().path()); - 3664
} - 3665
StatusCode::UNAUTHORIZED.into_response() - 3666
} - 3667
- 3668
fn health_projection(state: &AppState) -> serde_json::Value { - 3669
let report = vak_core::health::collect(&state.core, None); - 3670
let checks: Vec<serde_json::Value> = report - 3671
.checks - 3672
.into_iter() - 3673
.map(|check| match check.detail { - 3674
Ok(detail) => { - 3675
serde_json::json!({ "label": check.label, "status": "pass", "detail": detail }) - 3676
} - 3677
Err(detail) => { - 3678
serde_json::json!({ "label": check.label, "status": "fail", "detail": detail }) - 3679
} - 3680
}) - 3681
.collect(); - 3682
let route = state.core.effective_route(); - 3683
let posture = if report.failures == 0 { - 3684
"healthy" - 3685
} else { - 3686
"degraded" - 3687
}; - 3688
serde_json::json!({ - 3689
// `status = ok` is retained for existing health clients; posture is - 3690
// the truthful operational signal and is what the Operations Center - 3691
// renders. This keeps the compatibility contract without hiding - 3692
// failed doctor checks. - 3693
"status": "ok", - 3694
"posture": posture, - 3695
"provider": route.provider, - 3696
"model": route.model, - 3697
"provider_source": route.provider_source, - 3698
"model_source": route.model_source, - 3699
"route_revision": route.revision, - 3700
"permission_mode": format!("{:?}", state.core.effective_permission_mode()), - 3701
"approval_mode": state.core.effective_approval_mode().as_str(), - 3702
"sandbox": state.core.effective_sandbox_name(), - 3703
"context_window": state.core.config().context_window, - 3704
"voice": { - 3705
"enabled": state.core.effective_voice().enabled, - 3706
"provider": state.core.effective_voice().provider, - 3707
"transcription_model": state.core.effective_voice().transcription_model, - 3708
"synthesis_model": state.core.effective_voice().synthesis_model, - 3709
"max_session_secs": state.core.effective_voice().max_session_secs, - 3710
"max_concurrent": state.core.effective_voice().max_concurrent, - 3711
"max_audio_bytes": state.core.effective_voice().max_audio_bytes, - 3712
"source": "effective", - 3713
"active_sessions": state.voice_active.load(std::sync::atomic::Ordering::Relaxed), - 3714
"capacity_remaining": state.core.effective_voice().max_concurrent.saturating_sub( - 3715
state.voice_active.load(std::sync::atomic::Ordering::Relaxed), - 3716
), - 3717
"quota": { - 3718
"session_seconds": state.core.effective_voice().max_session_secs, - 3719
"concurrent_sessions": state.core.effective_voice().max_concurrent, - 3720
"inbound_audio_bytes": state.core.effective_voice().max_audio_bytes, - 3721
"scope": "workspace", - 3722
"source": "effective", - 3723
}, - 3724
// Historical voice telemetry is intentionally unavailable until - 3725
// it is derived from persisted evidence; never fabricate zeros. - 3726
"historical": { - 3727
"available": false, - 3728
"reason": "No persisted voice latency, error, or cost aggregates are available" - 3729
}, - 3730
}, - 3731
"cwd": state.core.cwd(), - 3732
"warnings": state.core.config().warnings, - 3733
"checks": checks, - 3734
"facts": report.facts, - 3735
"failures": report.failures, - 3736
}) - 3737
} - 3738
- 3739
async fn health(State(state): State<AppState>) -> Json<serde_json::Value> { - 3740
refresh_control_plane(&state); - 3741
Json(health_projection(&state)) - 3742
} - 3743
- 3744
#[allow(clippy::too_many_arguments)] - 3745
/// Build the one canonical presentation projection used when a live handle is - 3746
/// created or a run settles. Keeping both boundaries on this path prevents a - 3747
/// completion rebase from silently dropping plugin renderers, adaptive - 3748
/// library selection, or sandbox-artifact sidecars until the next reload. - 3749
fn live_presentation_snapshot( - 3750
core: &Core, - 3751
session_id: &str, - 3752
session: &SessionLog, - 3753
) -> vak_delivery::OutputTimeline { - 3754
let planner = delivery::merged_presentation_planner(core); - 3755
let adaptive_store = vak_store::presentation::PresentationStore::new( - 3756
core.sessions_home().join("presentations.json"), - 3757
); - 3758
let mut timeline = match adaptive_store.load() { - 3759
Ok(library) => { - 3760
let effective = effective_presentation_library(&library, &core.cwd().to_string_lossy()); - 3761
crate::projection::snapshot_with_planner_and_library( - 3762
session_id, session, &planner, &effective, - 3763
) - 3764
} - 3765
Err(_) => crate::projection::snapshot_with_planner(session_id, session, &planner), - 3766
}; - 3767
crate::projection::append_sandbox_artifacts(&mut timeline, &core.sessions_home(), session_id); - 3768
timeline - 3769
} - 3770
- 3771
pub(crate) fn register_handle( - 3772
state: &AppState, - 3773
id: String, - 3774
session: SessionLog, - 3775
cwd: PathBuf, - 3776
core: Core, - 3777
) -> Arc<SessionHandle> { - 3778
let core = core - 3779
.with_agent_identity(session.header().and_then(|header| header.agent.clone())) - 3780
.with_conversation_context( - 3781
session - 3782
.header() - 3783
.and_then(|header| header.conversation.clone()), - 3784
); - 3785
let durable_home = core.sessions_home(); - 3786
let latest_intent = session.chain_to_root().iter().rev().find_map(|entry| { - 3787
if let vak_session::EntryPayload::Intent(record) = &entry.payload { - 3788
Some((**record).clone()) - 3789
} else { - 3790
None - 3791
} - 3792
}); - 3793
let events_tx = events::EventBus::new(); - 3794
let side_events_tx = events::EventBus::new(); - 3795
let presentation_snapshot = live_presentation_snapshot(&core, &id, &session); - 3796
let presentation = Arc::new(Mutex::new(presentation_snapshot)); - 3797
// The handle's own projector, not a client: `external_subscribers()` - 3798
// must not count it, or "is anyone actually watching" (the /run attach - 3799
// wait, idle eviction) can never observe zero (finding 2/3). - 3800
let mut presentation_rx = events_tx.subscribe_internal(); - 3801
let presentation_state = presentation.clone(); - 3802
let presentation_activities = Arc::new(Mutex::new(Vec::new())); - 3803
let handle = Arc::new(SessionHandle { - 3804
id: id.clone(), - 3805
core, - 3806
cwd, - 3807
session: Arc::new(Mutex::new(Some(session))), - 3808
intent: Arc::new(Mutex::new(latest_intent)), - 3809
steering: Arc::new(SteeringQueues::new()), - 3810
cancel: Arc::new(std::sync::Mutex::new(CancellationToken::new())), - 3811
events_tx, - 3812
coworking_comments_tx: tokio::sync::broadcast::channel(32).0, - 3813
pending: Arc::new(Mutex::new(HashMap::new())), - 3814
activity_buffer: presentation_activities.clone(), - 3815
presentation, - 3816
subscribed: Arc::new(tokio::sync::Notify::new()), - 3817
side_events_tx, - 3818
last_touched: Mutex::new(std::time::Instant::now()), - 3819
side_cancel: Arc::new(std::sync::Mutex::new(CancellationToken::new())), - 3820
admissions: Arc::new(Mutex::new(HashSet::new())), - 3821
}); - 3822
if let Ok(runtime) = tokio::runtime::Handle::try_current() { - 3823
let durable_session_id = id.clone(); - 3824
runtime.spawn(async move { - 3825
loop { - 3826
match presentation_rx.recv().await { - 3827
Ok(framed) => { - 3828
let event = framed.event.clone(); - 3829
append_session_sandbox_event(&durable_home, &durable_session_id, &event); - 3830
crate::projection::project_frame( - 3831
&mut presentation_state - 3832
.lock() - 3833
.unwrap_or_else(std::sync::PoisonError::into_inner), - 3834
framed, - 3835
); - 3836
let activity = match event { - 3837
AgentEvent::WorkerStarted { label } => { - 3838
Some(vak_session::ActivityRecord { - 3839
activity_id: format!("worker-{label}"), - 3840
turn: None, - 3841
kind: vak_session::ActivityKind::Worker, - 3842
status: vak_session::ActivityStatus::Running, - 3843
label, - 3844
detail: Some("Worker started".into()), - 3845
data: std::collections::BTreeMap::new(), - 3846
}) - 3847
} - 3848
AgentEvent::WorkerFinished { - 3849
label, - 3850
is_error, - 3851
elapsed_ms, - 3852
} => Some(vak_session::ActivityRecord { - 3853
activity_id: format!("worker-{label}"), - 3854
turn: None, - 3855
kind: vak_session::ActivityKind::Worker, - 3856
status: if is_error { - 3857
vak_session::ActivityStatus::Failed - 3858
} else { - 3859
vak_session::ActivityStatus::Succeeded - 3860
}, - 3861
label, - 3862
detail: Some(format!("Completed in {elapsed_ms} ms")), - 3863
data: std::collections::BTreeMap::new(), - 3864
}), - 3865
_ => None, - 3866
}; - 3867
if let Some(activity) = activity { - 3868
presentation_activities - 3869
.lock() - 3870
.unwrap_or_else(std::sync::PoisonError::into_inner) - 3871
.push(activity); - 3872
} - 3873
} - 3874
Err(broadcast::error::RecvError::Lagged(_)) => continue, - 3875
Err(broadcast::error::RecvError::Closed) => break, - 3876
} - 3877
} - 3878
}); - 3879
} - 3880
state - 3881
.sessions - 3882
.lock() - 3883
.unwrap_or_else(std::sync::PoisonError::into_inner) - 3884
.insert(id, handle.clone()); - 3885
state.evict_idle_sessions(); - 3886
handle - 3887
} - 3888
- 3889
async fn create_session(State(state): State<AppState>) -> axum::response::Response { - 3890
use axum::response::IntoResponse; - 3891
refresh_control_plane(&state); - 3892
// The workspace the client currently has open, which on the web is - 3893
// switchable at runtime (docs/design/48-web-client.md §5). Existing - 3894
// sessions keep the `Core` they froze at creation (invariant 17); this - 3895
// only decides where the NEXT task lives. - 3896
let core = state.active_core(); - 3897
let session = match core.start_session().await { - 3898
Ok(s) => s, - 3899
Err(e) => { - 3900
return ( - 3901
StatusCode::INTERNAL_SERVER_ERROR, - 3902
Json(serde_json::json!({ "error": e.to_string() })), - 3903
) - 3904
.into_response(); - 3905
} - 3906
}; - 3907
let id = session - 3908
.header() - 3909
.map(|h| h.session_id.clone()) - 3910
.unwrap_or_default(); - 3911
register_handle( - 3912
&state, - 3913
id.clone(), - 3914
session, - 3915
core.cwd().clone(), - 3916
core.clone(), - 3917
); - 3918
- 3919
state.hub.emit_session_created(&id, ""); - 3920
index_session_later(state.store.clone(), state.core.sessions_home(), id.clone()); - 3921
- 3922
Json(serde_json::json!({ "session_id": id })).into_response() - 3923
} - 3924
- 3925
/// Re-index one session's JSONL in the background. Reading does not - 3926
/// conflict with the live handle's exclusive write lock. - 3927
pub(crate) fn index_session_later( - 3928
store: Option<vak_store::Store>, - 3929
home: std::path::PathBuf, - 3930
session_id: String, - 3931
) { - 3932
let Some(store) = store else { - 3933
return; - 3934
}; - 3935
tokio::spawn(async move { - 3936
import_session_sync(&store, &home, &session_id); - 3937
}); - 3938
} - 3939
- 3940
/// Locate `<home>/sessions/<hash>/<session>.jsonl` and import it into the - 3941
/// index synchronously. Idempotent; cheap when nothing changed. - 3942
pub(crate) fn import_session_sync( - 3943
store: &vak_store::Store, - 3944
home: &std::path::Path, - 3945
session_id: &str, - 3946
) -> bool { - 3947
let dir = home.join("sessions"); - 3948
if let Ok(read) = std::fs::read_dir(&dir) { - 3949
for project in read.flatten() { - 3950
let candidate = project.path().join(format!("{session_id}.jsonl")); - 3951
if candidate.is_file() - 3952
&& let Ok(stats) = store.import_session(home, &candidate) - 3953
{ - 3954
return stats.entries_indexed > 0 || stats.skipped > 0; - 3955
} - 3956
} - 3957
} - 3958
let shared = if home.join("agents").is_dir() { - 3959
home.to_path_buf() - 3960
} else if let Some(parent) = home.parent().and_then(|p| p.parent()) { - 3961
parent.to_path_buf() - 3962
} else { - 3963
home.to_path_buf() - 3964
}; - 3965
if let Ok(agents) = std::fs::read_dir(shared.join("agents")) { - 3966
for agent in agents.flatten() { - 3967
let agent_home = agent.path(); - 3968
let agent_sessions = agent_home.join("sessions"); - 3969
if let Ok(projects) = std::fs::read_dir(&agent_sessions) { - 3970
for project in projects.flatten() { - 3971
let candidate = project.path().join(format!("{session_id}.jsonl")); - 3972
if candidate.is_file() - 3973
&& let Ok(stats) = store.import_session(&agent_home, &candidate) - 3974
{ - 3975
return stats.entries_indexed > 0 || stats.skipped > 0; - 3976
} - 3977
} - 3978
} - 3979
} - 3980
} - 3981
false - 3982
} - 3983
- 3984
#[derive(serde::Deserialize)] - 3985
struct AttachBody { - 3986
session_id: String, - 3987
} - 3988
- 3989
async fn attach_session( - 3990
State(state): State<AppState>, - 3991
Json(body): Json<AttachBody>, - 3992
) -> axum::response::Response { - 3993
match ensure_session_handle(&state, &body.session_id).await { - 3994
Ok((id, _)) => ( - 3995
StatusCode::OK, - 3996
Json(serde_json::json!({ "session_id": id })), - 3997
) - 3998
.into_response(), - 3999
Err(e) => ( - 4000
StatusCode::NOT_FOUND, - 4001
Json(serde_json::json!({ "error": e.to_string() })), - 4002
) - 4003
.into_response(), - 4004
} - 4005
} - 4006
- 4007
/// Resolve a durable conversation into the live handle map at an admission - 4008
/// boundary. A browser may keep its page and EventSources across a server - 4009
/// restart; neither a follow-up run nor a reconnected stream can assume an - 4010
/// earlier explicit `/attach` call still exists in this process. - 4011
async fn ensure_session_handle( - 4012
state: &AppState, - 4013
session_id: &str, - 4014
) -> Result<(String, Arc<SessionHandle>), vak_core::CoreError> { - 4015
// Already attached? Return before touching the file. - 4016
// - 4017
// The live handle owns an exclusive lock on the session JSONL for its - 4018
// whole lifetime, and the lock is per open-file-description: opening the - 4019
// same path again from THIS process conflicts with our own handle just - 4020
// as it would with a stranger's. Re-attaching is routine — the desktop - 4021
// calls it on every task switch, and mid-run the handle's session is - 4022
// temporarily owned by the agent — so this must be a no-op, not a - 4023
// second open. - 4024
if vak_core::trash::is_trashed(&state.core.shared_data_home(), session_id) { - 4025
return Err(vak_core::CoreError::Session(vak_session::SessionError::Io( - 4026
std::io::Error::new( - 4027
std::io::ErrorKind::NotFound, - 4028
format!("session is in the trash: {session_id}"), - 4029
), - 4030
))); - 4031
} - 4032
if let Some(handle) = state.get(session_id) { - 4033
return Ok((session_id.to_owned(), handle)); - 4034
} - 4035
let session = if let Ok(s) = state.active_core().open_session(session_id).await { - 4036
Ok(s) - 4037
} else if let Ok(s) = state.core.open_session(session_id).await { - 4038
Ok(s) - 4039
} else if let Ok(s) = state.active_core().open_session_read_only(session_id).await { - 4040
Ok(s) - 4041
} else if let Ok(s) = state.core.open_session_read_only(session_id).await { - 4042
Ok(s) - 4043
} else { - 4044
find_session_on_disk(&state.core, session_id).ok_or_else(|| { - 4045
vak_core::CoreError::Session(vak_session::SessionError::Io(std::io::Error::new( - 4046
std::io::ErrorKind::NotFound, - 4047
format!("session not found: {session_id}"), - 4048
))) - 4049
}) - 4050
}; - 4051
session.map(|session| { - 4052
let header = session.header(); - 4053
let id = header - 4054
.map(|h| h.session_id.clone()) - 4055
.unwrap_or_else(|| session_id.to_owned()); - 4056
let session_cwd = header - 4057
.map(|h| h.cwd.clone()) - 4058
.unwrap_or_else(|| state.core.cwd().clone()); - 4059
let handle_core = match header { - 4060
Some(header) => resolve_core_for_header(state, header), - 4061
None => resolve_process_core_for_cwd(state, &state.active_core(), &session_cwd), - 4062
}; - 4063
// The header id can differ from the requested one; if that handle - 4064
// is already live, keep it rather than replacing it. - 4065
let handle = state.get(&id).unwrap_or_else(|| { - 4066
register_handle(state, id.clone(), session, session_cwd, handle_core) - 4067
}); - 4068
(id, handle) - 4069
}) - 4070
} - 4071
- 4072
/// The one derivation of "which `Core` should this session's next turn run - 4073
/// under," from its own recorded header (finding 6). A custom Agent's - 4074
/// session gets the exact pins `agent_chats::resolve_agent_core` applies — - 4075
/// permission-mode cap, sandbox backend override, provider instance - 4076
/// override, and shared sessions_home — via `agent_chats:: - 4077
/// pinned_core_for_workspace`, keyed by the session's OWN recorded - 4078
/// `header.cwd` rather than a workspace path re-derived from whatever - 4079
/// happens to be the CURRENT active workspace (which can disagree once the - 4080
/// active workspace has moved on since the session was created — see that - 4081
/// function's doc comment). Before this, `/attach`, `/run` and the SSE - 4082
/// endpoints resolved a plain pooled `Core` for the cwd with none of those - 4083
/// pins, so the same session's security ceiling depended on which endpoint - 4084
/// happened to touch it first. The built-in `vak` identity (and a - 4085
/// pre-Agent ledger with no `agent` on its header at all) keeps the plain - 4086
/// active/process core resolution. - 4087
fn resolve_core_for_header(state: &AppState, header: &vak_session::types::SessionHeader) -> Core { - 4088
let active = state.active_core(); - 4089
match header.agent.as_ref() { - 4090
Some(identity) if identity.id != "vak" => { - 4091
agent_chats::pinned_core_for_workspace(state, &active, identity, &header.cwd) - 4092
.unwrap_or_else(|_| resolve_process_core_for_cwd(state, &active, &header.cwd)) - 4093
} - 4094
_ => resolve_process_core_for_cwd(state, &active, &header.cwd), - 4095
} - 4096
} - 4097
- 4098
/// Plain cwd-keyed pooled `Core` resolution with no Agent-specific pins — - 4099
/// the built-in `vak` identity's own path, and `resolve_core_for_header`'s - 4100
/// fallback when a custom Agent's pinned resolution itself fails. - 4101
fn resolve_process_core_for_cwd(state: &AppState, active: &Core, cwd: &std::path::Path) -> Core { - 4102
if cwd == active.cwd().as_path() { - 4103
return active.clone(); - 4104
} - 4105
if cwd == state.core.cwd().as_path() { - 4106
return state.core.clone(); - 4107
} - 4108
if let Ok(c) = state - 4109
.gateway - 4110
.core_pool - 4111
.resolve_at(cwd, None, std::time::Instant::now()) - 4112
{ - 4113
return c; - 4114
} - 4115
// `resolve_at` failing here is not a trust decision — an unconditional - 4116
// `true` would let a workspace whose trust prompt an operator declined - 4117
// have its hooks/MCP servers/secret scope applied anyway. Recompute - 4118
// trust the same way `resolve_at` does rather than assuming it. - 4119
vak_core::Core::new_with_trust(cwd.to_path_buf(), vak_core::trust::is_trusted(cwd)) - 4120
.unwrap_or_else(|_| state.core.clone()) - 4121
} - 4122
- 4123
#[derive(serde::Deserialize, Default)] - 4124
struct ListSessionsQuery { - 4125
/// Lists the trash instead: only the sessions moved there, so a person - 4126
/// can restore one. Nothing else reads a trashed session. - 4127
#[serde(default)] - 4128
trash: bool, - 4129
} - 4130
- 4131
/// Sidebar projection over the persisted store: one summary per JSONL file. - 4132
async fn list_sessions( - 4133
State(state): State<AppState>, - 4134
axum::extract::Query(query): axum::extract::Query<ListSessionsQuery>, - 4135
) -> Json<serde_json::Value> { - 4136
// Sessions are stored per workspace, so this follows the workspace the - 4137
// client has open rather than the one the process started in. - 4138
let active = state.active_core(); - 4139
let dir = vak_session::SessionPath::sessions_dir(&state.core.sessions_home(), active.cwd()); - 4140
let active_cwd = active.cwd().to_string_lossy().into_owned(); - 4141
let archive_map = read_archive(&state.core); - 4142
let trashed = vak_core::trash::trashed(&state.core.shared_data_home()); - 4143
let mut sessions = Vec::new(); - 4144
let mut entries: Vec<std::fs::DirEntry> = std::fs::read_dir(&dir) - 4145
.map(|read| read.flatten().collect()) - 4146
.unwrap_or_default(); - 4147
// A workspace can be renamed or canonicalized between runs (notably - 4148
// `/var` vs `/private/var` on macOS). Recover sessions by their durable - 4149
// header cwd when the hashed directory no longer matches, while still - 4150
// filtering strictly to the active workspace. - 4151
if let Ok(projects) = std::fs::read_dir(state.core.sessions_home().join("sessions")) { - 4152
for project in projects.flatten() { - 4153
if let Ok(files) = std::fs::read_dir(project.path()) { - 4154
for file in files.flatten() { - 4155
let duplicate = entries - 4156
.iter() - 4157
.any(|existing| existing.path() == file.path()); - 4158
if !duplicate { - 4159
entries.push(file); - 4160
} - 4161
} - 4162
} - 4163
}
Indexing the workspace…
Vakyartha documentation is discovering safe artifacts, anchors, and source references.