agentmux_srv\server/
muxbus_handlers.rs

1// Copyright 2026, AgentMux Corp.
2// SPDX-License-Identifier: Apache-2.0
3
4//! MuxBus cloud connectivity RPC handlers.
5//!
6//! Three commands:
7//!   * `muxbus.login`      — PKCE browser flow (blocks until complete or timeout)
8//!   * `muxbus.status`     — current credential status
9//!   * `muxbus.disconnect` — clear stored credentials
10
11use std::sync::Arc;
12
13use serde::{Deserialize, Serialize};
14
15use crate::backend::rpc::engine::WshRpcEngine;
16
17use super::AppState;
18
19pub const COMMAND_MUXBUS_LOGIN: &str = "muxbus.login";
20pub const COMMAND_MUXBUS_STATUS: &str = "muxbus.status";
21pub const COMMAND_MUXBUS_DISCONNECT: &str = "muxbus.disconnect";
22
23#[derive(Debug, Deserialize)]
24#[serde(rename_all = "camelCase")]
25struct MuxBusLoginReq {
26    cognito_domain: String,
27    client_id: String,
28}
29
30#[derive(Debug, Serialize)]
31#[serde(rename_all = "camelCase")]
32struct MuxBusLoginResp {
33    success: bool,
34    email: String,
35    #[serde(skip_serializing_if = "Option::is_none")]
36    error: Option<String>,
37}
38
39#[derive(Debug, Serialize)]
40#[serde(rename_all = "camelCase")]
41struct MuxBusStatusResp {
42    connected: bool,
43    email: String,
44    cognito_domain: String,
45    expires_at: i64,
46    valid: bool,
47}
48
49pub fn register_muxbus_handlers(engine: &Arc<WshRpcEngine>, state: &AppState) {
50    // muxbus.login — PKCE browser flow, returns when browser login completes
51    let wstore_login = state.wstore.clone();
52    let http_client_login = state.http_client.clone();
53    engine.register_handler(
54        COMMAND_MUXBUS_LOGIN,
55        Box::new(move |data, _ctx| {
56            let wstore = wstore_login.clone();
57            let http = http_client_login.clone();
58            Box::pin(async move {
59                let req: MuxBusLoginReq = serde_json::from_value(data)
60                    .map_err(|e| format!("muxbus.login: {e}"))?;
61
62                match crate::muxbus::pkce::run_pkce_login(
63                    &req.cognito_domain,
64                    &req.client_id,
65                    &http,
66                )
67                .await
68                {
69                    Ok(result) => {
70                        if let Err(e) = wstore.muxbus_save(&result.credentials) {
71                            tracing::warn!(error = %e, "muxbus.login: failed to save credentials");
72                        }
73                        // Kick the cloud subscriber to open a WS with the new token
74                        if let Some(sub) = crate::muxbus::cloud_subscriber::get_global_subscriber() {
75                            sub.reload_token();
76                        }
77                        let resp = MuxBusLoginResp {
78                            success: true,
79                            email: result.credentials.user_email,
80                            error: None,
81                        };
82                        Ok(Some(serde_json::to_value(resp).unwrap()))
83                    }
84                    Err(e) => {
85                        tracing::warn!(error = %e, "muxbus.login: PKCE flow failed");
86                        let resp = MuxBusLoginResp {
87                            success: false,
88                            email: String::new(),
89                            error: Some(e),
90                        };
91                        Ok(Some(serde_json::to_value(resp).unwrap()))
92                    }
93                }
94            })
95        }),
96    );
97
98    // muxbus.status — return current credential state
99    let wstore_status = state.wstore.clone();
100    engine.register_handler(
101        COMMAND_MUXBUS_STATUS,
102        Box::new(move |_data, _ctx| {
103            let wstore = wstore_status.clone();
104            Box::pin(async move {
105                match wstore.muxbus_load() {
106                    Ok(Some(creds)) => {
107                        let valid = creds.is_valid();
108                        let resp = MuxBusStatusResp {
109                            connected: !creds.access_token.is_empty(),
110                            email: creds.user_email,
111                            cognito_domain: creds.cognito_domain,
112                            expires_at: creds.expires_at,
113                            valid,
114                        };
115                        Ok(Some(serde_json::to_value(resp).unwrap()))
116                    }
117                    Ok(None) => {
118                        let resp = MuxBusStatusResp {
119                            connected: false,
120                            email: String::new(),
121                            cognito_domain: String::new(),
122                            expires_at: 0,
123                            valid: false,
124                        };
125                        Ok(Some(serde_json::to_value(resp).unwrap()))
126                    }
127                    Err(e) => Err(format!("muxbus.status: {e}")),
128                }
129            })
130        }),
131    );
132
133    // muxbus.disconnect — clear credentials
134    let wstore_disconnect = state.wstore.clone();
135    engine.register_handler(
136        COMMAND_MUXBUS_DISCONNECT,
137        Box::new(move |_data, _ctx| {
138            let wstore = wstore_disconnect.clone();
139            Box::pin(async move {
140                wstore
141                    .muxbus_clear()
142                    .map_err(|e| format!("muxbus.disconnect: {e}"))?;
143                Ok(Some(serde_json::json!({})))
144            })
145        }),
146    );
147}
148
149/// Inject MUXBUS_TOKEN into spawn env if credentials are stored and valid.
150/// Token refresh is async — this path just injects whatever is currently stored.
151/// Agents should re-spawn after the user refreshes via muxbus.login if the token expires.
152pub fn inject_muxbus_env(
153    wstore: &crate::backend::storage::store::Store,
154    env_vars: &mut std::collections::HashMap<String, String>,
155) {
156    let creds = match wstore.muxbus_load() {
157        Ok(Some(c)) => c,
158        Ok(None) => return,
159        Err(e) => {
160            tracing::warn!(error = %e, "muxbus inject: failed to load credentials");
161            return;
162        }
163    };
164
165    if creds.access_token.is_empty() {
166        return;
167    }
168
169    if creds.is_valid() {
170        env_vars.insert("MUXBUS_TOKEN".to_string(), creds.access_token.clone());
171        env_vars.insert("MUXBUS_COGNITO_DOMAIN".to_string(), creds.cognito_domain.clone());
172        tracing::debug!(email = creds.user_email, "muxbus: injected MUXBUS_TOKEN into spawn env");
173    } else {
174        tracing::warn!(
175            email = creds.user_email,
176            expires_at = creds.expires_at,
177            "muxbus: token expired, skipping injection — user should reconnect via muxbus.login"
178        );
179    }
180}