1use std::sync::Arc;
5use std::time::{SystemTime, UNIX_EPOCH};
6
7use axum::{
8 extract::{
9 ws::{Message, WebSocket},
10 State, WebSocketUpgrade,
11 },
12 response::Response,
13};
14use base64::Engine as _;
15use serde::Deserialize;
16use serde_json::json;
17
18use crate::backend::blockcontroller;
19use crate::backend::rpc::engine::WshRpcEngine;
20use crate::backend::rpc_types::{
21 CommandBlockInputData, CommandControllerResyncData, CommandEventReadHistoryData,
22 CommandGetMetaData, CommandSetMetaData, CommandToolDecisionData,
23 RpcMessage, COMMAND_CONTROLLER_INPUT,
24 COMMAND_CONTROLLER_RESYNC, COMMAND_EVENT_READ_HISTORY, COMMAND_EVENT_SUB, COMMAND_EVENT_UNSUB,
25 COMMAND_EVENT_UNSUB_ALL, COMMAND_GET_FULL_CONFIG, COMMAND_GET_META,
26 COMMAND_GET_AI_RATE_LIMIT, COMMAND_ROUTE_ANNOUNCE, COMMAND_ROUTE_UNANNOUNCE,
27 COMMAND_SET_META, COMMAND_SET_CONFIG, COMMAND_APP_INFO,
28 COMMAND_TOOL_DECISION, COMMAND_AGENT_ANSWER,
29 CommandAgentAnswerData,
30};
31use crate::backend::obj::{Block, TermSize, WaveObjUpdate, wave_obj_to_value};
32use super::service::update_object_meta;
33
34use super::AppState;
35
36#[derive(Deserialize)]
39struct WSIncoming {
40 #[serde(rename = "type")]
41 msg_type: Option<String>,
42 #[allow(dead_code)]
43 stime: Option<i64>,
44 wscommand: Option<String>,
45 message: Option<RpcMessage>,
46 blockid: Option<String>,
48 inputdata64: Option<String>,
49 termsize: Option<serde_json::Value>,
50 agent_id: Option<String>,
52 from: Option<String>,
53 to: Option<String>,
54 target: Option<String>,
55 payload: Option<String>,
56 #[serde(rename = "bus_message")]
57 bus_message_text: Option<String>,
58 priority: Option<String>,
59}
60
61pub(super) async fn handle_ws(
62 State(state): State<AppState>,
63 ws: WebSocketUpgrade,
64) -> Response {
65 ws.on_upgrade(move |socket| handle_ws_connection(socket, state))
66}
67
68async fn handle_ws_connection(mut socket: WebSocket, state: AppState) {
69 let ws_start = std::time::Instant::now();
70 let conn_id = uuid::Uuid::new_v4().to_string();
71 let tab_id = String::new();
72
73 tracing::info!(conn_id = %conn_id, "WebSocket client connected");
74
75 let crate::backend::eventbus::WsReceivers {
81 priority: mut priority_rx,
82 background: mut background_rx,
83 } = state.event_bus.register_ws(&conn_id, &tab_id);
84 tracing::info!("[ws-perf] register_ws: {:.2}ms", ws_start.elapsed().as_secs_f64() * 1000.0);
85
86 let mut bus_rx: Option<tokio::sync::mpsc::UnboundedReceiver<crate::backend::messagebus::BusMessage>> = None;
88 let mut bus_agent_id: Option<String> = None;
89
90 {
94 let t = std::time::Instant::now();
95 let config = state.config_watcher.get_full_config();
96 if let Ok(config_val) = serde_json::to_value(config.as_ref()) {
97 let config_event = json!({
98 "eventtype": "rpc",
99 "data": {
100 "command": "eventrecv",
101 "data": {
102 "event": "config",
103 "data": { "fullconfig": config_val }
104 }
105 }
106 });
107 if let Ok(msg) = serde_json::to_string(&config_event) {
108 let _ = socket.send(Message::Text(msg.into())).await;
109 }
110 }
111 tracing::info!("[ws-perf] send_initial_config: {:.2}ms", t.elapsed().as_secs_f64() * 1000.0);
112 }
113
114 let t = std::time::Instant::now();
116 let (engine, mut rpc_output_rx) = WshRpcEngine::new();
117
118 register_handlers(&engine, state.clone(), conn_id.clone());
120 tracing::info!("[ws-perf] create_engine+register_handlers: {:.2}ms", t.elapsed().as_secs_f64() * 1000.0);
121 tracing::info!("[ws-perf] TOTAL ws_setup: {:.2}ms", ws_start.elapsed().as_secs_f64() * 1000.0);
122
123 let mut ping_interval = tokio::time::interval(std::time::Duration::from_secs(10));
125 ping_interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
126
127 loop {
128 tokio::select! {
129 biased;
135
136 msg = socket.recv() => {
138 match msg {
139 Some(Ok(Message::Close(_))) | None => break,
140 Some(Ok(Message::Ping(data))) => {
141 let _ = socket.send(Message::Pong(data)).await;
142 }
143 Some(Ok(Message::Text(text))) => {
144 match handle_incoming_text(&text, &engine, &state, &mut socket).await {
145 Err(true) => break,
146 Ok(Some((new_rx, agent_id))) => {
147 bus_rx = Some(new_rx);
149 bus_agent_id = Some(agent_id);
150 }
151 _ => {}
152 }
153 }
154 Some(Ok(_)) => {
155 }
157 Some(Err(_)) => break,
158 }
159 }
160
161 Some(rpc_msg) = rpc_output_rx.recv() => {
163 let wrapped = json!({
164 "eventtype": "rpc",
165 "data": rpc_msg,
166 });
167 let msg = serde_json::to_string(&wrapped).unwrap_or_default();
168 if socket.send(Message::Text(msg.into())).await.is_err() {
169 break;
170 }
171 }
172
173 Some(bus_msg) = async {
175 match bus_rx.as_mut() {
176 Some(rx) => rx.recv().await,
177 None => std::future::pending().await,
178 }
179 } => {
180 let wrapped = json!({
181 "type": "bus:message",
182 "data": bus_msg,
183 });
184 let msg = serde_json::to_string(&wrapped).unwrap_or_default();
185 if socket.send(Message::Text(msg.into())).await.is_err() {
186 break;
187 }
188 }
189
190 Some(event) = priority_rx.recv() => {
197 if forward_event(&mut socket, event).await {
198 break;
199 }
200 }
201
202 Some(event) = background_rx.recv() => {
219 if forward_event(&mut socket, event).await {
220 break;
221 }
222 }
223
224 _ = ping_interval.tick() => {
226 let now = SystemTime::now()
227 .duration_since(UNIX_EPOCH)
228 .unwrap_or_default()
229 .as_millis() as i64;
230 let ping = json!({ "type": "ping", "stime": now });
231 let msg = serde_json::to_string(&ping).unwrap_or_default();
232 if socket.send(Message::Text(msg.into())).await.is_err() {
233 break;
234 }
235 }
236 }
237 }
238
239 tracing::info!(conn_id = %conn_id, "WebSocket client disconnected");
240 state.event_bus.unregister_ws(&conn_id);
241 state.broker.unsubscribe_all(&conn_id);
242
243 if let Some(ref agent_id) = bus_agent_id {
245 state.messagebus.unregister(agent_id);
246 }
247}
248
249async fn forward_event(socket: &mut WebSocket, event: serde_json::Value) -> bool {
258 let msg = if event["eventtype"] == "rpc" {
259 serde_json::to_string(&event).unwrap_or_default()
261 } else {
262 let wave_event = json!({
264 "event": event["eventtype"],
265 "scopes": [event["oref"]],
266 "data": event["data"],
267 });
268 let wrapped = json!({
269 "eventtype": "rpc",
270 "data": {
271 "command": "eventrecv",
272 "data": wave_event,
273 },
274 });
275 serde_json::to_string(&wrapped).unwrap_or_default()
276 };
277 socket.send(Message::Text(msg.into())).await.is_err()
278}
279
280async fn handle_incoming_text(
284 text: &str,
285 engine: &Arc<WshRpcEngine>,
286 state: &AppState,
287 socket: &mut WebSocket,
288) -> Result<Option<(tokio::sync::mpsc::UnboundedReceiver<crate::backend::messagebus::BusMessage>, String)>, bool> {
289 let incoming: WSIncoming = match serde_json::from_str(text) {
290 Ok(v) => v,
291 Err(e) => {
292 tracing::warn!("ws: invalid JSON: {}", e);
293 return Ok(None);
294 }
295 };
296
297 if let Some(ref msg_type) = incoming.msg_type {
299 match msg_type.as_str() {
300 "ping" => {
301 let now = SystemTime::now()
302 .duration_since(UNIX_EPOCH)
303 .unwrap_or_default()
304 .as_millis() as i64;
305 let pong = json!({ "type": "pong", "stime": now });
306 let msg = serde_json::to_string(&pong).unwrap_or_default();
307 if socket.send(Message::Text(msg.into())).await.is_err() {
308 return Err(true);
309 }
310 return Ok(None);
311 }
312 "pong" => {
313 return Ok(None);
314 }
315 "bus:register" => {
316 if let Some(ref agent_id) = incoming.agent_id {
317 let rx = state.messagebus.register(agent_id, "websocket");
318 let ack = json!({ "type": "bus:registered", "agent_id": agent_id });
319 let msg = serde_json::to_string(&ack).unwrap_or_default();
320 if socket.send(Message::Text(msg.into())).await.is_err() {
321 return Err(true);
322 }
323 return Ok(Some((rx, agent_id.clone())));
324 }
325 return Ok(None);
326 }
327 "bus:send" => {
328 if let (Some(ref from), Some(ref to), Some(ref payload)) =
329 (&incoming.from, &incoming.to, &incoming.payload)
330 {
331 let priority = match incoming.priority.as_deref() {
332 Some("high") => crate::backend::messagebus::Priority::High,
333 Some("urgent") => crate::backend::messagebus::Priority::Urgent,
334 _ => crate::backend::messagebus::Priority::Normal,
335 };
336 let bus_msg = crate::backend::messagebus::BusMessage::new(
337 from, to, crate::backend::messagebus::MessageType::Send, payload, priority,
338 );
339 let msg_id = bus_msg.id.clone();
340 let _ = state.messagebus.send(bus_msg);
341 let ack = json!({ "type": "bus:sent", "message_id": msg_id });
342 let msg = serde_json::to_string(&ack).unwrap_or_default();
343 if socket.send(Message::Text(msg.into())).await.is_err() {
344 return Err(true);
345 }
346 }
347 return Ok(None);
348 }
349 "bus:inject" => {
350 let from = incoming.from.as_deref().unwrap_or("unknown");
351 if let (Some(ref target), Some(ref message)) =
352 (&incoming.target, &incoming.bus_message_text)
353 {
354 let reactive_req = crate::backend::reactive::InjectionRequest {
356 target_agent: target.clone(),
357 message: message.clone(),
358 source_agent: Some(from.to_string()),
359 request_id: None,
360 priority: incoming.priority.clone(),
361 wait_for_idle: false,
362 };
363 let resp = state.reactive_handler.inject_message(reactive_req);
364 if resp.success {
365 let ack = json!({ "type": "bus:injected", "via": "pty", "block_id": resp.block_id });
366 let msg = serde_json::to_string(&ack).unwrap_or_default();
367 if socket.send(Message::Text(msg.into())).await.is_err() {
368 return Err(true);
369 }
370 return Ok(None);
371 }
372
373 let is_not_found = resp.error.as_deref().map(|e| e.contains("not found")).unwrap_or(false);
375 if !is_not_found {
376 let err = json!({ "type": "bus:error", "error": resp.error });
377 let msg = serde_json::to_string(&err).unwrap_or_default();
378 if socket.send(Message::Text(msg.into())).await.is_err() {
379 return Err(true);
380 }
381 return Ok(None);
382 }
383
384 let priority = match incoming.priority.as_deref() {
386 Some("high") => crate::backend::messagebus::Priority::High,
387 Some("urgent") => crate::backend::messagebus::Priority::Urgent,
388 _ => crate::backend::messagebus::Priority::Normal,
389 };
390 match state.messagebus.inject(from, target, message, priority) {
391 Ok(msg_id) => {
392 let ack = json!({ "type": "bus:injected", "via": "messagebus", "message_id": msg_id });
393 let msg = serde_json::to_string(&ack).unwrap_or_default();
394 if socket.send(Message::Text(msg.into())).await.is_err() {
395 return Err(true);
396 }
397 }
398 Err(e) => {
399 let err = json!({ "type": "bus:error", "error": e });
400 let msg = serde_json::to_string(&err).unwrap_or_default();
401 if socket.send(Message::Text(msg.into())).await.is_err() {
402 return Err(true);
403 }
404 }
405 }
406 }
407 return Ok(None);
408 }
409 "bus:broadcast" => {
410 let from = incoming.from.as_deref().unwrap_or("unknown");
411 if let Some(ref payload) = incoming.payload {
412 let priority = match incoming.priority.as_deref() {
413 Some("high") => crate::backend::messagebus::Priority::High,
414 Some("urgent") => crate::backend::messagebus::Priority::Urgent,
415 _ => crate::backend::messagebus::Priority::Normal,
416 };
417 let _ = state.messagebus.broadcast(from, payload, priority);
418 }
419 return Ok(None);
420 }
421 _ => {}
422 }
423 }
424
425 if let Some(ref wscommand) = incoming.wscommand {
427 match wscommand.as_str() {
428 "rpc" => {
429 if let Some(rpc_msg) = incoming.message {
430 engine.handle_message(rpc_msg);
431 } else {
432 tracing::warn!("ws: rpc wscommand missing message field");
433 }
434 }
435 "blockinput" => {
436 if let Some(ref block_id) = incoming.blockid {
437 if let Some(ref data64) = incoming.inputdata64 {
438 if !data64.is_empty() {
439 match base64::engine::general_purpose::STANDARD.decode(data64) {
440 Ok(data) => {
441 let input = blockcontroller::BlockInputUnion::data(data);
442 if let Err(e) = blockcontroller::send_input(block_id, input, None) {
443 tracing::debug!("ws: blockinput error: {}", e);
444 }
445 }
446 Err(e) => {
447 tracing::warn!("ws: blockinput base64 decode error: {}", e);
448 }
449 }
450 }
451 }
452 }
453 }
454 "setblocktermsize" => {
455 if let Some(ref block_id) = incoming.blockid {
456 if let Some(ref ts_val) = incoming.termsize {
457 match serde_json::from_value::<TermSize>(ts_val.clone()) {
458 Ok(ts) => {
459 let input = blockcontroller::BlockInputUnion::resize(ts);
460 if let Err(e) = blockcontroller::send_input(block_id, input, None) {
461 tracing::debug!("ws: setblocktermsize error: {}", e);
462 }
463 }
464 Err(e) => {
465 tracing::warn!("ws: setblocktermsize parse error: {}", e);
466 }
467 }
468 }
469 }
470 }
471 other => {
472 tracing::warn!("ws: unknown wscommand: {}", other);
473 }
474 }
475 }
476
477 Ok(None)
478}
479
480fn register_handlers(engine: &Arc<WshRpcEngine>, state: AppState, conn_id: String) {
481 let config_watcher = state.config_watcher.clone();
483 engine.register_handler(
484 COMMAND_GET_FULL_CONFIG,
485 Box::new(move |_data, _ctx| {
486 let cw = config_watcher.clone();
487 Box::pin(async move {
488 let config = cw.get_full_config();
489 match serde_json::to_value(config.as_ref()) {
490 Ok(v) => Ok(Some(v)),
491 Err(e) => Err(format!("failed to serialize config: {}", e)),
492 }
493 })
494 }),
495 );
496
497 engine.register_handler(
499 COMMAND_ROUTE_ANNOUNCE,
500 Box::new(|data, _ctx| {
501 Box::pin(async move {
502 tracing::debug!("routeannounce: {:?}", data);
503 Ok(None)
504 })
505 }),
506 );
507
508 engine.register_handler(
510 COMMAND_ROUTE_UNANNOUNCE,
511 Box::new(|_data, _ctx| Box::pin(async move { Ok(None) })),
512 );
513
514 let broker_sub = state.broker.clone();
516 let conn_id_sub = conn_id.clone();
517 engine.register_handler(
518 COMMAND_EVENT_SUB,
519 Box::new(move |data, _ctx| {
520 let broker = broker_sub.clone();
521 let conn_id = conn_id_sub.clone();
522 Box::pin(async move {
523 let sub: crate::backend::wps::SubscriptionRequest =
524 serde_json::from_value(data).map_err(|e| format!("eventsub: {e}"))?;
525 tracing::debug!("eventsub: event={} scopes={:?} allscopes={}", sub.event, sub.scopes, sub.allscopes);
526 broker.subscribe(&conn_id, sub);
527 Ok(None)
528 })
529 }),
530 );
531
532 let broker_unsub = state.broker.clone();
534 let conn_id_unsub = conn_id.clone();
535 engine.register_handler(
536 COMMAND_EVENT_UNSUB,
537 Box::new(move |data, _ctx| {
538 let broker = broker_unsub.clone();
539 let conn_id = conn_id_unsub.clone();
540 Box::pin(async move {
541 let event_name = data.as_str().unwrap_or("").to_string();
542 if !event_name.is_empty() {
543 broker.unsubscribe(&conn_id, &event_name);
544 }
545 Ok(None)
546 })
547 }),
548 );
549
550 let broker_unsub_all = state.broker.clone();
552 let conn_id_unsub_all = conn_id.clone();
553 engine.register_handler(
554 COMMAND_EVENT_UNSUB_ALL,
555 Box::new(move |_data, _ctx| {
556 let broker = broker_unsub_all.clone();
557 let conn_id = conn_id_unsub_all.clone();
558 Box::pin(async move {
559 broker.unsubscribe_all(&conn_id);
560 Ok(None)
561 })
562 }),
563 );
564
565 let wstore_sm = state.wstore.clone();
567 let event_bus_sm = state.event_bus.clone();
568 engine.register_handler(
569 COMMAND_SET_META,
570 Box::new(move |data, _ctx| {
571 let wstore = wstore_sm.clone();
572 let event_bus = event_bus_sm.clone();
573 Box::pin(async move {
574 let cmd: CommandSetMetaData =
575 serde_json::from_value(data).map_err(|e| format!("setmeta: {e}"))?;
576 let oref_str = cmd.oref.to_string();
577 let meta_keys: Vec<&String> = cmd.meta.keys().collect();
578 tracing::info!(oref = %oref_str, keys = ?meta_keys, "SetMeta");
579 update_object_meta(&wstore, &oref_str, &cmd.meta)?;
580 let oref = crate::backend::ORef::parse(&oref_str)
583 .map_err(|e| e.to_string())?;
584 let update_data = if oref.otype == "block" {
585 if let Ok(block) = wstore.must_get::<Block>(&oref.oid) {
586 Some(serde_json::to_value(&WaveObjUpdate {
587 updatetype: "update".into(),
588 otype: oref.otype.clone(),
589 oid: oref.oid.clone(),
590 obj: Some(wave_obj_to_value(&block)),
591 }).unwrap_or_default())
592 } else { None }
593 } else { None };
594 event_bus.broadcast_event(&crate::backend::eventbus::WSEventType {
595 eventtype: "waveobj:update".to_string(),
596 oref: oref_str,
597 data: update_data,
598 });
599 Ok(None)
600 })
601 }),
602 );
603
604 let wstore_gm = state.wstore.clone();
606 engine.register_handler(
607 COMMAND_GET_META,
608 Box::new(move |data, _ctx| {
609 let wstore = wstore_gm.clone();
610 Box::pin(async move {
611 let cmd: CommandGetMetaData =
612 serde_json::from_value(data).map_err(|e| format!("getmeta: {e}"))?;
613 let obj: Option<serde_json::Value> = wstore
614 .get_raw(&cmd.oref.otype, &cmd.oref.oid)
615 .map_err(|e| format!("getmeta: {e}"))?;
616 match obj {
617 Some(val) => {
618 let meta = val.get("meta").cloned().unwrap_or(val);
620 Ok(Some(meta))
621 }
622 None => Err(format!("getmeta: object {} not found", cmd.oref)),
623 }
624 })
625 }),
626 );
627
628 let version_info = state.version.clone();
630 engine.register_handler(
631 COMMAND_APP_INFO,
632 Box::new(move |_data, _ctx| {
633 let version = version_info.clone();
634 Box::pin(async move {
635 Ok(Some(serde_json::json!({
636 "version": version,
637 })))
638 })
639 }),
640 );
641
642 engine.register_handler(
644 COMMAND_GET_AI_RATE_LIMIT,
645 Box::new(|_data, _ctx| {
646 Box::pin(async move {
647 Ok(Some(serde_json::json!({
648 "req": 9999,
649 "reqlimit": 9999,
650 "preq": 9999,
651 "preqlimit": 9999,
652 "resetepoch": 0,
653 "unknown": true
654 })))
655 })
656 }),
657 );
658
659 let wstore_resync = state.wstore.clone();
661 let broker_resync = state.broker.clone();
662 let event_bus_resync = state.event_bus.clone();
663 let filestore_resync = state.filestore.clone();
664 engine.register_handler(
665 COMMAND_CONTROLLER_RESYNC,
666 Box::new(move |data, _ctx| {
667 let wstore = wstore_resync.clone();
668 let broker = broker_resync.clone();
669 let event_bus = event_bus_resync.clone();
670 let filestore = filestore_resync.clone();
671 Box::pin(async move {
672 let cmd: CommandControllerResyncData = serde_json::from_value(data)
673 .map_err(|e| format!("controllerresync: {e}"))?;
674 tracing::info!(
675 block_id = %cmd.blockid,
676 tab_id = %cmd.tabid,
677 forcerestart = cmd.forcerestart,
678 "ControllerResync"
679 );
680 let block: Block = wstore
681 .get(&cmd.blockid)
682 .map_err(|e| format!("controllerresync: load block: {e}"))?
683 .ok_or_else(|| format!("controllerresync: block {} not found", cmd.blockid))?;
684 blockcontroller::resync_controller(
685 &block,
686 &cmd.tabid,
687 cmd.rtopts,
688 cmd.forcerestart,
689 Some(broker),
690 Some(event_bus),
691 Some(wstore),
692 Some(filestore),
693 )?;
694 Ok(None)
695 })
696 }),
697 );
698
699 engine.register_handler(
701 COMMAND_CONTROLLER_INPUT,
702 Box::new(|data, _ctx| {
703 Box::pin(async move {
704 let cmd: CommandBlockInputData = serde_json::from_value(data)
705 .map_err(|e| format!("controllerinput: {e}"))?;
706 let input = parse_block_input(&cmd)?;
707 blockcontroller::send_input(&cmd.blockid, input, cmd.seq)?;
708 Ok(None)
709 })
710 }),
711 );
712
713 engine.register_handler(
733 COMMAND_TOOL_DECISION,
734 Box::new(|data, _ctx| {
735 Box::pin(async move {
736 let cmd: CommandToolDecisionData = serde_json::from_value(data)
737 .map_err(|e| format!("tooldecision: {e}"))?;
738 match cmd.outcome.as_str() {
739 "allow" | "deny" => {}
740 other => {
741 return Err(format!(
742 "tooldecision: invalid outcome '{}' (expected 'allow' or 'deny')",
743 other
744 ));
745 }
746 }
747 match cmd.scope.as_str() {
751 "once" | "session" | "project" | "global" => {}
752 other => {
753 return Err(format!(
754 "tooldecision: invalid scope '{}' (expected 'once'/'session'/'project'/'global')",
755 other
756 ));
757 }
758 }
759 tracing::info!(
760 block_id = %cmd.blockid,
761 request_id = %cmd.request_id,
762 outcome = %cmd.outcome,
763 scope = %cmd.scope,
764 has_feedback = cmd.feedback.is_some(),
765 "[tooldecision] received (delivery mechanism deferred to PR-3b/PR-4)"
766 );
767 Ok(None)
768 })
769 }),
770 );
771
772 engine.register_handler(
781 COMMAND_AGENT_ANSWER,
782 Box::new(|data, _ctx| {
783 Box::pin(async move {
784 let cmd: CommandAgentAnswerData = serde_json::from_value(data)
785 .map_err(|e| format!("agent.answer: {e}"))?;
786 if cmd.tool_use_id.is_empty() {
787 return Err("agent.answer: MISSING_ARG: tool_use_id".to_string());
788 }
789 let ctrl = blockcontroller::get_controller(&cmd.blockid)
790 .ok_or_else(|| format!("agent.answer: no controller for block {}", cmd.blockid))?;
791 if let Some(persistent_ctrl) = ctrl
792 .as_any()
793 .downcast_ref::<blockcontroller::persistent::PersistentSubprocessController>()
794 {
795 persistent_ctrl.answer_question(cmd.tool_use_id.clone(), cmd.answers)?;
796 tracing::info!(
797 block_id = %cmd.blockid,
798 tool_use_id = %cmd.tool_use_id,
799 "[agent.answer] control_response delivered to persistent stdin"
800 );
801 Ok(None)
802 } else {
803 Err("agent.answer: UNSUPPORTED_CONTROLLER: answering AskUserQuestion \
804 requires a persistent (host) agent; container/one-shot agents are \
805 not yet supported (Phase 2)".to_string())
806 }
807 })
808 }),
809 );
810
811 super::agent_handlers::register_agent_input_handlers(engine, &state);
813
814 super::shell_handlers::register_shell_handlers(engine, &state);
816
817 super::editor_handlers::register_editor_handlers(engine, &state);
819
820 super::lsp_handlers::register_lsp_handlers(engine, &state);
822
823 super::cli_handlers::register_cli_handlers(engine, &state);
825
826 super::tool_handlers::register_tool_handlers(engine, &state);
828
829 let broker_history = state.broker.clone();
831 engine.register_handler(
832 COMMAND_EVENT_READ_HISTORY,
833 Box::new(move |data, _ctx| {
834 let broker = broker_history.clone();
835 Box::pin(async move {
836 let cmd: CommandEventReadHistoryData = serde_json::from_value(data)
837 .map_err(|e| format!("eventreadhistory: {e}"))?;
838 let max_items = if cmd.maxitems == 0 { 1024 } else { cmd.maxitems };
839 let events = broker.read_event_history(&cmd.event, &cmd.scope, max_items);
840 Ok(Some(serde_json::to_value(&events).unwrap_or_default()))
841 })
842 }),
843 );
844
845 let config_watcher_setconfig = state.config_watcher.clone();
850 let event_bus_setconfig = state.event_bus.clone();
851 let lan_discovery_setconfig = state.lan_discovery.clone();
852 engine.register_handler(
853 COMMAND_SET_CONFIG,
854 Box::new(move |data, _ctx| {
855 let cw = config_watcher_setconfig.clone();
856 let eb = event_bus_setconfig.clone();
857 let lan = lan_discovery_setconfig.clone();
858 Box::pin(async move {
859 let new_keys: serde_json::Map<String, serde_json::Value> =
860 serde_json::from_value(data).map_err(|e| format!("setconfig: {e}"))?;
861
862 crate::backend::config_watcher_fs::merge_settings_to_disk(new_keys.clone())
864 .map_err(|e| format!("setconfig write: {e}"))?;
865
866 let merged_settings = crate::backend::config_watcher_fs::merge_settings_into_current(&cw, new_keys);
868 let lan_enabled = merged_settings.network_lan_discovery;
869 cw.update_settings(merged_settings);
870
871 lan.apply(lan_enabled);
876
877 let config = cw.get_full_config();
879 if let Ok(config_val) = serde_json::to_value(config.as_ref()) {
880 let event = crate::backend::eventbus::WSEventType {
881 eventtype: crate::backend::eventbus::WS_EVENT_RPC.to_string(),
882 oref: String::new(),
883 data: Some(serde_json::json!({
884 "command": "eventrecv",
885 "data": {
886 "event": "config",
887 "data": { "fullconfig": config_val }
888 }
889 })),
890 };
891 eb.broadcast_event(&event);
892 }
893 Ok(None)
894 })
895 }),
896 );
897
898 super::agent_handlers::register_agent_handlers(engine, &state);
900
901 super::drone_handlers::register_drone_handlers(engine, &state);
903
904 super::identity_handlers::register_identity_handlers(engine, &state);
907
908 super::install_handlers::register_install_handlers(engine, &state);
911
912 super::app_api::register_app_api_handlers(engine, &state);
914
915 super::muxbus_handlers::register_muxbus_handlers(engine, &state);
917
918 super::native_memory_handlers::register_native_memory_handlers(engine, &state);
920}
921
922
923fn parse_block_input(
925 cmd: &CommandBlockInputData,
926) -> Result<blockcontroller::BlockInputUnion, String> {
927 if !cmd.inputdata64.is_empty() {
928 let data = base64::engine::general_purpose::STANDARD
929 .decode(&cmd.inputdata64)
930 .map_err(|e| format!("controllerinput: base64 decode: {e}"))?;
931 return Ok(blockcontroller::BlockInputUnion::data(data));
932 }
933 if !cmd.signame.is_empty() {
934 return Ok(blockcontroller::BlockInputUnion::signal(&cmd.signame));
935 }
936 if let Some(ref ts_val) = cmd.termsize {
937 let ts: TermSize =
938 serde_json::from_value(ts_val.clone()).map_err(|e| format!("controllerinput: {e}"))?;
939 return Ok(blockcontroller::BlockInputUnion::resize(ts));
940 }
941 Err("controllerinput: no input data, signal, or termsize".to_string())
942}