agentmux_srv\server/
websocket.rs

1// Copyright 2025-2026, AgentMux Corp.
2// SPDX-License-Identifier: Apache-2.0
3
4use 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/// Incoming WebSocket message envelope.
37/// Supports both ping/pong messages and wscommand-based RPC.
38#[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    // Fields for setblocktermsize / blockinput
47    blockid: Option<String>,
48    inputdata64: Option<String>,
49    termsize: Option<serde_json::Value>,
50    // Fields for bus:* commands
51    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    // Two egress lanes per connection: interactive (terminal echo, RPC-routed
76    // wave events, obj updates) and background (droppable perf telemetry —
77    // sysinfo/blockstats). The select! below drains priority before background
78    // so typing never waits behind a perf tick. See
79    // docs/specs/SPEC_TERMINAL_INPUT_PRIORITY_OVER_SYSINFO_2026_06_16.md.
80    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    // Optional messagebus receiver — activated when pane sends bus:register
87    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    // Send initial "config" wave event via the RPC eventrecv path so the frontend
91    // populates fullConfigAtom (and shows the widget bar).
92    // Frontend only processes events via: {"eventtype":"rpc","data":{"command":"eventrecv","data":{"event":"config","data":{...}}}}
93    {
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    // Create RPC engine for this connection
115    let t = std::time::Instant::now();
116    let (engine, mut rpc_output_rx) = WshRpcEngine::new();
117
118    // Register handlers
119    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    // Periodic ping interval (10 seconds)
124    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: poll branches top-to-bottom so interactive terminal I/O
130            // always wins over droppable perf telemetry. Order = incoming
131            // keystrokes → RPC replies → agent bus → priority events (terminal
132            // echo) → background events (sysinfo/blockstats) → keepalive ping.
133            // See docs/specs/SPEC_TERMINAL_INPUT_PRIORITY_OVER_SYSINFO_2026_06_16.md.
134            biased;
135
136            // Incoming WebSocket messages → parse & dispatch (keystrokes first)
137            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:register returned a new receiver
148                                bus_rx = Some(new_rx);
149                                bus_agent_id = Some(agent_id);
150                            }
151                            _ => {}
152                        }
153                    }
154                    Some(Ok(_)) => {
155                        // Binary or other message types — ignore
156                    }
157                    Some(Err(_)) => break,
158                }
159            }
160
161            // Forward RPC engine output → WebSocket (wrapped as eventtype:rpc)
162            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            // Forward MessageBus messages → WebSocket (if registered as agent)
174            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            // Priority event lane → WebSocket. Two sources feed it:
191            //   1. WPS Broker (via EventBusBridge) — already wrapped as
192            //      { eventtype: "rpc", data: { command: "eventrecv", data: WaveEvent } }
193            //   2. Direct broadcasts (e.g., SetMeta's obj:update) — raw
194            //      { eventtype: "waveobj:update", oref: "block:xxx", data: ... }
195            // This carries terminal echo output and all interactive events.
196            Some(event) = priority_rx.recv() => {
197                if forward_event(&mut socket, event).await {
198                    break;
199                }
200            }
201
202            // Background event lane → WebSocket. Droppable perf telemetry
203            // (sysinfo/blockstats) only; serviced when the priority lanes above
204            // are momentarily idle, so it can never delay a keystroke echo.
205            //
206            // Tradeoff of `biased;`: a SUSTAINED priority-lane flood (e.g.
207            // `! yes` pumping terminal output continuously) keeps the priority
208            // branches ready, so this lane is starved and its unbounded receiver
209            // accumulates sysinfo/blockstats until the flood pauses, then
210            // flushes as a stale burst. Acceptable for now because telemetry
211            // ingress is low-rate (sysinfo ~1/s) and a flood also saturates the
212            // socket writes — during which the priority lanes do drain, so this
213            // lane gets serviced — making true never-empty starvation the rare
214            // worst case. Phase 2 bounds it properly by coalescing this lane to
215            // the latest reading per (event, scope), so it cannot grow and the
216            // post-flood flush is one current frame, not a backlog. See
217            // SPEC_TERMINAL_INPUT_PRIORITY_OVER_SYSINFO_2026_06_16.md §4.2.
218            Some(event) = background_rx.recv() => {
219                if forward_event(&mut socket, event).await {
220                    break;
221                }
222            }
223
224            // Periodic ping
225            _ = 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    // Unregister from messagebus if this connection was an agent
244    if let Some(ref agent_id) = bus_agent_id {
245        state.messagebus.unregister(agent_id);
246    }
247}
248
249/// Forward a queued event-bus value to the WebSocket. Returns `true` if the
250/// send failed (the caller should break the loop).
251///
252/// Two shapes arrive: already-RPC-wrapped values (from the WPS broker via
253/// EventBusBridge) are forwarded as-is; raw event-bus values (e.g. SetMeta's
254/// `waveobj:update`) are wrapped as an RPC `eventrecv` so the frontend
255/// WshRouter routes them to handleWaveEvent → updateWaveObject. Shared by both
256/// the priority and background egress lanes.
257async fn forward_event(socket: &mut WebSocket, event: serde_json::Value) -> bool {
258    let msg = if event["eventtype"] == "rpc" {
259        // Already an RPC message (from WPS broker via EventBusBridge)
260        serde_json::to_string(&event).unwrap_or_default()
261    } else {
262        // Raw event bus event — wrap as RPC eventrecv
263        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
280/// Handle an incoming text message.
281/// Returns Err(true) if the socket send failed.
282/// Returns Ok(Some((rx, agent_id))) if a bus:register was processed.
283async 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    // Handle ping/pong by type field
298    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                    // Try direct PTY injection via ReactiveHandler first
355                    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                    // Non-"agent not found" error — report it
374                    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                    // Fall back to MessageBus WebSocket push
385                    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    // Handle wscommand-based messages
426    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    // getfullconfig → return full config as JSON
482    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    // routeannounce → log + no-op (fire-and-forget, may have no reqid)
498    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    // routeunannounce → no-op
509    engine.register_handler(
510        COMMAND_ROUTE_UNANNOUNCE,
511        Box::new(|_data, _ctx| Box::pin(async move { Ok(None) })),
512    );
513
514    // eventsub → register subscription with the WPS broker
515    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    // eventunsub → unsubscribe from the WPS broker
533    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    // eventunsuball → unsubscribe all from the WPS broker
551    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    // setmeta → update object metadata in the DB, broadcast update event
566    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                // Read the updated object and broadcast a proper WaveObjUpdate
581                // so all WS clients refresh their atoms with the new data.
582                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    // getmeta → return metadata for a wave object
605    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                        // Return the "meta" field if present, otherwise the full object
619                        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    // waveinfo → return version and build info
629    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    // getwaveairatelimit → AgentMux has no rate limits; return unlimited/unknown
643    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    // controllerresync → load block from DB, create/restart controller with PTY
660    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    // controllerinput → route keyboard input / signals / resize to block controller
700    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    // tooldecision → reply to a per-tool-call permission gate.
714    //
715    // The original PR-3a draft tried to write `y\n` / `n\n` to the
716    // subprocess's stdin via `blockcontroller::send_input`. Codex P1
717    // on PR #557 caught that this would fail: `SubprocessController::
718    // send_input` (and `PersistentSubprocessController::send_input`)
719    // both reject raw `input_data`, returning `Err("...use
720    // AgentInputCommand")`. The deeper truth is that AgentMux runs
721    // the agent CLI in non-interactive `--print` mode — the CLI
722    // never reads stdin and a y/n write would be a no-op even if
723    // the controller accepted it. See SPEC_DECISION_PROMPT
724    // _2026_04_24.md §9.1.
725    //
726    // For now this handler accepts the decision, validates the
727    // payload, logs it (audit trail via `~/.agentmux/logs/`), and
728    // returns Ok. The actual delivery mechanism (rule-persistence
729    // for next-turn application, or interactive-mode subprocess
730    // launch with stdin write) is decided in PR-3b / PR-4 once we
731    // pick a CLI integration strategy.
732    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                // Validate scope so PR-3b's rules-persistence layer
748                // can trust the value without re-checking. Reagent P1
749                // round-3 on PR #557.
750                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    // agent.answer → deliver an AskUserQuestion answer to the running agent CLI
773    // via the Agent SDK **control protocol**: the persistent controller replies
774    // to the CLI's parked `can_use_tool` control_request with a control_response
775    // carrying `updatedInput.answers`. (Delivering a `tool_result` on stdin does
776    // NOT work — the CLI auto-rejects AskUserQuestion within the turn; see spec
777    // §2/§3.) Only the persistent controller speaks the control protocol;
778    // container/one-shot subprocess agents return UNSUPPORTED_CONTROLLER (Phase 2).
779    // Spec: docs/specs/SPEC_AGENT_CONTROL_PROTOCOL_2026_06_15.md.
780    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    // Agent input/stop + subprocess spawn handlers
812    super::agent_handlers::register_agent_input_handlers(engine, &state);
813
814    // Shell exec/stop handlers
815    super::shell_handlers::register_shell_handlers(engine, &state);
816
817    // Editor/file-ops + write_agent_config handlers
818    super::editor_handlers::register_editor_handlers(engine, &state);
819
820    // LSP handlers (lspstart, lspsend, lspstop)
821    super::lsp_handlers::register_lsp_handlers(engine, &state);
822
823    // CLI handlers (resolvecli, checkcliauth, runclilogin)
824    super::cli_handlers::register_cli_handlers(engine, &state);
825
826    // Tool store handlers (gettoolstatus, installtool)
827    super::tool_handlers::register_tool_handlers(engine, &state);
828
829    // eventreadhistory → read persisted event history from the WPS broker
830    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    // setconfig → merge settings keys into settings.json AND update in-memory config immediately.
846    // Writing to disk + broadcasting directly gives instant UI response without waiting for
847    // the fs watcher (which has a ~300-800ms debounce + polling delay on Windows).
848    // The fs watcher's subsequent reload is a no-op (settings already up to date).
849    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                // 1. Write to disk (fs watcher will re-broadcast, harmlessly)
863                crate::backend::config_watcher_fs::merge_settings_to_disk(new_keys.clone())
864                    .map_err(|e| format!("setconfig write: {e}"))?;
865
866                // 2. Update in-memory config immediately
867                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                // 3. Live-toggle LAN discovery if the key changed. `apply` is
872                //    idempotent so it's safe to call unconditionally — when the
873                //    daemon is already in the requested state, this is a no-op.
874                //    See specs/lan-discovery-toggle.md.
875                lan.apply(lan_enabled);
876
877                // 4. Broadcast updated config now — no waiting for fs watcher
878                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    // Agent handlers (definitions, content, skills, history, import, reseed)
899    super::agent_handlers::register_agent_handlers(engine, &state);
900
901    // Drone handlers (issue #753 — Drone pane DAG executor)
902    super::drone_handlers::register_drone_handlers(engine, &state);
903
904    // Pre-launch OAuth handlers (auth.start / poll / submitcallback /
905    // cancel / submitapikey — see docs/specs/SPEC_PRE_LAUNCH_OAUTH_FLOW_2026_05_14.md)
906    super::identity_handlers::register_identity_handlers(engine, &state);
907
908    // Install handlers (install.start / install.cancel — see
909    // docs/specs/SPEC_AGENT_INSTALL_STAGE_2026_05_17.md)
910    super::install_handlers::register_install_handlers(engine, &state);
911
912    // App API handlers (agent.open, agent.send, agent.stop, agent.status, agent.list, agent.output)
913    super::app_api::register_app_api_handlers(engine, &state);
914
915    // MuxBus cloud connectivity (muxbus.login / muxbus.status / muxbus.disconnect)
916    super::muxbus_handlers::register_muxbus_handlers(engine, &state);
917
918    // Native memory file browser (agent:memory:list / read_file / write_file)
919    super::native_memory_handlers::register_native_memory_handlers(engine, &state);
920}
921
922
923/// Parse a CommandBlockInputData into a BlockInputUnion.
924fn 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}