- 7464
.into_response(); - 7465
}; - 7466
let agent = serde_json::json!({ - 7467
"id": agent.id, - 7468
"name": agent.name, - 7469
"character": agent.character, - 7470
"animation": agent.animation, - 7471
"voice": agent.voice, - 7472
"revision": agent.revision, - 7473
}); - 7474
match principal { - 7475
AuthenticatedPrincipal::Operator => { - 7476
Json(serde_json::json!({ "principal_id": "operator", "display_name": "You", "capabilities": ["owner"], "agent": agent })).into_response() - 7477
} - 7478
AuthenticatedPrincipal::Participant(participant) - 7479
if participant.conversation_id == conversation_id => - 7480
{ - 7481
Json(serde_json::json!({ - 7482
"principal_id": participant.principal_id, - 7483
"display_name": participant.display_name, - 7484
"capabilities": participant.capabilities, - 7485
"agent": agent, - 7486
})) - 7487
.into_response() - 7488
} - 7489
AuthenticatedPrincipal::Participant(_) => StatusCode::FORBIDDEN.into_response(), - 7490
} - 7491
} - 7492
- 7493
/// Add a verified human contribution to the shared ledger. This endpoint - 7494
/// never dispatches the Agent or carries approval, tool, or control authority. - 7495
async fn create_coworking_message( - 7496
State(state): State<AppState>, - 7497
Path(conversation_id): Path<String>, - 7498
axum::Extension(principal): axum::Extension<AuthenticatedPrincipal>, - 7499
Json(body): Json<CoworkingMessageBody>, - 7500
) -> axum::response::Response { - 7501
use axum::response::IntoResponse; - 7502
let AuthenticatedPrincipal::Participant(participant) = principal else { - 7503
return StatusCode::FORBIDDEN.into_response(); - 7504
}; - 7505
if participant.conversation_id != conversation_id - 7506
|| !participant.capabilities.iter().any(|value| value == "read") - 7507
|| !participant - 7508
.capabilities - 7509
.iter() - 7510
.any(|value| value == "message") - 7511
{ - 7512
return StatusCode::FORBIDDEN.into_response(); - 7513
} - 7514
let text = body.text.trim(); - 7515
let request_id = body.request_id.trim(); - 7516
if text.is_empty() - 7517
|| text.chars().count() > 32_768 - 7518
|| request_id.is_empty() - 7519
|| request_id.chars().count() > 120 - 7520
|| !request_id - 7521
.chars() - 7522
.all(|ch| ch.is_ascii_alphanumeric() || matches!(ch, '-' | '_')) - 7523
{ - 7524
return StatusCode::BAD_REQUEST.into_response(); - 7525
} - 7526
let append = |session: &mut vak_session::SessionLog| { - 7527
let duplicate = session.chain_to_root().iter().any(|entry| { - 7528
matches!(&entry.payload, vak_session::EntryPayload::Message(existing) - 7529
if existing.meta.as_ref().is_some_and(|meta| - 7530
meta.author_id.as_deref() == Some(participant.principal_id.as_str()) - 7531
&& meta.request_id.as_deref() == Some(request_id))) - 7532
}); - 7533
if duplicate { - 7534
return Ok(false); - 7535
} - 7536
session.append_message(vak_session::MessageRecord { - 7537
// Authorship is present both structurally and in the model-visible - 7538
// text. Later turns therefore know who contributed the message; - 7539
// human transcript projections remove this exact display prefix. - 7540
message: vak_llm::Message::user_text(format!("{}: {text}", participant.display_name)), - 7541
meta: Some(vak_session::MessageMeta { - 7542
author_id: Some(participant.principal_id.clone()), - 7543
author_name: Some(participant.display_name.clone()), - 7544
request_id: Some(request_id.to_string()), - 7545
..Default::default() - 7546
}), - 7547
})?; - 7548
Ok(true) - 7549
}; - 7550
let result = if let Some(handle) = state.get(&conversation_id) { - 7551
let result = match handle.session.lock() { - 7552
Ok(mut guard) => match guard.as_mut() { - 7553
Some(session) => append(session), - 7554
None => { - 7555
return ( - 7556
StatusCode::CONFLICT, - 7557
Json(serde_json::json!({"error":"Agent is working; send after this turn settles"})), - 7558
) - 7559
.into_response(); - 7560
} - 7561
}, - 7562
Err(_) => return StatusCode::INTERNAL_SERVER_ERROR.into_response(), - 7563
}; - 7564
if result.is_ok() { - 7565
let _ = handle.coworking_comments_tx.send(()); - 7566
} - 7567
result - 7568
} else if let Some(read_only) = open_historical_session(&state, &conversation_id) { - 7569
let path = read_only.path().to_path_buf(); - 7570
drop(read_only); - 7571
vak_session::SessionLog::open(path).and_then(|mut session| append(&mut session)) - 7572
} else { - 7573
return StatusCode::NOT_FOUND.into_response(); - 7574
}; - 7575
match result { - 7576
Ok(created) => ( - 7577
if created { - 7578
StatusCode::CREATED - 7579
} else { - 7580
StatusCode::OK - 7581
}, - 7582
Json(serde_json::json!({ - 7583
"request_id": request_id, - 7584
"created": created, - 7585
"agent_run_started": false, - 7586
})), - 7587
) - 7588
.into_response(), - 7589
Err(_) => StatusCode::INTERNAL_SERVER_ERROR.into_response(), - 7590
} - 7591
} - 7592
- 7593
async fn list_coworking_approvals( - 7594
State(state): State<AppState>, - 7595
Path(conversation_id): Path<String>, - 7596
axum::Extension(principal): axum::Extension<AuthenticatedPrincipal>, - 7597
) -> axum::response::Response { - 7598
use axum::response::IntoResponse; - 7599
let AuthenticatedPrincipal::Participant(participant) = principal else { - 7600
return StatusCode::FORBIDDEN.into_response(); - 7601
}; - 7602
if participant.conversation_id != conversation_id { - 7603
return StatusCode::FORBIDDEN.into_response(); - 7604
} - 7605
let Some(handle) = state.get(&conversation_id) else { - 7606
return StatusCode::NOT_FOUND.into_response(); - 7607
}; - 7608
let approvals: Vec<_> = handle - 7609
.pending - 7610
.lock() - 7611
.unwrap_or_else(std::sync::PoisonError::into_inner) - 7612
.values() - 7613
.filter(|request| { - 7614
request - 7615
.delegated_to - 7616
.lock() - 7617
.unwrap_or_else(std::sync::PoisonError::into_inner) - 7618
.as_deref() - 7619
== Some(participant.grant_id.as_str()) - 7620
}) - 7621
.map(|request| { - 7622
serde_json::json!({ - 7623
"request_id": request.id, - 7624
"tool": request.tool, - 7625
"args_json": request.args_json, - 7626
"reason": request.reason, - 7627
"requested_at": request.requested_at, - 7628
}) - 7629
}) - 7630
.collect(); - 7631
Json(serde_json::json!({ "approvals": approvals })).into_response() - 7632
} - 7633
- 7634
async fn answer_coworking_approval( - 7635
State(state): State<AppState>, - 7636
Path((conversation_id, request_id)): Path<(String, String)>, - 7637
axum::Extension(principal): axum::Extension<AuthenticatedPrincipal>, - 7638
Json(body): Json<CoworkingApprovalBody>, - 7639
) -> axum::response::Response { - 7640
use axum::response::IntoResponse; - 7641
let AuthenticatedPrincipal::Participant(participant) = principal else { - 7642
return StatusCode::FORBIDDEN.into_response(); - 7643
}; - 7644
if participant.conversation_id != conversation_id { - 7645
return StatusCode::FORBIDDEN.into_response(); - 7646
} - 7647
let Some(handle) = state.get(&conversation_id) else { - 7648
return StatusCode::NOT_FOUND.into_response(); - 7649
}; - 7650
let mut pending = handle - 7651
.pending - 7652
.lock() - 7653
.unwrap_or_else(std::sync::PoisonError::into_inner); - 7654
let Some(request) = pending.get(&request_id) else { - 7655
return StatusCode::NOT_FOUND.into_response(); - 7656
}; - 7657
if request - 7658
.delegated_to - 7659
.lock() - 7660
.unwrap_or_else(std::sync::PoisonError::into_inner) - 7661
.as_deref() - 7662
!= Some(participant.grant_id.as_str()) - 7663
{ - 7664
return StatusCode::FORBIDDEN.into_response(); - 7665
} - 7666
let Some(request) = pending.remove(&request_id) else { - 7667
return StatusCode::NOT_FOUND.into_response(); - 7668
}; - 7669
drop(pending); - 7670
*request - 7671
.answered_by - 7672
.lock() - 7673
.unwrap_or_else(std::sync::PoisonError::into_inner) = Some(( - 7674
participant.principal_id.clone(), - 7675
participant.display_name.clone(), - 7676
)); - 7677
request.respond(body.approve); - 7678
let _ = handle.coworking_comments_tx.send(()); - 7679
Json(serde_json::json!({ - 7680
"request_id": request_id, - 7681
"approved": body.approve, - 7682
"actor_id": participant.principal_id, - 7683
"actor_name": participant.display_name, - 7684
"remembered": false, - 7685
})) - 7686
.into_response() - 7687
} - 7688
- 7689
/// The owner delegates this pending gate to one currently active invitation. - 7690
/// The grant is bound to this request id; the participant cannot answer another gate. - 7691
async fn delegate_coworking_approval( - 7692
State(state): State<AppState>, - 7693
Path((conversation_id, request_id)): Path<(String, String)>, - 7694
axum::Extension(principal): axum::Extension<AuthenticatedPrincipal>, - 7695
Json(body): Json<CoworkingDelegationBody>, - 7696
) -> axum::response::Response { - 7697
use axum::response::IntoResponse; - 7698
if operator_only(&principal).is_err() { - 7699
return StatusCode::FORBIDDEN.into_response(); - 7700
} - 7701
let Some(audience_id) = conversation_audience(&state, &conversation_id) else { - 7702
return StatusCode::NOT_FOUND.into_response(); - 7703
}; - 7704
let path = coworking::store_path(&state.core.sessions_home()); - 7705
let Ok(grants) = coworking::list(&path, &conversation_id, chrono::Utc::now()) else { - 7706
return StatusCode::INTERNAL_SERVER_ERROR.into_response(); - 7707
}; - 7708
let Some(grant) = grants.iter().find(|grant| { - 7709
grant.grant_id == body.grant_id - 7710
&& grant.audience_id == audience_id - 7711
&& grant.status == coworking::GrantStatus::Active - 7712
}) else { - 7713
return StatusCode::FORBIDDEN.into_response(); - 7714
}; - 7715
let Some(handle) = state.get(&conversation_id) else { - 7716
return StatusCode::NOT_FOUND.into_response(); - 7717
}; - 7718
let pending = handle - 7719
.pending - 7720
.lock() - 7721
.unwrap_or_else(std::sync::PoisonError::into_inner); - 7722
let Some(request) = pending.get(&request_id) else { - 7723
return StatusCode::NOT_FOUND.into_response(); - 7724
}; - 7725
let mut assignment = request - 7726
.delegated_to - 7727
.lock() - 7728
.unwrap_or_else(std::sync::PoisonError::into_inner); - 7729
if assignment.as_deref() == Some(grant.grant_id.as_str()) { - 7730
return Json( - 7731
serde_json::json!({"request_id": request_id, "delegated_to": grant.display_name}), - 7732
) - 7733
.into_response(); - 7734
} - 7735
if assignment.is_some() { - 7736
return StatusCode::CONFLICT.into_response(); - 7737
} - 7738
record_activity_or_buffer( - 7739
&handle, - 7740
vak_session::ActivityRecord { - 7741
activity_id: format!("approval-delegation-{request_id}"), - 7742
turn: None, - 7743
kind: vak_session::ActivityKind::Approval, - 7744
status: vak_session::ActivityStatus::Succeeded, - 7745
label: format!("Approval assigned to {}", grant.display_name), - 7746
detail: Some("The owner assigned this pending decision to one invited person".into()), - 7747
data: std::collections::BTreeMap::from([ - 7748
("request_id".into(), request_id.clone()), - 7749
("actor_id".into(), "operator".into()), - 7750
("grant_id".into(), grant.grant_id.clone()), - 7751
("delegate_id".into(), grant.principal_id.clone()), - 7752
("delegate_name".into(), grant.display_name.clone()), - 7753
]), - 7754
}, - 7755
); - 7756
*assignment = Some(grant.grant_id.clone()); - 7757
drop(assignment); - 7758
drop(pending); - 7759
let _ = handle.coworking_comments_tx.send(()); - 7760
Json(serde_json::json!({"request_id": request_id, "delegated_to": grant.display_name})) - 7761
.into_response() - 7762
} - 7763
- 7764
const COWORKING_PRESENCE_TTL: Duration = Duration::from_secs(7); - 7765
- 7766
fn touch_coworking_presence( - 7767
state: &AppState, - 7768
conversation_id: &str, - 7769
principal_id: &str, - 7770
display_name: &str, - 7771
) { - 7772
let mut all = state - 7773
.coworking_presence - 7774
.lock() - 7775
.unwrap_or_else(std::sync::PoisonError::into_inner); - 7776
let presence = all - 7777
.entry(conversation_id.to_string()) - 7778
.or_default() - 7779
.entry(principal_id.to_string()) - 7780
.or_insert_with(|| CoworkingPresence { - 7781
display_name: display_name.to_string(), - 7782
seen_at: Instant::now(), - 7783
office_room_id: None, - 7784
office_anchor: None, - 7785
}); - 7786
presence.display_name = display_name.to_string(); - 7787
presence.seen_at = Instant::now(); - 7788
} - 7789
- 7790
fn coworking_presence_snapshot(state: &AppState, conversation_id: &str) -> serde_json::Value { - 7791
let now = Instant::now(); - 7792
let mut all = state - 7793
.coworking_presence - 7794
.lock() - 7795
.unwrap_or_else(std::sync::PoisonError::into_inner); - 7796
let Some(conversation) = all.get_mut(conversation_id) else { - 7797
return serde_json::json!({ "participants": [] }); - 7798
}; - 7799
conversation - 7800
.retain(|_, presence| now.duration_since(presence.seen_at) <= COWORKING_PRESENCE_TTL); - 7801
let mut participants: Vec<_> = conversation - 7802
.iter() - 7803
.map(|(principal_id, presence)| { - 7804
serde_json::json!({ - 7805
"principal_id": principal_id, - 7806
"display_name": presence.display_name, - 7807
"office_room_id": presence.office_room_id, - 7808
"office_anchor": presence.office_anchor, - 7809
}) - 7810
}) - 7811
.collect(); - 7812
participants.sort_by(|left, right| { - 7813
left["display_name"] - 7814
.as_str() - 7815
.cmp(&right["display_name"].as_str()) - 7816
}); - 7817
serde_json::json!({ "participants": participants }) - 7818
} - 7819
- 7820
async fn coworking_presence( - 7821
State(state): State<AppState>, - 7822
Path(conversation_id): Path<String>, - 7823
axum::Extension(principal): axum::Extension<AuthenticatedPrincipal>, - 7824
) -> axum::response::Response { - 7825
use axum::response::IntoResponse; - 7826
if !conversation_exists(&state, &conversation_id) { - 7827
return StatusCode::NOT_FOUND.into_response(); - 7828
} - 7829
match principal { - 7830
AuthenticatedPrincipal::Operator => {} - 7831
AuthenticatedPrincipal::Participant(participant) => { - 7832
if participant.conversation_id != conversation_id { - 7833
return StatusCode::FORBIDDEN.into_response(); - 7834
} - 7835
touch_coworking_presence( - 7836
&state, - 7837
&conversation_id, - 7838
&participant.principal_id, - 7839
&participant.display_name, - 7840
); - 7841
} - 7842
} - 7843
Json(coworking_presence_snapshot(&state, &conversation_id)).into_response() - 7844
} - 7845
- 7846
/// A content-free refresh signal for both sides of a shared conversation. - 7847
/// Participant credentials never reach the general Agent SSE route, and each - 7848
/// participant signal rechecks the durable grant. - 7849
async fn coworking_updates( - 7850
State(state): State<AppState>, - 7851
Path(conversation_id): Path<String>, - 7852
axum::Extension(principal): axum::Extension<AuthenticatedPrincipal>, - 7853
headers: axum::http::HeaderMap, - 7854
) -> axum::response::Response { - 7855
use axum::response::IntoResponse; - 7856
let Some(handle) = state.get(&conversation_id) else { - 7857
return StatusCode::NOT_FOUND.into_response(); - 7858
}; - 7859
let grant = match principal { - 7860
AuthenticatedPrincipal::Operator => None, - 7861
AuthenticatedPrincipal::Participant(participant) => { - 7862
if participant.conversation_id != conversation_id - 7863
|| conversation_audience(&state, &conversation_id).as_deref() - 7864
!= Some(participant.audience_id.as_str()) - 7865
{ - 7866
return StatusCode::FORBIDDEN.into_response(); - 7867
} - 7868
let Some(token) = headers - 7869
.get(axum::http::header::AUTHORIZATION) - 7870
.and_then(|value| value.to_str().ok()) - 7871
.and_then(|value| value.strip_prefix("Bearer ")) - 7872
.map(ToOwned::to_owned) - 7873
else { - 7874
return StatusCode::UNAUTHORIZED.into_response(); - 7875
}; - 7876
Some(( - 7877
token, - 7878
participant.grant_id, - 7879
participant.principal_id, - 7880
participant.display_name, - 7881
)) - 7882
} - 7883
}; - 7884
let grant_path = coworking::store_path(&state.core.sessions_home()); - 7885
let stream = futures::stream::unfold( - 7886
( - 7887
handle.events_tx.subscribe(), - 7888
handle.coworking_comments_tx.subscribe(), - 7889
tokio::time::interval(std::time::Duration::from_secs(2)), - 7890
true, - 7891
), - 7892
move |(mut events, mut comments, mut tick, active)| { - 7893
let grant_path = grant_path.clone(); - 7894
let grant = grant.clone(); - 7895
let state = state.clone(); - 7896
let conversation_id = conversation_id.clone(); - 7897
async move { - 7898
if !active { - 7899
return None; - 7900
} - 7901
let changed = tokio::select! { - 7902
_ = tick.tick() => false, - 7903
_ = events.recv() => true, - 7904
_ = comments.recv() => true, - 7905
}; - 7906
let valid = grant.as_ref().is_none_or(|(token, grant_id, _, _)| { - 7907
matches!( - 7908
coworking::verify(&grant_path, token, chrono::Utc::now()), - 7909
Ok(Some(current)) if current.grant_id == *grant_id - 7910
) - 7911
}); - 7912
if valid && let Some((_, _, principal_id, display_name)) = &grant { - 7913
touch_coworking_presence(&state, &conversation_id, principal_id, display_name); - 7914
} - 7915
let event = if valid { - 7916
Event::default() - 7917
.event(if changed { "refresh" } else { "heartbeat" }) - 7918
.data(coworking_presence_snapshot(&state, &conversation_id).to_string()) - 7919
} else { - 7920
Event::default().event("revoked").data("{}") - 7921
}; - 7922
Some(( - 7923
Ok::<_, std::convert::Infallible>(event), - 7924
(events, comments, tick, valid), - 7925
)) - 7926
} - 7927
}, - 7928
); - 7929
Sse::new(stream).into_response() - 7930
} - 7931
- 7932
async fn list_coworking_invitations( - 7933
State(state): State<AppState>, - 7934
Path(conversation_id): Path<String>, - 7935
axum::Extension(principal): axum::Extension<AuthenticatedPrincipal>, - 7936
) -> axum::response::Response { - 7937
use axum::response::IntoResponse; - 7938
if let Err(status) = operator_only(&principal) { - 7939
return status.into_response(); - 7940
} - 7941
if !conversation_exists(&state, &conversation_id) { - 7942
return StatusCode::NOT_FOUND.into_response(); - 7943
} - 7944
match coworking::list( - 7945
&coworking::store_path(&state.core.sessions_home()), - 7946
&conversation_id, - 7947
chrono::Utc::now(), - 7948
) { - 7949
Ok(invitations) => Json(serde_json::json!({ "invitations": invitations })).into_response(), - 7950
Err(_) => StatusCode::INTERNAL_SERVER_ERROR.into_response(), - 7951
} - 7952
} - 7953
- 7954
async fn create_coworking_invitation( - 7955
State(state): State<AppState>, - 7956
Path(conversation_id): Path<String>, - 7957
axum::Extension(principal): axum::Extension<AuthenticatedPrincipal>, - 7958
Json(body): Json<CoworkingInvitationBody>, - 7959
) -> axum::response::Response { - 7960
use axum::response::IntoResponse; - 7961
if let Err(status) = operator_only(&principal) { - 7962
return status.into_response(); - 7963
} - 7964
let Some(audience_id) = conversation_audience(&state, &conversation_id) else { - 7965
return StatusCode::NOT_FOUND.into_response(); - 7966
}; - 7967
let display_name = body.display_name.trim(); - 7968
if display_name.is_empty() - 7969
|| display_name.chars().count() > 120 - 7970
|| display_name.chars().any(char::is_control) - 7971
|| !(1..=30 * 24).contains(&body.expires_in_hours) - 7972
{ - 7973
return StatusCode::BAD_REQUEST.into_response(); - 7974
} - 7975
let now = chrono::Utc::now(); - 7976
let token = coworking::generate_token(); - 7977
let grant = coworking::AudienceGrant { - 7978
grant_id: uuid::Uuid::now_v7().to_string(), - 7979
principal_id: uuid::Uuid::now_v7().to_string(), - 7980
display_name: display_name.to_string(), - 7981
conversation_id: conversation_id.clone(), - 7982
audience_id, - 7983
capabilities: { - 7984
let mut capabilities = vec!["read".into()]; - 7985
if body.can_message { - 7986
capabilities.push("message".into()); - 7987
} - 7988
if body.can_comment { - 7989
capabilities.push("comment".into()); - 7990
} - 7991
if body.can_edit { - 7992
capabilities.push("edit".into()); - 7993
} - 7994
capabilities - 7995
}, - 7996
token_hash: coworking::token_hash(&token), - 7997
created_at: now.to_rfc3339(), - 7998
expires_at: (now + chrono::Duration::hours(i64::from(body.expires_in_hours))).to_rfc3339(), - 7999
}; - 8000
match coworking::invite( - 8001
&coworking::store_path(&state.core.sessions_home()), - 8002
grant.clone(), - 8003
) { - 8004
Ok(()) => ( - 8005
StatusCode::CREATED, - 8006
Json(serde_json::json!({ - 8007
"invitation": { - 8008
"grant_id": grant.grant_id, - 8009
"principal_id": grant.principal_id, - 8010
"display_name": grant.display_name, - 8011
"conversation_id": grant.conversation_id, - 8012
"audience_id": grant.audience_id, - 8013
"capabilities": grant.capabilities, - 8014
"created_at": grant.created_at, - 8015
"expires_at": grant.expires_at, - 8016
"status": "active", - 8017
}, - 8018
"token": token, - 8019
})), - 8020
) - 8021
.into_response(), - 8022
Err(_) => StatusCode::INTERNAL_SERVER_ERROR.into_response(), - 8023
} - 8024
} - 8025
- 8026
async fn revoke_coworking_invitation( - 8027
State(state): State<AppState>, - 8028
Path((conversation_id, grant_id)): Path<(String, String)>, - 8029
axum::Extension(principal): axum::Extension<AuthenticatedPrincipal>, - 8030
) -> axum::response::Response { - 8031
use axum::response::IntoResponse; - 8032
if let Err(status) = operator_only(&principal) { - 8033
return status.into_response(); - 8034
} - 8035
let path = coworking::store_path(&state.core.sessions_home()); - 8036
let invitations = match coworking::list(&path, &conversation_id, chrono::Utc::now()) { - 8037
Ok(invitations) => invitations, - 8038
Err(_) => return StatusCode::INTERNAL_SERVER_ERROR.into_response(), - 8039
}; - 8040
let Some(invitation) = invitations.iter().find(|item| item.grant_id == grant_id) else { - 8041
return StatusCode::NOT_FOUND.into_response(); - 8042
}; - 8043
if invitation.status != coworking::GrantStatus::Revoked - 8044
&& coworking::revoke(&path, &grant_id, "operator").is_err() - 8045
{ - 8046
return StatusCode::INTERNAL_SERVER_ERROR.into_response(); - 8047
} - 8048
StatusCode::NO_CONTENT.into_response() - 8049
} - 8050
- 8051
fn html_response(html: String) -> axum::response::Response { - 8052
use axum::response::IntoResponse; - 8053
( - 8054
[( - 8055
axum::http::header::CONTENT_TYPE, - 8056
axum::http::HeaderValue::from_static("text/html; charset=utf-8"), - 8057
)], - 8058
html, - 8059
) - 8060
.into_response() - 8061
} - 8062
- 8063
fn markdown_response(md: String) -> axum::response::Response { - 8064
( - 8065
[( - 8066
axum::http::header::CONTENT_TYPE, - 8067
axum::http::HeaderValue::from_static("text/markdown; charset=utf-8"), - 8068
)], - 8069
md, - 8070
) - 8071
.into_response() - 8072
} - 8073
- 8074
/// Find a session log on disk across current sessions_home and all agent directories. - 8075
fn find_session_on_disk(core: &Core, id: &str) -> Option<vak_session::SessionLog> { - 8076
let home = core.sessions_home(); - 8077
let path = home - 8078
.join("sessions") - 8079
.join(vak_core::memory::hash_cwd(core.cwd())) - 8080
.join(format!("{id}.jsonl")); - 8081
if let Ok(s) = vak_session::SessionLog::open_read_only(path) { - 8082
return Some(s); - 8083
} - 8084
if let Ok(entries) = std::fs::read_dir(home.join("sessions")) { - 8085
for entry in entries.flatten() { - 8086
let candidate = entry.path().join(format!("{id}.jsonl")); - 8087
if let Ok(s) = vak_session::SessionLog::open_read_only(candidate) { - 8088
return Some(s); - 8089
} - 8090
} - 8091
} - 8092
let shared = core.shared_data_home(); - 8093
if let Ok(entries) = std::fs::read_dir(shared.join("sessions")) { - 8094
for entry in entries.flatten() { - 8095
let candidate = entry.path().join(format!("{id}.jsonl")); - 8096
if let Ok(s) = vak_session::SessionLog::open_read_only(candidate) { - 8097
return Some(s); - 8098
} - 8099
} - 8100
} - 8101
if let Ok(agents) = std::fs::read_dir(shared.join("agents")) { - 8102
for agent in agents.flatten() { - 8103
if let Ok(projects) = std::fs::read_dir(agent.path().join("sessions")) { - 8104
for project in projects.flatten() { - 8105
let candidate = project.path().join(format!("{id}.jsonl")); - 8106
if let Ok(s) = vak_session::SessionLog::open_read_only(candidate) { - 8107
return Some(s); - 8108
} - 8109
} - 8110
} - 8111
} - 8112
} - 8113
None - 8114
} - 8115
- 8116
/// Historical sessions live on disk but not in the in-memory handle map - 8117
/// (a fresh server process starts with an empty map). Open read-only for - 8118
/// export/inspection without mutating run bookkeeping. A trashed session is - 8119
/// not opened: it is hidden everywhere (`vak_core::trash`). - 8120
fn open_historical_session(state: &AppState, id: &str) -> Option<vak_session::SessionLog> { - 8121
if vak_core::trash::is_trashed(&state.core.shared_data_home(), id) { - 8122
return None; - 8123
} - 8124
find_session_on_disk(&state.core, id) - 8125
} - 8126
- 8127
/// Read only the immutable first header entry without acquiring the session's - 8128
/// writer lock. Admin forensics must still work when another local process is - 8129
/// actively serving the channel. - 8130
pub(crate) fn read_historical_header( - 8131
state: &AppState, - 8132
id: &str, - 8133
workspace: Option<&std::path::Path>, - 8134
) -> Option<vak_session::types::SessionHeader> { - 8135
fn read(path: &std::path::Path) -> Option<vak_session::types::SessionHeader> { - 8136
use std::io::BufRead; - 8137
let file = std::fs::File::open(path).ok()?; - 8138
for line in std::io::BufReader::new(file).lines().take(4) { - 8139
let entry: vak_session::types::Entry = serde_json::from_str(&line.ok()?).ok()?; - 8140
if let vak_session::types::EntryPayload::Header(header) = entry.payload { - 8141
return Some(header); - 8142
} - 8143
} - 8144
None - 8145
} - 8146
- 8147
if let Some(workspace) = workspace { - 8148
let path = - 8149
vak_session::SessionPath::new_session_file(&state.core.sessions_home(), workspace, id); - 8150
if let Some(header) = read(&path) { - 8151
return Some(header); - 8152
} - 8153
} - 8154
if let Ok(entries) = std::fs::read_dir(state.core.sessions_home().join("sessions")) { - 8155
for project in entries.flatten().filter(|entry| entry.path().is_dir()) { - 8156
let path = project.path().join(format!("{id}.jsonl")); - 8157
if let Some(header) = read(&path) { - 8158
return Some(header); - 8159
} - 8160
} - 8161
} - 8162
let shared = state.core.shared_data_home(); - 8163
if let Ok(entries) = std::fs::read_dir(shared.join("sessions")) { - 8164
for project in entries.flatten().filter(|entry| entry.path().is_dir()) { - 8165
let path = project.path().join(format!("{id}.jsonl")); - 8166
if let Some(header) = read(&path) { - 8167
return Some(header); - 8168
} - 8169
} - 8170
} - 8171
if let Ok(agents) = std::fs::read_dir(shared.join("agents")) { - 8172
for agent in agents.flatten().filter(|entry| entry.path().is_dir()) { - 8173
if let Ok(projects) = std::fs::read_dir(agent.path().join("sessions")) { - 8174
for project in projects.flatten().filter(|entry| entry.path().is_dir()) { - 8175
let path = project.path().join(format!("{id}.jsonl")); - 8176
if let Some(header) = read(&path) { - 8177
return Some(header); - 8178
} - 8179
} - 8180
} - 8181
} - 8182
} - 8183
None - 8184
} - 8185
- 8186
/// Reopen a session whose in-memory handle was consumed by a turn that - 8187
/// then failed: `run_turn_with` returns `Err(CoreError)` without the log, - 8188
/// but the append-only ledger file is durable — restore from it so the - 8189
/// session does not stay wedged as "run in progress" forever. - 8190
fn reopen_ledger(core: &vak_core::Core, id: &str) -> Option<vak_session::SessionLog> { - 8191
let path = core - 8192
.sessions_home() - 8193
.join("sessions") - 8194
.join(vak_core::memory::hash_cwd(core.cwd())) - 8195
.join(format!("{id}.jsonl")); - 8196
vak_session::SessionLog::open(path).ok() - 8197
} - 8198
- 8199
// ---- Personal-OS surfaces (docs/design/29-personal-os.md P1–P4) ------------- - 8200
- 8201
#[derive(serde::Deserialize)] - 8202
struct DoctorQuery { - 8203
#[serde(default)] - 8204
session: Option<String>, - 8205
} - 8206
- 8207
fn health_report_json(report: vak_core::health::HealthReport) -> serde_json::Value { - 8208
let ok = report.failures == 0; - 8209
let report_text = report - 8210
.checks - 8211
.iter() - 8212
.map(|check| match &check.detail { - 8213
Ok(detail) => format!("✓ {} — {detail}", check.label), - 8214
Err(detail) => format!("✗ {} — {detail}", check.label), - 8215
}) - 8216
.chain(report.facts.iter().map(|fact| format!("· {fact}"))) - 8217
.collect::<Vec<_>>() - 8218
.join("\n"); - 8219
serde_json::json!({ - 8220
"ok": ok, - 8221
"report": report_text, - 8222
"failures": report.failures, - 8223
"checks": report.checks.iter().map(|c| serde_json::json!({ - 8224
"label": c.label, - 8225
"ok": c.detail.is_ok(), - 8226
"detail": match &c.detail { Ok(d) => d, Err(e) => e }, - 8227
})).collect::<Vec<_>>(), - 8228
"facts": report.facts, - 8229
"ladder": report.ladder.map(|l| serde_json::json!({ - 8230
"legs": l.legs, - 8231
"rendered": l.rendered, - 8232
"objective": l.objective, - 8233
"fallback_legs": l.fallback_legs, - 8234
"annotations": l.annotations, - 8235
})), - 8236
}) - 8237
} - 8238
- 8239
/// `GET /onboarding` — the derived setup projection - 8240
/// (`docs/design/46-stabilization-install-and-onboarding.md` Part III). - 8241
/// - 8242
/// The same `vak_core::onboarding::derive` the CLI renders, so web, - 8243
/// desktop, and terminal cannot disagree about what is configured. The - 8244
/// install manifest is probed by the CLI (which owns install layout) and - 8245
/// is therefore reported here as not-probed; services are probed from the - 8246
/// same service manager the ops endpoints already use. - 8247
async fn onboarding_state(State(state): State<AppState>) -> axum::response::Response { - 8248
use axum::response::IntoResponse; - 8249
// The service manager probe shells out, so it belongs on a blocking - 8250
// worker rather than inside an async handler (invariant 26). - 8251
let services = tokio::task::spawn_blocking(|| { - 8252
let cfg = vak_ops::OpsConfig::detect(); - 8253
let probed: Vec<(String, bool)> = [vak_ops::Service::Gateway, vak_ops::Service::Bridges] - 8254
.into_iter() - 8255
.filter_map(|service| { - 8256
let st = vak_ops::status(service, &cfg); - 8257
(st != vak_ops::State::NotInstalled) - 8258
.then(|| (service.label().to_string(), st == vak_ops::State::Running)) - 8259
}) - 8260
.collect(); - 8261
(!probed.is_empty()).then_some(probed) - 8262
}) - 8263
.await - 8264
.unwrap_or(None); - 8265
- 8266
let awaiting_activation = service_control::activation_drift(&state.core) - 8267
.await - 8268
.map(|drift| drift.awaiting_activation) - 8269
.unwrap_or_default(); - 8270
let projection = vak_core::onboarding::derive( - 8271
&state.core, - 8272
&vak_core::onboarding::ProbedFacts { - 8273
services, - 8274
install: None, - 8275
awaiting_activation, - 8276
}, - 8277
); - 8278
Json(projection).into_response() - 8279
} - 8280
- 8281
/// `POST /onboarding/seed` — install the Shared starter capabilities. - 8282
/// - 8283
/// Idempotent: new standard skills/plugins are added, untouched shipped - 8284
/// content may advance on update, and edited or independently installed - 8285
/// content is preserved. Hooks and network defaults are seeded only when - 8286
/// their configuration layer is empty. Explicit because seeding is a setup - 8287
/// action, never an install side effect (doc 46 D6). - 8288
async fn onboarding_seed(State(state): State<AppState>) -> axum::response::Response { - 8289
use axum::response::IntoResponse; - 8290
// Touches the filesystem and the plugin store; not an async handler's - 8291
// work (invariant 26). - 8292
let outcome = tokio::task::spawn_blocking(vak_core::seed::seed_shared_capabilities).await; - 8293
match outcome { - 8294
Ok(Ok(())) => { - 8295
state.hub.emit_config_changed("capabilities_seeded", ""); - 8296
Json(serde_json::json!({ "ok": true })).into_response() - 8297
} - 8298
Ok(Err(e)) => ( - 8299
StatusCode::INTERNAL_SERVER_ERROR, - 8300
Json(serde_json::json!({ "error": e })), - 8301
) - 8302
.into_response(), - 8303
Err(e) => ( - 8304
StatusCode::INTERNAL_SERVER_ERROR, - 8305
Json(serde_json::json!({ "error": format!("seeding did not complete: {e}") })), - 8306
) - 8307
.into_response(), - 8308
} - 8309
} - 8310
- 8311
#[derive(serde::Deserialize)] - 8312
struct WorkspacePath { - 8313
path: String, - 8314
} - 8315
- 8316
/// `POST /onboarding/workspace-review` — what a folder would ask for. - 8317
/// - 8318
/// Reports privileged sections **without loading them** (doc 46, Step 2). - 8319
/// Describing a project's config by parsing it through the normal loader - 8320
/// would activate the very thing the operator is being asked about. - 8321
async fn onboarding_workspace_review(Json(body): Json<WorkspacePath>) -> axum::response::Response { - 8322
use axum::response::IntoResponse; - 8323
let path = std::path::PathBuf::from(body.path.trim()); - 8324
if std::fs::read_dir(&path).is_err() { - 8325
return ( - 8326
StatusCode::BAD_REQUEST, - 8327
Json(serde_json::json!({ "error": format!("cannot read {}", path.display()) })), - 8328
) - 8329
.into_response(); - 8330
} - 8331
Json(serde_json::json!({ - 8332
"path": path, - 8333
"git": path.join(".git").exists(), - 8334
"requests_privilege": vak_core::trust::requests_privilege(&path), - 8335
"privileges": vak_core::trust::requested_privileges(&path), - 8336
"trusted": vak_core::trust::is_trusted(&path), - 8337
})) - 8338
.into_response() - 8339
} - 8340
- 8341
/// `POST /onboarding/trust` — record an explicit trust decision. - 8342
/// - 8343
/// Only ever *grants*: opening safely is the absence of a decision, and is - 8344
/// already the default, so there is nothing to write for it. Selecting a - 8345
/// folder is never itself consent (doc 46 security invariant 2). - 8346
async fn onboarding_trust( - 8347
State(state): State<AppState>, - 8348
Json(body): Json<WorkspacePath>, - 8349
) -> axum::response::Response { - 8350
use axum::response::IntoResponse; - 8351
let path = std::path::PathBuf::from(body.path.trim()); - 8352
match vak_core::trust::record(&path) { - 8353
Ok(()) => { - 8354
vak_core::security_events::record( - 8355
&state.core.sessions_home(), - 8356
vak_core::security_events::EventKind::ConfigChange, - 8357
"workspace_trusted", - 8358
&path.display().to_string(), - 8359
None, - 8360
); - 8361
Json(serde_json::json!({ "trusted": true, "path": path })).into_response() - 8362
} - 8363
Err(e) => ( - 8364
StatusCode::INTERNAL_SERVER_ERROR, - 8365
Json(serde_json::json!({ "error": e.to_string() })), - 8366
) - 8367
.into_response(), - 8368
} - 8369
} - 8370
- 8371
/// The starter task. Read-only by construction, and deliberately not - 8372
/// something the caller supplies: a prompt this endpoint accepted would be - 8373
/// a way to run arbitrary work under the onboarding path. - 8374
const FIRST_TASK_PROMPT: &str = "Map this codebase and explain its architecture, key flows, \ - 8375
and highest-risk areas. Do not modify files or run any destructive command."; - 8376
- 8377
/// `POST /onboarding/first-task` — create the guided starter session. - 8378
/// - 8379
/// Capped to read-only **regardless of the workspace's configured mode** - 8380
/// (doc 46 security invariant 5). The cap is not advisory and not the - 8381
/// caller's to choose: the session is created against a `Core` resolved - 8382
/// through `CorePool` with a read-only override, which `capped_by` folds - 8383
/// against the workspace ceiling as a `min` — so the result is provably - 8384
/// never more permissive than the workspace, and never less strict than - 8385
/// read-only. The handle carries that `Core`, and every run path uses the - 8386
/// handle's `Core`, so the cap holds for the actual dispatch rather than - 8387
/// only at creation. - 8388
/// - 8389
/// Returns the session id; the caller drives it through the normal run and - 8390
/// SSE endpoints, which is what makes its receipt an ordinary receipt. - 8391
async fn onboarding_first_task(State(state): State<AppState>) -> axum::response::Response { - 8392
use axum::response::IntoResponse; - 8393
let workspace = state.core.cwd().clone(); - 8394
let capped = match state.gateway.core_pool.resolve_at( - 8395
&workspace, - 8396
Some(vak_config::PermissionMode::ReadOnly), - 8397
std::time::Instant::now(), - 8398
) { - 8399
Ok(core) => core, - 8400
Err(e) => { - 8401
return ( - 8402
StatusCode::INTERNAL_SERVER_ERROR, - 8403
Json(serde_json::json!({ "error": e })), - 8404
) - 8405
.into_response(); - 8406
} - 8407
}; - 8408
- 8409
let session = match capped.start_session().await { - 8410
Ok(s) => s, - 8411
Err(e) => { - 8412
return ( - 8413
StatusCode::INTERNAL_SERVER_ERROR, - 8414
Json(serde_json::json!({ "error": e.to_string() })), - 8415
) - 8416
.into_response(); - 8417
} - 8418
}; - 8419
let id = session - 8420
.header() - 8421
.map(|h| h.session_id.clone()) - 8422
.unwrap_or_default(); - 8423
register_handle(&state, id.clone(), session, workspace, capped.clone()); - 8424
state.hub.emit_session_created(&id, ""); - 8425
index_session_later(state.store.clone(), state.core.sessions_home(), id.clone()); - 8426
- 8427
Json(serde_json::json!({ - 8428
"session_id": id, - 8429
"prompt": FIRST_TASK_PROMPT, - 8430
"permission_mode": format!("{:?}", capped.effective_permission_mode()), - 8431
})) - 8432
.into_response() - 8433
} - 8434
- 8435
async fn doctor_report( - 8436
State(state): State<AppState>, - 8437
axum::extract::Query(q): axum::extract::Query<DoctorQuery>, - 8438
) -> axum::response::Response { - 8439
use axum::response::IntoResponse; - 8440
// The optional session adds its frozen-ladder section; a live run owns - 8441
// the ledger, in which case doctor reports without that section rather - 8442
// than failing. - 8443
let session_handle = q.session.and_then(|sid| state.get(&sid)); - 8444
let session_guard = session_handle.as_deref().map(|h| { - 8445
h.session - 8446
.lock() - 8447
.unwrap_or_else(std::sync::PoisonError::into_inner) - 8448
}); - 8449
let report = - 8450
vak_core::health::collect(&state.core, session_guard.as_ref().and_then(|g| g.as_ref())); - 8451
(StatusCode::OK, Json(health_report_json(report))).into_response() - 8452
} - 8453
- 8454
#[derive(serde::Deserialize)] - 8455
struct BackupExportBody { - 8456
dest_dir: String, - 8457
#[serde(default)] - 8458
include_secrets: bool, - 8459
} - 8460
- 8461
#[derive(serde::Deserialize)] - 8462
struct BackupImportBody { - 8463
src_dir: String,
Indexing the workspace…
Vakyartha documentation is discovering safe artifacts, anchors, and source references.