- 6001
} - 6002
if let Some(value) = routing.target_work_id.as_deref() { - 6003
data.insert("target_work_id".into(), value.into()); - 6004
} - 6005
if let Some(value) = routing.target_result_id.as_deref() { - 6006
data.insert("target_result_id".into(), value.into()); - 6007
} - 6008
if let Some(value) = routing.relation.as_deref() { - 6009
data.insert("relation".into(), value.into()); - 6010
} - 6011
if let Some(value) = routing.outcome_revision { - 6012
data.insert("outcome_revision".into(), value.to_string()); - 6013
} - 6014
if let Some(value) = routing.provenance.as_deref() { - 6015
data.insert("routing_provenance".into(), value.into()); - 6016
} - 6017
} - 6018
- 6019
/// The session that dispatched this one, from its own header: what "an agent - 6020
/// may control only its own children" is checked against. Unknown while the - 6021
/// ledger is out on a running turn, and unknown refuses an agent's control. - 6022
fn target_parent_of(handle: &SessionHandle) -> Option<String> { - 6023
handle.session.lock().ok().and_then(|guard| { - 6024
guard.as_ref().and_then(|log| { - 6025
log.header() - 6026
.and_then(|header| header.parent_session_id.clone()) - 6027
}) - 6028
}) - 6029
} - 6030
- 6031
fn routing_revision_is_current(handle: &SessionHandle, expected: u64) -> bool { - 6032
handle.intent.lock().ok().is_some_and(|record| { - 6033
record - 6034
.as_ref() - 6035
.and_then(|intent| intent.outcome.as_ref()) - 6036
.is_some_and(|outcome| outcome.revision == expected) - 6037
}) - 6038
} - 6039
- 6040
fn routing_summary(routing: Option<&RoutingEnvelope>) -> String { - 6041
let Some(routing) = routing else { - 6042
return String::new(); - 6043
}; - 6044
let mut data = std::collections::BTreeMap::new(); - 6045
record_routing_data(&mut data, routing); - 6046
serde_json::to_string(&data).unwrap_or_default() - 6047
} - 6048
- 6049
fn deny_pending_approvals(handle: &SessionHandle) { - 6050
let requests: Vec<ApprovalRequest> = handle - 6051
.pending - 6052
.lock() - 6053
.unwrap_or_else(std::sync::PoisonError::into_inner) - 6054
.drain() - 6055
.map(|(_, request)| request) - 6056
.collect(); - 6057
for request in requests { - 6058
request.respond(false); - 6059
} - 6060
} - 6061
- 6062
// ---- Worker control plane ------------------------------------------------- - 6063
// - 6064
// Children already stream lifecycle/tool events into the parent session's - 6065
// SSE channel; these endpoints add the missing half: listing, steering, and - 6066
// stopping from a remote surface. Scope-checked against the parent so one - 6067
// session can never touch another's child. - 6068
- 6069
async fn list_workers( - 6070
State(state): State<AppState>, - 6071
Path(id): Path<String>, - 6072
) -> Json<serde_json::Value> { - 6073
let children = state.core.workers().active_for(&id); - 6074
Json(serde_json::json!({ "workers": children })) - 6075
} - 6076
- 6077
#[derive(serde::Deserialize)] - 6078
struct WorkerSteerBody { - 6079
text: String, - 6080
} - 6081
- 6082
async fn steer_worker( - 6083
State(state): State<AppState>, - 6084
Path((id, child)): Path<(String, String)>, - 6085
Json(body): Json<WorkerSteerBody>, - 6086
) -> StatusCode { - 6087
if state.core.workers().parent_of(&child).as_deref() != Some(id.as_str()) { - 6088
return StatusCode::NOT_FOUND; - 6089
} - 6090
if state.core.workers().steer(&child, &body.text) { - 6091
StatusCode::ACCEPTED - 6092
} else { - 6093
StatusCode::CONFLICT - 6094
} - 6095
} - 6096
- 6097
async fn stop_worker( - 6098
State(state): State<AppState>, - 6099
Path((id, child)): Path<(String, String)>, - 6100
) -> StatusCode { - 6101
if state.core.workers().parent_of(&child).as_deref() != Some(id.as_str()) { - 6102
return StatusCode::NOT_FOUND; - 6103
} - 6104
if state.core.workers().stop(&child) { - 6105
StatusCode::ACCEPTED - 6106
} else { - 6107
StatusCode::CONFLICT - 6108
} - 6109
} - 6110
- 6111
#[derive(serde::Deserialize)] - 6112
struct ApprovalBody { - 6113
approve: bool, - 6114
/// "…and don't ask again for calls like this one". Derives the narrowest - 6115
/// rule that covers this call and persists it to - 6116
/// `.vak/permissions.local.toml`. Only meaningful alongside - 6117
/// `approve: true` — remembering a refusal would be a deny rule, which - 6118
/// is a different and much heavier decision than answering one gate. - 6119
#[serde(default)] - 6120
remember: bool, - 6121
} - 6122
- 6123
#[derive(serde::Deserialize)] - 6124
struct OutcomeReviewBody { - 6125
verdict: String, - 6126
#[serde(default)] - 6127
turn: Option<usize>, - 6128
#[serde(default)] - 6129
note: Option<String>, - 6130
} - 6131
- 6132
/// Record an operator's review of a computed outcome. This is deliberately - 6133
/// separate from approval: reviewing a result never authorizes a tool call. - 6134
async fn record_outcome_review( - 6135
State(state): State<AppState>, - 6136
Path(id): Path<String>, - 6137
Json(body): Json<OutcomeReviewBody>, - 6138
) -> axum::response::Response { - 6139
use axum::response::IntoResponse; - 6140
if !matches!( - 6141
body.verdict.as_str(), - 6142
"accepted" | "needs_work" | "rejected" - 6143
) { - 6144
return (StatusCode::BAD_REQUEST, "invalid outcome review verdict").into_response(); - 6145
} - 6146
let Some(handle) = state.get(&id) else { - 6147
return StatusCode::NOT_FOUND.into_response(); - 6148
}; - 6149
let Ok(mut guard) = handle.session.lock() else { - 6150
return StatusCode::CONFLICT.into_response(); - 6151
}; - 6152
let Some(session) = guard.as_mut() else { - 6153
return StatusCode::NOT_FOUND.into_response(); - 6154
}; - 6155
let reviewed_turn = body.turn.or_else(|| { - 6156
session - 6157
.chain_to_root() - 6158
.iter() - 6159
.filter_map(|entry| match &entry.payload { - 6160
vak_session::EntryPayload::Activity(activity) - 6161
if activity.label == "Outcome evaluation" => - 6162
{ - 6163
activity.turn - 6164
} - 6165
_ => None, - 6166
}) - 6167
.next_back() - 6168
}); - 6169
if body.turn.is_some() && reviewed_turn.is_none() { - 6170
return (StatusCode::NOT_FOUND, "outcome evaluation turn not found").into_response(); - 6171
} - 6172
let activity = vak_session::ActivityRecord { - 6173
activity_id: format!("outcome-review-{}", uuid::Uuid::now_v7()), - 6174
turn: reviewed_turn, - 6175
kind: vak_session::ActivityKind::Diagnostic, - 6176
status: vak_session::ActivityStatus::Succeeded, - 6177
label: "Outcome review".into(), - 6178
detail: body.note, - 6179
data: std::collections::BTreeMap::from([("verdict".into(), body.verdict)]), - 6180
}; - 6181
if let Err(error) = session.append_activity(activity) { - 6182
return (StatusCode::INTERNAL_SERVER_ERROR, error.to_string()).into_response(); - 6183
} - 6184
Json(serde_json::json!({ "recorded": true })).into_response() - 6185
} - 6186
- 6187
/// `POST /sessions/{id}/approvals/{req_id}` — answer one gate. - 6188
/// - 6189
/// `remember` is the mechanism `docs/design/08-permissions.md` has described - 6190
/// since the engine shipped and nothing ever called: `Core::learn_allow_rule` - 6191
/// existed, was tested, wrote a validated file, and had no caller on any - 6192
/// surface, because the interactive approver its comment referred to was - 6193
/// never built. This is that caller. - 6194
/// - 6195
/// Remembering never blocks the answer. The gate is resolved first; a - 6196
/// failure to derive or persist a rule is reported alongside a successful - 6197
/// approval, because the run is already waiting and a bookkeeping problem - 6198
/// must not become a denial. - 6199
async fn answer_approval( - 6200
State(state): State<AppState>, - 6201
Path((id, req_id)): Path<(String, String)>, - 6202
Json(body): Json<ApprovalBody>, - 6203
) -> axum::response::Response { - 6204
use axum::response::IntoResponse; - 6205
// Look the request up in THIS session's pending map only: approvals - 6206
// are never resolvable across sessions. - 6207
let Some(handle) = state.get(&id) else { - 6208
return StatusCode::NOT_FOUND.into_response(); - 6209
}; - 6210
let Some(req) = handle - 6211
.pending - 6212
.lock() - 6213
.unwrap_or_else(std::sync::PoisonError::into_inner) - 6214
.remove(&req_id) - 6215
else { - 6216
return StatusCode::NOT_FOUND.into_response(); - 6217
}; - 6218
let tool = req.tool.clone(); - 6219
let args_json = req.args_json.clone(); - 6220
*req.answered_by - 6221
.lock() - 6222
.unwrap_or_else(std::sync::PoisonError::into_inner) = - 6223
Some(("operator".into(), "You".into())); - 6224
req.respond(body.approve); - 6225
- 6226
let mut learned: Option<String> = None; - 6227
let mut learn_error: Option<String> = None; - 6228
if body.remember && body.approve { - 6229
match serde_json::from_str::<serde_json::Value>(&args_json) { - 6230
// The gate's own core, not the server's: a learned rule belongs - 6231
// to the workspace whose call raised it. - 6232
Ok(args) => match handle.core.learn_from_call(&tool, &args) { - 6233
Ok(spec) => { - 6234
vak_core::security_events::record( - 6235
&state.core.sessions_home(), - 6236
vak_core::security_events::EventKind::ConfigChange, - 6237
"permission_rule_learned", - 6238
&format!("session={id} rule={spec}"), - 6239
None, - 6240
); - 6241
learned = Some(spec); - 6242
} - 6243
Err(error) => learn_error = Some(error.to_string()), - 6244
}, - 6245
Err(error) => learn_error = Some(format!("unreadable tool arguments: {error}")), - 6246
} - 6247
} - 6248
( - 6249
StatusCode::OK, - 6250
Json(serde_json::json!({ - 6251
"approved": body.approve, - 6252
"learned_rule": learned, - 6253
"learn_error": learn_error, - 6254
})), - 6255
) - 6256
.into_response() - 6257
} - 6258
- 6259
/// Dispatch forensics (docs/design/42-managed-work-contracts.md): the session's work - 6260
/// receipts, newest last. - 6261
async fn session_receipts( - 6262
State(state): State<AppState>, - 6263
Path(id): Path<String>, - 6264
) -> Result<Json<Vec<vak_llm::WorkReceipt>>, StatusCode> { - 6265
let Some(handle) = state.get(&id) else { - 6266
return Err(StatusCode::NOT_FOUND); - 6267
}; - 6268
let Ok(session) = handle.session.lock() else { - 6269
return Err(StatusCode::NOT_FOUND); - 6270
}; - 6271
match session.as_ref() { - 6272
Some(log) => Ok(Json(log.receipts().into_iter().cloned().collect())), - 6273
None => Err(StatusCode::NOT_FOUND), - 6274
} - 6275
} - 6276
- 6277
async fn session_work( - 6278
State(state): State<AppState>, - 6279
Path(id): Path<String>, - 6280
) -> Result<Json<Option<vak_session::work::WorkProjection>>, StatusCode> { - 6281
if let Some(handle) = state.get(&id) { - 6282
let guard = handle - 6283
.session - 6284
.lock() - 6285
.unwrap_or_else(std::sync::PoisonError::into_inner); - 6286
let Some(session) = guard.as_ref() else { - 6287
return Err(StatusCode::NOT_FOUND); - 6288
}; - 6289
return session - 6290
.work_projection() - 6291
.map(Json) - 6292
.map_err(|_| StatusCode::INTERNAL_SERVER_ERROR); - 6293
} - 6294
let Some(session) = open_historical_session(&state, &id) else { - 6295
return Err(StatusCode::NOT_FOUND); - 6296
}; - 6297
session - 6298
.work_projection() - 6299
.map(Json) - 6300
.map_err(|_| StatusCode::INTERNAL_SERVER_ERROR) - 6301
} - 6302
- 6303
#[derive(serde::Deserialize)] - 6304
#[serde(tag = "operation", rename_all = "snake_case")] - 6305
enum WorkCommand { - 6306
Transition { - 6307
item_id: String, - 6308
to: vak_session::types::WorkItemStatus, - 6309
#[serde(default)] - 6310
reason: String, - 6311
}, - 6312
Cancel { - 6313
#[serde(default)] - 6314
reason: String, - 6315
}, - 6316
Resume { - 6317
#[serde(default)] - 6318
reason: String, - 6319
}, - 6320
Retry { - 6321
item_id: String, - 6322
#[serde(default)] - 6323
reason: String, - 6324
}, - 6325
ResolveAssumption { - 6326
assumption_id: String, - 6327
resolution: String, - 6328
}, - 6329
AttachEvidence { - 6330
item_id: String, - 6331
evidence: vak_session::types::EvidenceRef, - 6332
}, - 6333
Assign { - 6334
item_id: String, - 6335
owner: vak_session::types::WorkOwner, - 6336
child_session_id: Option<String>, - 6337
}, - 6338
Revise { - 6339
contract: vak_session::types::WorkContract, - 6340
#[serde(default)] - 6341
reason: String, - 6342
}, - 6343
} - 6344
- 6345
#[derive(serde::Deserialize)] - 6346
struct WorkReason { - 6347
#[serde(default)] - 6348
reason: String, - 6349
} - 6350
- 6351
#[derive(serde::Deserialize)] - 6352
struct WorkReassign { - 6353
owner: vak_session::types::WorkOwner, - 6354
#[serde(default)] - 6355
child_session_id: Option<String>, - 6356
} - 6357
- 6358
async fn session_work_confirm( - 6359
state: State<AppState>, - 6360
path: Path<String>, - 6361
body: Option<Json<WorkReason>>, - 6362
) -> axum::response::Response { - 6363
session_work_command( - 6364
state, - 6365
path, - 6366
Json(WorkCommand::Resume { - 6367
reason: body - 6368
.map(|body| body.0.reason) - 6369
.unwrap_or_else(|| "confirmed by operator".into()), - 6370
}), - 6371
) - 6372
.await - 6373
} - 6374
- 6375
async fn session_work_revise( - 6376
state: State<AppState>, - 6377
path: Path<String>, - 6378
Json(body): Json<WorkCommand>, - 6379
) -> axum::response::Response { - 6380
session_work_command(state, path, Json(body)).await - 6381
} - 6382
- 6383
async fn session_work_retry( - 6384
state: State<AppState>, - 6385
Path((id, item_id)): Path<(String, String)>, - 6386
body: Option<Json<WorkReason>>, - 6387
) -> axum::response::Response { - 6388
session_work_command( - 6389
state, - 6390
Path(id), - 6391
Json(WorkCommand::Retry { - 6392
item_id, - 6393
reason: body - 6394
.map(|body| body.0.reason) - 6395
.unwrap_or_else(|| "retry requested by operator".into()), - 6396
}), - 6397
) - 6398
.await - 6399
} - 6400
- 6401
async fn session_work_cancel_item( - 6402
state: State<AppState>, - 6403
Path((id, item_id)): Path<(String, String)>, - 6404
body: Option<Json<WorkReason>>, - 6405
) -> axum::response::Response { - 6406
session_work_command( - 6407
state, - 6408
Path(id), - 6409
Json(WorkCommand::Transition { - 6410
item_id, - 6411
to: vak_session::types::WorkItemStatus::Cancelled, - 6412
reason: body - 6413
.map(|body| body.0.reason) - 6414
.unwrap_or_else(|| "cancelled by operator".into()), - 6415
}), - 6416
) - 6417
.await - 6418
} - 6419
- 6420
async fn session_work_reassign( - 6421
state: State<AppState>, - 6422
Path((id, item_id)): Path<(String, String)>, - 6423
Json(body): Json<WorkReassign>, - 6424
) -> axum::response::Response { - 6425
session_work_command( - 6426
state, - 6427
Path(id), - 6428
Json(WorkCommand::Assign { - 6429
item_id, - 6430
owner: body.owner, - 6431
child_session_id: body.child_session_id, - 6432
}), - 6433
) - 6434
.await - 6435
} - 6436
- 6437
/// Apply an operator/user work action to the same append-only ledger used by - 6438
/// the agent. The current projection supplies both the expected state and - 6439
/// revision, so stale clients receive a conflict instead of silently - 6440
/// overwriting a newer action. - 6441
async fn session_work_command( - 6442
State(state): State<AppState>, - 6443
Path(id): Path<String>, - 6444
Json(command): Json<WorkCommand>, - 6445
) -> axum::response::Response { - 6446
use axum::response::IntoResponse; - 6447
let Some(handle) = state.get(&id) else { - 6448
return StatusCode::NOT_FOUND.into_response(); - 6449
}; - 6450
let mut guard = handle - 6451
.session - 6452
.lock() - 6453
.unwrap_or_else(std::sync::PoisonError::into_inner); - 6454
let Some(session) = guard.as_mut() else { - 6455
return ( - 6456
StatusCode::CONFLICT, - 6457
Json(serde_json::json!({ "error": "session is still running" })), - 6458
) - 6459
.into_response(); - 6460
}; - 6461
let Ok(Some(projection)) = session.work_projection() else { - 6462
return ( - 6463
StatusCode::CONFLICT, - 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
Indexing the workspace…
Vakyartha documentation is discovering safe artifacts, anchors, and source references.