- 6464
Json(serde_json::json!({ "error": "session has no valid work contract" })), - 6465
) - 6466
.into_response(); - 6467
}; - 6468
let contract_id = projection.contract.contract_id.clone(); - 6469
let revision = projection.contract.revision; - 6470
let event_kind = match command { - 6471
WorkCommand::Transition { - 6472
item_id, - 6473
to, - 6474
reason, - 6475
} => { - 6476
let Some(item) = projection.items.get(&item_id) else { - 6477
return ( - 6478
StatusCode::BAD_REQUEST, - 6479
Json(serde_json::json!({ "error": format!("unknown work item '{item_id}'") })), - 6480
) - 6481
.into_response(); - 6482
}; - 6483
if matches!(to, vak_session::types::WorkItemStatus::Succeeded) { - 6484
return ( - 6485
StatusCode::BAD_REQUEST, - 6486
Json(serde_json::json!({ "error": "succeeded requires independent verification" })), - 6487
) - 6488
.into_response(); - 6489
} - 6490
vak_session::types::WorkEventKind::ItemStatusChanged { - 6491
item_id, - 6492
from: item.status.clone(), - 6493
to, - 6494
attempt: item.attempt, - 6495
reason, - 6496
} - 6497
} - 6498
WorkCommand::Cancel { reason } => { - 6499
vak_session::types::WorkEventKind::ContractStatusChanged { - 6500
from: projection.status, - 6501
to: vak_session::types::WorkContractStatus::Cancelled, - 6502
reason, - 6503
} - 6504
} - 6505
WorkCommand::Resume { reason } => { - 6506
if projection.status == vak_session::types::WorkContractStatus::AwaitingInput - 6507
&& projection.contract.assumptions.iter().any(|assumption| { - 6508
assumption.requires_confirmation && assumption.resolution.is_none() - 6509
}) - 6510
{ - 6511
return ( - 6512
StatusCode::BAD_REQUEST, - 6513
Json(serde_json::json!({ - 6514
"error": "required assumptions must be resolved before resuming" - 6515
})), - 6516
) - 6517
.into_response(); - 6518
} - 6519
vak_session::types::WorkEventKind::ContractStatusChanged { - 6520
from: projection.status, - 6521
to: vak_session::types::WorkContractStatus::Active, - 6522
reason, - 6523
} - 6524
} - 6525
WorkCommand::Retry { item_id, reason } => { - 6526
let Some(item) = projection.items.get(&item_id) else { - 6527
return ( - 6528
StatusCode::BAD_REQUEST, - 6529
Json(serde_json::json!({ "error": format!("unknown work item '{item_id}'") })), - 6530
) - 6531
.into_response(); - 6532
}; - 6533
vak_session::types::WorkEventKind::ItemStatusChanged { - 6534
item_id, - 6535
from: item.status.clone(), - 6536
to: vak_session::types::WorkItemStatus::Ready, - 6537
attempt: item.attempt.saturating_add(1), - 6538
reason, - 6539
} - 6540
} - 6541
WorkCommand::ResolveAssumption { - 6542
assumption_id, - 6543
resolution, - 6544
} => { - 6545
if resolution.trim().is_empty() - 6546
|| !projection - 6547
.contract - 6548
.assumptions - 6549
.iter() - 6550
.any(|assumption| assumption.assumption_id == assumption_id) - 6551
{ - 6552
return ( - 6553
StatusCode::BAD_REQUEST, - 6554
Json(serde_json::json!({ "error": "unknown assumption or empty resolution" })), - 6555
) - 6556
.into_response(); - 6557
} - 6558
vak_session::types::WorkEventKind::AssumptionResolved { - 6559
assumption_id, - 6560
resolution, - 6561
} - 6562
} - 6563
WorkCommand::AttachEvidence { item_id, evidence } => { - 6564
if !projection.items.contains_key(&item_id) { - 6565
return ( - 6566
StatusCode::BAD_REQUEST, - 6567
Json(serde_json::json!({ "error": format!("unknown work item '{item_id}'") })), - 6568
) - 6569
.into_response(); - 6570
} - 6571
if matches!( - 6572
&evidence, - 6573
vak_session::types::EvidenceRef::FlowNode { .. } - 6574
| vak_session::types::EvidenceRef::ChildSession { .. } - 6575
| vak_session::types::EvidenceRef::ExternalOperation { .. } - 6576
) { - 6577
return ( - 6578
StatusCode::BAD_REQUEST, - 6579
Json(serde_json::json!({ - 6580
"error": "runtime-owned evidence must be produced by its integration" - 6581
})), - 6582
) - 6583
.into_response(); - 6584
} - 6585
vak_session::types::WorkEventKind::EvidenceAttached { item_id, evidence } - 6586
} - 6587
WorkCommand::Assign { - 6588
item_id, - 6589
owner, - 6590
child_session_id, - 6591
} => { - 6592
if !projection.items.contains_key(&item_id) { - 6593
return ( - 6594
StatusCode::BAD_REQUEST, - 6595
Json(serde_json::json!({ "error": format!("unknown work item '{item_id}'") })), - 6596
) - 6597
.into_response(); - 6598
} - 6599
vak_session::types::WorkEventKind::ItemAssigned { - 6600
item_id, - 6601
owner, - 6602
child_session_id, - 6603
} - 6604
} - 6605
WorkCommand::Revise { contract, reason } => { - 6606
if contract.contract_id != projection.contract.contract_id - 6607
|| contract.revision != revision.saturating_add(1) - 6608
{ - 6609
return ( - 6610
StatusCode::BAD_REQUEST, - 6611
Json(serde_json::json!({ "error": "revision must target the active contract and be exactly one greater" })), - 6612
) - 6613
.into_response(); - 6614
} - 6615
if revision >= state.core.effective_work().max_revisions { - 6616
return ( - 6617
StatusCode::CONFLICT, - 6618
Json(serde_json::json!({ "error": "maximum work revisions reached" })), - 6619
) - 6620
.into_response(); - 6621
} - 6622
if let Err(error) = vak_session::validate_contract_for_admission(&contract) { - 6623
return ( - 6624
StatusCode::BAD_REQUEST, - 6625
Json(serde_json::json!({ "error": error.to_string() })), - 6626
) - 6627
.into_response(); - 6628
} - 6629
if let Err(error) = vak_agent::validate_work_paths(&contract) { - 6630
return ( - 6631
StatusCode::BAD_REQUEST, - 6632
Json(serde_json::json!({ "error": error })), - 6633
) - 6634
.into_response(); - 6635
} - 6636
vak_session::types::WorkEventKind::ContractRevised { - 6637
previous_revision: revision, - 6638
contract, - 6639
reason, - 6640
} - 6641
} - 6642
}; - 6643
let event = vak_session::types::WorkEvent { - 6644
contract_id, - 6645
revision: if matches!( - 6646
&event_kind, - 6647
vak_session::types::WorkEventKind::ContractRevised { .. } - 6648
) { - 6649
revision.saturating_add(1) - 6650
} else { - 6651
revision - 6652
}, - 6653
kind: event_kind, - 6654
}; - 6655
match session.append_work(event) { - 6656
Ok(_) => { - 6657
if let Ok(Some(updated)) = session.work_projection() - 6658
&& updated.status == vak_session::types::WorkContractStatus::AwaitingInput - 6659
&& updated - 6660
.contract - 6661
.assumptions - 6662
.iter() - 6663
.filter(|assumption| assumption.requires_confirmation) - 6664
.all(|assumption| assumption.resolution.is_some()) - 6665
{ - 6666
let _ = session.append_work(vak_session::types::WorkEvent { - 6667
contract_id: updated.contract.contract_id.clone(), - 6668
revision: updated.contract.revision, - 6669
kind: vak_session::types::WorkEventKind::ContractStatusChanged { - 6670
from: vak_session::types::WorkContractStatus::AwaitingInput, - 6671
to: vak_session::types::WorkContractStatus::Active, - 6672
reason: "all required assumptions resolved".into(), - 6673
}, - 6674
}); - 6675
} - 6676
match session.work_projection() { - 6677
Ok(Some(updated)) => Json(updated).into_response(), - 6678
_ => StatusCode::INTERNAL_SERVER_ERROR.into_response(), - 6679
} - 6680
} - 6681
Err(error) => ( - 6682
StatusCode::CONFLICT, - 6683
Json(serde_json::json!({ "error": error.to_string() })), - 6684
) - 6685
.into_response(), - 6686
} - 6687
} - 6688
- 6689
/// Flow names discovered under `<sessions_home>/flow-runs` (docs/design/42-managed-work-contracts.mdG). - 6690
async fn flows_list(State(state): State<AppState>) -> Json<Vec<String>> { - 6691
let root = state.core.sessions_home().join("flow-runs"); - 6692
let mut out = Vec::new(); - 6693
if let Ok(entries) = std::fs::read_dir(&root) { - 6694
for e in entries.flatten() { - 6695
if e.path().is_dir() { - 6696
out.push(e.file_name().to_string_lossy().into_owned()); - 6697
} - 6698
} - 6699
} - 6700
out.sort(); - 6701
Json(out) - 6702
} - 6703
- 6704
/// Run ledger filenames for one flow, oldest first. - 6705
async fn flow_runs_list( - 6706
State(state): State<AppState>, - 6707
Path(name): Path<String>, - 6708
) -> Result<Json<Vec<String>>, StatusCode> { - 6709
let dir = state.core.sessions_home().join("flow-runs").join(&name); - 6710
let mut out = Vec::new(); - 6711
match std::fs::read_dir(&dir) { - 6712
Ok(entries) => { - 6713
for e in entries.flatten() { - 6714
if e.path().extension().map(|x| x == "json").unwrap_or(false) { - 6715
out.push(e.file_name().to_string_lossy().into_owned()); - 6716
} - 6717
} - 6718
out.sort(); - 6719
Ok(Json(out)) - 6720
} - 6721
Err(_) => Err(StatusCode::NOT_FOUND), - 6722
} - 6723
} - 6724
- 6725
/// Typed run-graph snapshot (delta+snapshot invariant 4): projection of a - 6726
/// single run ledger — statuses, layers, counts. No rendering opinions. - 6727
async fn flow_run_graph( - 6728
State(state): State<AppState>, - 6729
Path((name, run)): Path<(String, String)>, - 6730
) -> Result<Json<vak_flow::graph::RunGraph>, StatusCode> { - 6731
// `run` is either the ledger filename or its stem. - 6732
let run_file = if run.ends_with(".json") { - 6733
run.clone() - 6734
} else { - 6735
format!("{run}.json") - 6736
}; - 6737
let path = state - 6738
.core - 6739
.sessions_home() - 6740
.join("flow-runs") - 6741
.join(&name) - 6742
.join(&run_file); - 6743
match std::fs::read_to_string(&path) { - 6744
Ok(body) => { - 6745
let state: vak_flow::FlowState = - 6746
serde_json::from_str(&body).map_err(|_| StatusCode::NOT_FOUND)?; - 6747
Ok(Json(vak_flow::graph::graph_snapshot(&state))) - 6748
} - 6749
Err(_) => Err(StatusCode::NOT_FOUND), - 6750
} - 6751
} - 6752
- 6753
/// The `Last-Event-ID` a reconnecting client sent, if any. - 6754
/// - 6755
/// Browsers resend it automatically on their own reconnect; the client also - 6756
/// passes it explicitly when it reopens a stream it tore down itself. - 6757
fn resume_from(headers: &axum::http::HeaderMap, uri: &axum::http::Uri) -> Option<u64> { - 6758
headers - 6759
.get("last-event-id") - 6760
.and_then(|v| v.to_str().ok()) - 6761
.and_then(|v| v.trim().parse::<u64>().ok()) - 6762
.or_else(|| { - 6763
uri.query() - 6764
.and_then(|query| { - 6765
query.split('&').find_map(|part| { - 6766
let (key, value) = part.split_once('=')?; - 6767
(key == "last_event_id").then_some(value) - 6768
}) - 6769
}) - 6770
.and_then(|value| value.parse::<u64>().ok()) - 6771
}) - 6772
} - 6773
- 6774
#[cfg(test)] - 6775
#[allow(clippy::unwrap_used, clippy::expect_used)] - 6776
mod resume_cursor_tests { - 6777
use super::resume_from; - 6778
- 6779
#[test] - 6780
fn accepts_native_header_and_manual_reconnect_query_cursor() { - 6781
let mut headers = axum::http::HeaderMap::new(); - 6782
let uri = "/sessions/s/events?last_event_id=17".parse().unwrap(); - 6783
assert_eq!(resume_from(&headers, &uri), Some(17)); - 6784
- 6785
headers.insert("last-event-id", "23".parse().unwrap()); - 6786
assert_eq!(resume_from(&headers, &uri), Some(23)); - 6787
} - 6788
} - 6789
- 6790
async fn events_sse( - 6791
State(state): State<AppState>, - 6792
Path(id): Path<String>, - 6793
headers: axum::http::HeaderMap, - 6794
uri: axum::http::Uri, - 6795
) -> Sse<impl tokio_stream::Stream<Item = Result<Event, std::convert::Infallible>>> { - 6796
use tokio_stream::StreamExt; - 6797
- 6798
let resume = resume_from(&headers, &uri); - 6799
let stream: std::pin::Pin< - 6800
Box<dyn tokio_stream::Stream<Item = Result<Event, std::convert::Infallible>> + Send>, - 6801
> = match ensure_session_handle(&state, &id) - 6802
.await - 6803
.ok() - 6804
.map(|(_, h)| h) - 6805
{ - 6806
Some(h) => Box::pin(stream::agent_frames(&h, resume).map(|frame| { - 6807
Ok(match frame { - 6808
stream::AgentFrame::Event { seq, event } => { - 6809
let data = serde_json::to_string(&event).unwrap_or_else(|error| { - 6810
serde_json::json!({ - 6811
"error": "event serialization failed", - 6812
"detail": error.to_string(), - 6813
}) - 6814
.to_string() - 6815
}); - 6816
Event::default().id(seq.to_string()).data(data) - 6817
} - 6818
stream::AgentFrame::Resync(reason) => Event::default() - 6819
.event("resync") - 6820
.data(serde_json::json!({ "reason": reason }).to_string()), - 6821
}) - 6822
})), - 6823
None => Box::pin(tokio_stream::once(Ok( - 6824
Event::default().data("{\"error\":\"unknown session\"}") - 6825
))), - 6826
}; - 6827
Sse::new(stream).keep_alive(KeepAlive::default()) - 6828
} - 6829
- 6830
async fn presentation_snapshot( - 6831
State(state): State<AppState>, - 6832
Path(id): Path<String>, - 6833
) -> axum::response::Response { - 6834
use axum::response::IntoResponse; - 6835
if let Some(handle) = state.get(&id) { - 6836
let guard = handle - 6837
.session - 6838
.lock() - 6839
.unwrap_or_else(std::sync::PoisonError::into_inner); - 6840
if let Some(session) = guard.as_ref() { - 6841
let planner = delivery::merged_presentation_planner(&handle.core); - 6842
let mut timeline = match presentation_store(&state).load() { - 6843
Ok(library) => { - 6844
let effective = effective_presentation_library( - 6845
&library, - 6846
&handle.core.cwd().to_string_lossy(), - 6847
); - 6848
crate::projection::snapshot_with_planner_and_library( - 6849
&id, session, &planner, &effective, - 6850
) - 6851
} - 6852
Err(_) => crate::projection::snapshot_with_planner(&id, session, &planner), - 6853
}; - 6854
crate::projection::append_sandbox_artifacts( - 6855
&mut timeline, - 6856
&handle.core.sessions_home(), - 6857
&id, - 6858
); - 6859
return Json(timeline).into_response(); - 6860
} - 6861
let mut timeline = handle - 6862
.presentation - 6863
.lock() - 6864
.unwrap_or_else(std::sync::PoisonError::into_inner) - 6865
.clone(); - 6866
timeline.diagnostics.push("run in progress".into()); - 6867
return Json(timeline).into_response(); - 6868
} - 6869
match open_historical_session(&state, &id) { - 6870
Some(session) => { - 6871
let mut timeline = match presentation_store(&state).load() { - 6872
Ok(library) => { - 6873
let owner = session - 6874
.header() - 6875
.map(|header| header.contract_cwd().to_string_lossy().into_owned()) - 6876
.unwrap_or_default(); - 6877
let effective = effective_presentation_library(&library, &owner); - 6878
let planner = delivery::merged_presentation_planner(&state.active_core()); - 6879
crate::projection::snapshot_with_planner_and_library( - 6880
&id, &session, &planner, &effective, - 6881
) - 6882
} - 6883
Err(_) => crate::projection::snapshot(&id, &session), - 6884
}; - 6885
crate::projection::append_sandbox_artifacts( - 6886
&mut timeline, - 6887
&state.core.sessions_home(), - 6888
&id, - 6889
); - 6890
Json(timeline).into_response() - 6891
} - 6892
None => ( - 6893
StatusCode::NOT_FOUND, - 6894
Json(serde_json::json!({ "error": "unknown session" })), - 6895
) - 6896
.into_response(), - 6897
} - 6898
} - 6899
- 6900
/// Fetch one immutable result from the typed presentation projection. This - 6901
/// keeps background notifications addressable without exposing transcript or - 6902
/// transport internals to the client. - 6903
async fn session_result( - 6904
State(state): State<AppState>, - 6905
Path((id, result_id)): Path<(String, String)>, - 6906
) -> axum::response::Response { - 6907
use axum::response::IntoResponse; - 6908
let timeline = if let Some(handle) = state.get(&id) { - 6909
let guard = handle - 6910
.session - 6911
.lock() - 6912
.unwrap_or_else(std::sync::PoisonError::into_inner); - 6913
if let Some(session) = guard.as_ref() { - 6914
let planner = delivery::merged_presentation_planner(&handle.core); - 6915
match presentation_store(&state).load() { - 6916
Ok(library) => { - 6917
let effective = effective_presentation_library( - 6918
&library, - 6919
&handle.core.cwd().to_string_lossy(), - 6920
); - 6921
crate::projection::snapshot_with_planner_and_library( - 6922
&id, session, &planner, &effective, - 6923
) - 6924
} - 6925
Err(_) => crate::projection::snapshot_with_planner(&id, session, &planner), - 6926
} - 6927
} else { - 6928
handle - 6929
.presentation - 6930
.lock() - 6931
.unwrap_or_else(std::sync::PoisonError::into_inner) - 6932
.clone() - 6933
} - 6934
} else if let Some(session) = open_historical_session(&state, &id) { - 6935
match presentation_store(&state).load() { - 6936
Ok(library) => { - 6937
let owner = session - 6938
.header() - 6939
.map(|header| header.contract_cwd().to_string_lossy().into_owned()) - 6940
.unwrap_or_default(); - 6941
let effective = effective_presentation_library(&library, &owner); - 6942
let planner = delivery::merged_presentation_planner(&state.active_core()); - 6943
crate::projection::snapshot_with_planner_and_library( - 6944
&id, &session, &planner, &effective, - 6945
) - 6946
} - 6947
Err(_) => crate::projection::snapshot(&id, &session), - 6948
} - 6949
} else { - 6950
return ( - 6951
StatusCode::NOT_FOUND, - 6952
Json(serde_json::json!({ "error": "unknown session" })), - 6953
) - 6954
.into_response(); - 6955
}; - 6956
match timeline.items.into_iter().find(|item| { - 6957
item.id == result_id - 6958
|| item - 6959
.outcome - 6960
.as_ref() - 6961
.is_some_and(|outcome| outcome.result_id == result_id) - 6962
}) { - 6963
Some(item) => Json(item).into_response(), - 6964
None => ( - 6965
StatusCode::NOT_FOUND, - 6966
Json(serde_json::json!({ "error": "unknown result" })), - 6967
) - 6968
.into_response(), - 6969
} - 6970
} - 6971
- 6972
#[derive(Debug, serde::Deserialize)] - 6973
struct PresentationFeedbackBody { - 6974
choice: String, - 6975
#[serde(default)] - 6976
feedback: Option<String>, - 6977
#[serde(default)] - 6978
chain_id: Option<String>, - 6979
/// The `Presentation` ledger entry id this feedback is about - 6980
/// (docs/design/68-context-engine.md §10: "the user dismissed this - 6981
/// card" is an event about a ledger fact, so feedback keys on the - 6982
/// entry, not on an unscoped choice string). Greenfield: required, no - 6983
/// compatibility path for feedback that names no card. - 6984
presentation_id: String, - 6985
} - 6986
- 6987
#[derive(Debug, serde::Deserialize)] - 6988
struct PresentationSelectionBody { - 6989
#[serde(default)] - 6990
spec_id: String, - 6991
revision: u64, - 6992
#[serde(default)] - 6993
semantic_type: Option<String>, - 6994
#[serde(default = "default_presentation_selection")] - 6995
lifetime: String, - 6996
#[serde(default)] - 6997
scope: Option<vak_presentation::LibraryScope>, - 6998
#[serde(default)] - 6999
owner: Option<String>, - 7000
/// The `Presentation` ledger entry id the user was looking at when they - 7001
/// picked this renderer (docs/design/68-context-engine.md §10). - 7002
presentation_id: String, - 7003
} - 7004
- 7005
fn default_presentation_selection() -> String { - 7006
"use_once".into() - 7007
} - 7008
- 7009
/// Whether `presentation_id` names a real `Presentation` entry in this - 7010
/// session's ledger — live if the runner currently owns the log, otherwise - 7011
/// a read-only reopen from disk (mirrors `session_result`'s fallback). - 7012
fn presentation_entry_exists( - 7013
state: &AppState, - 7014
handle: &SessionHandle, - 7015
id: &str, - 7016
presentation_id: &str, - 7017
) -> bool { - 7018
if let Ok(guard) = handle.session.lock() - 7019
&& let Some(session) = guard.as_ref() - 7020
{ - 7021
return session - 7022
.presentations() - 7023
.into_iter() - 7024
.any(|(entry_id, _)| entry_id == presentation_id); - 7025
} - 7026
open_historical_session(state, id).is_some_and(|session| { - 7027
session - 7028
.presentations() - 7029
.into_iter() - 7030
.any(|(entry_id, _)| entry_id == presentation_id) - 7031
}) - 7032
} - 7033
- 7034
async fn select_presentation_for_session( - 7035
State(state): State<AppState>, - 7036
Path(id): Path<String>, - 7037
Json(body): Json<PresentationSelectionBody>, - 7038
) -> axum::response::Response { - 7039
if body.spec_id.len() > 256 - 7040
|| body - 7041
.semantic_type - 7042
.as_deref() - 7043
.is_some_and(|value| value.len() > 256) - 7044
|| !matches!(body.lifetime.as_str(), "use_once" | "remember") - 7045
|| (body.lifetime == "remember" - 7046
&& (body.scope.is_none() - 7047
|| body.owner.as_deref().unwrap_or_default().trim().is_empty())) - 7048
|| body.presentation_id.trim().is_empty() - 7049
|| body.presentation_id.len() > 256 - 7050
{ - 7051
return StatusCode::BAD_REQUEST.into_response(); - 7052
} - 7053
let Some(handle) = state.get(&id) else { - 7054
return StatusCode::NOT_FOUND.into_response(); - 7055
}; - 7056
if !presentation_entry_exists(&state, &handle, &id, &body.presentation_id) { - 7057
return StatusCode::NOT_FOUND.into_response(); - 7058
} - 7059
let store = presentation_store(&state); - 7060
let library = match store.load() { - 7061
Ok(library) => library, - 7062
Err(error) => { - 7063
return ( - 7064
StatusCode::INTERNAL_SERVER_ERROR, - 7065
Json(serde_json::json!({ "error": error.to_string() })), - 7066
) - 7067
.into_response(); - 7068
} - 7069
}; - 7070
let (spec_id, revision) = if body.spec_id.trim().is_empty() { - 7071
let Some(semantic_type) = body - 7072
.semantic_type - 7073
.as_deref() - 7074
.filter(|value| !value.trim().is_empty()) - 7075
else { - 7076
return StatusCode::BAD_REQUEST.into_response(); - 7077
}; - 7078
let owner = state.core.cwd().to_string_lossy().into_owned(); - 7079
let Some(definition) = library.select_preferred(semantic_type, "builtin", &owner) else { - 7080
return StatusCode::NOT_FOUND.into_response(); - 7081
}; - 7082
(definition.spec.id.clone(), definition.spec.revision) - 7083
} else { - 7084
(body.spec_id.clone(), body.revision) - 7085
}; - 7086
let Some(definition) = library.get(&spec_id, revision) else { - 7087
return StatusCode::NOT_FOUND.into_response(); - 7088
}; - 7089
if body.lifetime == "remember" { - 7090
let Some(scope) = body.scope else { - 7091
return StatusCode::BAD_REQUEST.into_response(); - 7092
}; - 7093
let owner = body.owner.as_deref().unwrap_or_default(); - 7094
if definition.origin.scope != scope || definition.origin.owner != owner { - 7095
return StatusCode::FORBIDDEN.into_response(); - 7096
} - 7097
let result = store.load().and_then(|mut current| { - 7098
current - 7099
.activate(&spec_id, revision, scope, owner) - 7100
.map_err(|error| { - 7101
vak_store::presentation::PresentationStoreError::Invalid(error.to_string()) - 7102
})?; - 7103
store.save(¤t) - 7104
}); - 7105
if result.is_err() { - 7106
return StatusCode::BAD_REQUEST.into_response(); - 7107
} - 7108
} - 7109
let mut data = std::collections::BTreeMap::new(); - 7110
data.insert("spec_id".into(), spec_id); - 7111
data.insert("revision".into(), revision.to_string()); - 7112
data.insert("lifetime".into(), body.lifetime); - 7113
data.insert("presentation_id".into(), body.presentation_id); - 7114
let activity = vak_session::ActivityRecord { - 7115
activity_id: format!("presentation-select-{}", uuid::Uuid::now_v7()), - 7116
turn: None, - 7117
kind: vak_session::ActivityKind::PresentationSelection, - 7118
status: vak_session::ActivityStatus::Succeeded, - 7119
label: "Presentation selected".into(), - 7120
detail: None, - 7121
data, - 7122
}; - 7123
handle - 7124
.activity_buffer - 7125
.lock() - 7126
.unwrap_or_else(std::sync::PoisonError::into_inner) - 7127
.push(activity); - 7128
Json(serde_json::json!({ "selected": true })).into_response() - 7129
} - 7130
- 7131
async fn presentation_feedback( - 7132
State(state): State<AppState>, - 7133
Path(id): Path<String>, - 7134
Json(body): Json<PresentationFeedbackBody>, - 7135
) -> StatusCode { - 7136
if body.choice.trim().is_empty() || body.choice.len() > 128 { - 7137
return StatusCode::BAD_REQUEST; - 7138
} - 7139
if body - 7140
.feedback - 7141
.as_deref() - 7142
.is_some_and(|text| text.len() > 32 * 1024) - 7143
{ - 7144
return StatusCode::BAD_REQUEST; - 7145
} - 7146
if body - 7147
.chain_id - 7148
.as_deref() - 7149
.is_some_and(|chain| chain.trim().is_empty() || chain.len() > 256) - 7150
{ - 7151
return StatusCode::BAD_REQUEST; - 7152
} - 7153
if body.presentation_id.trim().is_empty() || body.presentation_id.len() > 256 { - 7154
return StatusCode::BAD_REQUEST; - 7155
} - 7156
let Some(handle) = state.get(&id) else { - 7157
return StatusCode::NOT_FOUND; - 7158
}; - 7159
if !presentation_entry_exists(&state, &handle, &id, &body.presentation_id) { - 7160
return StatusCode::NOT_FOUND; - 7161
} - 7162
let feedback_denied = matches!( - 7163
body.choice.trim().to_ascii_lowercase().as_str(), - 7164
"keep_original" | "reject" | "dismiss" - 7165
); - 7166
let mut data = std::collections::BTreeMap::new(); - 7167
data.insert("choice".into(), body.choice); - 7168
data.insert("presentation_id".into(), body.presentation_id); - 7169
if let Some(chain_id) = body.chain_id { - 7170
data.insert("chain_id".into(), chain_id); - 7171
} - 7172
if let Some(feedback) = body.feedback { - 7173
data.insert("feedback".into(), feedback); - 7174
} - 7175
let activity = vak_session::ActivityRecord { - 7176
activity_id: uuid::Uuid::now_v7().to_string(), - 7177
turn: None, - 7178
kind: vak_session::ActivityKind::PresentationFeedback, - 7179
status: if feedback_denied { - 7180
vak_session::ActivityStatus::Denied - 7181
} else { - 7182
vak_session::ActivityStatus::Succeeded - 7183
}, - 7184
label: "Presentation feedback".into(), - 7185
detail: None, - 7186
data, - 7187
}; - 7188
handle - 7189
.activity_buffer - 7190
.lock() - 7191
.unwrap_or_else(std::sync::PoisonError::into_inner) - 7192
.push(activity); - 7193
StatusCode::ACCEPTED - 7194
} - 7195
- 7196
async fn presentation_events_sse( - 7197
State(state): State<AppState>, - 7198
Path(id): Path<String>, - 7199
) -> Sse<impl tokio_stream::Stream<Item = Result<Event, std::convert::Infallible>>> { - 7200
use tokio_stream::StreamExt; - 7201
- 7202
// A reconnect's `Last-Event-ID` is deliberately not read: the stream - 7203
// opens on an authoritative snapshot whose id is the new cursor. - 7204
let handle = ensure_session_handle(&state, &id) - 7205
.await - 7206
.ok() - 7207
.map(|(_, h)| h); - 7208
let stream: std::pin::Pin< - 7209
Box<dyn tokio_stream::Stream<Item = Result<Event, std::convert::Infallible>> + Send>, - 7210
> = match stream::presentation_frames(&state, &id, handle) { - 7211
Some(frames) => Box::pin(frames.map(|frame| { - 7212
let event = Event::default().data(frame.json); - 7213
Ok(match frame.sequence { - 7214
Some(sequence) => event.id(sequence.to_string()), - 7215
None => event, - 7216
}) - 7217
})), - 7218
None => Box::pin(tokio_stream::once(Ok( - 7219
Event::default().data("{\"error\":\"unknown session\"}") - 7220
))), - 7221
}; - 7222
Sse::new(stream).keep_alive(KeepAlive::default()) - 7223
} - 7224
- 7225
async fn transcript( - 7226
State(state): State<AppState>, - 7227
Path(id): Path<String>, - 7228
) -> axum::response::Response { - 7229
use axum::response::IntoResponse; - 7230
// `count` is the number of model-visible messages in the derived - 7231
// projection (`derive_messages().len()`), not a raw ledger-entry count. - 7232
// Since v3.0.17 the projection may legitimately exceed the exchanged - 7233
// messages: the multi-turn continuity layer injects a synthetic - 7234
// `<conversation_thread>` (and context compaction may add entries), so a - 7235
// two-message exchange can project as five. Both numbers are - 7236
// reconstructable from the append-only ledger; `count` describes exactly - 7237
// what the model consumed. - 7238
if let Some(handle) = state.get(&id) { - 7239
let guard = handle - 7240
.session - 7241
.lock() - 7242
.unwrap_or_else(std::sync::PoisonError::into_inner); - 7243
let Some(s) = guard.as_ref() else { - 7244
return ( - 7245
StatusCode::CONFLICT, - 7246
Json(serde_json::json!({ "error": "run in progress" })), - 7247
) - 7248
.into_response(); - 7249
}; - 7250
return Json(transcript_json(s)).into_response(); - 7251
} - 7252
match open_historical_session(&state, &id) { - 7253
Some(s) => Json(transcript_json(&s)).into_response(), - 7254
None => ( - 7255
StatusCode::NOT_FOUND, - 7256
Json(serde_json::json!({ "error": "unknown session" })), - 7257
) - 7258
.into_response(), - 7259
} - 7260
} - 7261
- 7262
/// Ledger entry ids of assistant drafts the runtime itself discarded: a - 7263
/// text-only assistant message (no tool call) immediately followed by a - 7264
/// runtime control message that asks for a redo - 7265
/// (`ControlKind::retries_answer`, e.g. a presentation/grounding/freshness - 7266
/// check). The model's real answer is whatever came after the redo; this one - 7267
/// was never shown to the user as final and must not appear as if it were - 7268
/// (docs/audits Finding 3). - 7269
/// - 7270
/// This belongs in `vak-session`'s own transcript derivation so every - 7271
/// consumer gets it for free; it is implemented here for now because - 7272
/// `vak-server` is the only place that currently projects a transcript. - 7273
fn rejected_draft_entry_ids( - 7274
transcript: &[vak_session::TranscriptMessage], - 7275
) -> std::collections::HashSet<String> { - 7276
let mut ids = std::collections::HashSet::new(); - 7277
for pair in transcript.windows(2) { - 7278
let prev = &pair[0]; - 7279
let next = &pair[1]; - 7280
let is_text_only_assistant = prev.message.role == vak_llm::Role::Assistant - 7281
&& !prev - 7282
.message - 7283
.content - 7284
.iter() - 7285
.any(|block| matches!(block, vak_llm::ContentBlock::ToolUse { .. })); - 7286
let next_asks_for_redo = next - 7287
.control - 7288
.is_some_and(vak_intent::control::ControlKind::retries_answer); - 7289
if is_text_only_assistant && next_asks_for_redo { - 7290
ids.insert(prev.entry_id.clone()); - 7291
} - 7292
} - 7293
ids - 7294
} - 7295
- 7296
/// The JSON transcript, built once for the live and the historical path. - 7297
/// - 7298
/// It carries what a person can see and nothing else. The model-visible - 7299
/// projection also holds runtime-authored nudges and derived context blocks - 7300
/// (`<context_summary>`, `<intent>`, …); those are not output, so they are not - 7301
/// sent — the client used to receive them only to strip them again. Likewise - 7302
/// the frozen contract (system prompt and prompt layers), which no client of - 7303
/// this endpoint reads. A discarded draft (see `rejected_draft_entry_ids`) is - 7304
/// excluded the same way. - 7305
/// - 7306
/// `count` is still the model-visible total (`derive_messages().len()`); the - 7307
/// multi-turn continuity layer legitimately makes it exceed `messages`. - 7308
/// `entries` runs parallel to `messages` and gives each one's ledger entry id, - 7309
/// the stable identity the client pairs with the projection's - 7310
/// `provenance.entry_id` instead of counting turns. - 7311
pub(crate) fn transcript_json(s: &SessionLog) -> serde_json::Value { - 7312
let transcript = s.derive_transcript(); - 7313
let rejected_drafts = rejected_draft_entry_ids(&transcript); - 7314
let visible: Vec<&vak_session::TranscriptMessage> = transcript - 7315
.iter() - 7316
.filter(|item| { - 7317
item.control.is_none() - 7318
&& !item.context - 7319
&& !rejected_drafts.contains(item.entry_id.as_str()) - 7320
}) - 7321
.collect(); - 7322
let entries: Vec<serde_json::Value> = visible - 7323
.iter() - 7324
.map(|item| { - 7325
serde_json::json!({ - 7326
"entry_id": item.entry_id, - 7327
"author_id": item.author_id, - 7328
"author_name": item.author_name, - 7329
"attachments": item.attachments, - 7330
}) - 7331
}) - 7332
.collect(); - 7333
let messages: Vec<&vak_llm::Message> = visible.iter().map(|item| &item.message).collect(); - 7334
serde_json::json!({ - 7335
"count": transcript.len(), - 7336
"usage": s.total_usage(), - 7337
"messages": messages, - 7338
"entries": entries, - 7339
}) - 7340
} - 7341
- 7342
/// `SessionLog::derive_conversation` minus discarded drafts (see - 7343
/// `rejected_draft_entry_ids`) — the same exclusion `transcript_json` applies, - 7344
/// kept here rather than in `vak-session` for the reason given there. - 7345
pub(crate) fn conversation_messages(s: &SessionLog) -> Vec<vak_llm::Message> { - 7346
let transcript = s.derive_transcript(); - 7347
let rejected_drafts = rejected_draft_entry_ids(&transcript); - 7348
transcript - 7349
.into_iter() - 7350
.filter(|item| item.control.is_none() && !rejected_drafts.contains(item.entry_id.as_str())) - 7351
.map(|item| item.message) - 7352
.collect() - 7353
} - 7354
- 7355
/// Markdown export over the same projection the JSON transcript serves. - 7356
/// One shared renderer with the TUI export — byte-identical output for the - 7357
/// same session (docs/design/29-personal-os.md P4). - 7358
async fn transcript_markdown( - 7359
State(state): State<AppState>, - 7360
Path(id): Path<String>, - 7361
axum::extract::Query(query): axum::extract::Query<std::collections::HashMap<String, String>>, - 7362
) -> axum::response::Response { - 7363
use axum::response::IntoResponse; - 7364
let format_html = query.get("format").map(|s| s.as_str()) == Some("html"); - 7365
- 7366
let md_opt = if let Some(handle) = state.get(&id) { - 7367
let guard = handle - 7368
.session - 7369
.lock() - 7370
.unwrap_or_else(std::sync::PoisonError::into_inner); - 7371
let Some(s) = guard.as_ref() else { - 7372
return Json(serde_json::json!({ "error": "run in progress" })).into_response(); - 7373
}; - 7374
Some(vak_core::transcript_md::render_markdown( - 7375
&conversation_messages(s), - 7376
)) - 7377
} else { - 7378
open_historical_session(&state, &id) - 7379
.map(|s| vak_core::transcript_md::render_markdown(&conversation_messages(&s))) - 7380
}; - 7381
- 7382
match md_opt { - 7383
Some(md) => { - 7384
if format_html { - 7385
html_response(vak_presentation::transcode_to_html( - 7386
&format!("Session {id}"), - 7387
&md, - 7388
)) - 7389
} else { - 7390
markdown_response(md) - 7391
} - 7392
} - 7393
None => ( - 7394
StatusCode::NOT_FOUND, - 7395
Json(serde_json::json!({ "error": "unknown session" })), - 7396
) - 7397
.into_response(), - 7398
} - 7399
} - 7400
- 7401
#[derive(Debug, serde::Deserialize)] - 7402
struct CoworkingInvitationBody { - 7403
display_name: String, - 7404
#[serde(default = "default_coworking_invitation_hours")] - 7405
expires_in_hours: u32, - 7406
#[serde(default)] - 7407
can_comment: bool, - 7408
#[serde(default)] - 7409
can_message: bool, - 7410
#[serde(default)] - 7411
can_edit: bool, - 7412
} - 7413
- 7414
#[derive(Debug, serde::Deserialize)] - 7415
struct CoworkingMessageBody { - 7416
text: String, - 7417
request_id: String, - 7418
} - 7419
- 7420
#[derive(Debug, serde::Deserialize)] - 7421
struct CoworkingApprovalBody { - 7422
approve: bool, - 7423
} - 7424
- 7425
#[derive(Debug, serde::Deserialize)] - 7426
struct CoworkingDelegationBody { - 7427
grant_id: String, - 7428
} - 7429
- 7430
fn default_coworking_invitation_hours() -> u32 { - 7431
7 * 24 - 7432
} - 7433
- 7434
fn operator_only(principal: &AuthenticatedPrincipal) -> Result<(), StatusCode> { - 7435
match principal { - 7436
AuthenticatedPrincipal::Operator => Ok(()), - 7437
AuthenticatedPrincipal::Participant(_) => Err(StatusCode::FORBIDDEN), - 7438
} - 7439
} - 7440
- 7441
fn conversation_exists(state: &AppState, id: &str) -> bool { - 7442
state.get(id).is_some() || open_historical_session(state, id).is_some() - 7443
} - 7444
- 7445
fn conversation_audience(state: &AppState, id: &str) -> Option<String> { - 7446
read_historical_header(state, id, None)? - 7447
.conversation - 7448
.map(|context| context.audience_id) - 7449
} - 7450
- 7451
async fn coworking_me( - 7452
State(state): State<AppState>, - 7453
Path(conversation_id): Path<String>, - 7454
axum::Extension(principal): axum::Extension<AuthenticatedPrincipal>, - 7455
) -> axum::response::Response { - 7456
use axum::response::IntoResponse; - 7457
let Some(agent) = - 7458
read_historical_header(&state, &conversation_id, None).and_then(|header| header.agent) - 7459
else { - 7460
return ( - 7461
StatusCode::CONFLICT, - 7462
Json(serde_json::json!({ "error": "conversation has no Agent identity" })), - 7463
)
Indexing the workspace…
Vakyartha documentation is discovering safe artifacts, anchors, and source references.