- 15464
| "stop" - 15465
) { - 15466
return Err(format!("unknown hook event '{}'", hook.event)); - 15467
} - 15468
if hook.enabled && hook.command.trim().is_empty() { - 15469
return Err("enabled hooks need a command".into()); - 15470
} - 15471
let timeout_ms = hook.timeout_ms.unwrap_or(vak_hooks::DEFAULT_TIMEOUT_MS); - 15472
if timeout_ms == 0 { - 15473
return Err("hook timeout must be greater than zero".into()); - 15474
} - 15475
let failure_mode = hook.failure_mode.as_deref().unwrap_or("open").trim(); - 15476
if !matches!(failure_mode, "open" | "closed") { - 15477
return Err(format!("unknown hook failure mode '{failure_mode}'")); - 15478
} - 15479
Ok(vak_config::HookConfig { - 15480
event: hook.event.clone(), - 15481
matcher: hook - 15482
.matcher - 15483
.clone() - 15484
.filter(|matcher| !matcher.trim().is_empty()), - 15485
command: hook.command.trim().to_string(), - 15486
timeout_ms: Some(timeout_ms), - 15487
enabled: hook.enabled, - 15488
failure_mode: Some(failure_mode.to_string()), - 15489
}) - 15490
}) - 15491
.collect() - 15492
} - 15493
- 15494
async fn put_hooks( - 15495
State(state): State<AppState>, - 15496
Json(body): Json<HooksPutBody>, - 15497
) -> axum::response::Response { - 15498
use axum::response::IntoResponse; - 15499
let core = scoped_core!(&state, None, body.agent.as_deref()); - 15500
let hooks = match validated_hook_configs(&body.hooks) { - 15501
Ok(hooks) => hooks, - 15502
Err(message) => { - 15503
return ( - 15504
StatusCode::BAD_REQUEST, - 15505
Json(serde_json::json!({ "error": message })), - 15506
) - 15507
.into_response(); - 15508
} - 15509
}; - 15510
// A disabled hook is kept in config, not dropped — round-tripping the - 15511
// toggle used to delete the definition outright (there was nowhere in - 15512
// `[[hooks]]` to record "off"), which is not what a checkbox should do. - 15513
match vak_config::persist_hooks(&vak_config::project_path(core.cwd()), &hooks) { - 15514
Ok(()) => {} - 15515
Err(vak_config::ConfigError::Parse { .. }) => { - 15516
return ( - 15517
StatusCode::BAD_REQUEST, - 15518
Json(serde_json::json!({ "error": "workspace config is invalid" })), - 15519
) - 15520
.into_response(); - 15521
} - 15522
Err(error) => { - 15523
return ( - 15524
StatusCode::INTERNAL_SERVER_ERROR, - 15525
Json(serde_json::json!({ "error": format!("write config: {error}") })), - 15526
) - 15527
.into_response(); - 15528
} - 15529
} - 15530
// Disabled hooks are still handed to Core — `build_hooks_from` is what - 15531
// skips them when it builds the live `HookDef` list — so the effective - 15532
// set stays correct without this endpoint duplicating that filter. - 15533
core.apply_persisted_hooks(hooks); - 15534
let enabled_count = body.hooks.iter().filter(|h| h.enabled).count(); - 15535
vak_core::security_events::record( - 15536
&core.sessions_home(), - 15537
vak_core::security_events::EventKind::ConfigChange, - 15538
"hooks_updated", - 15539
&format!("enabled={enabled_count} total={}", body.hooks.len()), - 15540
None, - 15541
); - 15542
state.hub.emit_config_changed( - 15543
"hooks_updated", - 15544
&format!("enabled={enabled_count} total={}", body.hooks.len()), - 15545
); - 15546
( - 15547
StatusCode::OK, - 15548
Json(serde_json::json!({ "saved": true, "count": enabled_count })), - 15549
) - 15550
.into_response() - 15551
} - 15552
- 15553
#[derive(serde::Deserialize, Clone)] - 15554
struct McpServerInput { - 15555
command: String, - 15556
#[serde(default)] - 15557
args: Vec<String>, - 15558
#[serde(default)] - 15559
env: std::collections::BTreeMap<String, String>, - 15560
#[serde(default)] - 15561
network: bool, - 15562
} - 15563
- 15564
#[derive(serde::Deserialize)] - 15565
struct McpPutBody { - 15566
servers: std::collections::BTreeMap<String, McpServerInput>, - 15567
#[serde(default)] - 15568
agent: Option<String>, - 15569
} - 15570
- 15571
fn valid_server_name(name: &str) -> bool { - 15572
!name.is_empty() - 15573
&& name.len() <= 64 - 15574
&& name - 15575
.chars() - 15576
.all(|c| c.is_ascii_alphanumeric() || matches!(c, '-' | '_' | '.')) - 15577
} - 15578
- 15579
fn persist_mcp_to_project_config( - 15580
cwd: &std::path::Path, - 15581
servers: &std::collections::BTreeMap<String, McpServerInput>, - 15582
) -> Result<std::path::PathBuf, String> { - 15583
let path = vak_config::project_path(cwd); - 15584
let config = mcp_config_from_input(servers); - 15585
vak_config::persist_mcp_servers(&path, &config.servers).map_err(|error| error.to_string())?; - 15586
Ok(path) - 15587
} - 15588
- 15589
fn mcp_config_from_input( - 15590
servers: &std::collections::BTreeMap<String, McpServerInput>, - 15591
) -> vak_config::McpConfig { - 15592
vak_config::McpConfig { - 15593
servers: servers - 15594
.iter() - 15595
.map(|(name, server)| { - 15596
( - 15597
name.clone(), - 15598
vak_config::McpServerConfig { - 15599
command: server.command.trim().to_string(), - 15600
args: server.args.clone(), - 15601
env: server.env.clone(), - 15602
network: server.network, - 15603
// Undeclared: the per-turn slice never narrows a - 15604
// capability that has not classified itself, so a - 15605
// server configured through the API stays reachable - 15606
// without the operator having to know about domains. - 15607
serves: Vec::new(), - 15608
}, - 15609
) - 15610
}) - 15611
.collect(), - 15612
} - 15613
} - 15614
- 15615
fn read_mcp_config(path: &std::path::Path) -> Result<vak_config::McpConfig, String> { - 15616
if !path.is_file() { - 15617
return Ok(vak_config::McpConfig::default()); - 15618
} - 15619
let raw = std::fs::read_to_string(path).map_err(|error| error.to_string())?; - 15620
toml::from_str::<vak_config::FileConfig>(&raw) - 15621
.map(|config| config.mcp) - 15622
.map_err(|error| error.to_string()) - 15623
} - 15624
- 15625
/// Which layer a setting belongs to. - 15626
/// - 15627
/// One word for one concept: the config section is `workspace_roots`, the - 15628
/// path helper is `default_workspace`, the API is `/workspaces` — and this - 15629
/// said "project". Two names for the same thing is how a UI ends up - 15630
/// labelling one panel "This project" and its own store `workspace`. - 15631
#[derive(Clone, Copy, serde::Deserialize)] - 15632
#[serde(rename_all = "lowercase")] - 15633
enum ConfigScope { - 15634
User, - 15635
Workspace, - 15636
} - 15637
- 15638
impl ConfigScope { - 15639
fn is_workspace(self) -> bool { - 15640
matches!(self, Self::Workspace) - 15641
} - 15642
- 15643
fn label(self) -> &'static str { - 15644
match self { - 15645
Self::User => "user", - 15646
Self::Workspace => "workspace", - 15647
} - 15648
} - 15649
- 15650
/// Root whose `.vak/prompts` directory this scope edits. - 15651
fn prompt_root(self, core: &vak_core::Core) -> std::path::PathBuf { - 15652
match self { - 15653
Self::User => vak_config::paths::default_workspace(), - 15654
Self::Workspace => core.cwd().clone(), - 15655
} - 15656
} - 15657
- 15658
fn config_path(self, core: &vak_core::Core) -> Result<std::path::PathBuf, String> { - 15659
match self { - 15660
Self::User => vak_config::global_path().ok_or_else(|| "user home unavailable".into()), - 15661
Self::Workspace => Ok(vak_config::project_path(core.cwd())), - 15662
} - 15663
} - 15664
} - 15665
- 15666
#[derive(serde::Deserialize)] - 15667
struct ScopeQuery { - 15668
scope: ConfigScope, - 15669
/// Which user-facing Agent's isolated workspace this config layer is - 15670
/// rooted under. Absent means the built-in "vak" Agent, resolved - 15671
/// through `resolve_scoped_core` exactly like the memory/proposal - 15672
/// endpoints (see commit 15c9c256). - 15673
#[serde(default)] - 15674
agent: Option<String>, - 15675
} - 15676
- 15677
/// A scope query where omitting the parameter is legal and means "workspace". - 15678
/// Kept separate from [`ScopeQuery`] so the endpoints that genuinely - 15679
/// require an explicit scope keep rejecting a request without one. - 15680
#[derive(serde::Deserialize)] - 15681
struct OptionalScopeQuery { - 15682
#[serde(default)] - 15683
scope: Option<ConfigScope>, - 15684
#[serde(default)] - 15685
agent: Option<String>, - 15686
} - 15687
- 15688
#[derive(Clone, Copy)] - 15689
struct IntegrationCatalogEntry { - 15690
id: &'static str, - 15691
label: &'static str, - 15692
description: &'static str, - 15693
command: &'static str, - 15694
args: &'static [&'static str], - 15695
env_var: Option<&'static str>, - 15696
key_required: bool, - 15697
documentation_url: &'static str, - 15698
} - 15699
- 15700
/// The curated integrations, **alphabetically**. - 15701
/// - 15702
/// Order is not cosmetic here. Whatever sits first reads as the default, - 15703
/// and this list led with Tavily — which is how one connector came to look - 15704
/// like the real one and the others like extras (doc 46 D5). They are - 15705
/// peers: same shape, same status projection, same scoped read/write path. - 15706
const INTEGRATION_CATALOG: &[IntegrationCatalogEntry] = &[ - 15707
IntegrationCatalogEntry { - 15708
id: "context7", - 15709
label: "Context7", - 15710
description: "Current library documentation and version-specific code examples.", - 15711
command: "npx", - 15712
args: &["-y", "@upstash/context7-mcp@latest"], - 15713
env_var: Some("CONTEXT7_API_KEY"), - 15714
key_required: false, - 15715
documentation_url: "https://github.com/upstash/context7", - 15716
}, - 15717
IntegrationCatalogEntry { - 15718
id: "exa", - 15719
label: "Exa", - 15720
description: "Web, code, company, and research search with page retrieval.", - 15721
command: "npx", - 15722
args: &["-y", "exa-mcp-server"], - 15723
env_var: Some("EXA_API_KEY"), - 15724
key_required: true, - 15725
documentation_url: "https://github.com/exa-labs/exa-mcp-server", - 15726
}, - 15727
IntegrationCatalogEntry { - 15728
id: "firecrawl", - 15729
label: "Firecrawl", - 15730
description: "Search, scrape, crawl, extract, and operate cloud browser sessions.", - 15731
command: "npx", - 15732
args: &["-y", "firecrawl-mcp"], - 15733
env_var: Some("FIRECRAWL_API_KEY"), - 15734
key_required: true, - 15735
documentation_url: "https://github.com/firecrawl/firecrawl-mcp-server", - 15736
}, - 15737
IntegrationCatalogEntry { - 15738
id: "tavily", - 15739
label: "Tavily", - 15740
description: "Real-time web search, extraction, site maps, and crawling.", - 15741
command: "npx", - 15742
args: &["-y", "tavily-mcp@latest"], - 15743
env_var: Some("TAVILY_API_KEY"), - 15744
key_required: true, - 15745
documentation_url: "https://github.com/tavily-ai/tavily-mcp", - 15746
}, - 15747
]; - 15748
- 15749
fn catalog_entry(id: &str) -> Option<IntegrationCatalogEntry> { - 15750
INTEGRATION_CATALOG - 15751
.iter() - 15752
.copied() - 15753
.find(|entry| entry.id == id) - 15754
} - 15755
- 15756
fn catalog_server(entry: IntegrationCatalogEntry) -> vak_config::McpServerConfig { - 15757
vak_config::McpServerConfig { - 15758
command: entry.command.into(), - 15759
args: entry.args.iter().map(|arg| (*arg).to_string()).collect(), - 15760
env: entry - 15761
.env_var - 15762
.map(|name| { - 15763
[(name.to_string(), format!("${{{name}}}"))] - 15764
.into_iter() - 15765
.collect() - 15766
}) - 15767
.unwrap_or_default(), - 15768
network: true, - 15769
// Catalog entries stay undeclared for the same reason: a server the - 15770
// operator just installed must be reachable immediately, and - 15771
// declaring domains only ever narrows. - 15772
serves: Vec::new(), - 15773
} - 15774
} - 15775
- 15776
fn integration_status( - 15777
core: &vak_core::Core, - 15778
scope: ConfigScope, - 15779
entry: IntegrationCatalogEntry, - 15780
) -> Result<serde_json::Value, String> { - 15781
let selected = read_mcp_config(&scope.config_path(core)?)?; - 15782
let user = match vak_config::global_path() { - 15783
Some(path) => read_mcp_config(&path)?, - 15784
None => vak_config::McpConfig::default(), - 15785
}; - 15786
let configured_here = selected.servers.contains_key(entry.id); - 15787
let inherited = scope.is_workspace() && !configured_here && user.servers.contains_key(entry.id); - 15788
let effective = core.effective_mcp().servers.contains_key(entry.id); - 15789
let key_here = entry - 15790
.env_var - 15791
.is_some_and(|name| core.mcp_secret_at_scope(name, scope.is_workspace())); - 15792
let key_inherited = scope.is_workspace() - 15793
&& !key_here - 15794
&& entry - 15795
.env_var - 15796
.is_some_and(|name| core.mcp_secret(name).is_some()); - 15797
Ok(serde_json::json!({ - 15798
"id": entry.id, - 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
Indexing the workspace…
Vakyartha documentation is discovering safe artifacts, anchors, and source references.