- 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
} - 4164
} - 4165
let shared_home = state.core.shared_data_home(); - 4166
if let Ok(agents) = std::fs::read_dir(shared_home.join("agents")) { - 4167
for agent in agents.flatten() { - 4168
if let Ok(projects) = std::fs::read_dir(agent.path().join("sessions")) { - 4169
for project in projects.flatten() { - 4170
if let Ok(files) = std::fs::read_dir(project.path()) { - 4171
for file in files.flatten() { - 4172
let duplicate = entries - 4173
.iter() - 4174
.any(|existing| existing.path() == file.path()); - 4175
if !duplicate { - 4176
entries.push(file); - 4177
} - 4178
} - 4179
} - 4180
} - 4181
} - 4182
} - 4183
} - 4184
for entry in entries { - 4185
let path = entry.path(); - 4186
if path.extension().and_then(|e| e.to_str()) != Some("jsonl") { - 4187
continue; - 4188
} - 4189
let Some(session_id) = path.file_stem().and_then(|s| s.to_str()).map(String::from) else { - 4190
continue; - 4191
}; - 4192
if trashed.contains(&session_id) != query.trash { - 4193
continue; - 4194
} - 4195
let updated_at = std::fs::metadata(&path) - 4196
.ok() - 4197
.and_then(|m| m.modified().ok()) - 4198
.map(|t| chrono::DateTime::<chrono::Utc>::from(t).to_rfc3339()); - 4199
let (created_at, title, entry_count, cwd, agent) = summarize_jsonl(&path); - 4200
// A user-created Agent lives in its own isolated workspace (see - 4201
// agent_chats::open / agent_workspace), independent of whichever - 4202
// default workspace the browser client currently has open — that - 4203
// switch only ever applied to the built-in "vak" identity, which - 4204
// still shares the process's default workspace. So only the "vak" - 4205
// (or header-less/legacy) sessions are filtered by `active_cwd`; - 4206
// every other Agent's sessions are always its own to show. - 4207
let is_default_agent = agent.as_ref().is_none_or(|a| a.id == "vak"); - 4208
if is_default_agent && cwd.as_deref() != Some(active_cwd.as_str()) { - 4209
continue; - 4210
} - 4211
// Header-only sessions are abandoned drafts (for example, creating a - 4212
// task and immediately switching away). Keep the ledger append-only, - 4213
// but do not let empty drafts accumulate in the task switcher. This - 4214
// applies equally to built-in and user-created Agents, but actively - 4215
// registered sessions must remain discoverable. - 4216
let is_active = state.get(&session_id).is_some(); - 4217
if entry_count <= 1 && !is_active { - 4218
continue; - 4219
} - 4220
let running = state.get(&session_id).is_some_and(|handle| { - 4221
handle - 4222
.session - 4223
.lock() - 4224
.unwrap_or_else(std::sync::PoisonError::into_inner) - 4225
.is_none() - 4226
}); - 4227
let archived = archive_map.get(&session_id).copied().unwrap_or(false); - 4228
sessions.push(serde_json::json!({ - 4229
"session_id": session_id, - 4230
"cwd": cwd.unwrap_or_else(|| state.core.cwd().to_string_lossy().into_owned()), - 4231
"created_at": created_at, - 4232
"updated_at": updated_at, - 4233
"entries": entry_count, - 4234
"agent": agent, - 4235
"title": title, - 4236
"running": running, - 4237
"archived": archived, - 4238
})); - 4239
} - 4240
sessions.sort_by_key(|s| s["updated_at"].as_str().unwrap_or("").to_string()); - 4241
sessions.reverse(); - 4242
Json(serde_json::json!({ "sessions": sessions })) - 4243
} - 4244
- 4245
/// Bounded scan: header line for created_at + first user message as title. - 4246
fn summarize_jsonl( - 4247
path: &std::path::Path, - 4248
) -> ( - 4249
Option<String>, - 4250
Option<String>, - 4251
u64, - 4252
Option<String>, - 4253
Option<vak_session::types::AgentIdentity>, - 4254
) { - 4255
use std::io::BufRead; - 4256
let Ok(file) = std::fs::File::open(path) else { - 4257
return (None, None, 0, None, None); - 4258
}; - 4259
let mut reader = std::io::BufReader::new(file); - 4260
let mut created_at = None; - 4261
let mut title = None; - 4262
let mut cwd = None; - 4263
let mut agent = None; - 4264
let mut entries = 0u64; - 4265
let mut line = String::new(); - 4266
loop { - 4267
line.clear(); - 4268
match reader.read_line(&mut line) { - 4269
Ok(0) => break, - 4270
Ok(_) => { - 4271
entries += 1; - 4272
if let Ok(entry) = serde_json::from_str::<vak_session::Entry>(line.trim()) { - 4273
match entry.payload { - 4274
vak_session::EntryPayload::Header(h) => { - 4275
created_at = Some(h.created_at.to_rfc3339()); - 4276
cwd = Some(h.cwd.to_string_lossy().into_owned()); - 4277
agent = h.agent; - 4278
} - 4279
vak_session::EntryPayload::Message(rec) => { - 4280
if title.is_none() - 4281
&& rec.message.role == vak_llm::Role::User - 4282
&& rec.control_kind().is_none() - 4283
{ - 4284
let text = rec.message.text_content(); - 4285
let text = text.trim(); - 4286
if !text.is_empty() { - 4287
let first_line = text.lines().next().unwrap_or(text).trim(); - 4288
let mut snippet: String = first_line.chars().take(72).collect(); - 4289
if first_line.chars().count() > 72 { - 4290
snippet.push('…'); - 4291
} - 4292
title = Some(snippet); - 4293
} - 4294
} - 4295
} - 4296
vak_session::EntryPayload::Compaction(_) => {} - 4297
vak_session::EntryPayload::Receipt(_) => {} - 4298
vak_session::EntryPayload::Goal(_) => {} - 4299
vak_session::EntryPayload::GoalUpdate(_) => {} - 4300
vak_session::EntryPayload::Activity(_) => {} - 4301
vak_session::EntryPayload::Work(_) => {} - 4302
vak_session::EntryPayload::Intent(_) => {} - 4303
vak_session::EntryPayload::TurnCapabilitiesBound(_) - 4304
| vak_session::EntryPayload::TurnCapabilitiesRef(_) => {} - 4305
vak_session::EntryPayload::ChildRun { .. } => {} - 4306
vak_session::EntryPayload::Presentation(_) => {} - 4307
vak_session::EntryPayload::TurnCard(_) => {} - 4308
vak_session::EntryPayload::EvidenceBody(_) => {} - 4309
} - 4310
} - 4311
if title.is_some() && entries > 400 { - 4312
break; - 4313
} - 4314
} - 4315
Err(_) => break, - 4316
} - 4317
} - 4318
(created_at, title, entries, cwd, agent) - 4319
} - 4320
- 4321
#[derive(serde::Deserialize)] - 4322
struct RoutingEnvelope { - 4323
/// The user message that caused this admission. This is metadata, not a - 4324
/// capability or an instruction to the model. - 4325
#[serde(default)] - 4326
message_id: Option<String>, - 4327
#[serde(default)] - 4328
conversation_id: Option<String>, - 4329
#[serde(default)] - 4330
target_work_id: Option<String>, - 4331
#[serde(default)] - 4332
target_result_id: Option<String>, - 4333
/// `independent`, `follow_up`, `correction`, `status`, `cancel`, or - 4334
/// `schedule`; unknown values are retained as provenance but never used - 4335
/// to authorize work. - 4336
#[serde(default)] - 4337
relation: Option<String>, - 4338
#[serde(default)] - 4339
outcome_revision: Option<u64>, - 4340
#[serde(default)] - 4341
provenance: Option<String>, - 4342
} - 4343
- 4344
#[derive(serde::Deserialize)] - 4345
struct RunBody { - 4346
prompt: String, - 4347
/// Stable client identity used to make network retries idempotent. - 4348
#[serde(default)] - 4349
request_id: Option<String>, - 4350
#[serde(default)] - 4351
routing: Option<RoutingEnvelope>, - 4352
/// Optional run-scoped work profile. `managed` creates and persists a - 4353
/// work contract before the agent can execute tools. - 4354
#[serde(default)] - 4355
work_mode: Option<String>, - 4356
/// Optional base64 images appended to the prompt as vision content - 4357
/// (docs/design/22-gateway.md media passthrough). - 4358
#[serde(default)] - 4359
attachments: Vec<RunAttachment>, - 4360
/// Files already saved in the workspace inbox (`POST /fs/inbox`), named - 4361
/// by the paths that route returned. Each reaches the model as a note - 4362
/// saying where it is, never as its bytes (docs/design/72, F1). - 4363
#[serde(default)] - 4364
files: Vec<String>, - 4365
/// Goal mode (docs/design/42-managed-work-contracts.md): durable objective; completion - 4366
/// is audited against `criteria`, never self-reported. - 4367
#[serde(default)] - 4368
goal: Option<String>, - 4369
/// Acceptance criteria for goal mode (`verify:` prefixed criteria run - 4370
/// as brokered shell commands; others are judged from evidence). - 4371
#[serde(default)] - 4372
criteria: Vec<String>, - 4373
} - 4374
- 4375
#[derive(serde::Deserialize)] - 4376
struct RunAttachment { - 4377
#[serde(default = "default_image_mime")] - 4378
mime: String, - 4379
data: String, - 4380
} - 4381
- 4382
fn default_image_mime() -> String { - 4383
"image/png".into() - 4384
} - 4385
- 4386
pub(crate) fn mpsc_to_broadcast(tx: events::EventBus) -> mpsc::Sender<AgentEvent> { - 4387
let (tx_in, mut rx) = mpsc::channel::<AgentEvent>(512); - 4388
tokio::spawn(async move { - 4389
// Forward into the BROADCAST channel (sync send). Forwarding into - 4390
// tx_in would feed the channel back into itself. - 4391
// - 4392
// Headless consumers (gateway turns, cron routines) legitimately run - 4393
// with zero broadcast subscribers; send errors must NEVER tear the - 4394
// pump down — the agent treats a dropped mpsc receiver as a lost - 4395
// consumer and cancels the run mid-flight. - 4396
while let Some(ev) = rx.recv().await { - 4397
let _ = tx.send(ev); - 4398
} - 4399
}); - 4400
tx_in - 4401
} - 4402
- 4403
/// A run cannot start without a working provider credential. - 4404
/// - 4405
/// Every one of these three call sites used to return a bare 503 with no - 4406
/// body and nothing logged, which made "the agent never replied" a - 4407
/// symptom with no server-side trail: a client saw an empty response, an - 4408
/// operator reading gateway.log saw nothing at all, and diagnosing it - 4409
/// meant reading this file. `Core::provider()` already carries a precise - 4410
/// `CoreError::MissingAuth { env, provider }` — this puts it where an - 4411
/// operator and a client can both actually see it. - 4412
/// - 4413
/// The body is `{"error": <message>}`, matching every other handler in - 4414
/// this file. A `{"error": <code>, "detail": <message>}` shape was tried - 4415
/// first and reverted: the desktop frontend's error handling already - 4416
/// reads `.error` as the human-readable string every other endpoint puts - 4417
/// there, so a two-field body would have shown the user the machine code - 4418
/// ("provider_unavailable") instead of the message that says what to fix. - 4419
/// - 4420
/// A missing credential also carries `"kind": "no_ai_service"` (see - 4421
/// `provider_error_body`). - 4422
fn provider_unavailable(err: vak_core::CoreError) -> axum::response::Response { - 4423
use axum::response::IntoResponse; - 4424
eprintln!("[run] refused: {err}"); - 4425
( - 4426
StatusCode::SERVICE_UNAVAILABLE, - 4427
axum::Json(provider_error_body(&err)), - 4428
) - 4429
.into_response() - 4430
} - 4431
- 4432
/// `{"error": <message>}` for a provider failure, plus `"kind": - 4433
/// "no_ai_service"` when the cause is a missing credential, so a client can - 4434
/// say so in plain words without matching on the message text; `error` - 4435
/// stays the precise message an operator or the CLI needs. The one place - 4436
/// that decides the kind, for a refused turn and a model catalogue alike. - 4437
fn provider_error_body(err: &vak_core::CoreError) -> serde_json::Value { - 4438
match err { - 4439
vak_core::CoreError::MissingAuth { .. } => { - 4440
serde_json::json!({ "error": err.to_string(), "kind": "no_ai_service" }) - 4441
} - 4442
_ => serde_json::json!({ "error": err.to_string() }), - 4443
} - 4444
} - 4445
- 4446
/// `run_prompt`, `side_chat`, and `start_bestofn` all fall back to
Indexing the workspace…
Vakyartha documentation is discovering safe artifacts, anchors, and source references.