- 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 - 4447
/// `provider_unavailable` when `Core::provider()` fails; this pins the - 4448
/// response it produces so a regression — an empty body, or the - 4449
/// `{"error": <code>, "detail": <message>}` shape tried and reverted - 4450
/// above — fails a fast unit test instead of surfacing as "the agent - 4451
/// never replied" with nothing in gateway.log to explain why. - 4452
#[cfg(test)] - 4453
#[allow(clippy::unwrap_used, clippy::expect_used)] - 4454
mod provider_unavailable_tests { - 4455
use super::provider_unavailable; - 4456
use axum::response::IntoResponse as _; - 4457
use http_body_util::BodyExt as _; - 4458
- 4459
#[tokio::test] - 4460
async fn reports_status_and_a_body_naming_the_missing_credential() { - 4461
let err = vak_core::CoreError::MissingAuth { - 4462
env: "ANTHROPIC_API_KEY".into(), - 4463
provider: "anthropic".into(), - 4464
}; - 4465
let response = provider_unavailable(err).into_response(); - 4466
assert_eq!( - 4467
response.status(), - 4468
axum::http::StatusCode::SERVICE_UNAVAILABLE - 4469
); - 4470
- 4471
let bytes = response - 4472
.into_body() - 4473
.collect() - 4474
.await - 4475
.expect("body readable") - 4476
.to_bytes(); - 4477
assert!( - 4478
!bytes.is_empty(), - 4479
"body must not be empty — that was the original bug" - 4480
); - 4481
- 4482
let body: serde_json::Value = serde_json::from_slice(&bytes).expect("body is JSON"); - 4483
// `error` carries the human-readable message, same as every other - 4484
// handler in this file — not a machine code with the message hidden - 4485
// in a `detail` the frontend never reads — and `kind` types it. - 4486
let fields: Vec<&String> = body.as_object().expect("object body").keys().collect(); - 4487
assert_eq!( - 4488
fields, - 4489
vec!["error", "kind"], - 4490
"body must have the `error` message and its `kind`" - 4491
); - 4492
assert_eq!(body["kind"], "no_ai_service"); - 4493
let message = body["error"].as_str().expect("error is a string"); - 4494
assert!( - 4495
message.contains("ANTHROPIC_API_KEY"), - 4496
"message must name the env var to set, got: {message}" - 4497
); - 4498
assert!( - 4499
message.contains("anthropic"), - 4500
"message must name the provider, got: {message}" - 4501
); - 4502
} - 4503
- 4504
#[tokio::test] - 4505
async fn only_a_missing_credential_is_typed_as_no_ai_service() { - 4506
let err = vak_core::CoreError::InvalidConfig("bad route".into()); - 4507
let bytes = provider_unavailable(err) - 4508
.into_response() - 4509
.into_body() - 4510
.collect() - 4511
.await - 4512
.expect("body readable") - 4513
.to_bytes(); - 4514
let body: serde_json::Value = serde_json::from_slice(&bytes).expect("body is JSON"); - 4515
let fields: Vec<&String> = body.as_object().expect("object body").keys().collect(); - 4516
assert_eq!(fields, vec!["error"], "other refusals carry no kind"); - 4517
} - 4518
} - 4519
- 4520
// ---- Shared turn-chain executor (invariant 30; docs/design/ - 4521
// 64-agent-owned-platform.md, "Request durability and delivery") ---------- - 4522
// - 4523
// `run_prompt`, `send_steering`, and `gateway::execute_turn_chain` all - 4524
// admit a prompt, run it, and — if more input arrived while the run was - 4525
// settling — keep going rather than silently stranding it. Before this, - 4526
// each surface implemented that loop separately: the HTTP path did not - 4527
// implement it at all (steering queued after a run's last internal drain - 4528
// was never picked back up), and the gateway's own version restored - 4529
// `handle.session` before draining, leaving a race window where a - 4530
// concurrent admission could steal the ledger. `admit_or_queue` and - 4531
// `continue_or_release` are the one busy/idle decision, in both - 4532
// directions; `run_turn_chain` is the one loop that runs a leg and decides - 4533
// whether to continue, parameterized by approver and run kind so each - 4534
// surface keeps its own settle bookkeeping (durable activity records vs. - 4535
// reply channel + rendered text) without duplicating the loop mechanics. - 4536
- 4537
/// What a turn chain's FIRST leg runs. Every leg after the first is always - 4538
/// a plain message turn: draining `handle.steering` only ever produces a - 4539
/// `vak_llm::Message` via `SteeringQueues::merge_prompt`, never a fresh - 4540
/// goal/managed/auto request — that is `/run`'s own admission, which a - 4541
/// queued steering message never claims to be. - 4542
enum TurnStart { - 4543
/// The person's message as recorded, with its metadata (attached files). - 4544
Message(vak_session::MessageRecord), - 4545
Managed(String), - 4546
Auto(String), - 4547
Goal { - 4548
prompt: String, - 4549
objective: String, - 4550
criteria: Vec<String>, - 4551
}, - 4552
} - 4553
- 4554
impl TurnStart { - 4555
fn message(message: vak_llm::Message) -> Self { - 4556
TurnStart::Message(vak_session::MessageRecord { - 4557
message, - 4558
meta: None, - 4559
}) - 4560
} - 4561
- 4562
/// The message this leg would present — used both to seed the preview - 4563
/// intent before a run starts and, on the busy path, as the queued - 4564
/// steering entry (attachments and all; invariant 1, model-visible - 4565
/// input is never degraded to bare text). - 4566
fn preview_message(&self) -> vak_llm::Message { - 4567
match self { - 4568
TurnStart::Message(m) => m.message.clone(), - 4569
TurnStart::Managed(p) | TurnStart::Auto(p) => vak_llm::Message::user_text(p), - 4570
TurnStart::Goal { prompt, .. } => vak_llm::Message::user_text(prompt), - 4571
} - 4572
} - 4573
- 4574
#[allow(clippy::too_many_arguments)] - 4575
async fn run( - 4576
self, - 4577
core: &Core, - 4578
session: SessionLog, - 4579
cancel: CancellationToken, - 4580
approver: Arc<dyn Approver>, - 4581
steering: Arc<SteeringQueues>, - 4582
events: mpsc::Sender<AgentEvent>, - 4583
) -> Result<(vak_agent::TurnOutcome, SessionLog), vak_core::CoreError> { - 4584
match self { - 4585
TurnStart::Message(m) => { - 4586
core.run_turn_with_message( - 4587
session, - 4588
m, - 4589
cancel, - 4590
Some(approver), - 4591
None, - 4592
Some(steering), - 4593
events, - 4594
) - 4595
.await - 4596
} - 4597
TurnStart::Managed(prompt) => { - 4598
core.run_managed_turn_with( - 4599
session, - 4600
&prompt, - 4601
cancel, - 4602
Some(approver), - 4603
None, - 4604
Some(steering), - 4605
events, - 4606
) - 4607
.await - 4608
} - 4609
TurnStart::Auto(prompt) => { - 4610
core.run_auto_turn_with( - 4611
session, - 4612
&prompt, - 4613
cancel, - 4614
Some(approver), - 4615
None, - 4616
Some(steering), - 4617
events, - 4618
) - 4619
.await - 4620
} - 4621
TurnStart::Goal { - 4622
prompt, - 4623
objective, - 4624
criteria, - 4625
} => { - 4626
core.run_goal_turn_with( - 4627
session, - 4628
&prompt, - 4629
&objective, - 4630
criteria, - 4631
cancel, - 4632
Some(approver), - 4633
None, - 4634
Some(steering), - 4635
events, - 4636
) - 4637
.await - 4638
} - 4639
} - 4640
} - 4641
} - 4642
- 4643
/// Result of admitting input at the busy boundary. - 4644
enum Admission { - 4645
/// The ledger was idle; the caller now owns it and must run a chain.
Indexing the workspace…
Vakyartha documentation is discovering safe artifacts, anchors, and source references.