agentmux_srv\server/
muxbus_handlers.rs1use 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 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 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 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 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
149pub 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}