- 2447
vak_core::security_events::record( - 2448
&home, - 2449
vak_core::security_events::EventKind::ConfigChange, - 2450
"memory_append", - 2451
&format!("scope={scope_str} note_id={}", note.id), - 2452
None, - 2453
); - 2454
(StatusCode::CREATED, Json(note_payload(¬e, scope_str))).into_response() - 2455
} - 2456
Err(e) => ( - 2457
StatusCode::BAD_REQUEST, - 2458
Json(serde_json::json!({ "error": e })), - 2459
) - 2460
.into_response(), - 2461
} - 2462
} - 2463
- 2464
fn memory_store_path(core: &vak_core::Core, scope: MemoryScope) -> PathBuf { - 2465
let home = core.sessions_home(); - 2466
match scope { - 2467
MemoryScope::Workspace => home - 2468
.join("memory") - 2469
.join(vak_core::memory::hash_cwd(core.cwd())) - 2470
.join("MEMORY.md"), - 2471
MemoryScope::Profile => vak_core::memory::profile_path(&home), - 2472
} - 2473
} - 2474
- 2475
#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Deserialize, Default)] - 2476
#[serde(rename_all = "lowercase")] - 2477
enum MemoryScope { - 2478
#[default] - 2479
Workspace, - 2480
Profile, - 2481
} - 2482
- 2483
async fn forget_memory_note( - 2484
State(state): State<AppState>, - 2485
Path(note_id): Path<String>, - 2486
axum::extract::Query(q): axum::extract::Query<MemoryScopeQuery>, - 2487
) -> axum::response::Response { - 2488
use axum::response::IntoResponse; - 2489
let core = scoped_core!(&state, None, q.agent.as_deref()); - 2490
let path = memory_store_path(&core, q.scope.unwrap_or_default()); - 2491
let result = vak_core::memory::forget_note(&path, ¬e_id); - 2492
// A note written before this Agent's data moved to its own sessions_home - 2493
// subfolder still lives at the legacy, un-scoped path — fall back to it - 2494
// the same way `promote_proposal`/`reject_proposal` already do, so an - 2495
// old note surfaced by `list_memory`'s merged view can still be forgotten. - 2496
let legacy_path = memory_store_path(&state.core, q.scope.unwrap_or_default()); - 2497
let result = match result { - 2498
Err(_) if legacy_path != path => vak_core::memory::forget_note(&legacy_path, ¬e_id), - 2499
other => other, - 2500
}; - 2501
match result { - 2502
Ok(bytes) => { - 2503
vak_core::security_events::record( - 2504
&core.sessions_home(), - 2505
vak_core::security_events::EventKind::ConfigChange, - 2506
"memory_forget", - 2507
&format!("scope={:?} note_id={note_id}", q.scope.unwrap_or_default()), - 2508
None, - 2509
); - 2510
( - 2511
StatusCode::OK, - 2512
Json(serde_json::json!({ "forgotten": note_id, "bytes": bytes })), - 2513
) - 2514
.into_response() - 2515
} - 2516
Err(e) => ( - 2517
StatusCode::NOT_FOUND, - 2518
Json(serde_json::json!({ "error": e })), - 2519
) - 2520
.into_response(), - 2521
} - 2522
} - 2523
- 2524
#[derive(serde::Deserialize)] - 2525
struct MemoryAmendBody { - 2526
text: String, - 2527
#[serde(default)] - 2528
scope: Option<MemoryScope>, - 2529
/// See `AgentScopeQuery`. - 2530
#[serde(default)] - 2531
agent: Option<String>, - 2532
} - 2533
- 2534
#[derive(serde::Deserialize, Default)] - 2535
struct MemoryScopeQuery { - 2536
#[serde(default)] - 2537
scope: Option<MemoryScope>, - 2538
/// See `AgentScopeQuery`. - 2539
#[serde(default)] - 2540
agent: Option<String>, - 2541
} - 2542
- 2543
async fn amend_memory_note( - 2544
State(state): State<AppState>, - 2545
Path(note_id): Path<String>, - 2546
Json(body): Json<MemoryAmendBody>, - 2547
) -> axum::response::Response { - 2548
use axum::response::IntoResponse; - 2549
if body.text.trim().is_empty() { - 2550
return ( - 2551
StatusCode::BAD_REQUEST, - 2552
Json(serde_json::json!({ "error": "note must not be empty" })), - 2553
) - 2554
.into_response(); - 2555
} - 2556
let core = scoped_core!(&state, None, body.agent.as_deref()); - 2557
let path = memory_store_path(&core, body.scope.unwrap_or_default()); - 2558
let result = vak_core::memory::amend_note(&path, ¬e_id, &body.text); - 2559
// See the identical fallback in `forget_memory_note`. - 2560
let legacy_path = memory_store_path(&state.core, body.scope.unwrap_or_default()); - 2561
let result = match result { - 2562
Err(_) if legacy_path != path => { - 2563
vak_core::memory::amend_note(&legacy_path, ¬e_id, &body.text) - 2564
} - 2565
other => other, - 2566
}; - 2567
match result { - 2568
Ok(()) => { - 2569
vak_core::security_events::record( - 2570
&core.sessions_home(), - 2571
vak_core::security_events::EventKind::ConfigChange, - 2572
"memory_amend", - 2573
&format!( - 2574
"scope={:?} note_id={note_id}", - 2575
body.scope.unwrap_or_default() - 2576
), - 2577
None, - 2578
); - 2579
( - 2580
StatusCode::OK, - 2581
Json(serde_json::json!({ "amended": note_id })), - 2582
) - 2583
.into_response() - 2584
} - 2585
Err(e) => ( - 2586
StatusCode::NOT_FOUND, - 2587
Json(serde_json::json!({ "error": e })), - 2588
) - 2589
.into_response(), - 2590
} - 2591
} - 2592
- 2593
fn proposals_payload(core: &Core) -> Vec<serde_json::Value> { - 2594
let mut proposals = vak_core::learning::list_proposals(&core.sessions_home(), core.cwd()); - 2595
if proposals.is_empty() && core.shared_data_home() != core.sessions_home() { - 2596
proposals = vak_core::learning::list_proposals(&core.shared_data_home(), core.cwd()); - 2597
} - 2598
proposals - 2599
.iter() - 2600
.map(|p| { - 2601
serde_json::json!({ - 2602
"id": p.id, - 2603
"name": p.name, - 2604
"description": p.description, - 2605
}) - 2606
}) - 2607
.collect() - 2608
} - 2609
- 2610
async fn list_proposals_route( - 2611
State(state): State<AppState>, - 2612
axum::extract::Query(q): axum::extract::Query<AgentScopeQuery>, - 2613
) -> axum::response::Response { - 2614
use axum::response::IntoResponse; - 2615
let core = scoped_core!(&state, None, q.agent.as_deref()); - 2616
Json(serde_json::json!({ "proposals": proposals_payload(&core) })).into_response() - 2617
} - 2618
- 2619
async fn promote_proposal( - 2620
State(state): State<AppState>, - 2621
Path(id): Path<String>, - 2622
axum::extract::Query(q): axum::extract::Query<AgentScopeQuery>, - 2623
) -> axum::response::Response { - 2624
use axum::response::IntoResponse; - 2625
let core = scoped_core!(&state, None, q.agent.as_deref()); - 2626
let res = vak_core::learning::promote(&core.sessions_home(), core.cwd(), &id); - 2627
let res = match res { - 2628
Err(_) if state.core.shared_data_home() != core.sessions_home() => { - 2629
vak_core::learning::promote(&state.core.shared_data_home(), core.cwd(), &id) - 2630
} - 2631
other => other, - 2632
}; - 2633
match res { - 2634
Ok(name) => ( - 2635
StatusCode::OK, - 2636
Json(serde_json::json!({ "promoted": name })), - 2637
) - 2638
.into_response(), - 2639
Err(e) => ( - 2640
StatusCode::NOT_FOUND, - 2641
Json(serde_json::json!({ "error": e })), - 2642
) - 2643
.into_response(), - 2644
} - 2645
} - 2646
- 2647
async fn reject_proposal( - 2648
State(state): State<AppState>, - 2649
Path(id): Path<String>, - 2650
axum::extract::Query(q): axum::extract::Query<AgentScopeQuery>, - 2651
) -> axum::response::Response { - 2652
use axum::response::IntoResponse; - 2653
let core = scoped_core!(&state, None, q.agent.as_deref()); - 2654
let res = vak_core::learning::reject(&core.sessions_home(), core.cwd(), &id); - 2655
let res = match res { - 2656
Err(_) if state.core.shared_data_home() != core.sessions_home() => { - 2657
vak_core::learning::reject(&state.core.shared_data_home(), core.cwd(), &id) - 2658
} - 2659
other => other, - 2660
}; - 2661
match res { - 2662
Ok(()) => (StatusCode::OK, Json(serde_json::json!({ "rejected": id }))).into_response(), - 2663
Err(e) => ( - 2664
StatusCode::NOT_FOUND, - 2665
Json(serde_json::json!({ "error": e })), - 2666
) - 2667
.into_response(), - 2668
} - 2669
} - 2670
- 2671
#[derive(serde::Deserialize)] - 2672
struct SearchQuery { - 2673
q: String, - 2674
#[serde(default)] - 2675
limit: Option<usize>, - 2676
/// Session id whose (already-in-context) content should be skipped. - 2677
#[serde(default)] - 2678
exclude: Option<String>, - 2679
/// Cross-project recall: search every project's ledgers under the - 2680
/// sessions home (docs/design/29-personal-os.md P1), not just this cwd. - 2681
#[serde(default)] - 2682
all: bool, - 2683
/// The Agent whose ledgers and memory are searched (default `vak`). - 2684
/// Each Agent's memory is private (AGENTS.md invariant 37), so a search - 2685
/// resolves one Agent the way `/memory` does and reads only its home. - 2686
#[serde(default)] - 2687
agent: Option<String>, - 2688
} - 2689
- 2690
async fn search_sessions( - 2691
State(state): State<AppState>, - 2692
axum::extract::Query(q): axum::extract::Query<SearchQuery>, - 2693
) -> axum::response::Response { - 2694
use axum::response::IntoResponse; - 2695
let core = scoped_core!(&state, None, q.agent.as_deref()); - 2696
let home = core.sessions_home(); - 2697
let cwd = core.cwd().clone(); - 2698
let query = q.q.clone(); - 2699
let limit = q.limit.unwrap_or(vak_session::DEFAULT_LIMIT); - 2700
let excluded = - 2701
vak_core::trash::search_exclusions(&core.shared_data_home(), q.exclude.as_deref()); - 2702
let all = q.all; - 2703
let mut extras = Vec::new(); - 2704
let mut workspace_notes = vak_core::memory::list_notes(&home, &cwd); - 2705
if all { - 2706
let root = home.join("memory"); - 2707
if let Ok(entries) = std::fs::read_dir(&root) { - 2708
workspace_notes.clear(); - 2709
for entry in entries.flatten() { - 2710
let path = entry.path().join("MEMORY.md"); - 2711
if let Ok(raw) = std::fs::read_to_string(path) { - 2712
workspace_notes.extend(vak_core::memory::parse_blocks(&raw)); - 2713
} - 2714
} - 2715
} - 2716
} - 2717
for note in workspace_notes { - 2718
let id = if note.tag.is_empty() { - 2719
note.id.clone() - 2720
} else { - 2721
note.tag.clone() - 2722
}; - 2723
extras.push(vak_session::ExternalDoc { - 2724
id, - 2725
text: format!("[{}] {}", note.kind, note.text), - 2726
ts: Some(note.ts), - 2727
role: Some("memory".into()), - 2728
}); - 2729
} - 2730
for note in vak_core::memory::list_profile_notes(&home) { - 2731
let id = format!( - 2732
"profile/{}", - 2733
if note.tag.is_empty() { - 2734
note.id.clone() - 2735
} else { - 2736
note.tag.clone() - 2737
} - 2738
); - 2739
extras.push(vak_session::ExternalDoc { - 2740
id, - 2741
text: format!("[{}] {}", note.kind, note.text), - 2742
ts: Some(note.ts), - 2743
role: Some("profile".into()), - 2744
}); - 2745
} - 2746
match tokio::task::spawn_blocking(move || { - 2747
// Both hit shapes are Serialize; the workspace path keeps its flat - 2748
// SessionHit wire shape, cross-project adds the project_hash wrapper. - 2749
let searched = if all { - 2750
vak_session::search_all_extended(&home, &query, limit, &excluded, &extras) - 2751
.map(|hits| serde_json::to_value(&hits).map_err(|e| e.to_string())) - 2752
} else { - 2753
vak_session::search_extended(&home, &cwd, &query, limit, &excluded, &extras) - 2754
.map(|hits| serde_json::to_value(&hits).map_err(|e| e.to_string())) - 2755
}; - 2756
match searched { - 2757
Ok(inner) => inner, - 2758
Err(e) => Err(e.to_string()), - 2759
} - 2760
}) - 2761
.await - 2762
{ - 2763
Ok(Ok(hits)) => Json(serde_json::json!({ "all": all, "hits": hits })).into_response(), - 2764
Ok(Err(e)) => ( - 2765
StatusCode::INTERNAL_SERVER_ERROR, - 2766
Json(serde_json::json!({ "error": e })), - 2767
) - 2768
.into_response(), - 2769
Err(e) => ( - 2770
StatusCode::INTERNAL_SERVER_ERROR, - 2771
Json(serde_json::json!({ "error": e.to_string() })), - 2772
) - 2773
.into_response(), - 2774
} - 2775
} - 2776
- 2777
/// The bearer token lives for the life of the process; embedders (desktop - 2778
/// shell, tests) need it to hand to their webview, so build the secured - 2779
/// stack here instead of inside `serve()`. - 2780
pub fn secured_router(core: Core) -> (Router, String) { - 2781
secured_router_with(core, false) - 2782
} - 2783
- 2784
/// Same stack with a CLI-level gateway override (`serve --gateway`). - 2785
/// - 2786
/// Token selection: when `VAK_GATEWAY_TOKEN` is set in the - 2787
/// environment, it is used verbatim so service-managed bridges and other - 2788
/// long-lived clients can survive process restarts. Otherwise a fresh - 2789
/// per-process token is minted as before. The variable is never logged. - 2790
pub fn secured_router_with(core: Core, force_gateway: bool) -> (Router, String) { - 2791
secured_router_with_port(core, force_gateway, vak_ops::OpsConfig::detect().port) - 2792
} - 2793
- 2794
/// Same secured stack with the actual listener port carried into operational - 2795
/// probes. `serve_with` uses this so a non-default `--port` cannot make the - 2796
/// console probe a different process. - 2797
pub fn secured_router_with_port(core: Core, force_gateway: bool, port: u16) -> (Router, String) { - 2798
// Tauri can use either its custom scheme or the loopback-style origin, - 2799
// depending on the platform and WebView runtime, plus vite dev servers. - 2800
let origins = [ - 2801
"tauri://localhost", - 2802
"http://tauri.localhost", - 2803
"https://tauri.localhost", - 2804
"http://localhost:1420", - 2805
"http://127.0.0.1:1420", - 2806
"http://localhost:5173", - 2807
"http://127.0.0.1:5173", - 2808
] - 2809
.into_iter() - 2810
.filter_map(|o| o.parse::<axum::http::HeaderValue>().ok()) - 2811
.collect::<Vec<_>>(); - 2812
let cors = tower_http::cors::CorsLayer::new() - 2813
.allow_origin(origins) - 2814
// Must cover every method the router exposes: PATCH (/config, - 2815
// /sessions/:id/config) and DELETE are preflighted, so omitting them - 2816
// makes the browser reject the request before it is ever sent. - 2817
.allow_methods([ - 2818
axum::http::Method::GET, - 2819
axum::http::Method::POST, - 2820
axum::http::Method::PUT, - 2821
axum::http::Method::PATCH, - 2822
axum::http::Method::DELETE, - 2823
]) - 2824
.allow_headers([ - 2825
axum::http::header::AUTHORIZATION, - 2826
axum::http::header::CONTENT_TYPE, - 2827
]); - 2828
let mut state = AppState::new(core); - 2829
state.ops_port = port; - 2830
if force_gateway { - 2831
state.enable_gateway(); - 2832
} - 2833
let token = (*state.auth_token).clone(); - 2834
let rl_settings = state.core.config().gateway.rate_limit.clone(); - 2835
let rl_config = rate_limit::RateLimitConfig::from_settings(rl_settings); - 2836
let limiter = rate_limit::RateLimiter::new(rl_config, state.core.sessions_home()); - 2837
let app = router_with_state(state.clone()) - 2838
.layer(axum::middleware::from_fn_with_state( - 2839
limiter, - 2840
rate_limit::rate_limit_layer, - 2841
)) - 2842
.layer(axum::middleware::from_fn_with_state( - 2843
state.clone(), - 2844
enforce_participant_audience, - 2845
)) - 2846
.layer(axum::middleware::from_fn_with_state( - 2847
AuthPolicy { - 2848
token: token.clone(), - 2849
home: state.core.sessions_home(), - 2850
trusted_hosts: state.core.config().server.trusted_hosts.clone(), - 2851
}, - 2852
require_bearer, - 2853
)) - 2854
.layer(cors); - 2855
#[cfg(unix)] - 2856
{ - 2857
let broker = state.core.agent_network_broker(); - 2858
let policy_file = state - 2859
.core - 2860
.sessions_home() - 2861
.join("agent-network/policies.json"); - 2862
if let Err(error) = broker.load_policies(&policy_file) { - 2863
eprintln!("[agent-network] policy load failed: {error}"); - 2864
} - 2865
let socket = - 2866
vak_core::agent_network::AgentNetworkBroker::socket_path(&state.core.sessions_home()); - 2867
tokio::spawn(async move { - 2868
if let Err(error) = broker.serve_unix(&socket).await { - 2869
eprintln!("[agent-network] broker stopped: {error}"); - 2870
} - 2871
}); - 2872
} - 2873
// Local routines: fires due scheduled tasks while this server lives. - 2874
start_scheduler(&state); - 2875
delivery::start_replay(&state.core); - 2876
// Capability discovery is NOT started here. - 2877
// - 2878
// It used to be, and the CLI did its own bounded wait, and the desktop - 2879
// did neither — three surfaces answering "when may a prompt be frozen?" - 2880
// three different ways, which is how an admitted, working MCP server - 2881
// still produced a session that had never seen its catalog. - 2882
// `Core::admitted_capabilities` owns that decision now, so every surface - 2883
// gets the same packet whether it was reached from a terminal, this - 2884
// server, or the desktop app. - 2885
// Background index sync: keeps the admin console populated from the - 2886
// very first boot. Idempotent; never blocks request handling. - 2887
if let Some(store) = state.store.clone() { - 2888
let home = state.core.sessions_home(); - 2889
tokio::spawn(async move { - 2890
match store.rebuild(&home) { - 2891
Ok(s) if s.files_scanned > 0 => eprintln!( - 2892
"[store] indexed {} files / {} entries", - 2893
s.files_scanned, s.entries_indexed - 2894
), - 2895
Ok(_) => {} - 2896
Err(e) => eprintln!("[store] startup rebuild failed: {e}"), - 2897
} - 2898
}); - 2899
} - 2900
(app, token) - 2901
} - 2902
- 2903
fn participant_matches_conversation_audience( - 2904
participant: &coworking::VerifiedPrincipal, - 2905
conversation_id: &str, - 2906
observed_audience: Option<&str>, - 2907
) -> bool { - 2908
participant.conversation_id == conversation_id - 2909
&& observed_audience == Some(participant.audience_id.as_str()) - 2910
} - 2911
- 2912
/// Recheck the live conversation audience after bearer authentication and - 2913
/// before any scoped handler runs. The route whitelist limits *where* a - 2914
/// participant may go; this boundary also proves the durable grant still - 2915
/// names the audience owned by that conversation. - 2916
async fn enforce_participant_audience( - 2917
State(state): State<AppState>, - 2918
req: axum::extract::Request, - 2919
next: axum::middleware::Next, - 2920
) -> axum::response::Response { - 2921
if let Some(AuthenticatedPrincipal::Participant(participant)) = - 2922
req.extensions().get::<AuthenticatedPrincipal>() - 2923
{ - 2924
let conversation_id = req - 2925
.uri() - 2926
.path() - 2927
.trim_matches('/') - 2928
.split('/') - 2929
.nth(1) - 2930
.unwrap_or_default(); - 2931
let audience = conversation_audience(&state, conversation_id); - 2932
if !participant_matches_conversation_audience( - 2933
participant, - 2934
conversation_id, - 2935
audience.as_deref(), - 2936
) { - 2937
return StatusCode::FORBIDDEN.into_response(); - 2938
} - 2939
} - 2940
next.run(req).await - 2941
} - 2942
- 2943
/// Initialize the distributed event bus from config and attach it to the - 2944
/// global `EventHub`. Called from `serve_with` (async) so NATS connection - 2945
/// retries don't block request handling. - 2946
pub async fn init_server_bus(core: &Core) { - 2947
let config = core.config().server.bus.clone(); - 2948
if config.nats_url.is_none() { - 2949
return; - 2950
} - 2951
install_server_bus(core, &config).await; - 2952
} - 2953
- 2954
/// Builds the bus `config` describes, with its secrets from the secrets - 2955
/// chain, and makes it the live one. - 2956
async fn install_server_bus(core: &Core, config: &vak_config::BusResolved) { - 2957
let sessions_home = core.sessions_home(); - 2958
let workspace_id = { - 2959
use sha2::{Digest, Sha256}; - 2960
let mut hasher = Sha256::new(); - 2961
hasher.update(sessions_home.to_string_lossy().as_bytes()); - 2962
let result = hasher.finalize(); - 2963
// First 6 bytes → 12 hex chars, a compact workspace-scoped subject prefix. - 2964
let prefix = &result[..6]; - 2965
let hex: String = prefix.iter().map(|b| format!("{b:02x}")).collect(); - 2966
format!("ws_{}", hex) - 2967
}; - 2968
let bus = crate::bus::ServerBus::from_resolved( - 2969
&workspace_id, - 2970
config, - 2971
&crate::bus::BusCredentials::of(core), - 2972
) - 2973
.await; - 2974
if let Some(mut hub) = crate::events::global() { - 2975
hub.set_server_bus(std::sync::Arc::new(bus)); - 2976
} - 2977
} - 2978
- 2979
pub async fn serve(core: Core, addr: std::net::SocketAddr) -> std::io::Result<()> { - 2980
serve_with(core, addr, false).await - 2981
} - 2982
- 2983
/// `force_gateway` mirrors `serve --gateway`: enable routing regardless of - 2984
/// the (untrusted-stripped) project config. - 2985
/// Serve an already-bound listener with an already-built router. - 2986
/// - 2987
/// The pieces `secured_router` returns, joined. `vak setup` binds its own - 2988
/// ephemeral loopback port so it can print the URL *before* serving, and - 2989
/// needs the token from the same call — which `serve_with` cannot give it, - 2990
/// because that mints and consumes the token internally. Exposed here so - 2991
/// axum stays a dependency of this crate rather than leaking into the CLI. - 2992
pub async fn serve_router(listener: tokio::net::TcpListener, app: Router) -> std::io::Result<()> { - 2993
axum::serve(listener, app).await - 2994
} - 2995
- 2996
pub async fn serve_with( - 2997
core: Core, - 2998
addr: std::net::SocketAddr, - 2999
force_gateway: bool, - 3000
) -> std::io::Result<()> { - 3001
// Presentation seeds are versioned, additive previews. Reconcile them at - 3002
// process startup so a newly installed binary reaches existing workspaces - 3003
// even when the admin presentation list is never opened. User revisions - 3004
// and activations remain untouched by register(). - 3005
reconcile_builtin_presentations(&core) - 3006
.map_err(|error| std::io::Error::other(format!("presentation seed failed: {error}")))?; - 3007
// Local-only does not mean safe-by-default: any local process could - 3008
// reach an unauthenticated agent and drive arbitrary tool execution - 3009
// plus self-approval. Every serve() instance gets a per-process - 3010
// bearer token; /health stays open. - 3011
let listener = tokio::net::TcpListener::bind(addr).await?; - 3012
let actual_addr = listener.local_addr()?; - 3013
let (app, token) = secured_router_with_port(core.clone(), force_gateway, actual_addr.port()); - 3014
// Initialize the distributed event bus (vak-bus, docs/design/53). - 3015
// Falls back to InMemoryBus when NATS is absent or unreachable. - 3016
init_server_bus(&core).await; - 3017
eprintln!("Vakyartha server listening on http://{actual_addr}"); - 3018
// Same source the real token-selection logic above (auth_token, in - 3019
// AppState::new) already checks: `vak_config::get_var` also sees a - 3020
// value that only reached the process through the credential store - 3021
// (never a real `std::env` var), so a plain `std::env::var` check - 3022
// here — which is all this ever did — reported "generated" for - 3023
// every service-managed deployment, since none of them export - 3024
// VAK_GATEWAY_TOKEN into the actual process environment; they rely - 3025
// on the Shared secret scope main() already loads unconditionally at - 3026
// startup. The token itself was always correctly pinned; only this - 3027
// log line was wrong, in exactly the deployment shape (a durable - 3028
// service reading from the credential store) where getting it right - 3029
// matters most for debugging a stale-cookie/token mismatch after a - 3030
// restart. - 3031
if vak_config::get_var("VAK_GATEWAY_TOKEN").is_some_and(|t| !t.trim().is_empty()) { - 3032
eprintln!("auth token: (pinned via VAK_GATEWAY_TOKEN)"); - 3033
} else if std::io::IsTerminal::is_terminal(&std::io::stderr()) { - 3034
eprintln!("auth token: {token}"); - 3035
eprintln!("clients must send 'Authorization: Bearer {token}' (or ?token=)"); - 3036
} else { - 3037
eprintln!("auth token: generated for this process (suppressed in non-interactive output)"); - 3038
} - 3039
if force_gateway { - 3040
eprintln!("gateway: ENABLED (--gateway overrides config)"); - 3041
} - 3042
let (draining_tx, draining_rx) = oneshot::channel(); - 3043
let server = std::future::IntoFuture::into_future( - 3044
axum::serve(listener, app).with_graceful_shutdown(async { - 3045
let _ = tokio::signal::ctrl_c().await; - 3046
eprintln!("\n[shutting down: draining connections]"); - 3047
let _ = draining_tx.send(()); - 3048
}), - 3049
); - 3050
tokio::pin!(server); - 3051
tokio::select! { - 3052
result = &mut server => result, - 3053
_ = draining_rx => { - 3054
// EventSource streams can remain open indefinitely. A restart - 3055
// must release session-ledger locks even when a browser keeps - 3056
// those connections alive, or the replacement server can only - 3057
// attach read-only and every follow-up conflicts. Give ordinary - 3058
// requests a short drain window, then end this process. - 3059
match tokio::time::timeout(Duration::from_secs(5), &mut server).await { - 3060
Ok(result) => result, - 3061
Err(_) => { - 3062
eprintln!("[shutdown: closing long-lived connections]"); - 3063
Ok(()) - 3064
} - 3065
} - 3066
} - 3067
} - 3068
} - 3069
- 3070
fn reconcile_builtin_presentations(core: &Core) -> Result<(), String> { - 3071
let store = vak_store::presentation::PresentationStore::new( - 3072
core.sessions_home().join("presentations.json"), - 3073
); - 3074
let mut library = store.load().map_err(|error| error.to_string())?; - 3075
let before = library.definitions().count(); - 3076
let mut changed = false; - 3077
for seed in vak_presentation::seeds::built_in_seed_pack() { - 3078
if let Some(existing) = library.get(&seed.spec.id, seed.spec.revision) - 3079
&& existing.digest != seed.digest - 3080
{ - 3081
changed = true; - 3082
} - 3083
library.register(seed).map_err(|error| error.to_string())?; - 3084
} - 3085
if changed || library.definitions().count() != before { - 3086
store.save(&library).map_err(|error| error.to_string())?; - 3087
} - 3088
Ok(()) - 3089
} - 3090
- 3091
/// Built-in recipes are available in every build without writing a standing - 3092
/// user preference. An explicit user or workspace activation still wins. - 3093
fn effective_presentation_library( - 3094
library: &vak_presentation::PresentationLibrary, - 3095
workspace_owner: &str, - 3096
) -> vak_presentation::PresentationLibrary { - 3097
let mut effective = library.clone(); - 3098
for seed in vak_presentation::seeds::built_in_seed_pack() { - 3099
let id = seed.spec.id.clone(); - 3100
let revision = seed.spec.revision; - 3101
let accepts = seed.spec.accepts.clone(); - 3102
if library.is_suppressed(&id, vak_presentation::LibraryScope::User, "user") - 3103
|| library.is_suppressed( - 3104
&id, - 3105
vak_presentation::LibraryScope::Workspace, - 3106
workspace_owner, - 3107
) - 3108
{ - 3109
continue; - 3110
} - 3111
if let Err(error) = effective.register(seed) { - 3112
eprintln!("[presentation] built-in pack {id} unavailable: {error}"); - 3113
continue; - 3114
} - 3115
if accepts.iter().any(|semantic_type| { - 3116
effective - 3117
.select_preferred(semantic_type, "user", workspace_owner) - 3118
.is_some_and(|selected| { - 3119
selected.spec.metadata.get("seed").map(String::as_str) != Some("true") - 3120
}) - 3121
}) { - 3122
continue; - 3123
} - 3124
if let Err(error) = effective.activate( - 3125
&id, - 3126
revision, - 3127
vak_presentation::LibraryScope::Workspace, - 3128
workspace_owner, - 3129
) { - 3130
eprintln!("[presentation] built-in pack {id} could not activate: {error}"); - 3131
} - 3132
} - 3133
effective - 3134
} - 3135
- 3136
#[cfg(test)] - 3137
#[allow(clippy::expect_used)] - 3138
mod built_in_presentation_tests { - 3139
#[test] - 3140
fn every_built_in_pack_is_selected_without_a_saved_activation() { - 3141
let library = super::effective_presentation_library( - 3142
&vak_presentation::PresentationLibrary::default(), - 3143
"/tmp/presentation-selection", - 3144
); - 3145
for seed in vak_presentation::seeds::built_in_seed_pack() { - 3146
let selected = library - 3147
.select_preferred(&seed.spec.accepts[0], "user", "/tmp/presentation-selection") - 3148
.expect("built-in pack is selected"); - 3149
assert_eq!(selected.spec.id, seed.spec.id, "{}", seed.spec.accepts[0]); - 3150
} - 3151
} - 3152
- 3153
#[test] - 3154
fn deactivated_builtin_stays_off_until_explicitly_activated() { - 3155
let owner = "/tmp/presentation-selection"; - 3156
let mut library = vak_presentation::PresentationLibrary::default(); - 3157
library.deactivate( - 3158
"seed.timeline", - 3159
vak_presentation::LibraryScope::Workspace, - 3160
owner, - 3161
); - 3162
let effective = super::effective_presentation_library(&library, owner); - 3163
assert!( - 3164
effective - 3165
.select_preferred("timeline", "user", owner) - 3166
.is_none() - 3167
); - 3168
assert!(effective.select_preferred("table", "user", owner).is_some()); - 3169
let dir = tempfile::tempdir().expect("tempdir"); - 3170
let store = - 3171
vak_store::presentation::PresentationStore::new(dir.path().join("presentations.json")); - 3172
store.save(&library).expect("save suppression"); - 3173
library = store.load().expect("reload suppression"); - 3174
assert!( - 3175
super::effective_presentation_library(&library, owner) - 3176
.select_preferred("timeline", "user", owner) - 3177
.is_none() - 3178
); - 3179
let seed = vak_presentation::seeds::built_in_seed_pack() - 3180
.into_iter() - 3181
.find(|record| record.spec.id == "seed.timeline") - 3182
.expect("timeline seed"); - 3183
library.register(seed).expect("register seed"); - 3184
library - 3185
.activate( - 3186
"seed.timeline", - 3187
6, - 3188
vak_presentation::LibraryScope::Workspace, - 3189
owner, - 3190
) - 3191
.expect("activate seed"); - 3192
let effective = super::effective_presentation_library(&library, owner); - 3193
assert_eq!( - 3194
effective - 3195
.select_preferred("timeline", "user", owner) - 3196
.map(|record| record.spec.id.as_str()), - 3197
Some("seed.timeline") - 3198
); - 3199
} - 3200
} - 3201
- 3202
/// Paths that must be reachable without a token: health probe, the SPA - 3203
/// shell (static assets carry no data), and the login endpoint itself. - 3204
fn auth_exempt_path(path: &str) -> bool { - 3205
path == "/health" - 3206
|| path == "/admin" - 3207
|| path == "/admin/" - 3208
|| path == "/admin/favicon.svg" - 3209
|| path == "/admin/vak-icon.png" - 3210
|| path.starts_with("/admin/assets/") - 3211
// The workspace client's shell and its hashed assets carry no data - 3212
// and must load before a session exists — the login form is part of - 3213
// the bundle. Every route it then calls is authenticated. - 3214
|| path == "/app" - 3215
|| path == "/app/" - 3216
|| path == "/app/vak-icon.png" - 3217
|| path == "/app/manifest.webmanifest" - 3218
|| path == "/app/sw.js" - 3219
|| path.starts_with("/app/assets/") - 3220
|| path.starts_with("/app/characters/") - 3221
// The login exchange itself, and the probe that decides whether to - 3222
// show it. `/auth/session` answers `{authenticated:false}` rather - 3223
// than 401 so an unauthenticated client can tell "no session" from - 3224
// "server unreachable". - 3225
|| path == "/auth/login" - 3226
|| path == "/auth/session" - 3227
// The public site and its build stamp. Every page answers an - 3228
// unauthenticated stranger by design — a blank 401 at `/` told a - 3229
// visitor nothing at all, not even that anything was listening — - 3230
// and `site.rs` owns what those pages may say. - 3231
|| site::ROUTES.iter().any(|(uri, _)| { - 3232
*uri == path || (*uri != "/" && path.len() == uri.len() + 1 && path.starts_with(uri) && path.ends_with('/')) - 3233
}) - 3234
|| path.starts_with("/site/") - 3235
|| path == "/version" - 3236
|| path == "/favicon.ico" - 3237
|| path == "/favicon.svg" - 3238
} - 3239
- 3240
#[cfg(test)] - 3241
mod auth_exempt_path_tests { - 3242
use super::{auth_exempt_path, host_is_loopback}; - 3243
- 3244
/// A page that resolves its own domain to 127.0.0.1 becomes same-origin - 3245
/// with this server, and `SameSite=Strict` does not help — after - 3246
/// rebinding the request is not cross-site. The Host name it carries is - 3247
/// still the attacker's, which is what this rejects. - 3248
#[test] - 3249
fn only_loopback_hostnames_are_answered() { - 3250
for host in [ - 3251
"localhost", - 3252
"localhost:8901", - 3253
"127.0.0.1", - 3254
"127.0.0.1:8901", - 3255
"[::1]:8901", - 3256
"::1", - 3257
] { - 3258
assert!(host_is_loopback(Some(host)), "{host} is loopback"); - 3259
} - 3260
for host in [ - 3261
"rebind.attacker.example", - 3262
"rebind.attacker.example:8901", - 3263
"vak.internal", - 3264
"127.0.0.1.attacker.example", - 3265
"evil.com:80", - 3266
] { - 3267
assert!(!host_is_loopback(Some(host)), "{host} must be rejected"); - 3268
} - 3269
} - 3270
- 3271
/// Rebinding is a browser attack and browsers always send Host, so an - 3272
/// absent value is a non-browser client rather than something to defend - 3273
/// against. - 3274
#[test] - 3275
fn an_absent_host_is_not_treated_as_an_attack() { - 3276
assert!(host_is_loopback(None)); - 3277
} - 3278
- 3279
/// The login screen's own logo must be reachable before a cookie - 3280
/// exists to authenticate the request that would fetch it — the same - 3281
/// reasoning that already exempts favicon.svg. This regressed once - 3282
/// already: the admin console shipped a login-screen `<img>` pointing - 3283
/// at a root-level dist file, and only /admin/assets/* (the hashed - 3284
/// JS/CSS bundle) was exempt, so the logo 401'd on every fresh login. - 3285
#[test] - 3286
fn login_screen_logo_is_exempt() { - 3287
assert!(auth_exempt_path("/admin/vak-icon.png")); - 3288
} - 3289
- 3290
#[test] - 3291
fn admin_api_routes_still_require_auth() { - 3292
assert!(!auth_exempt_path("/admin/api/config")); - 3293
assert!(!auth_exempt_path("/admin/api/gateway/status")); - 3294
} - 3295
} - 3296
- 3297
/// Whether a `Host` header names this server legitimately. - 3298
/// - 3299
/// Loopback names always pass. Anything else must be listed verbatim in - 3300
/// `[server] trusted_hosts`, which is a privileged setting an untrusted - 3301
/// project cannot write (docs/design/48-web-client.md §4.2). - 3302
/// - 3303
/// This is the DNS-rebinding defence: binding 127.0.0.1 does not stop a - 3304
/// page the user visits from resolving its *own* domain to 127.0.0.1 and - 3305
/// becoming same-origin with this server, and `SameSite=Strict` does not - 3306
/// help once that has happened, because the request is then not cross-site. - 3307
/// Pinning the `Host` name is the check that does — and the reason - 3308
/// `trusted_hosts` takes exact names rather than patterns: a wildcard here - 3309
/// re-opens exactly the hole the list closes. - 3310
fn host_is_trusted(host: Option<&str>, trusted: &[String]) -> bool { - 3311
if host_is_loopback(host) { - 3312
return true; - 3313
} - 3314
let Some(host) = host else { return true }; - 3315
let name = strip_port(host).to_ascii_lowercase(); - 3316
trusted.iter().any(|allowed| allowed == &name) - 3317
} - 3318
- 3319
/// Host name without its port, leaving a bare IPv6 literal intact. - 3320
fn strip_port(host: &str) -> &str { - 3321
if let Some(rest) = host.strip_prefix('[') { - 3322
rest.split_once(']').map_or(rest, |(name, _)| name) - 3323
} else if host.matches(':').count() > 1 { - 3324
host - 3325
} else { - 3326
host.split_once(':').map_or(host, |(name, _)| name) - 3327
} - 3328
} - 3329
- 3330
/// Whether `Host` names this machine's own loopback interface. - 3331
fn host_is_loopback(host: Option<&str>) -> bool { - 3332
let Some(host) = host else { - 3333
// Absent entirely: not a browser. Rebinding is a browser attack and - 3334
// every browser sends Host (or HTTP/2 `:authority`, which axum - 3335
// surfaces as the URI authority), so an absent value carries no - 3336
// attacker-chosen name to defend against. Rejecting here would only - 3337
// break well-behaved non-browser clients that omit it. - 3338
return true; - 3339
}; - 3340
// Strip the port without mangling a bare IPv6 literal, which is all - 3341
// colons: `"::1".rsplit_once(':')` yields `"::"`. - 3342
matches!( - 3343
strip_port(host), - 3344
"localhost" | "127.0.0.1" | "::1" | "0:0:0:0:0:0:0:1" - 3345
) - 3346
} - 3347
- 3348
/// Whether a state-changing request's `Origin` is one of ours. - 3349
/// - 3350
/// Cookies alone are not enough to authorize a mutation: a cookie is - 3351
/// attached by the browser to whoever asks, and `SameSite=Strict` covers - 3352
/// the common cases but not a same-site subdomain or a rebound name. So a - 3353
/// mutation carrying a cookie must ALSO carry an `Origin` we recognise. - 3354
/// - 3355
/// A request with no `Origin` at all is not a browser form post — browsers - 3356
/// always send one on cross-origin mutations — so it is allowed through - 3357
/// here and still has to satisfy `require_bearer` with a real header - 3358
/// token. That is what keeps curl, the CLI, and the bridges working - 3359
/// without giving a page any new power. - 3360
fn origin_is_trusted(origin: Option<&str>, trusted: &[String]) -> bool { - 3361
let Some(origin) = origin else { return true }; - 3362
// The Tauri webview's own origins: the desktop is a first-party client - 3363
// and its scheme is not something an attacker can mint. - 3364
if matches!( - 3365
origin, - 3366
"tauri://localhost" | "http://tauri.localhost" | "https://tauri.localhost" - 3367
) { - 3368
return true; - 3369
} - 3370
let authority = origin - 3371
.split_once("://") - 3372
.map(|(_, rest)| rest) - 3373
.unwrap_or(origin); - 3374
host_is_trusted(Some(authority), trusted) - 3375
} - 3376
- 3377
/// What the auth layer needs to know about this deployment's exposure. - 3378
#[derive(Clone)] - 3379
pub(crate) struct AuthPolicy { - 3380
pub(crate) token: String, - 3381
pub(crate) home: std::path::PathBuf, - 3382
/// Non-loopback `Host` names this server answers to (`[server] - 3383
/// trusted_hosts`). Empty on a default install. - 3384
pub(crate) trusted_hosts: Vec<String>, - 3385
} - 3386
- 3387
/// Identity established by the HTTP boundary. Participant grants are kept - 3388
/// distinct from the operator credential so downstream handlers can add - 3389
/// attribution and narrower decisions without ever treating a collaborator - 3390
/// as the workspace owner. - 3391
#[derive(Clone, Debug)] - 3392
pub(crate) enum AuthenticatedPrincipal { - 3393
Operator, - 3394
Participant(coworking::VerifiedPrincipal), - 3395
} - 3396
- 3397
fn participant_read_route_allowed( - 3398
method: &axum::http::Method, - 3399
path: &str, - 3400
principal: &coworking::VerifiedPrincipal, - 3401
) -> bool { - 3402
let segments: Vec<&str> = path.trim_matches('/').split('/').collect(); - 3403
let ["sessions", conversation_id, rest @ ..] = segments.as_slice() else { - 3404
return false; - 3405
}; - 3406
if *conversation_id != principal.conversation_id { - 3407
return false; - 3408
} - 3409
if method == axum::http::Method::POST { - 3410
let can_read = principal - 3411
.capabilities - 3412
.iter() - 3413
.any(|capability| capability == "read"); - 3414
if !can_read { - 3415
return false; - 3416
} - 3417
if matches!(*rest, ["coworking", "messages"]) { - 3418
return principal - 3419
.capabilities - 3420
.iter() - 3421
.any(|capability| capability == "message"); - 3422
} - 3423
if matches!(*rest, ["coworking", "approvals", _]) { - 3424
return true; - 3425
} - 3426
if matches!(*rest, ["office-workspaces", _, "presence"]) { - 3427
return true; - 3428
} - 3429
if matches!(*rest, ["office-workspaces", _]) { - 3430
return principal - 3431
.capabilities - 3432
.iter() - 3433
.any(|capability| capability == "edit"); - 3434
} - 3435
return principal - 3436
.capabilities - 3437
.iter() - 3438
.any(|capability| capability == "comment") - 3439
&& matches!(*rest, ["sandbox", "candidates", _, "comments"]); - 3440
} - 3441
if method != axum::http::Method::GET - 3442
|| !principal - 3443
.capabilities - 3444
.iter() - 3445
.any(|capability| capability == "read") - 3446
{
Indexing the workspace…
Vakyartha documentation is discovering safe artifacts, anchors, and source references.