- 15799
"label": entry.label, - 15800
"description": entry.description, - 15801
"command": entry.command, - 15802
"args": entry.args, - 15803
"network": true, - 15804
"env_var": entry.env_var, - 15805
"key_required": entry.key_required, - 15806
"documentation_url": entry.documentation_url, - 15807
"scope": scope.label(), - 15808
"configured_here": configured_here, - 15809
"inherited": inherited, - 15810
"effective": effective, - 15811
"key_here": key_here, - 15812
"key_inherited": key_inherited, - 15813
"key_effective": entry.env_var.is_none_or(|name| core.mcp_secret(name).is_some()), - 15814
})) - 15815
} - 15816
- 15817
async fn get_integration_catalog( - 15818
State(state): State<AppState>, - 15819
Query(query): Query<ScopeQuery>, - 15820
) -> axum::response::Response { - 15821
use axum::response::IntoResponse; - 15822
let core = scoped_core!(&state, None, query.agent.as_deref()); - 15823
match INTEGRATION_CATALOG - 15824
.iter() - 15825
.copied() - 15826
.map(|entry| integration_status(&core, query.scope, entry)) - 15827
.collect::<Result<Vec<_>, _>>() - 15828
{ - 15829
Ok(integrations) => Json(serde_json::json!({ - 15830
"scope": query.scope.label(), - 15831
"integrations": integrations, - 15832
})) - 15833
.into_response(), - 15834
Err(error) => ( - 15835
StatusCode::INTERNAL_SERVER_ERROR, - 15836
Json(serde_json::json!({ "error": error })), - 15837
) - 15838
.into_response(), - 15839
} - 15840
} - 15841
- 15842
async fn get_scoped_integration( - 15843
State(state): State<AppState>, - 15844
Path(id): Path<String>, - 15845
Query(query): Query<ScopeQuery>, - 15846
) -> axum::response::Response { - 15847
use axum::response::IntoResponse; - 15848
let Some(entry) = catalog_entry(&id) else { - 15849
return StatusCode::NOT_FOUND.into_response(); - 15850
}; - 15851
let core = scoped_core!(&state, None, query.agent.as_deref()); - 15852
match integration_status(&core, query.scope, entry) { - 15853
Ok(status) => Json(status).into_response(), - 15854
Err(error) => ( - 15855
StatusCode::INTERNAL_SERVER_ERROR, - 15856
Json(serde_json::json!({ "error": error })), - 15857
) - 15858
.into_response(), - 15859
} - 15860
} - 15861
- 15862
#[derive(serde::Deserialize)] - 15863
struct IntegrationPutBody { - 15864
scope: ConfigScope, - 15865
key: Option<String>, - 15866
#[serde(default)] - 15867
agent: Option<String>, - 15868
} - 15869
- 15870
fn apply_scoped_mcp_change( - 15871
core: &vak_core::Core, - 15872
scope: ConfigScope, - 15873
id: &str, - 15874
server: Option<vak_config::McpServerConfig>, - 15875
) -> Result<(), String> { - 15876
let path = scope.config_path(core)?; - 15877
vak_config::persist_mcp_server(&path, id, server.as_ref()) - 15878
.map_err(|error| error.to_string())?; - 15879
let effective = vak_config::load_with_trust(core.cwd(), core.project_config_trusted()) - 15880
.map_err(|error| error.to_string())?; - 15881
core.apply_persisted_mcp_servers(effective.mcp); - 15882
Ok(()) - 15883
} - 15884
- 15885
async fn put_scoped_integration( - 15886
State(state): State<AppState>, - 15887
Path(id): Path<String>, - 15888
Json(body): Json<IntegrationPutBody>, - 15889
) -> axum::response::Response { - 15890
use axum::response::IntoResponse; - 15891
let Some(entry) = catalog_entry(&id) else { - 15892
return StatusCode::NOT_FOUND.into_response(); - 15893
}; - 15894
let core = scoped_core!(&state, None, body.agent.as_deref()); - 15895
if let Some(key) = body.key.as_deref() - 15896
&& let Err(error) = core.set_mcp_secret_scoped( - 15897
entry.env_var.unwrap_or_default(), - 15898
key, - 15899
body.scope.is_workspace(), - 15900
) - 15901
{ - 15902
return ( - 15903
StatusCode::BAD_REQUEST, - 15904
Json(serde_json::json!({ "error": error.to_string() })), - 15905
) - 15906
.into_response(); - 15907
} - 15908
let key_available = entry - 15909
.env_var - 15910
.is_none_or(|name| core.mcp_secret(name).is_some()); - 15911
if entry.key_required && !key_available { - 15912
return ( - 15913
StatusCode::BAD_REQUEST, - 15914
Json(serde_json::json!({ - 15915
"error": format!("{} requires {} at this scope or an inherited scope", entry.label, entry.env_var.unwrap_or("a key")) - 15916
})), - 15917
) - 15918
.into_response(); - 15919
} - 15920
if let Err(error) = - 15921
apply_scoped_mcp_change(&core, body.scope, entry.id, Some(catalog_server(entry))) - 15922
{ - 15923
return ( - 15924
StatusCode::INTERNAL_SERVER_ERROR, - 15925
Json(serde_json::json!({ "error": error })), - 15926
) - 15927
.into_response(); - 15928
} - 15929
state.hub.emit_config_changed( - 15930
"integration_enabled", - 15931
&format!("scope={} integration={}", body.scope.label(), entry.id), - 15932
); - 15933
match integration_status(&core, body.scope, entry) { - 15934
Ok(status) => Json(status).into_response(), - 15935
Err(error) => ( - 15936
StatusCode::INTERNAL_SERVER_ERROR, - 15937
Json(serde_json::json!({ "error": error })), - 15938
) - 15939
.into_response(), - 15940
} - 15941
} - 15942
- 15943
async fn delete_scoped_integration( - 15944
State(state): State<AppState>, - 15945
Path(id): Path<String>, - 15946
Query(query): Query<ScopeQuery>, - 15947
) -> axum::response::Response { - 15948
use axum::response::IntoResponse; - 15949
let Some(entry) = catalog_entry(&id) else { - 15950
return StatusCode::NOT_FOUND.into_response(); - 15951
}; - 15952
let core = scoped_core!(&state, None, query.agent.as_deref()); - 15953
if let Some(env_var) = entry.env_var - 15954
&& let Err(error) = core.remove_mcp_secret_scoped(env_var, query.scope.is_workspace()) - 15955
{ - 15956
return ( - 15957
StatusCode::INTERNAL_SERVER_ERROR, - 15958
Json(serde_json::json!({ "error": error.to_string() })), - 15959
) - 15960
.into_response(); - 15961
} - 15962
if let Err(error) = apply_scoped_mcp_change(&core, query.scope, entry.id, None) { - 15963
return ( - 15964
StatusCode::INTERNAL_SERVER_ERROR, - 15965
Json(serde_json::json!({ "error": error })), - 15966
) - 15967
.into_response(); - 15968
} - 15969
state.hub.emit_config_changed( - 15970
"integration_removed", - 15971
&format!("scope={} integration={}", query.scope.label(), entry.id), - 15972
); - 15973
match integration_status(&core, query.scope, entry) { - 15974
Ok(status) => Json(status).into_response(), - 15975
Err(error) => ( - 15976
StatusCode::INTERNAL_SERVER_ERROR, - 15977
Json(serde_json::json!({ "error": error })), - 15978
) - 15979
.into_response(), - 15980
} - 15981
} - 15982
- 15983
async fn get_global_mcp_servers() -> axum::response::Response { - 15984
use axum::response::IntoResponse; - 15985
let Some(path) = vak_config::global_path() else { - 15986
return (StatusCode::INTERNAL_SERVER_ERROR, "user home unavailable").into_response(); - 15987
}; - 15988
match read_mcp_config(&path) { - 15989
Ok(mcp) => { - 15990
Json(serde_json::json!({ "scope": "global", "path": path, "servers": mcp.servers })) - 15991
.into_response() - 15992
} - 15993
Err(error) => ( - 15994
StatusCode::INTERNAL_SERVER_ERROR, - 15995
Json(serde_json::json!({ "error": error })), - 15996
) - 15997
.into_response(), - 15998
} - 15999
} - 16000
- 16001
async fn put_global_mcp_servers( - 16002
State(state): State<AppState>, - 16003
Json(body): Json<McpPutBody>, - 16004
) -> axum::response::Response { - 16005
use axum::response::IntoResponse; - 16006
if let Err(error) = validate_mcp_servers(&body.servers) { - 16007
return ( - 16008
StatusCode::BAD_REQUEST, - 16009
Json(serde_json::json!({ "error": error })), - 16010
) - 16011
.into_response(); - 16012
} - 16013
let Some(path) = vak_config::global_path() else { - 16014
return (StatusCode::INTERNAL_SERVER_ERROR, "user home unavailable").into_response(); - 16015
}; - 16016
let config = mcp_config_from_input(&body.servers); - 16017
if let Err(error) = vak_config::persist_mcp_servers(&path, &config.servers) { - 16018
return ( - 16019
StatusCode::INTERNAL_SERVER_ERROR, - 16020
Json(serde_json::json!({ "error": error.to_string() })), - 16021
) - 16022
.into_response(); - 16023
} - 16024
if state.core.refresh_persisted_preferences().is_err() { - 16025
return ( - 16026
StatusCode::INTERNAL_SERVER_ERROR, - 16027
"could not apply user capability configuration", - 16028
) - 16029
.into_response(); - 16030
} - 16031
state.hub.emit_config_changed( - 16032
"global_mcp_servers_updated", - 16033
&format!("count={}", body.servers.len()), - 16034
); - 16035
Json(serde_json::json!({ "saved": true, "scope": "global", "count": body.servers.len() })) - 16036
.into_response() - 16037
} - 16038
- 16039
fn validate_mcp_servers( - 16040
servers: &std::collections::BTreeMap<String, McpServerInput>, - 16041
) -> Result<(), String> { - 16042
for name in servers.keys() { - 16043
if !valid_server_name(name) { - 16044
return Err(format!("invalid server name '{name}'")); - 16045
} - 16046
} - 16047
for (name, server) in servers { - 16048
if server.command.trim().is_empty() { - 16049
return Err(format!("server '{name}' needs a command")); - 16050
} - 16051
} - 16052
Ok(()) - 16053
} - 16054
- 16055
async fn put_mcp_servers( - 16056
State(state): State<AppState>, - 16057
Json(body): Json<McpPutBody>, - 16058
) -> axum::response::Response { - 16059
use axum::response::IntoResponse; - 16060
if let Err(error) = validate_mcp_servers(&body.servers) { - 16061
return ( - 16062
StatusCode::BAD_REQUEST, - 16063
Json(serde_json::json!({ "error": error })), - 16064
) - 16065
.into_response(); - 16066
} - 16067
let core = scoped_core!(&state, None, body.agent.as_deref()); - 16068
match persist_mcp_to_project_config(core.cwd(), &body.servers) { - 16069
Ok(_) => {} - 16070
Err(e) => { - 16071
return ( - 16072
StatusCode::INTERNAL_SERVER_ERROR, - 16073
Json(serde_json::json!({ "error": e })), - 16074
) - 16075
.into_response(); - 16076
} - 16077
} - 16078
let cfg = vak_config::load_with_trust(core.cwd(), core.project_config_trusted()) - 16079
.map(|config| config.mcp); - 16080
let Ok(cfg) = cfg else { - 16081
return StatusCode::INTERNAL_SERVER_ERROR.into_response(); - 16082
}; - 16083
core.apply_persisted_mcp_servers(cfg); - 16084
vak_core::security_events::record( - 16085
&core.sessions_home(), - 16086
vak_core::security_events::EventKind::ConfigChange, - 16087
"mcp_servers_updated", - 16088
&format!("count={}", body.servers.len()), - 16089
None, - 16090
); - 16091
state.hub.emit_config_changed( - 16092
"mcp_servers_updated", - 16093
&format!("count={}", body.servers.len()), - 16094
); - 16095
( - 16096
StatusCode::OK, - 16097
Json(serde_json::json!({ "saved": true, "count": body.servers.len() })), - 16098
) - 16099
.into_response() - 16100
} - 16101
- 16102
#[derive(serde::Deserialize)] - 16103
struct TreeQuery { - 16104
path: Option<String>, - 16105
limit: Option<usize>, - 16106
} - 16107
- 16108
/// Bounded recursive listing for @-mention autocomplete. Vendored/build - 16109
/// directories are skipped; results are cwd-relative and capped. - 16110
async fn fs_tree( - 16111
State(state): State<AppState>, - 16112
axum::extract::Query(q): axum::extract::Query<TreeQuery>, - 16113
) -> axum::response::Response { - 16114
use axum::response::IntoResponse; - 16115
use walkdir::WalkDir; - 16116
- 16117
let limit = q.limit.unwrap_or(400).min(2000); - 16118
let base = match confined_path(state.core.cwd(), q.path.as_deref().unwrap_or(".")) { - 16119
Some(p) => p, - 16120
None => return (StatusCode::FORBIDDEN, "path outside workspace").into_response(), - 16121
}; - 16122
const SKIP: &[&str] = &[ - 16123
".git", - 16124
"target", - 16125
"node_modules", - 16126
"dist", - 16127
"build", - 16128
".venv", - 16129
"venv", - 16130
"__pycache__", - 16131
".vak", - 16132
".next", - 16133
".cache", - 16134
"coverage", - 16135
]; - 16136
let mut files: Vec<String> = Vec::new(); - 16137
for entry in WalkDir::new(&base) - 16138
.max_depth(8) - 16139
.follow_links(false) - 16140
.into_iter() - 16141
.filter_entry(|e| { - 16142
e.file_name() - 16143
.to_str() - 16144
.map(|n| !SKIP.contains(&n) || e.depth() == 0) - 16145
.unwrap_or(true) - 16146
}) - 16147
{ - 16148
let Ok(entry) = entry else { continue }; - 16149
if !entry.file_type().is_file() { - 16150
continue; - 16151
} - 16152
// Strip against `base` (canonicalized): on macOS /var is a symlink - 16153
// to /private/var, so prefixes against raw cwd never match. - 16154
let Ok(rel) = entry.path().strip_prefix(&base) else { - 16155
continue; - 16156
}; - 16157
let mut text = rel.to_string_lossy().replace('\\', "/"); - 16158
if let Some(sub) = q - 16159
.path - 16160
.as_deref() - 16161
.map(str::trim) - 16162
.filter(|s| !s.is_empty() && *s != ".") - 16163
{ - 16164
text = format!("{}/{}", sub.trim_end_matches('/'), text); - 16165
} - 16166
files.push(text); - 16167
if files.len() >= limit { - 16168
break; - 16169
} - 16170
} - 16171
files.sort_unstable(); - 16172
( - 16173
StatusCode::OK, - 16174
Json(serde_json::json!({ "files": files, "truncated": files.len() >= limit })), - 16175
) - 16176
.into_response() - 16177
} - 16178
- 16179
#[derive(serde::Deserialize)] - 16180
struct SideBody { - 16181
question: String, - 16182
} - 16183
- 16184
/// `/btw`: ask a question using the session's context WITHOUT landing it on - 16185
/// the main chain. Mechanics: append the Q + run the turn as a sibling - 16186
/// branch (parent = current main tail), then restore the tail so future - 16187
/// main turns continue exactly where they were. The side entries stay in - 16188
/// the ledger — reconstructable, never deleted. - 16189
async fn side_chat( - 16190
State(state): State<AppState>, - 16191
Path(id): Path<String>, - 16192
Json(body): Json<SideBody>, - 16193
) -> axum::response::Response { - 16194
use axum::response::IntoResponse; - 16195
let Some(handle) = state.get(&id) else { - 16196
return StatusCode::NOT_FOUND.into_response(); - 16197
}; - 16198
let Some(mut taken) = handle - 16199
.session - 16200
.lock() - 16201
.unwrap_or_else(std::sync::PoisonError::into_inner) - 16202
.take() - 16203
else { - 16204
return StatusCode::CONFLICT.into_response(); // main run active - 16205
}; - 16206
if let Err(e) = state.core.provider() { - 16207
*handle - 16208
.session - 16209
.lock() - 16210
.unwrap_or_else(std::sync::PoisonError::into_inner) = Some(taken); - 16211
return provider_unavailable(e); - 16212
} - 16213
- 16214
let tail_main = taken.tail_id().cloned(); - 16215
if let Err(_e) = taken.append_message(vak_session::MessageRecord { - 16216
message: vak_llm::Message::user_text(body.question.clone()), - 16217
meta: None, - 16218
}) { - 16219
*handle - 16220
.session - 16221
.lock() - 16222
.unwrap_or_else(std::sync::PoisonError::into_inner) = Some(taken); - 16223
return StatusCode::INTERNAL_SERVER_ERROR.into_response(); - 16224
} - 16225
let side_tx = handle.side_events_tx.clone(); - 16226
let side_activity = Arc::new(Mutex::new(Vec::new())); - 16227
let approver: Arc<dyn Approver> = Arc::new(HttpApprover { - 16228
events_tx: handle.events_tx.clone(), - 16229
pending: handle.pending.clone(), - 16230
session_id: handle.id.clone(), - 16231
activity_buffer: side_activity.clone(), - 16232
answerable: true, - 16233
}); - 16234
let events = mpsc_to_broadcast(side_tx.clone()); - 16235
let cancel = handle - 16236
.side_cancel - 16237
.lock() - 16238
.unwrap_or_else(std::sync::PoisonError::into_inner) - 16239
.clone(); - 16240
// A side chat reads this session's context, so it runs under the same - 16241
// ceiling the session was created with. - 16242
let core = handle.core.clone(); - 16243
- 16244
tokio::spawn(async move { - 16245
// No steering on side chats by design: they are read-only Q&A over - 16246
// the session context, not a second control surface. - 16247
let outcome = core - 16248
.run_turn_with( - 16249
taken, - 16250
&body.question, - 16251
cancel, - 16252
Some(approver), - 16253
None, - 16254
None, - 16255
events, - 16256
) - 16257
.await; - 16258
let (summary, is_error) = match &outcome { - 16259
Ok((vak_agent::TurnOutcome::Completed { .. }, _)) => ("completed".to_string(), false), - 16260
Ok((vak_agent::TurnOutcome::Aborted { .. }, _)) => ("aborted".to_string(), false), - 16261
Ok((_, _)) => ("ended".to_string(), false), - 16262
Err(e) => (format!("error: {e}"), true), - 16263
}; - 16264
let turn_ok = matches!(&outcome, Ok((_, _))); - 16265
if let Ok((_, mut restored)) = outcome { - 16266
for activity in std::mem::take( - 16267
&mut *side_activity - 16268
.lock() - 16269
.unwrap_or_else(std::sync::PoisonError::into_inner), - 16270
) { - 16271
let _ = restored.append_activity(activity); - 16272
} - 16273
// Rewind the branch pointer to the main line: the side entries - 16274
// remain in the ledger as a sibling branch — reconstructable via - 16275
// their parent chain, invisible to derive_messages(). - 16276
if let Some(main_tail) = &tail_main { - 16277
let _ = restored.branch_at(main_tail); - 16278
} - 16279
*handle - 16280
.presentation - 16281
.lock() - 16282
.unwrap_or_else(std::sync::PoisonError::into_inner) = - 16283
crate::projection::snapshot(&id, &restored); - 16284
*handle - 16285
.session - 16286
.lock() - 16287
.unwrap_or_else(std::sync::PoisonError::into_inner) = Some(restored); - 16288
} - 16289
let _ = side_tx.send(AgentEvent::RunFinished { summary, is_error }); - 16290
// On a failed turn the taken log is gone with the Err — reopen the - 16291
// durable ledger so the session does not stay wedged as - 16292
// "run in progress" forever (found by the v0.6 deployment gate). - 16293
if !turn_ok && let Some(log) = reopen_ledger(&core, &id) { - 16294
*handle - 16295
.session - 16296
.lock() - 16297
.unwrap_or_else(std::sync::PoisonError::into_inner) = Some(log); - 16298
} - 16299
}); - 16300
- 16301
StatusCode::ACCEPTED.into_response() - 16302
} - 16303
- 16304
async fn side_cancel_run(State(state): State<AppState>, Path(id): Path<String>) -> StatusCode { - 16305
let Some(handle) = state.get(&id) else { - 16306
return StatusCode::NOT_FOUND; - 16307
}; - 16308
handle - 16309
.side_cancel - 16310
.lock() - 16311
.unwrap_or_else(std::sync::PoisonError::into_inner) - 16312
.cancel(); - 16313
let _ = handle.side_events_tx.send(AgentEvent::RunFinished { - 16314
summary: "cancelled by client".into(), - 16315
is_error: false, - 16316
}); - 16317
StatusCode::ACCEPTED - 16318
} - 16319
- 16320
#[derive(serde::Deserialize)] - 16321
struct BestBody { - 16322
prompt: String, - 16323
n: Option<usize>, - 16324
} - 16325
- 16326
/// Best-of-N: fan the same prompt across N isolated git worktrees, each with - 16327
/// its own session + event stream. Candidates are compared by diff; `keep` - 16328
/// merges a branch, `discard` drops it. Ledger-native: every run is a normal - 16329
/// session under the shared store. - 16330
async fn start_bestofn( - 16331
State(state): State<AppState>, - 16332
Path(id): Path<String>, - 16333
Json(body): Json<BestBody>, - 16334
) -> axum::response::Response { - 16335
use axum::response::IntoResponse; - 16336
- 16337
let Some(anchor) = state.get(&id) else { - 16338
return StatusCode::NOT_FOUND.into_response(); - 16339
}; - 16340
if anchor - 16341
.session - 16342
.lock() - 16343
.unwrap_or_else(std::sync::PoisonError::into_inner) - 16344
.is_none() - 16345
{ - 16346
return StatusCode::CONFLICT.into_response(); - 16347
} - 16348
let provider = match state.core.provider() { - 16349
Ok(p) => p, - 16350
Err(e) => return provider_unavailable(e), - 16351
}; - 16352
- 16353
let n = body.n.unwrap_or(2).clamp(1, 4); - 16354
let repo = state.core.cwd().clone(); - 16355
if !vak_core::worktree::is_git_repo(&repo) { - 16356
return StatusCode::CONFLICT.into_response(); - 16357
} - 16358
- 16359
// Create worktrees first; roll back everything on partial failure. - 16360
let mut created: Vec<(String, vak_core::worktree::Worktree)> = Vec::new(); - 16361
for i in 0..n { - 16362
// v7 shares its leading chars within one millisecond; disambiguate. - 16363
let rid = format!("{}-{i}", uuid::Uuid::now_v7().simple()); - 16364
match vak_core::worktree::create(&repo, &rid) { - 16365
Ok(wt) => created.push((rid, wt)), - 16366
Err(e) => { - 16367
for (_, wt) in &created { - 16368
let _ = vak_core::worktree::remove(&repo, wt); - 16369
} - 16370
return ( - 16371
StatusCode::INTERNAL_SERVER_ERROR, - 16372
Json(serde_json::json!({ "error": format!("worktree create failed: {e}") })), - 16373
) - 16374
.into_response(); - 16375
} - 16376
} - 16377
} - 16378
- 16379
let mut runs = Vec::new(); - 16380
for (_, wt) in &created { - 16381
match spawn_isolated_run( - 16382
&state, - 16383
provider.clone(), - 16384
wt, - 16385
&body.prompt, - 16386
None, - 16387
None, - 16388
None, - 16389
true, - 16390
) - 16391
.await - 16392
{ - 16393
Ok(child_id) => { - 16394
state - 16395
.best_runs - 16396
.lock() - 16397
.unwrap_or_else(std::sync::PoisonError::into_inner) - 16398
.insert( - 16399
child_id.clone(), - 16400
BestRunMeta { - 16401
repo: repo.clone(), - 16402
wt_path: wt.path.clone(), - 16403
branch: wt.branch.clone(), - 16404
}, - 16405
); - 16406
runs.push(serde_json::json!({ - 16407
"session_id": child_id, - 16408
"branch": wt.branch, - 16409
"path": wt.path, - 16410
})); - 16411
} - 16412
Err(e) => { - 16413
for (_, w) in &created { - 16414
let _ = vak_core::worktree::remove(&repo, w); - 16415
} - 16416
return ( - 16417
StatusCode::INTERNAL_SERVER_ERROR, - 16418
Json(serde_json::json!({ "error": e })), - 16419
) - 16420
.into_response(); - 16421
} - 16422
} - 16423
} - 16424
- 16425
(StatusCode::OK, Json(serde_json::json!({ "runs": runs }))).into_response() - 16426
} - 16427
- 16428
/// One isolated run inside `wt`: child Core + session + registered handle + - 16429
/// fired turn. Shared by best-of-N and the task scheduler. `model_pin` - 16430
/// (docs/design/29-personal-os.md P2) overrides the child's provider/model - 16431
/// so BOTH main dispatches and any receipts carry the pinned id only — a - 16432
/// pinned task never escalates to another model. - 16433
#[allow(clippy::too_many_arguments)] - 16434
async fn spawn_isolated_run( - 16435
state: &AppState, - 16436
provider: Arc<dyn Provider>, - 16437
wt: &vak_core::worktree::Worktree, - 16438
prompt: &str, - 16439
model_pin: Option<&str>, - 16440
agent_id: Option<&str>, - 16441
agent_revision: Option<u64>, - 16442
start_turn: bool, - 16443
) -> Result<String, String> { - 16444
let identity = if let Some(agent_id) = agent_id { - 16445
let profiles = agents::effective(&state.active_core())?; - 16446
let profile = profiles - 16447
.iter() - 16448
.find(|profile| profile.id == agent_id) - 16449
.ok_or_else(|| format!("Agent '{agent_id}' no longer exists"))?; - 16450
if !profile.is_admissible() { - 16451
return Err(format!("Agent '{agent_id}' is paused or archived")); - 16452
} - 16453
if let Some(expected) = agent_revision - 16454
&& expected != profile.revision - 16455
{ - 16456
return Err(format!( - 16457
"Agent '{agent_id}' changed from revision {expected} to {}", - 16458
profile.revision - 16459
)); - 16460
} - 16461
Some(profile.identity()) - 16462
} else { - 16463
None - 16464
}; - 16465
let child_core = vak_core::Core::new_with_trust(wt.path.clone(), true) - 16466
.map(|c| { - 16467
c.with_agent_identity(identity) - 16468
.with_surface(vak_core::Surface::Background) - 16469
// Unattended, and stamped BEFORE `start_session` composes and - 16470
// freezes the prompt. Stamping afterwards would be too late: - 16471
// the prompt would already have advertised a gated capability - 16472
// that this run can only ever be refused, which is the exact - 16473
// mismatch `vak_core::reach` exists to remove. `begin_turn` - 16474
// installs the matching approver. - 16475
.with_approver_answerable(false) - 16476
}) - 16477
.map_err(|e| format!("child core failed: {e}"))?; - 16478
child_core.set_provider_instance(provider); - 16479
// The shared root: the child resolves its own Agent's home beneath it, - 16480
// as every Core does. Seeding it with this Core's (already Agent-scoped) - 16481
// home nested one Agent's home inside another's. - 16482
child_core.set_sessions_home(state.core.shared_data_home()); - 16483
if let Some(pin) = model_pin.map(str::trim).filter(|p| !p.is_empty()) { - 16484
let (pin_provider, pin_model) = split_model_pin(pin, &child_core.effective_provider()); - 16485
child_core.set_route(pin_provider, pin_model); - 16486
} - 16487
- 16488
let child_log = child_core - 16489
.start_session() - 16490
.await - 16491
.map_err(|_| "child session failed to start".to_string())?; - 16492
let Some(child_header) = child_log.header() else { - 16493
return Err("child session has no header".to_string()); - 16494
}; - 16495
// The handle is the ledger's own id, so a run recorded on a task - 16496
// (`last_session_id`) opens from disk after a restart. - 16497
let child_id = child_header.session_id.clone(); - 16498
let handle = register_handle( - 16499
state, - 16500
child_id.clone(), - 16501
child_log, - 16502
wt.path.clone(), - 16503
child_core.clone(), - 16504
); - 16505
if start_turn { - 16506
begin_turn(&handle, &child_core, prompt, false); - 16507
} - 16508
Ok(child_id) - 16509
} - 16510
- 16511
/// Fire a single-turn agent run on a (usually fresh) session handle. - 16512
/// - 16513
/// `attended` says whether anyone is watching this run's event stream. It is - 16514
/// not cosmetic: a scheduled routine and a best-of-N leg both arrive here, - 16515
/// nobody is subscribed to either, and an approval gate raised on one used - 16516
/// to emit an SSE event into the void and then block the run until the - 16517
/// process restarted. An unattended run gets an approver that says so, and - 16518
/// `Core::with_approver` carries that fact into the prompt so the model is - 16519
/// never offered a capability whose gate can only ever be refused. - 16520
fn begin_turn(handle: &Arc<SessionHandle>, core: &Core, prompt: &str, attended: bool) { - 16521
let approver: Arc<dyn Approver> = Arc::new(HttpApprover { - 16522
events_tx: handle.events_tx.clone(), - 16523
pending: handle.pending.clone(), - 16524
session_id: handle.id.clone(), - 16525
activity_buffer: handle.activity_buffer.clone(), - 16526
answerable: attended, - 16527
}); - 16528
let core = core.clone().with_approver(approver.as_ref()); - 16529
let events = mpsc_to_broadcast(handle.events_tx.clone()); - 16530
let steering = Arc::new(SteeringQueues::new()); - 16531
let cancel = handle - 16532
.cancel - 16533
.lock() - 16534
.unwrap_or_else(std::sync::PoisonError::into_inner) - 16535
.clone(); - 16536
let Some(log) = handle - 16537
.session - 16538
.lock() - 16539
.unwrap_or_else(std::sync::PoisonError::into_inner) - 16540
.take() - 16541
else { - 16542
return; // busy — caller should have checked - 16543
}; - 16544
let turn_session_id = log - 16545
.header() - 16546
.map(|h| h.session_id.clone()) - 16547
.unwrap_or_default(); - 16548
let core = core.clone(); - 16549
let prompt = prompt.to_string(); - 16550
let h2 = handle.clone(); - 16551
tokio::spawn(async move { - 16552
let outcome = core - 16553
.run_turn_with( - 16554
log, - 16555
&prompt, - 16556
cancel, - 16557
Some(approver), - 16558
None, - 16559
Some(steering.clone()), - 16560
events, - 16561
) - 16562
.await; - 16563
*h2.cancel - 16564
.lock() - 16565
.unwrap_or_else(std::sync::PoisonError::into_inner) = CancellationToken::new(); - 16566
match outcome { - 16567
Ok((_, mut restored)) => { - 16568
for activity in std::mem::take( - 16569
&mut *h2 - 16570
.activity_buffer - 16571
.lock() - 16572
.unwrap_or_else(std::sync::PoisonError::into_inner), - 16573
) { - 16574
let _ = restored.append_activity(activity); - 16575
} - 16576
*h2.presentation - 16577
.lock() - 16578
.unwrap_or_else(std::sync::PoisonError::into_inner) = - 16579
crate::projection::snapshot(&turn_session_id, &restored); - 16580
*h2.session - 16581
.lock() - 16582
.unwrap_or_else(std::sync::PoisonError::into_inner) = Some(restored); - 16583
let _ = h2.events_tx.send(AgentEvent::RunFinished { - 16584
summary: "completed".into(), - 16585
is_error: false, - 16586
}); - 16587
} - 16588
Err(e) => { - 16589
// Same leak class as side chats: restore from the durable - 16590
// ledger so the handle is not wedged on "run in progress". - 16591
if let Some(log) = reopen_ledger(&core, &turn_session_id) { - 16592
*h2.session - 16593
.lock() - 16594
.unwrap_or_else(std::sync::PoisonError::into_inner) = Some(log); - 16595
} - 16596
let _ = h2.events_tx.send(AgentEvent::RunFinished { - 16597
summary: format!("error: {e}"), - 16598
is_error: true, - 16599
}); - 16600
} - 16601
} - 16602
drop(steering); - 16603
}); - 16604
} - 16605
- 16606
async fn keep_best_run( - 16607
State(state): State<AppState>, - 16608
Path(id): Path<String>, - 16609
) -> axum::response::Response { - 16610
use axum::response::IntoResponse; - 16611
let meta = state - 16612
.best_runs - 16613
.lock() - 16614
.unwrap_or_else(std::sync::PoisonError::into_inner) - 16615
.get(&id) - 16616
.cloned(); - 16617
let Some(meta) = meta else { - 16618
return StatusCode::NOT_FOUND.into_response(); - 16619
}; - 16620
// Merge into the user's checkout. Conflicts/dirty trees surface as errors. - 16621
let out = tokio::process::Command::new("git") - 16622
.args([ - 16623
"merge", - 16624
"--no-ff", - 16625
"-m", - 16626
&format!("best-of-n: merge {}", meta.branch), - 16627
]) - 16628
.arg(&meta.branch) - 16629
.current_dir(&meta.repo) - 16630
.output() - 16631
.await; - 16632
match out { - 16633
Ok(o) if o.status.success() => { - 16634
cleanup_worktree(&state, &id, &meta); - 16635
(StatusCode::OK, Json(serde_json::json!({"kept": id}))).into_response() - 16636
} - 16637
Ok(o) => { - 16638
// Abort any conflicted merge so the tree is not left dirty. - 16639
let _ = tokio::process::Command::new("git") - 16640
.args(["merge", "--abort"]) - 16641
.current_dir(&meta.repo) - 16642
.output() - 16643
.await; - 16644
( - 16645
StatusCode::CONFLICT, - 16646
Json(serde_json::json!({ - 16647
"error": "merge failed", - 16648
"stderr": String::from_utf8_lossy(&o.stderr), - 16649
})), - 16650
) - 16651
.into_response() - 16652
} - 16653
Err(e) => ( - 16654
StatusCode::INTERNAL_SERVER_ERROR, - 16655
Json(serde_json::json!({ "error": e.to_string() })), - 16656
) - 16657
.into_response(), - 16658
} - 16659
} - 16660
- 16661
async fn discard_best_run( - 16662
State(state): State<AppState>, - 16663
Path(id): Path<String>, - 16664
) -> axum::response::Response { - 16665
use axum::response::IntoResponse; - 16666
let meta = state - 16667
.best_runs - 16668
.lock() - 16669
.unwrap_or_else(std::sync::PoisonError::into_inner) - 16670
.get(&id) - 16671
.cloned(); - 16672
let Some(meta) = meta else { - 16673
return StatusCode::NOT_FOUND.into_response(); - 16674
}; - 16675
cleanup_worktree(&state, &id, &meta); - 16676
(StatusCode::OK, Json(serde_json::json!({"discarded": id}))).into_response() - 16677
} - 16678
- 16679
fn cleanup_worktree(state: &AppState, child_id: &str, meta: &BestRunMeta) { - 16680
let wt = vak_core::worktree::Worktree { - 16681
path: meta.wt_path.clone(), - 16682
branch: meta.branch.clone(), - 16683
}; - 16684
let _ = vak_core::worktree::remove(&meta.repo, &wt); - 16685
state - 16686
.best_runs - 16687
.lock() - 16688
.unwrap_or_else(std::sync::PoisonError::into_inner) - 16689
.remove(child_id); - 16690
} - 16691
- 16692
// ---- PR monitoring (gh-backed) --------------------------------------------- - 16693
- 16694
async fn gh_output(cwd: &std::path::Path, args: &[&str]) -> Result<String, String> { - 16695
let out = tokio::process::Command::new("gh") - 16696
.args(args) - 16697
.current_dir(cwd) - 16698
.output() - 16699
.await - 16700
.map_err(|e| format!("gh not available: {e}"))?; - 16701
if !out.status.success() { - 16702
return Err(String::from_utf8_lossy(&out.stderr).trim().to_string()); - 16703
} - 16704
Ok(String::from_utf8_lossy(&out.stdout).into_owned()) - 16705
} - 16706
- 16707
/// One-shot PR status for the session workspace: current branch, the PR - 16708
/// attached to it (if any), and its check rollup. Tooling absence surfaces - 16709
/// as `{error}` — never a hang, never a panic. - 16710
async fn session_pr( - 16711
State(state): State<AppState>, - 16712
Path(id): Path<String>, - 16713
) -> Json<serde_json::Value> { - 16714
let Some(handle) = state.get(&id) else { - 16715
return Json(serde_json::json!({ "error": "unknown session" })); - 16716
}; - 16717
let cwd = handle.cwd.clone(); - 16718
let Some(branch) = git_output(&cwd, &["rev-parse", "--abbrev-ref", "HEAD"]).await else { - 16719
return Json(serde_json::json!({ "error": "not a git repository" })); - 16720
}; - 16721
let branch = branch.trim().to_string(); - 16722
if branch.is_empty() || branch == "HEAD" { - 16723
return Json(serde_json::json!({ "error": "detached HEAD" })); - 16724
} - 16725
- 16726
let raw = match gh_output( - 16727
&cwd, - 16728
&[ - 16729
"pr", - 16730
"view", - 16731
&branch, - 16732
"--json", - 16733
"number,title,url,state,mergeable,statusCheckRollup", - 16734
], - 16735
) - 16736
.await - 16737
{ - 16738
Ok(r) => r, - 16739
Err(e) => { - 16740
// No PR for this branch vs no gh at all — distinguish for UX. - 16741
let msg = e.to_lowercase(); - 16742
let kind = if msg.contains("no pull requests") || msg.contains("no merges requested") { - 16743
"no_pr" - 16744
} else { - 16745
"gh_unavailable" - 16746
}; - 16747
return Json(serde_json::json!({ - 16748
"branch": branch, - 16749
"pr": serde_json::Value::Null, - 16750
"reason": kind, - 16751
"error": e, - 16752
})); - 16753
} - 16754
}; - 16755
- 16756
let Ok(view) = serde_json::from_str::<serde_json::Value>(&raw) else { - 16757
return Json(serde_json::json!({ "error": "unparsable gh output", "branch": branch })); - 16758
}; - 16759
let mut pass = 0u32; - 16760
let mut fail = 0u32; - 16761
let mut pending = 0u32; - 16762
if let Some(rollup) = view["statusCheckRollup"].as_array() { - 16763
for c in rollup { - 16764
let status = c["status"].as_str().unwrap_or(""); - 16765
let conclusion = c["conclusion"].as_str().unwrap_or(""); - 16766
match (status, conclusion) { - 16767
(_, "SUCCESS") => pass += 1, - 16768
(_, "FAILURE") | (_, "CANCELLED") | (_, "TIMED_OUT") => fail += 1, - 16769
("COMPLETED", _) => {} - 16770
_ => pending += 1, - 16771
} - 16772
} - 16773
} - 16774
Json(serde_json::json!({ - 16775
"branch": branch, - 16776
"pr": { - 16777
"number": view["number"], - 16778
"title": view["title"], - 16779
"url": view["url"], - 16780
"state": view["state"], - 16781
"mergeable": view["mergeable"], - 16782
}, - 16783
"checks": view["statusCheckRollup"], - 16784
"summary": { "pass": pass, "fail": fail, "pending": pending }, - 16785
})) - 16786
} - 16787
- 16788
#[derive(serde::Deserialize)] - 16789
struct PrMergeBody { - 16790
number: u64, - 16791
#[serde(default)] - 16792
method: Option<String>, - 16793
} - 16794
- 16795
/// Merge an open PR via gh. `--auto` honors branch protection: gh merges - 16796
/// when checks go green. - 16797
async fn pr_merge( - 16798
State(state): State<AppState>,
Indexing the workspace…
Vakyartha documentation is discovering safe artifacts, anchors, and source references.