Oregami
Repositories/oxedyne/daimond

oxedyne/daimond/src/handler.rs

39.6 KiB, 1 run

created by r2519314175:949, which is this file's identity for as long as the history lasts, whatever it is later renamed to

download · who wrote it · its history

1//! Daimond WS handler — chat protocol over WebSocket.
2//!
3//! Connected via Steel's path-based WS dispatch in https.rs when the
4//! path is `/chat`. Handles the full chat lifecycle:
5//!
6//! - Receives user messages (syntax commands from o3db.js)
7//! - Runs the agent loop
8//! - Streams LLM response tokens back as binary WS messages
9
10use oxedyne_fe2o3_core::prelude::*;
11use oxedyne_fe2o3_jdat::{
12 prelude::*,
13 id::NumIdDat,
14};
15use oxedyne_fe2o3_net::ws::core::{WebSocket, WebSocketMessage};
16use oxedyne_fe2o3_iop_crypto::enc::Encrypter;
17use oxedyne_fe2o3_iop_db::api::Database;
18use oxedyne_fe2o3_iop_hash::api::Hasher;
19use oxedyne_fe2o3_syntax::{
20 SyntaxRef,
21 msg::{Msg, MsgCmd},
22};
23
24use std::sync::{Arc, RwLock};
25
26use crate::agent::Agent;
27use crate::executor::Executor;
28use crate::protocol::{AgentEvent, Session};
29use crate::session::SessionStore;
30use crate::tools::{Tool, ToolContext, ToolRegistry};
31use crate::workspace::Workspace;
32use crate::llm::datmap_to_json;
33use std::path::PathBuf;
34use std::pin::pin;
35
36
37// ┌───────────────────────────────────────────────────────────────┐
38// │ DaimondHandler │
39// └───────────────────────────────────────────────────────────────┘
40
41/// Per-vhost configuration for Daimond, parsed from Steel's config.jdat.
42#[derive(Clone, Debug, Eq, PartialEq)]
43pub struct DaimondConfig {
44 pub llm_host: String,
45 pub llm_port: u16,
46 pub llm_path: String,
47 pub llm_key: String,
48 pub llm_model: String,
49 pub max_tokens: u32,
50 pub system_prompt: String,
51 /// Base directory (relative to the app root, or absolute) under
52 /// which each user's workspace lives as `<workspace_root>/<user>`.
53 pub workspace_root: String,
54}
55
56impl DaimondConfig {
57
58 pub fn from_datmap(m: &DaticleMap) -> Outcome<Self> {
59 let llm_host = match m.get(&dat!("llm_host")) {
60 Some(Dat::Str(s)) => s.clone(),
61 _ => "api.fireworks.ai".to_string(),
62 };
63 let llm_port = match m.get(&dat!("llm_port")) {
64 Some(Dat::U16(n)) => *n,
65 Some(Dat::U64(n)) => *n as u16,
66 _ => 443,
67 };
68 let llm_path = match m.get(&dat!("llm_path")) {
69 Some(Dat::Str(s)) => s.clone(),
70 _ => "/inference/v1/chat/completions".to_string(),
71 };
72 let llm_key = match m.get(&dat!("llm_key")) {
73 Some(Dat::Str(s)) => s.clone(),
74 _ => return Err(err!("DaimondConfig: 'llm_key' is required."; Invalid, Input, Missing)),
75 };
76 let llm_model = match m.get(&dat!("llm_model")) {
77 Some(Dat::Str(s)) => s.clone(),
78 _ => "accounts/fireworks/models/glm-5p2".to_string(),
79 };
80 let max_tokens = match m.get(&dat!("max_tokens")) {
81 Some(Dat::U32(n)) => *n,
82 Some(Dat::U64(n)) => *n as u32,
83 Some(Dat::U16(n)) => *n as u32,
84 _ => 4096, // sensible default; caps runaway reasoning loops
85 };
86 let system_prompt = match m.get(&dat!("system_prompt")) {
87 Some(Dat::Str(s)) => s.clone(),
88 _ => "You are Daimond, an AI assistant.".to_string(),
89 };
90 let workspace_root = match m.get(&dat!("workspace_root")) {
91 Some(Dat::Str(s)) => s.clone(),
92 _ => "workspaces".to_string(),
93 };
94 Ok(Self {
95 llm_host, llm_port, llm_path, llm_key, llm_model, max_tokens,
96 system_prompt, workspace_root,
97 })
98 }
99}
100
101
102/// State for the Daimond WS handler, shared across connections on a vhost.
103#[derive(Clone)]
104pub struct DaimondState<
105 const UIDL: usize,
106 UID: NumIdDat<UIDL> + Clone + 'static,
107 ENC: Encrypter + 'static,
108 KH: Hasher + 'static,
109 DB: Database<UIDL, UID, ENC, KH> + 'static,
110> {
111 pub agent: Agent,
112 pub session_store: SessionStore<UIDL, UID, ENC, KH, DB>,
113 /// Absolute base directory under which per-user workspaces live.
114 pub workspace_base: PathBuf,
115 /// Command-execution backend for the shell tool.
116 pub executor: Executor,
117}
118
119// ┌───────────────────────────────────────────────────────────────┐
120// │ Chat WS handler │
121// └───────────────────────────────────────────────────────────────┘
122
123/// Handle a chat WebSocket connection.
124///
125/// This is a path-based handler (like the terminal bridge) that
126/// manages its own WS read/write loop. It receives syntax commands
127/// from the client, executes session/chat operations, and streams
128/// agent responses back.
129///
130/// The connection uses Steel's existing session cookie for auth —
131/// the `sid` parameter identifies the user's session.
132pub async fn handle_chat_websocket<
133 const UIDL: usize,
134 UID: NumIdDat<UIDL> + Clone + 'static,
135 ENC: Encrypter + 'static,
136 KH: Hasher + 'static,
137 DB: Database<UIDL, UID, ENC, KH> + 'static,
138 S: tokio::io::AsyncRead + tokio::io::AsyncWrite + Unpin,
139>(
140 mut stream: S,
141 state: DaimondState<UIDL, UID, ENC, KH, DB>,
142 sid: Option<String>,
143 vhost_db: Option<(Arc<RwLock<DB>>, UID)>,
144 request: oxedyne_fe2o3_net::http::msg::HttpMessage,
145 id: &String,
146)
147 -> Outcome<()>
148{
149 // ── WebSocket handshake ───────────────────────────────────
150 let chunk_size = 65536;
151 let chunk_thresh = 32768;
152 let mut ws: WebSocket<
153 '_,
154 UIDL, UID, ENC, KH, DB,
155 S,
156 oxedyne_fe2o3_net::ws::handler::WebSocketEchoHandler,
157 > = WebSocket::new_server(
158 &mut stream,
159 oxedyne_fe2o3_net::ws::handler::WebSocketEchoHandler,
160 chunk_size,
161 chunk_thresh,
162 );
163 match ws.connect_as_server(request).await {
164 Ok(()) => {
165 debug!("{}: Daimond WS handshake completed.", id);
166 }
167 Err(e) => return Err(err!(e,
168 "{}: Daimond WS handshake failed.", id;
169 IO, Network, Wire)),
170 }
171
172 // ── Build Daimond's own syntax ────────────────────────────────
173 let syntax = match crate::syntax::build_syntax() {
174 Ok(s) => s,
175 Err(e) => {
176 error!(e, "{}: Daimond: failed to build syntax.", id);
177 return Ok(());
178 }
179 };
180
181 // ── Determine the authenticated username ──────────────────
182 let username = match get_username(&vhost_db, &sid) {
183 Some(u) => {
184 info!("{}: Daimond WS authenticated as '{}'.", id, u);
185 u
186 }
187 None => {
188 info!("{}: Daimond WS not authenticated (sid={:?}).", id, sid);
189 let _ = ws.send(&text_msg(syntax.clone(), "error", vec![
190 dat!("Not authenticated. Please log in first."),
191 ])).await;
192 return Ok(());
193 }
194 };
195
196 // ── Chat loop ─────────────────────────────────────────────
197 //
198 // We read WS messages, parse them as syntax commands, and
199 // dispatch:
200 // session_new → create a new session
201 // session_list → list user's sessions
202 // session_switch → set current session
203 // session_close → delete a session
204 // chat → run agent turn (streams response)
205 // file_list → list sandbox files (Phase 2)
206
207 let mut current_session_id: Option<String> = None;
208 let mut current_session: Option<Session> = None;
209
210 loop {
211 let msg = match ws.read().await {
212 Ok(Some(WebSocketMessage::Text(txt))) => {
213 debug!("{}: Daimond WS received text: '{}'", id, txt);
214 txt
215 }
216 Ok(Some(WebSocketMessage::Binary(_))) => continue,
217 Ok(Some(WebSocketMessage::Close(_, _))) => {
218 info!("{}: Daimond WS client closed.", id);
219 break;
220 }
221 Ok(None) => {
222 info!("{}: Daimond WS client disconnected.", id);
223 break;
224 }
225 Ok(_) => continue,
226 Err(e) => {
227 warn!("{}: Daimond WS read error: {}", id, e);
228 break;
229 }
230 };
231
232 // Parse syntax command.
233 let msgrx = Msg::new(syntax.clone());
234 let msgrx = match msgrx.from_str(&msg, None) {
235 Ok(m) => m,
236 Err(e) => {
237 let _ = ws.send(&text_msg(syntax.clone(), "error", vec![
238 dat!(e.to_string()),
239 ])).await;
240 continue;
241 }
242 };
243 if msgrx.cmds.len() != 1 {
244 let _ = ws.send(&text_msg(syntax.clone(), "error", vec![
245 dat!("Expected one command."),
246 ])).await;
247 continue;
248 }
249 let (cmd_name, cmdrx) = match msgrx.cmds.into_iter().next() {
250 Some(v) => v,
251 None => continue, // length checked == 1 above
252 };
253
254 debug!("{}: Daimond WS command '{}'.", id, cmd_name);
255
256 match cmd_name.as_str() {
257 "session_new" => {
258 let name = match cmdrx.vals.get(0) {
259 Some(Dat::Str(s)) => s.clone(),
260 _ => fmt!("Session"),
261 };
262 let model = match cmdrx.vals.get(1) {
263 Some(Dat::Str(s)) => s.clone(),
264 _ => state.agent.llm.model.clone(),
265 };
266 match state.session_store.create_session(
267 &username, &name, &model,
268 ) {
269 Ok(session) => {
270 current_session_id = Some(session.id.clone());
271 current_session = Some(session.clone());
272 let _ = ws.send(&text_msg(syntax.clone(), "data", vec![
273 dat!(datmap_to_json(&session.to_meta_datmap())),
274 ])).await;
275 }
276 Err(e) => {
277 let _ = ws.send(&text_msg(syntax.clone(), "error", vec![
278 dat!(e.to_string()),
279 ])).await;
280 }
281 }
282 }
283 "session_list" => {
284 match state.session_store.list_sessions(&username) {
285 Ok(sessions) => {
286 let list: Vec<Dat> = sessions.iter()
287 .map(|s| Dat::Map(s.to_meta_datmap()))
288 .collect();
289 let mut m = DaticleMap::new();
290 m.insert(dat!("sessions"), Dat::List(list));
291 let json = crate::llm::datmap_to_json(&m);
292 let _ = ws.send(&text_msg(syntax.clone(), "data", vec![
293 dat!(json),
294 ])).await;
295 }
296 Err(e) => {
297 let _ = ws.send(&text_msg(syntax.clone(), "error", vec![
298 dat!(e.to_string()),
299 ])).await;
300 }
301 }
302 }
303 "session_switch" => {
304 let session_id = match cmdrx.vals.get(0) {
305 Some(Dat::Str(s)) => s.clone(),
306 _ => {
307 let _ = ws.send(&text_msg(syntax.clone(), "error", vec![
308 dat!("session_switch: missing session id."),
309 ])).await;
310 continue;
311 }
312 };
313 match state.session_store.get_session(&session_id) {
314 Ok(session) => {
315 current_session_id = Some(session.id.clone());
316 current_session = Some(session.clone());
317 let json = crate::llm::datmap_to_json(&session.to_datmap());
318 let _ = ws.send(&text_msg(syntax.clone(), "data", vec![
319 dat!(json),
320 ])).await;
321 }
322 Err(e) => {
323 let _ = ws.send(&text_msg(syntax.clone(), "error", vec![
324 dat!(e.to_string()),
325 ])).await;
326 }
327 }
328 }
329 "session_close" => {
330 let session_id = match cmdrx.vals.get(0) {
331 Some(Dat::Str(s)) => s.clone(),
332 _ => {
333 let _ = ws.send(&text_msg(syntax.clone(), "error", vec![
334 dat!("session_close: missing session id."),
335 ])).await;
336 continue;
337 }
338 };
339 match state.session_store.delete_session(&username, &session_id) {
340 Ok(()) => {
341 if current_session_id.as_deref() == Some(&session_id) {
342 current_session_id = None;
343 current_session = None;
344 }
345 let _ = ws.send(&text_msg(syntax.clone(), "info", vec![
346 dat!(fmt!("Session '{}' closed.", session_id)),
347 ])).await;
348 }
349 Err(e) => {
350 let _ = ws.send(&text_msg(syntax.clone(), "error", vec![
351 dat!(e.to_string()),
352 ])).await;
353 }
354 }
355 }
356 "session_rename" => {
357 let session_id = match cmdrx.vals.get(0) {
358 Some(Dat::Str(s)) => s.clone(),
359 _ => {
360 let _ = ws.send(&text_msg(syntax.clone(), "error", vec![
361 dat!("session_rename: missing session id."),
362 ])).await;
363 continue;
364 }
365 };
366 let new_name = match cmdrx.vals.get(1) {
367 Some(Dat::Str(s)) => s.clone(),
368 _ => {
369 let _ = ws.send(&text_msg(syntax.clone(), "error", vec![
370 dat!("session_rename: missing new name."),
371 ])).await;
372 continue;
373 }
374 };
375 match state.session_store.rename_session(&session_id, &new_name) {
376 Ok(()) => {
377 if current_session_id.as_deref() == Some(&session_id) {
378 if let Some(ref mut s) = current_session {
379 s.name = new_name.clone();
380 }
381 }
382 let _ = ws.send(&text_msg(syntax.clone(), "info", vec![
383 dat!(fmt!("Session renamed to '{}'.", new_name)),
384 ])).await;
385 }
386 Err(e) => {
387 let _ = ws.send(&text_msg(syntax.clone(), "error", vec![
388 dat!(e.to_string()),
389 ])).await;
390 }
391 }
392 }
393 "chat" => {
394 let content = match cmdrx.vals.get(0) {
395 Some(Dat::Str(s)) => s.clone(),
396 _ => {
397 let _ = ws.send(&text_msg(syntax.clone(), "error", vec![
398 dat!("chat: missing content."),
399 ])).await;
400 continue;
401 }
402 };
403
404 // Ensure we have a current session.
405 if current_session.is_none() {
406 // Auto-create a session.
407 match state.session_store.create_session(
408 &username, "Untitled", &state.agent.llm.model,
409 ) {
410 Ok(s) => {
411 current_session_id = Some(s.id.clone());
412 current_session = Some(s);
413 }
414 Err(e) => {
415 let _ = ws.send(&text_msg(syntax.clone(), "error", vec![
416 dat!(e.to_string()),
417 ])).await;
418 continue;
419 }
420 }
421 }
422
423 // Run the agent turn with streaming events.
424 //
425 // The on_event callback is synchronous (FnMut), but
426 // we need to send WS messages asynchronously. We
427 // use an mpsc channel and pin the agent future so
428 // we can concurrently drain events while the agent
429 // turn runs — giving the user true incremental
430 // streaming.
431 //
432 // The session is taken out of current_session so
433 // the pinned future can hold &mut without
434 // conflicting with the save after. The future is
435 // scoped inside a block so the borrow ends before
436 // we put the session back.
437 let mut session = match current_session.take() {
438 Some(s) => s,
439 None => continue, // presence checked above
440 };
441 // Use the session's selected model (the picker sets it);
442 // fall back to the configured default.
443 let mut agent = state.agent.clone();
444 if !session.model.is_empty() {
445 agent.llm.model = session.model.clone();
446 }
447 let syntax_ref = syntax.clone();
448
449 // Build the per-user tool registry: a workspace jailed
450 // to <workspace_base>/<user>, plus the executor.
451 let registry = {
452 let ws_dir = state.workspace_base.join(&username);
453 match Workspace::new(ws_dir) {
454 Ok(ws) => {
455 let ctx = ToolContext {
456 workspace: ws,
457 executor: state.executor.clone(),
458 cwd: String::new(),
459 path_prefix: String::new(),
460 root: crate::tools::FileRoot::Workspace,
461 read_seen: crate::tools::new_read_cache(),
462 no_write: Vec::new(),
463 // The native handler serves no Diamond.
464 daimon_of: String::new(),
465 };
466 ToolRegistry::new(Tool::defaults(), ctx)
467 }
468 Err(e) => {
469 let _ = ws.send(&text_msg(syntax.clone(), "error", vec![
470 dat!(fmt!("Workspace unavailable: {}", e)),
471 ])).await;
472 current_session = Some(session);
473 continue;
474 }
475 }
476 };
477
478 // Expand any skill invocation (<name ...>) in the message by injecting the matching
479 // skill's instructions -- and bound the turn by what that skill declared it needs.
480 //
481 // A skill is instructions, and instructions are as powerful as the agent reading
482 // them. Its text is therefore not the thing to trust; its declaration is. A skill
483 // that named `file_read` runs this turn against a registry holding `file_read` and
484 // nothing else, so it cannot send mail, spawn an agent, or drive a logged-in
485 // browser -- whatever its text says, and however cleverly it says it. The model is
486 // not even shown the tools it did not ask for.
487 //
488 // A skill that declares nothing runs unrestricted, which is right for one the user
489 // wrote themselves. When skills can be imported, an imported one that declares
490 // nothing must be refused rather than trusted.
491 let (content, registry) = match crate::skills::expand(&content, &registry.ctx.workspace) {
492 Ok(exp) => {
493 // A skill that ships a script and does not declare `shell` is refused, not
494 // narrowed: shipping code is asking to run it, and the user is owed the
495 // asking. Refused at expansion, so no turn runs at all.
496 if let Some(refusal) = &exp.refused {
497 let _ = ws.send(&text_msg(syntax.clone(), "error", vec![
498 dat!(refusal.clone()),
499 ])).await;
500 current_session = Some(session);
501 continue;
502 }
503 match exp.declared_tools() {
504 None => (exp.text, registry),
505 Some(names) => {
506 let unknown = registry.unknown_tools(&names);
507 if !unknown.is_empty() {
508 // A skill asking for a tool that does not exist has either a
509 // typo or a misdescription of itself. Say so; do not quietly
510 // hand it an empty toolbelt and let it fail obscurely.
511 let _ = ws.send(&text_msg(syntax.clone(), "error", vec![
512 dat!(fmt!(
513 "That skill asks for {}, which {} not a tool Daimond \
514 has. Fix its 'uses' line.",
515 unknown.join(", "),
516 if unknown.len() == 1 { "is" } else { "are" })),
517 ])).await;
518 current_session = Some(session);
519 continue;
520 }
521 let dirs = exp.skill_dirs();
522 let mut narrowed = registry.narrowed(&names);
523 // A bound is only a bound while the declaration cannot be edited.
524 // Skills live in the workspace, so a skill holding nothing but
525 // `file_write` could rewrite its own `uses` line -- or another
526 // skill's -- to ask for everything, and escape on the next
527 // invocation. One move, and the containment is theatre. So a turn
528 // running under a declaration is fenced out of Daimond's own
529 // directory, where the rules about what a skill may do are kept --
530 // and let back in to read one place only, its own folder, because
531 // the references a skill ships are part of the skill.
532 //
533 // COMPOSED and never assigned. A turn can already be bounded when
534 // it arrives here -- scoped to a Diamond, granted a toolkit -- and
535 // assigning over that discarded the Diamond's allow-list outright,
536 // silently, leaving a skill turn able to read every other Diamond
537 // (`hand/REVIEW.md` §1.12). `compose` gives what BOTH permit and is
538 // never able to widen either, so this line can only narrow.
539 narrowed.ctx.no_write = crate::tools::compose(
540 &narrowed.ctx.no_write,
541 &crate::tools::skill_bounds(&dirs));
542 info!(
543 "Skill turn bounded to {} of {} tools ({}), fenced out of \
544 .daimond/ but for reading {}",
545 narrowed.tools.len(),
546 registry.tools.len(),
547 narrowed.tool_names().join(", "),
548 if dirs.is_empty() {
549 fmt!("nothing (it ships no files)")
550 } else {
551 dirs.join(", ")
552 });
553 (exp.text, narrowed)
554 },
555 }
556 },
557 Err(_) => (content, registry),
558 };
559
560 let result = {
561 let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel::<AgentEvent>();
562
563 let mut on_event = |ev| { let _ = tx.send(ev); };
564 let mut agent_fut = pin!(agent.run_turn(
565 &mut session,
566 content,
567 &registry,
568 &mut on_event,
569 ));
570
571 let mut result = Ok(());
572 loop {
573 tokio::select! {
574 biased;
575 r = &mut agent_fut => {
576 result = r;
577 // Drain any remaining events.
578 while let Ok(ev) = rx.try_recv() {
579 if let Some((cmd_name, vals)) = event_to_ws(&ev) {
580 let _ = ws.send(&text_msg(
581 syntax_ref.clone(), cmd_name, vals,
582 )).await;
583 }
584 }
585 break;
586 }
587 ev = rx.recv() => {
588 match ev {
589 Some(ev) => {
590 if let Some((cmd_name, vals)) = event_to_ws(&ev) {
591 let _ = ws.send(&text_msg(
592 syntax_ref.clone(), cmd_name, vals,
593 )).await;
594 }
595 }
596 None => break,
597 }
598 }
599 }
600 }
601 result
602 };
603
604 // Put the session back and save.
605 current_session = Some(session);
606 if let Some(ref s) = current_session {
607 let _ = state.session_store.save_session(s);
608 }
609
610 if let Err(e) = result {
611 warn!("{}: Agent turn error: {}", id, e);
612 }
613 }
614 "fs_list" => {
615 let path = match cmdrx.vals.get(0) {
616 Some(Dat::Str(s)) => s.clone(),
617 _ => ".".to_string(),
618 };
619 match user_ws(&state.workspace_base, &username)
620 .and_then(|w| fs_list_json(&w, &path))
621 {
622 Ok(json) => { let _ = ws.send(&text_msg(syntax.clone(), "fs_tree", vec![dat!(json)])).await; }
623 Err(e) => { let _ = ws.send(&text_msg(syntax.clone(), "error", vec![dat!(e.to_string())])).await; }
624 }
625 }
626 "fs_read" => {
627 let path = match cmdrx.vals.get(0) {
628 Some(Dat::Str(s)) => s.clone(),
629 _ => { let _ = ws.send(&text_msg(syntax.clone(), "error", vec![dat!("fs_read: missing path.")])).await; continue; }
630 };
631 match user_ws(&state.workspace_base, &username)
632 .and_then(|w| fs_read_file(&w, &path))
633 {
634 Ok(content) => { let _ = ws.send(&text_msg(syntax.clone(), "fs_content", vec![dat!(path.clone()), dat!(content)])).await; }
635 Err(e) => { let _ = ws.send(&text_msg(syntax.clone(), "error", vec![dat!(e.to_string())])).await; }
636 }
637 }
638 "fs_delete" => {
639 let path = match cmdrx.vals.get(0) {
640 Some(Dat::Str(s)) => s.clone(),
641 _ => { let _ = ws.send(&text_msg(syntax.clone(), "error", vec![dat!("fs_delete: missing path.")])).await; continue; }
642 };
643 match user_ws(&state.workspace_base, &username)
644 .and_then(|w| fs_delete_file(&w, &path))
645 {
646 Ok(()) => { let _ = ws.send(&text_msg(syntax.clone(), "info", vec![dat!(fmt!("Deleted {}.", path))])).await; }
647 Err(e) => { let _ = ws.send(&text_msg(syntax.clone(), "error", vec![dat!(e.to_string())])).await; }
648 }
649 }
650 "fs_write" => {
651 let path = match cmdrx.vals.get(0) {
652 Some(Dat::Str(s)) => s.clone(),
653 _ => { let _ = ws.send(&text_msg(syntax.clone(), "error", vec![dat!("fs_write: missing path.")])).await; continue; }
654 };
655 let content_hex = match cmdrx.vals.get(1) {
656 Some(Dat::Str(s)) => s.clone(),
657 _ => String::new(),
658 };
659 let write_res = hex_decode(&content_hex).and_then(|bytes| {
660 user_ws(&state.workspace_base, &username)
661 .and_then(|w| fs_write_file(&w, &path, &bytes))
662 });
663 match write_res {
664 Ok(()) => { let _ = ws.send(&text_msg(syntax.clone(), "info", vec![dat!(fmt!("Wrote {}.", path))])).await; }
665 Err(e) => { let _ = ws.send(&text_msg(syntax.clone(), "error", vec![dat!(e.to_string())])).await; }
666 }
667 }
668 _ => {
669 let _ = ws.send(&text_msg(syntax.clone(), "error", vec![
670 dat!(fmt!("Unknown command '{}'.", cmd_name)),
671 ])).await;
672 }
673 }
674 }
675
676 info!("{}: Daimond WS closed.", id);
677 Ok(())
678}
679
680
681// ┌───────────────────────────────────────────────────────────────┐
682// │ Helpers │
683// └───────────────────────────────────────────────────────────────┘
684
685/// Map an `AgentEvent` to a WS command name and value list.
686/// Open a user's workspace under the shared base directory.
687fn user_ws(base: &PathBuf, user: &str) -> Outcome<Workspace> {
688 Workspace::new(base.join(user))
689}
690
691/// Build a JSON directory listing: `{"path":..,"entries":[{name,dir,size}]}`.
692fn fs_list_json(ws: &Workspace, path: &str) -> Outcome<String> {
693 let abs = res!(ws.resolve(path));
694 let rd = res!(std::fs::read_dir(&abs)
695 .map_err(|e| err!(e, "fs_list: cannot read directory '{}'.", path; IO, File, Read)));
696 let mut entries: Vec<(bool, String, u64)> = rd
697 .filter_map(|e| e.ok())
698 .map(|e| {
699 let name = e.file_name().to_string_lossy().to_string();
700 let is_dir = e.path().is_dir();
701 let size = e.metadata().map(|m| m.len()).unwrap_or(0);
702 (is_dir, name, size)
703 })
704 .collect();
705 entries.sort_by(|a, b| b.0.cmp(&a.0).then(a.1.cmp(&b.1)));
706 let disp = if path == "." { String::new() } else { path.to_string() };
707 let mut out = fmt!("{{\"path\":\"{}\",\"entries\":[", crate::llm::json_escape(&disp));
708 for (i, (is_dir, name, size)) in entries.iter().enumerate() {
709 if i > 0 { out.push(','); }
710 out.push_str(&fmt!("{{\"name\":\"{}\",\"dir\":{},\"size\":{}}}",
711 crate::llm::json_escape(name), is_dir, size));
712 }
713 out.push_str("]}");
714 Ok(out)
715}
716
717/// Read a workspace text file.
718fn fs_read_file(ws: &Workspace, path: &str) -> Outcome<String> {
719 let abs = res!(ws.resolve(path));
720 let data = res!(std::fs::read(&abs)
721 .map_err(|e| err!(e, "fs_read: cannot read '{}'.", path; IO, File, Read)));
722 Ok(String::from_utf8_lossy(&data).to_string())
723}
724
725/// Delete a workspace file.
726fn fs_delete_file(ws: &Workspace, path: &str) -> Outcome<()> {
727 let abs = res!(ws.resolve(path));
728 res!(std::fs::remove_file(&abs)
729 .map_err(|e| err!(e, "fs_delete: cannot delete '{}'.", path; IO, File)));
730 Ok(())
731}
732
733/// Create or overwrite a workspace file with raw bytes.
734fn fs_write_file(ws: &Workspace, path: &str, content: &[u8]) -> Outcome<()> {
735 let abs = res!(ws.resolve(path));
736 if let Some(parent) = abs.parent() {
737 res!(std::fs::create_dir_all(parent)
738 .map_err(|e| err!(e, "fs_write: cannot create parent dirs for '{}'.", path; IO, File)));
739 }
740 res!(std::fs::write(&abs, content)
741 .map_err(|e| err!(e, "fs_write: cannot write '{}'.", path; IO, File, Write)));
742 Ok(())
743}
744
745/// Decode a lowercase/uppercase hex string to bytes. Uploads send
746/// file content hex-encoded so arbitrary bytes (newlines, quotes,
747/// `$`, binary) pass through the WS syntax parser unharmed.
748fn hex_decode(s: &str) -> Outcome<Vec<u8>> {
749 let b = s.as_bytes();
750 if b.len() % 2 != 0 {
751 return Err(err!("fs_write: hex content has odd length."; Invalid, Input));
752 }
753 let val = |c: u8| -> Outcome<u8> {
754 match c {
755 b'0'..=b'9' => Ok(c - b'0'),
756 b'a'..=b'f' => Ok(c - b'a' + 10),
757 b'A'..=b'F' => Ok(c - b'A' + 10),
758 _ => Err(err!("fs_write: invalid hex digit."; Invalid, Input)),
759 }
760 };
761 let mut out = Vec::with_capacity(b.len() / 2);
762 let mut i = 0;
763 while i < b.len() {
764 out.push(res!(val(b[i])) * 16 + res!(val(b[i + 1])));
765 i += 2;
766 }
767 Ok(out)
768}
769
770/// The Steel WS command for one agent event, or `None` where this wire carries none.
771fn event_to_ws(ev: &AgentEvent) -> Option<(&'static str, Vec<Dat>)> {
772 Some(match ev {
773 AgentEvent::Text(t) => ("text", vec![dat!(t.clone())]),
774 AgentEvent::ToolCall { name, args, .. } =>
775 ("tool_call", vec![dat!(name.clone()), dat!(args.clone())]),
776 // The outcome is deliberately NOT carried here. This is the Steel WS wire, whose
777 // `tool_result` command is declared in `src/syntax.rs` as two values, and widening a
778 // declared command is a change to a protocol nobody in this round owns a reader for.
779 // The browser gets the outcome over `wasm::app::event_to_js`, which is the path the
780 // page actually runs on.
781 AgentEvent::ToolResult { name, result, .. } =>
782 ("tool_result", vec![dat!(name.clone()), dat!(result.clone())]),
783 AgentEvent::Interjected(text) => ("interjected", vec![dat!(text.clone())]),
784 // NOT carried on this wire, and now not sent at all. The WS commands are declared
785 // in `src/syntax.rs`, which declares no `thinking` -- so the empty frame this used
786 // to send named a command the syntax does not know, and `text_msg` could never
787 // build one. It was a no-op, once a round, and nobody noticed.
788 //
789 // Since 2026-08-28 the reasoning STREAMS, so the same no-op would happen for every
790 // delta: 1,378 of them in one measured round, each a syntax lookup and a discarded
791 // send. The old note argued that an empty frame at least said the act had happened;
792 // hundreds of identical empty frames say nothing a client can use, and this wire has
793 // no reader for them either way. The browser gets the text over
794 // `wasm::app::event_to_js`, which is the path the page runs on.
795 AgentEvent::Thinking(_) => return None,
796 // Not carried on this wire either, and for the reason given above: the WS commands are
797 // declared in `src/syntax.rs`. The browser reads it over `wasm::app::event_to_js`.
798 AgentEvent::Ended { .. } => ("ended", vec![]),
799 // The two counts before the prose, so a client can draw the act without
800 // parsing the sentence: what went, what is left, and then what it says.
801 AgentEvent::Compacted { folded, kept, note } =>
802 ("compacted", vec![Dat::U64(*folded as u64), Dat::U64(*kept as u64),
803 dat!(note.clone())]),
804 // The count before the model name, for the same reason as `compacted` above: how
805 // many pictures were left out is drawable without reading a sentence.
806 AgentEvent::Unseeable { images, model } =>
807 ("unseeable", vec![Dat::U64(*images as u64), dat!(model.clone())]),
808 // Not carried on this wire, for the reason `thinking` and `ended` are not: the WS
809 // commands are declared in `src/syntax.rs` and widening one is a protocol change with no
810 // reader in this round. The road ladder is a BROWSER mechanism in any case -- the mark it
811 // classifies on is set in `window.fetch` (www/js/gateway.js) -- so the native path can
812 // never raise this. Named rather than swept into a catch-all, so a future event cannot
813 // reach this wire unnoticed.
814 AgentEvent::Roading { .. } => ("roading", vec![]),
815 AgentEvent::Truncated => ("truncated", vec![]),
816 AgentEvent::Done => ("done", vec![]),
817 AgentEvent::Error(msg) => ("error", vec![dat!(msg.clone())]),
818 })
819}
820
821fn text_msg(syntax: SyntaxRef, cmd: &str, vals: Vec<Dat>)
822 -> WebSocketMessage
823{
824 let mut response = match MsgCmd::new(syntax.clone(), cmd) {
825 Ok(r) => r,
826 Err(e) => {
827 debug!("text_msg: MsgCmd::new('{}') failed: {}", cmd, e);
828 return WebSocketMessage::Text("error \"internal\"".to_string());
829 }
830 };
831 for val in &vals {
832 response = match response.add_cmd_val(val.clone()) {
833 Ok(r) => r,
834 Err(e) => {
835 debug!("text_msg: add_cmd_val for '{}' failed: {} (val kind={:?})", cmd, e, val.kind());
836 return WebSocketMessage::Text("error \"internal\"".to_string());
837 }
838 }
839 }
840 WebSocketMessage::Text(response.to_string())
841}
842
843/// Look up the authenticated username from the session metadata
844/// in O3db.
845fn get_username<
846 const UIDL: usize,
847 UID: NumIdDat<UIDL> + Clone,
848 ENC: Encrypter,
849 KH: Hasher,
850 DB: Database<UIDL, UID, ENC, KH>,
851>(
852 db: &Option<(Arc<RwLock<DB>>, UID)>,
853 sid: &Option<String>,
854) -> Option<String> {
855 let sid = match sid.as_ref() {
856 Some(s) => s,
857 None => return None,
858 };
859 let (db, _uid) = match db.as_ref() {
860 Some(v) => v,
861 None => return None,
862 };
863 let db_r = match db.read() {
864 Ok(v) => v,
865 Err(_) => return None,
866 };
867 let meta_key = Dat::Str(fmt!("sess_meta:{}", sid));
868 match db_r.get(&meta_key, None) {
869 Ok(Some((data, _))) => {
870 if let Dat::Map(m) = &data {
871 if let Some(Dat::Str(user)) = m.get(&dat!("user")) {
872 if !user.is_empty() {
873 return Some(user.clone());
874 }
875 }
876 }
877 None
878 }
879 _ => None,
880 }
881}