- 1001
use futures::StreamExt; - 1002
let mut stream = res.bytes_stream(); - 1003
let deadline = std::time::Instant::now() + Duration::from_secs(10); - 1004
while std::time::Instant::now() < deadline { - 1005
if let Some(Ok(chunk)) = stream.next().await { - 1006
let text = String::from_utf8_lossy(&chunk).into_owned(); - 1007
if text.contains("StreamOpened") { - 1008
if let Some(t) = opened_tx.take() { - 1009
let _ = t.send(()); - 1010
} - 1011
break; - 1012
} - 1013
} - 1014
} - 1015
}); - 1016
let _ = tokio::time::timeout(Duration::from_secs(5), opened_rx).await; - 1017
- 1018
let start = std::time::Instant::now(); - 1019
let res = client - 1020
.post(format!("{base}/sessions/{session_id}/run")) - 1021
.json(&serde_json::json!({"prompt": "go"})) - 1022
.send() - 1023
.await - 1024
.unwrap(); - 1025
let elapsed = start.elapsed(); - 1026
assert_eq!(res.status(), 202); - 1027
assert!( - 1028
elapsed < Duration::from_millis(1500), - 1029
"an already-attached subscriber must skip the 2s attach wait entirely, took {elapsed:?}" - 1030
); - 1031
} - 1032
- 1033
/// Finding 4: validation must happen BEFORE any side effect — a rejected - 1034
/// request must never write a durable admission activity, insert into the - 1035
/// admissions set, or broadcast a synthesized `RunFinished` for a run that - 1036
/// never started. - 1037
#[tokio::test(flavor = "multi_thread", worker_threads = 4)] - 1038
async fn rejected_request_writes_no_durable_admission() { - 1039
let provider = Arc::new(Scripted { - 1040
responses: Mutex::new(VecDeque::new()), - 1041
}); - 1042
let (base, _server) = spawn_server(provider, vak_config::PermissionMode::FullAccess).await; - 1043
let client = reqwest::Client::new(); - 1044
let session_id: String = client - 1045
.post(format!("{base}/sessions")) - 1046
.send() - 1047
.await - 1048
.unwrap() - 1049
.json::<serde_json::Value>() - 1050
.await - 1051
.unwrap()["session_id"] - 1052
.as_str() - 1053
.unwrap() - 1054
.to_string(); - 1055
- 1056
// Goal mode explicitly refuses attachments; this must be rejected - 1057
// before any admission is ever recorded. - 1058
let res = client - 1059
.post(format!("{base}/sessions/{session_id}/run")) - 1060
.json(&serde_json::json!({ - 1061
"prompt": "do a thing", - 1062
"request_id": "rejected-1", - 1063
"goal": "finish the thing", - 1064
"attachments": [{"mime": "image/png", "data": "AAAA"}], - 1065
})) - 1066
.send() - 1067
.await - 1068
.unwrap(); - 1069
assert_eq!(res.status(), 400); - 1070
- 1071
let transcript: serde_json::Value = client - 1072
.get(format!("{base}/sessions/{session_id}/transcript")) - 1073
.send() - 1074
.await - 1075
.unwrap() - 1076
.json() - 1077
.await - 1078
.unwrap(); - 1079
assert_eq!( - 1080
transcript["count"].as_u64(), - 1081
Some(0), - 1082
"a rejected request must leave no durable entry behind: {transcript}" - 1083
); - 1084
- 1085
// The same request_id must be admittable afterward — nothing was - 1086
// parked in the in-memory admissions guard either. - 1087
let retry = client - 1088
.post(format!("{base}/sessions/{session_id}/run")) - 1089
.json(&serde_json::json!({"prompt": "do a thing", "request_id": "rejected-1"})) - 1090
.send() - 1091
.await - 1092
.unwrap(); - 1093
assert_eq!(retry.status(), 202); - 1094
let retry_body: serde_json::Value = retry.json().await.unwrap(); - 1095
assert_eq!(retry_body["state"], "started"); - 1096
} - 1097
- 1098
/// A provider that answers once and keeps every request it was sent. - 1099
struct Recording { - 1100
requests: Arc<Mutex<Vec<ChatRequest>>>, - 1101
} - 1102
- 1103
#[async_trait::async_trait] - 1104
impl Provider for Recording { - 1105
fn name(&self) -> &str { - 1106
"scripted" - 1107
} - 1108
- 1109
async fn stream( - 1110
&self, - 1111
request: ChatRequest, - 1112
_cancel: CancellationToken, - 1113
) -> Result<EventStream, LlmError> { - 1114
self.requests.lock().unwrap().push(request); - 1115
let (mut sink, rx) = stream::channel(64); - 1116
let m = text("It is a three-slide deck."); - 1117
sink.push(stream::StreamEvent::Start { partial: m.clone() }); - 1118
sink.close_message(m).await; - 1119
Ok(rx) - 1120
} - 1121
} - 1122
- 1123
/// A file dropped on the conversation (docs/design/72, "File in") is saved - 1124
/// to the workspace inbox, reaches the model as a note naming where it is - 1125
/// and never as its bytes, and reaches the client as a typed attachment of - 1126
/// that message, so the chat draws a file rather than the note. - 1127
#[tokio::test(flavor = "multi_thread", worker_threads = 4)] - 1128
async fn a_dropped_file_reaches_the_model_as_a_note_and_the_chat_as_a_file() { - 1129
let requests = Arc::new(Mutex::new(Vec::new())); - 1130
let provider = Arc::new(Recording { - 1131
requests: requests.clone(), - 1132
}); - 1133
let (base, _server) = spawn_server(provider, vak_config::PermissionMode::WorkspaceWrite).await; - 1134
let client = reqwest::Client::new(); - 1135
let bytes = b"PK\x03\x04\x14\x00 deck bytes \x00\x01\x02".to_vec(); - 1136
let uploaded: serde_json::Value = client - 1137
.post(format!("{base}/fs/inbox?name=Q3%20deck.pptx")) - 1138
.body(bytes.clone()) - 1139
.send() - 1140
.await - 1141
.unwrap() - 1142
.json() - 1143
.await - 1144
.unwrap(); - 1145
let path = uploaded["path"].as_str().unwrap().to_string(); - 1146
assert!( - 1147
path.starts_with("inbox/") && path.ends_with("-Q3 deck.pptx"), - 1148
"{uploaded}" - 1149
); - 1150
assert_eq!(uploaded["name"], "Q3 deck.pptx"); - 1151
assert_eq!(uploaded["bytes"], bytes.len()); - 1152
- 1153
let session_id: String = client - 1154
.post(format!("{base}/sessions")) - 1155
.send() - 1156
.await - 1157
.unwrap() - 1158
.json::<serde_json::Value>() - 1159
.await - 1160
.unwrap()["session_id"] - 1161
.as_str() - 1162
.unwrap() - 1163
.to_string(); - 1164
for outside in [ - 1165
"../secret.pptx", - 1166
"deck.pptx", - 1167
"inbox/../deck.pptx", - 1168
"inbox/missing.pptx", - 1169
] { - 1170
let refused = client - 1171
.post(format!("{base}/sessions/{session_id}/run")) - 1172
.json(&serde_json::json!({"prompt": "read it", "files": [outside]})) - 1173
.send() - 1174
.await - 1175
.unwrap(); - 1176
assert_eq!(refused.status(), 400, "{outside} must be refused"); - 1177
} - 1178
let res = client - 1179
.post(format!("{base}/sessions/{session_id}/run")) - 1180
.json(&serde_json::json!({"prompt": "what is in it?", "files": [path]})) - 1181
.send() - 1182
.await - 1183
.unwrap(); - 1184
assert_eq!(res.status(), 202); - 1185
- 1186
let deadline = std::time::Instant::now() + Duration::from_secs(10); - 1187
let raw = poll_transcript_contains(&client, &base, &session_id, "three-slide", deadline).await; - 1188
let transcript: serde_json::Value = serde_json::from_str(&raw).unwrap(); - 1189
let attachment = &transcript["entries"][0]["attachments"][0]; - 1190
assert_eq!(attachment["path"], path.as_str(), "{transcript}"); - 1191
assert_eq!(attachment["name"], "Q3 deck.pptx"); - 1192
assert_eq!(attachment["block"], 1); - 1193
let note = transcript["messages"][0]["content"][1]["text"] - 1194
.as_str() - 1195
.unwrap(); - 1196
assert!( - 1197
note.contains(&format!("at path \"{path}\"")) && note.contains("doc_read"), - 1198
"{note}" - 1199
); - 1200
assert_eq!( - 1201
transcript["messages"][0]["content"][0]["text"], - 1202
"what is in it?" - 1203
); - 1204
- 1205
let requests = requests.lock().unwrap(); - 1206
let sent = serde_json::to_string(&requests.last().unwrap().messages).unwrap(); - 1207
assert!(sent.contains(&format!("at path \\\"{path}\\\"")), "{sent}"); - 1208
assert!( - 1209
!sent.contains("deck bytes"), - 1210
"the file's bytes reached the model" - 1211
); - 1212
} - 1213
Indexing the workspace…
Vakyartha documentation is discovering safe artifacts, anchors, and source references.