mlua_swarm_server/operator_ws/login.rs
1//! REST-like Operator session resource.
2//!
3//! Provides the `POST/GET/DELETE /v1/operators` + `WS /v1/operators/:sid/ws`
4//! route family — the sole WS Operator session route. `session.rs` /
5//! `protocol.rs` are unchanged by this module.
6//!
7//! ## Login flow
8//!
9//! ```text
10//! POST /v1/operators { roles?: ["main-ai"], capability_manifest?: {...} }
11//! → 409 if any role already owns a live entry (roles alias exclusivity,
12//! v1.md §Auth session flow)
13//! → { sid: "S-<hex>", token: "<10-hex>", roles: [...] }
14//! The manifest is pinned to this session and later resolved through the
15//! Core `AgentBindingProvider` interface before any Runner-backed spawn.
16//!
17//! WS /v1/operators/:sid/ws
18//! Authorization: Bearer <token> (mandatory — no empty-string default)
19//! → 401 missing/empty Bearer, 404 unknown sid, 401 token mismatch
20//! → registers a `WSOperatorSession` into the engine's 3 registries
21//! (senior_bridge / spawn_hook / operator) + role aliases, same pattern
22//! as `handler::handle_socket`. Reconnect (same sid, matching token)
23//! reuses the existing `WSOperatorSession` via `replace_tx`.
24//!
25//! DELETE /v1/operators/:sid (Bearer required)
26//! → unregisters the 3 registries + role aliases + `operator_sessions`
27//! entry + releases `roles_to_sid` ownership.
28//!
29//! GET /v1/operators/:sid (Bearer required)
30//! → { sid, roles, connected }
31//! ```
32//!
33//! `OperatorSessionEntry` is the login-flow record (`AppState.operator_sessions`),
34//! distinct from `mlua_swarm::OperatorSession` (the engine-side
35//! `attach`/session-token record) and from `WSOperatorSession` (the 3-trait WS
36//! session, `session.rs`) — this module owns the mapping `sid → (token, roles,
37//! Option<WSOperatorSession>)` that the login flow is built on.
38
39use axum::{
40 extract::{
41 ws::{Message, WebSocket, WebSocketUpgrade},
42 Path, State,
43 },
44 http::{HeaderMap, StatusCode},
45 response::{IntoResponse, Response},
46 Json,
47};
48use futures_util::{sink::SinkExt, stream::StreamExt};
49use mlua_swarm::{AgentProviderManifest, Operator, SeniorBridge, SessionId, SpawnHook};
50use serde::{Deserialize, Serialize};
51use serde_json::json;
52use std::sync::Arc;
53use tokio::sync::{mpsc, Mutex};
54
55use super::protocol::{ClientMsg, PendingReply, ServerMsg};
56use super::session::WSOperatorSession;
57use crate::AppState;
58
59/// Login-flow record for a minted Operator session. Held in
60/// `AppState.operator_sessions`, keyed by `sid`. `ws_session` starts `None`
61/// (login only mints sid+token) and is set on first successful WS connect;
62/// on reconnect the same `WSOperatorSession` is reused (`replace_tx`) rather
63/// than re-registered.
64pub struct OperatorSessionEntry {
65 /// Server-minted session id (typed [`SessionId`] since issue #14).
66 pub sid: SessionId,
67 /// Bearer auth token (10-hex-char) required on the WS upgrade and admin routes.
68 pub token: String,
69 /// Role aliases claimed by this session (roles-exclusivity set).
70 pub roles: Vec<String>,
71 /// Provider-owned effective capability manifest submitted at join.
72 pub capability_manifest: Option<AgentProviderManifest>,
73 /// GH #81 Layer 2: unix epoch seconds when `POST /v1/operators` minted
74 /// this entry. Surfaced by `GET /v1/operators` so a recovery driver
75 /// can pick the oldest stale session without probing each sid
76 /// individually.
77 pub joined_at_secs: u64,
78 /// The reusable 3-trait session object once a WS has connected at least
79 /// once; `None` before first connect. Its sender tracks current connectivity.
80 pub ws_session: Mutex<Option<Arc<WSOperatorSession>>>,
81}
82
83// ─── POST /v1/operators (mint) ──────────────────────────────────────────────
84
85/// Body for `POST /v1/operators`.
86#[derive(Debug, Deserialize, Default)]
87pub struct OperatorsCreateReq {
88 /// Role aliases to claim exclusively (empty = no exclusivity claimed).
89 #[serde(default)]
90 pub roles: Vec<String>,
91 /// Effective execution capabilities supplied by the Operator/MainAI.
92 #[serde(default)]
93 pub capability_manifest: Option<AgentProviderManifest>,
94}
95
96/// Response for `POST /v1/operators`.
97#[derive(Debug, Serialize)]
98pub struct OperatorsCreateResp {
99 /// Newly minted session id (typed [`SessionId`]; serializes as the
100 /// plain `S-<hex>` string — the wire shape is unchanged).
101 pub sid: SessionId,
102 /// Bearer auth token required on the WS upgrade and admin routes.
103 pub token: String,
104 /// Echoes the granted role aliases.
105 pub roles: Vec<String>,
106}
107
108/// `POST /v1/operators`. Mints `sid` (`S-<hex>` — the shared `SessionId`
109/// shape; issue #11) + a 10-hex-char token
110/// (`mlua_swarm::types::secure_hex(5)` — OS-RNG hex, unguessable across
111/// calls and restarts, which is the point: this token is the sole bearer
112/// secret on the short-handle path). When `roles` is non-empty, checks
113/// `AppState.roles_to_sid` for conflicts under a single lock (check + insert
114/// atomic w.r.t. concurrent mints) and returns `409 CONFLICT` with the
115/// conflicting role names on collision. Empty `roles` never conflicts (= no
116/// exclusivity is claimed).
117pub async fn operators_create(
118 State(state): State<AppState>,
119 Json(req): Json<OperatorsCreateReq>,
120) -> Response {
121 let roles = req.roles;
122 let capability_manifest = req.capability_manifest;
123 // The sid is the operator-session identity, so it mints in the same
124 // `SessionId` shape (`S-<hex>`) as the engine-side session id — one
125 // session-id form across the system (issue #11 observation 2; the old
126 // `op-<uuid>` shape collided with the operator-backend registry prefix).
127 // It is an identifier, not a secret: `token` (secure_hex) is the sole
128 // bearer credential on this path.
129 let sid = SessionId::new();
130 let token = mlua_swarm::types::secure_hex(5);
131
132 {
133 let mut map = state.roles_to_sid.lock().await;
134 let conflicts: Vec<String> = roles
135 .iter()
136 .filter(|r| map.contains_key(r.as_str()))
137 .cloned()
138 .collect();
139 if !conflicts.is_empty() {
140 // GH #81 Layer 2 (a): identify the holding session per
141 // conflicted role so a recovery driver knows which sid to
142 // release without probing. The pre-#81 `conflicts: [role]`
143 // array stays byte-identical for callers that already
144 // ignore unknown keys; the new `conflicts_detail: [{role,
145 // sid}]` array is an additive companion.
146 let conflicts_detail: Vec<serde_json::Value> = conflicts
147 .iter()
148 .map(|r| {
149 let holder = map.get(r.as_str()).map(|sid| sid.to_string());
150 json!({ "role": r, "sid": holder })
151 })
152 .collect();
153 return (
154 StatusCode::CONFLICT,
155 Json(json!({
156 "error": "roles conflict",
157 "conflicts": conflicts,
158 "conflicts_detail": conflicts_detail,
159 })),
160 )
161 .into_response();
162 }
163 for r in &roles {
164 map.insert(r.clone(), sid.clone());
165 }
166 }
167
168 let joined_at_secs = std::time::SystemTime::now()
169 .duration_since(std::time::UNIX_EPOCH)
170 .map(|d| d.as_secs())
171 .unwrap_or(0);
172 let entry = Arc::new(OperatorSessionEntry {
173 sid: sid.clone(),
174 token: token.clone(),
175 roles: roles.clone(),
176 capability_manifest,
177 joined_at_secs,
178 ws_session: Mutex::new(None),
179 });
180 state
181 .operator_sessions
182 .lock()
183 .await
184 .insert(sid.clone(), entry);
185
186 (
187 StatusCode::OK,
188 Json(OperatorsCreateResp { sid, token, roles }),
189 )
190 .into_response()
191}
192
193// ─── WS /v1/operators/:sid/ws (Bearer required) ─────────────────────────────
194
195/// Extracts `Authorization: Bearer <token>`; missing header, wrong scheme, or
196/// an empty token all resolve to a `401` response. `Authorization` is
197/// mandatory on the WS path — there is no empty-string default.
198fn extract_bearer_token_required(headers: &HeaderMap) -> Result<String, Box<Response>> {
199 let token = headers
200 .get(axum::http::header::AUTHORIZATION)
201 .and_then(|v| v.to_str().ok())
202 .and_then(|s| s.strip_prefix("Bearer "))
203 .map(|s| s.trim().to_string())
204 .filter(|s| !s.is_empty());
205 token.ok_or_else(|| {
206 Box::new((StatusCode::UNAUTHORIZED, "missing or empty Bearer token").into_response())
207 })
208}
209
210/// `GET /v1/operators/:sid/ws` (WS upgrade). Bearer mandatory. `404` on
211/// unknown sid, `401` on token mismatch. On successful upgrade, registers (or
212/// reuses, on reconnect) a `WSOperatorSession` under `sid` — same 3-registry
213/// pattern as `handler::handle_socket`, plus role-alias registration for
214/// every role minted alongside this sid.
215pub async fn operators_ws_connect(
216 State(state): State<AppState>,
217 Path(sid): Path<String>,
218 headers: HeaderMap,
219 ws: WebSocketUpgrade,
220) -> Response {
221 let bearer = match extract_bearer_token_required(&headers) {
222 Ok(t) => t,
223 Err(resp) => return *resp,
224 };
225 // A string that doesn't even parse as a SessionId can't be a known sid.
226 let Ok(sid) = SessionId::parse(sid) else {
227 return (StatusCode::NOT_FOUND, "unknown sid").into_response();
228 };
229
230 let entry = {
231 let map = state.operator_sessions.lock().await;
232 map.get(&sid).cloned()
233 };
234 let entry = match entry {
235 Some(e) => e,
236 None => return (StatusCode::NOT_FOUND, "unknown sid").into_response(),
237 };
238 if !mlua_swarm::types::ct_eq(entry.token.as_bytes(), bearer.as_bytes()) {
239 return (StatusCode::UNAUTHORIZED, "token mismatch").into_response();
240 }
241
242 ws.on_upgrade(move |socket| handle_operator_socket(socket, state, entry))
243}
244
245/// Bidirectional pump for a single WS connection, bound to an
246/// `OperatorSessionEntry`. Owns the full wire protocol pump (write task /
247/// read task / `ClientMsg` dispatch / disconnect) for this session.
248async fn handle_operator_socket(
249 socket: WebSocket,
250 state: AppState,
251 entry: Arc<OperatorSessionEntry>,
252) {
253 let (tx, mut rx) = mpsc::unbounded_channel::<ServerMsg>();
254
255 let existing_ws = entry.ws_session.lock().await.clone();
256 let session = match existing_ws {
257 Some(ws_session) => {
258 // Reconnect: reuse the existing WSOperatorSession on this entry; only swap out `tx`.
259 ws_session.replace_tx(tx.clone()).await;
260 ws_session
261 }
262 None => {
263 let ws_session = Arc::new(WSOperatorSession::new_with_base_url(
264 entry.sid.clone(),
265 tx.clone(),
266 state.base_url.clone(),
267 ));
268 state
269 .engine
270 .register_senior_bridge(
271 entry.sid.clone(),
272 ws_session.clone() as Arc<dyn SeniorBridge>,
273 )
274 .await;
275 state
276 .engine
277 .register_spawn_hook(entry.sid.clone(), ws_session.clone() as Arc<dyn SpawnHook>)
278 .await;
279 state
280 .engine
281 .register_operator(entry.sid.clone(), ws_session.clone() as Arc<dyn Operator>)
282 .await;
283 if let Some(factory) = &state.ws_operator_factory {
284 factory
285 .register_operator(entry.sid.clone(), ws_session.clone() as Arc<dyn Operator>);
286 }
287 // Role exclusivity was already resolved at login (POST) time. Here
288 // we just bind the same session into the three registries + factory
289 // under its role aliases (same shape as handler::handle_socket's
290 // ?roles= path).
291 for role in &entry.roles {
292 if let Some(factory) = &state.ws_operator_factory {
293 factory
294 .register_operator(role.clone(), ws_session.clone() as Arc<dyn Operator>);
295 }
296 state
297 .engine
298 .register_operator(role.clone(), ws_session.clone() as Arc<dyn Operator>)
299 .await;
300 }
301 *entry.ws_session.lock().await = Some(ws_session.clone());
302 ws_session
303 }
304 };
305
306 let (mut ws_sink, mut ws_stream) = socket.split();
307
308 // write task: mpsc → WebSocket
309 let write_task = tokio::spawn(async move {
310 while let Some(msg) = rx.recv().await {
311 let txt = match serde_json::to_string(&msg) {
312 Ok(s) => s,
313 Err(_) => continue,
314 };
315 if ws_sink.send(Message::Text(txt)).await.is_err() {
316 break;
317 }
318 }
319 let _ = ws_sink.close().await;
320 });
321
322 // read task: WS message → ClientMsg parse → session.resolve_pending
323 let session_for_read = session.clone();
324 let read_result: Result<(), String> = async {
325 while let Some(item) = ws_stream.next().await {
326 match item {
327 Ok(Message::Text(t)) => {
328 let parsed: ClientMsg = match serde_json::from_str(&t) {
329 Ok(p) => p,
330 Err(_) => continue,
331 };
332 match parsed {
333 ClientMsg::Answer { req_id, value } => {
334 session_for_read
335 .resolve_pending(&req_id, PendingReply::Answer(value))
336 .await;
337 }
338 ClientMsg::HookAck { req_id, ok, reason } => {
339 session_for_read
340 .resolve_pending(&req_id, PendingReply::HookAck { ok, reason })
341 .await;
342 }
343 ClientMsg::SpawnAck {
344 req_id,
345 value,
346 ok,
347 error,
348 stats,
349 } => {
350 session_for_read
351 .resolve_pending(
352 &req_id,
353 PendingReply::SpawnAck {
354 value,
355 ok,
356 error,
357 stats,
358 },
359 )
360 .await;
361 }
362 ClientMsg::SpawnHalt {
363 req_id,
364 value,
365 reason,
366 } => {
367 session_for_read
368 .resolve_pending(&req_id, PendingReply::SpawnHalt { value, reason })
369 .await;
370 }
371 }
372 }
373 Ok(Message::Ping(_)) | Ok(Message::Pong(_)) => {}
374 Ok(Message::Close(_)) | Err(_) => break,
375 _ => {}
376 }
377 }
378 Ok(())
379 }
380 .await;
381
382 // Clear only this socket's sender. A reconnect may already have installed
383 // a replacement while this older socket was unwinding.
384 session.clear_tx_if(&tx).await;
385 write_task.abort();
386 let _ = read_result;
387}
388
389// ─── DELETE /v1/operators/:sid (Bearer required) ────────────────────────────
390
391/// Shared teardown for `DELETE /v1/operators/:sid` (`operators_delete`) and
392/// `DELETE /v1/operators/by-role/:role` (`operators_delete_by_role` — GH #81
393/// Layer 2 (c)): drops the 3 engine registries + role aliases +
394/// `ws_operator_factory` bindings + `operator_sessions` entry, and releases
395/// the sid's ownership in `roles_to_sid`. Idempotent w.r.t. a concurrent
396/// delete — every `remove` / `unregister` is a no-op when the entry is
397/// already gone.
398async fn teardown_operator_session(
399 state: &AppState,
400 sid: &SessionId,
401 entry: &Arc<OperatorSessionEntry>,
402) {
403 state.engine.unregister_senior_bridge(sid.as_str()).await;
404 state.engine.unregister_spawn_hook(sid.as_str()).await;
405 state.engine.unregister_operator(sid.as_str()).await;
406 if let Some(factory) = &state.ws_operator_factory {
407 factory.unregister_operator(sid.as_str());
408 }
409 for role in &entry.roles {
410 state.engine.unregister_operator(role).await;
411 if let Some(factory) = &state.ws_operator_factory {
412 factory.unregister_operator(role);
413 }
414 }
415
416 if let Some(session) = entry.ws_session.lock().await.take() {
417 // B-2: fail every parked spawn/ask/hook_before on this session
418 // right away. Teardown removes the session from `operator_sessions`
419 // below (no reconnect can find it again), so unlike a plain WS
420 // disconnect there is no reconnect/resend contract to preserve —
421 // an in-flight spawn parked in `send_and_await` would otherwise
422 // orphan until the run's sync timeout (up to 300s) fires.
423 session.fail_pending("operator session torn down").await;
424 session.clear_tx().await;
425 }
426
427 state.operator_sessions.lock().await.remove(sid);
428
429 {
430 let mut map = state.roles_to_sid.lock().await;
431 for role in &entry.roles {
432 if map.get(role) == Some(sid) {
433 map.remove(role);
434 }
435 }
436 }
437}
438
439/// `DELETE /v1/operators/:sid`. Bearer mandatory. `404` on unknown sid, `401`
440/// on token mismatch. Drops the 3 engine registries + role aliases +
441/// `ws_operator_factory` bindings + `operator_sessions` entry, and releases
442/// this sid's ownership in `roles_to_sid` (re-opening the role names for a
443/// future mint).
444pub async fn operators_delete(
445 State(state): State<AppState>,
446 Path(sid): Path<String>,
447 headers: HeaderMap,
448) -> Response {
449 let bearer = match extract_bearer_token_required(&headers) {
450 Ok(t) => t,
451 Err(resp) => return *resp,
452 };
453 let Ok(sid) = SessionId::parse(sid) else {
454 return (StatusCode::NOT_FOUND, "unknown sid").into_response();
455 };
456
457 let entry = {
458 let map = state.operator_sessions.lock().await;
459 map.get(&sid).cloned()
460 };
461 let entry = match entry {
462 Some(e) => e,
463 None => return (StatusCode::NOT_FOUND, "unknown sid").into_response(),
464 };
465 if !mlua_swarm::types::ct_eq(entry.token.as_bytes(), bearer.as_bytes()) {
466 return (StatusCode::UNAUTHORIZED, "token mismatch").into_response();
467 }
468
469 teardown_operator_session(&state, &sid, &entry).await;
470
471 StatusCode::NO_CONTENT.into_response()
472}
473
474// ─── GH #81 Layer 2: GET /v1/operators + DELETE /v1/operators/by-role/:role
475
476/// GH #81 Layer 2 (b): one entry in the `GET /v1/operators` list response.
477/// Bare identity fields (no token, no capability manifest — those live
478/// behind Bearer on `GET /v1/operators/:sid`); this list surface is
479/// read-only observability, on the same trust tier as `GET /v1/status`.
480#[derive(Debug, Serialize)]
481pub struct OperatorsListEntry {
482 /// Session id (`S-<hex>`) — safe to expose; token is the sole bearer secret.
483 pub sid: SessionId,
484 /// Role aliases held by this session.
485 pub roles: Vec<String>,
486 /// Unix epoch seconds when the session minted (from
487 /// [`OperatorSessionEntry::joined_at_secs`]).
488 pub joined_at_secs: u64,
489 /// Whether a WS is currently attached to this session (matches the
490 /// `connected` field on `GET /v1/operators/:sid`).
491 pub connected: bool,
492}
493
494/// Response body for `GET /v1/operators` (GH #81 Layer 2 (b)).
495#[derive(Debug, Serialize)]
496pub struct OperatorsListResp {
497 /// One entry per live session, ordered by `sid` (deterministic —
498 /// callers can `.iter().find(...)` without probing the map order).
499 pub operators: Vec<OperatorsListEntry>,
500}
501
502/// `GET /v1/operators`. Read-only enumeration of every live session's
503/// `{sid, roles, joined_at_secs, connected}` (GH #81 Layer 2 (b)). Same
504/// trust tier as `GET /v1/status` — no Bearer required; sids are
505/// identifiers, not secrets. Answers "which sid holds `main-ai`?"
506/// without probing every sid individually via `GET /v1/operators/:sid`,
507/// which was the pre-#81 recovery gap.
508pub async fn operators_list(State(state): State<AppState>) -> Response {
509 let entries: Vec<(SessionId, Arc<OperatorSessionEntry>)> = {
510 let map = state.operator_sessions.lock().await;
511 map.iter().map(|(k, v)| (k.clone(), v.clone())).collect()
512 };
513 let mut operators = Vec::with_capacity(entries.len());
514 for (sid, entry) in entries {
515 let session = entry.ws_session.lock().await.clone();
516 let connected = match session {
517 Some(session) => session.is_connected().await,
518 None => false,
519 };
520 operators.push(OperatorsListEntry {
521 sid,
522 roles: entry.roles.clone(),
523 joined_at_secs: entry.joined_at_secs,
524 connected,
525 });
526 }
527 operators.sort_by(|a, b| a.sid.as_str().cmp(b.sid.as_str()));
528 (StatusCode::OK, Json(OperatorsListResp { operators })).into_response()
529}
530
531/// `DELETE /v1/operators/by-role/:role`. Releases the session currently
532/// holding `role` without requiring the caller to know the sid or its
533/// Bearer token (GH #81 Layer 2 (c)). Recovery route for a stale session
534/// whose driver crashed after minting the sid — pre-#81 the only reliable
535/// recovery was a full server restart, which also dropped every OTHER live
536/// session. Same trust tier as the server-shutdown surface
537/// (`mlua_swarm_server_shutdown`): admin observability, no Bearer.
538///
539/// `404` when no session holds the role, `204` on successful teardown. The
540/// response body on `204` is empty (`teardown_operator_session` performs
541/// the same cleanup as `operators_delete`).
542pub async fn operators_delete_by_role(
543 State(state): State<AppState>,
544 Path(role): Path<String>,
545) -> Response {
546 let sid = {
547 let map = state.roles_to_sid.lock().await;
548 match map.get(role.as_str()) {
549 Some(sid) => sid.clone(),
550 None => {
551 return (
552 StatusCode::NOT_FOUND,
553 Json(json!({"error": "no session holds this role", "role": role})),
554 )
555 .into_response();
556 }
557 }
558 };
559 let entry = {
560 let map = state.operator_sessions.lock().await;
561 map.get(&sid).cloned()
562 };
563 let entry = match entry {
564 Some(e) => e,
565 None => {
566 // The role was mapped to a sid that has no matching
567 // `operator_sessions` entry — a torn state that a mint-time
568 // atomic guard prevents in normal operation. Release the
569 // stale role mapping so a future mint can reclaim the
570 // name, then report NOT_FOUND.
571 let mut map = state.roles_to_sid.lock().await;
572 if map.get(role.as_str()) == Some(&sid) {
573 map.remove(role.as_str());
574 }
575 return (
576 StatusCode::NOT_FOUND,
577 Json(json!({
578 "error": "torn role mapping cleared; role now open",
579 "role": role,
580 })),
581 )
582 .into_response();
583 }
584 };
585 teardown_operator_session(&state, &sid, &entry).await;
586 StatusCode::NO_CONTENT.into_response()
587}
588
589// ─── GET /v1/operators/:sid (Bearer required) ───────────────────────────────
590
591/// Response for `GET /v1/operators/:sid`.
592#[derive(Debug, Serialize)]
593pub struct OperatorsInfoResp {
594 /// Echoes the requested session id.
595 pub sid: SessionId,
596 /// Role aliases held by this session.
597 pub roles: Vec<String>,
598 /// Capability manifest pinned when this session joined.
599 #[serde(skip_serializing_if = "Option::is_none")]
600 pub capability_manifest: Option<AgentProviderManifest>,
601 /// Whether a WS is currently attached (not merely that the session ever connected).
602 pub connected: bool,
603}
604
605/// `GET /v1/operators/:sid`. Bearer mandatory. `404` on unknown sid, `401` on
606/// token mismatch. `connected` reflects whether the reusable session currently
607/// owns a live sender, not merely whether it connected at least once.
608pub async fn operators_info(
609 State(state): State<AppState>,
610 Path(sid): Path<String>,
611 headers: HeaderMap,
612) -> Response {
613 let bearer = match extract_bearer_token_required(&headers) {
614 Ok(t) => t,
615 Err(resp) => return *resp,
616 };
617 let Ok(sid) = SessionId::parse(sid) else {
618 return (StatusCode::NOT_FOUND, "unknown sid").into_response();
619 };
620
621 let entry = {
622 let map = state.operator_sessions.lock().await;
623 map.get(&sid).cloned()
624 };
625 let entry = match entry {
626 Some(e) => e,
627 None => return (StatusCode::NOT_FOUND, "unknown sid").into_response(),
628 };
629 if !mlua_swarm::types::ct_eq(entry.token.as_bytes(), bearer.as_bytes()) {
630 return (StatusCode::UNAUTHORIZED, "token mismatch").into_response();
631 }
632
633 let session = entry.ws_session.lock().await.clone();
634 let connected = match session {
635 Some(session) => session.is_connected().await,
636 None => false,
637 };
638 (
639 StatusCode::OK,
640 Json(OperatorsInfoResp {
641 sid: entry.sid.clone(),
642 roles: entry.roles.clone(),
643 capability_manifest: entry.capability_manifest.clone(),
644 connected,
645 }),
646 )
647 .into_response()
648}
649
650#[cfg(test)]
651mod tests {
652 use super::*;
653 use axum::http::HeaderValue;
654
655 fn headers_with_bearer(token: &str) -> HeaderMap {
656 let mut h = HeaderMap::new();
657 h.insert(
658 axum::http::header::AUTHORIZATION,
659 HeaderValue::from_str(&format!("Bearer {token}")).unwrap(),
660 );
661 h
662 }
663
664 #[test]
665 fn extract_bearer_token_required_accepts_valid() {
666 let h = headers_with_bearer("abc123");
667 assert_eq!(extract_bearer_token_required(&h).unwrap(), "abc123");
668 }
669
670 #[test]
671 fn extract_bearer_token_required_rejects_missing_header() {
672 let h = HeaderMap::new();
673 assert!(extract_bearer_token_required(&h).is_err());
674 }
675
676 #[test]
677 fn extract_bearer_token_required_rejects_empty_token() {
678 let h = headers_with_bearer("");
679 assert!(extract_bearer_token_required(&h).is_err());
680 }
681
682 #[test]
683 fn extract_bearer_token_required_rejects_wrong_scheme() {
684 let mut h = HeaderMap::new();
685 h.insert(
686 axum::http::header::AUTHORIZATION,
687 HeaderValue::from_static("Basic dXNlcjpwYXNz"),
688 );
689 assert!(extract_bearer_token_required(&h).is_err());
690 }
691
692 #[test]
693 fn operators_create_request_accepts_capability_manifest() {
694 let req: OperatorsCreateReq = serde_json::from_value(serde_json::json!({
695 "roles": ["main-ai"],
696 "capability_manifest": {
697 "provider_id": "main-ai-self-report",
698 "capabilities": [{
699 "launch_variant": "mse-coder",
700 "resolved_model": "claude-sonnet-4",
701 "effective_tools": ["Read", "Edit"]
702 }]
703 }
704 }))
705 .unwrap();
706 assert_eq!(req.roles, ["main-ai"]);
707 assert_eq!(
708 req.capability_manifest.unwrap().provider_id,
709 "main-ai-self-report"
710 );
711 }
712
713 #[test]
714 fn operators_create_request_keeps_manifest_optional_on_wire() {
715 let req: OperatorsCreateReq =
716 serde_json::from_value(serde_json::json!({ "roles": [] })).unwrap();
717 assert!(req.capability_manifest.is_none());
718 }
719}