- 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
) - 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(
Indexing the workspace…
Vakyartha documentation is discovering safe artifacts, anchors, and source references.