1use std::sync::Arc;
16
17use dashmap::DashMap;
18use tokio::sync::{broadcast, mpsc, Mutex};
19use trustee_core::session::Session;
20use trustee_core::types::TuiMessage;
21
22use crate::auth::AuthState;
23
24pub struct UserSessionEntry {
30 pub session: Arc<Mutex<Session>>,
32 pub ws_tx: broadcast::Sender<String>,
34 pub created_at: chrono::DateTime<chrono::Utc>,
36 pub last_active: Arc<Mutex<chrono::DateTime<chrono::Utc>>>,
39}
40
41pub struct UserSessions {
43 pub sessions: DashMap<String, UserSessionEntry>,
45 pub token_store: Arc<pep::MemoryTokenStore>,
47 pub active_session_id: Mutex<String>,
49}
50
51#[derive(Debug, serde::Serialize)]
53pub struct SessionListItem {
54 pub session_id: String,
55 pub session_name: Option<String>,
56 pub workflow_state: String,
57 pub created_at: String,
58 pub last_active: String,
59 pub handoff_count: u32,
64}
65
66#[derive(Debug)]
68pub enum SessionError {
69 MaxSessionsReached(usize),
71 NotFound(String),
73 NotIdle(String),
75}
76
77impl std::fmt::Display for SessionError {
78 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
79 match self {
80 SessionError::MaxSessionsReached(n) => {
81 write!(f, "Maximum {} sessions per user reached", n)
82 }
83 SessionError::NotFound(id) => write!(f, "Session {} not found", id),
84 SessionError::NotIdle(state) => write!(f, "Session is not idle (state: {})", state),
85 }
86 }
87}
88
89impl std::error::Error for SessionError {}
90
91pub type SessionRegistry = Arc<DashMap<String, UserSessions>>;
93
94pub struct McpLoaderEntry {
105 pub loader: Option<std::sync::Arc<abk::agent::McpToolLoader>>,
107 pub fingerprint: u64,
110 pub built_at: chrono::DateTime<chrono::Utc>,
111 pub degraded: Option<String>,
114 pub failed_at: Option<chrono::DateTime<chrono::Utc>>,
117}
118
119pub const MCP_BUILD_RETRY_BACKOFF: std::time::Duration = std::time::Duration::from_secs(30);
121
122#[derive(Debug, Clone)]
129pub struct ThqDispatchEntry {
130 pub user_key: String,
135 pub service_token: Option<String>,
139 pub issuer_url: Option<String>,
145}
146
147pub fn service_issuers_from_config(config_toml: &str) -> Vec<String> {
151 let Ok(v) = config_toml.parse::<toml::Value>() else {
152 return Vec::new();
153 };
154 let Some(creds) = v.get("mcp").and_then(|m| m.get("credentials")).and_then(|c| c.as_table())
155 else {
156 return Vec::new();
157 };
158 let mut out = Vec::new();
159 for (_name, cred) in creds {
160 if cred.get("type").and_then(|t| t.as_str()) == Some("service-account") {
161 if let Some(issuer) = cred.get("issuer_url").and_then(|i| i.as_str()) {
162 if !out.iter().any(|e: &String| e == issuer) {
163 out.push(issuer.to_string());
164 }
165 }
166 }
167 }
168 out
169}
170
171#[derive(Clone)]
173pub struct ServerState {
174 pub sessions: SessionRegistry,
176 pub ws_tx: broadcast::Sender<String>,
178 pub auth: Option<Arc<AuthState>>,
180 pub config_toml: Option<String>,
182 pub secrets: Option<std::collections::HashMap<String, String>>,
184 pub build_info: Option<trustee_core::types::BuildInfo>,
186 pub workflow_semaphore: Arc<tokio::sync::Semaphore>,
189 pub max_sessions_per_user: usize,
191 pub allow_llm_overlay: bool,
194 pub mcp_loaders: Arc<DashMap<String, McpLoaderEntry>>,
196 mcp_build_locks: Arc<DashMap<String, Arc<tokio::sync::Mutex<()>>>>,
198 pub thq_dispatch: Arc<DashMap<String, ThqDispatchEntry>>,
200 pub agent_dispatch_tokens: Arc<DashMap<String, (String, std::time::Instant)>>,
203}
204
205impl ServerState {
206 pub fn new(
208 session: Session,
209 ws_tx: broadcast::Sender<String>,
210 auth: Option<Arc<AuthState>>,
211 ) -> Self {
212 let sessions = Arc::new(DashMap::new());
213
214 let token_store = Arc::new(pep::MemoryTokenStore::new());
216 let (ws_tx_entry, _) = broadcast::channel::<String>(256);
217
218 let now = chrono::Utc::now();
219 let initial_entry = UserSessionEntry {
220 session: Arc::new(Mutex::new(session)),
221 ws_tx: ws_tx_entry,
222 created_at: now,
223 last_active: Arc::new(Mutex::new(now)),
224 };
225
226 let user_sessions = UserSessions {
227 sessions: DashMap::new(),
228 token_store,
229 active_session_id: Mutex::new(String::new()),
230 };
231 user_sessions.sessions.insert("default".to_string(), initial_entry);
232
233 sessions.insert("default".to_string(), user_sessions);
234
235 Self {
236 sessions,
237 ws_tx,
238 auth,
239 config_toml: None,
240 secrets: None,
241 build_info: None,
242 workflow_semaphore: Arc::new(tokio::sync::Semaphore::new(8)),
243 max_sessions_per_user: 4,
244 allow_llm_overlay: false,
245 mcp_loaders: Arc::new(DashMap::new()),
246 mcp_build_locks: Arc::new(DashMap::new()),
247 thq_dispatch: Arc::new(DashMap::new()),
248 agent_dispatch_tokens: Arc::new(DashMap::new()),
249 }
250 }
251
252 pub fn with_config_toml(mut self, config_toml: String) -> Self {
253 self.config_toml = Some(config_toml);
254 self
255 }
256
257 pub fn service_issuers(&self) -> Vec<String> {
268 service_issuers_from_config(self.config_toml.as_deref().unwrap_or(""))
269 }
270
271 pub fn with_secrets(mut self, secrets: std::collections::HashMap<String, String>) -> Self {
272 self.secrets = Some(secrets);
273 self
274 }
275
276 pub fn with_build_info(mut self, build_info: trustee_core::types::BuildInfo) -> Self {
277 self.build_info = Some(build_info);
278 self
279 }
280
281 pub fn with_max_concurrent_workflows(mut self, max: usize) -> Self {
282 self.workflow_semaphore = Arc::new(tokio::sync::Semaphore::new(max));
283 self
284 }
285
286 pub fn with_max_sessions_per_user(mut self, max: usize) -> Self {
288 self.max_sessions_per_user = max;
289 self
290 }
291
292 pub fn with_allow_llm_overlay(mut self, allow: bool) -> Self {
295 self.allow_llm_overlay = allow;
296 self
297 }
298
299 pub async fn create_session(
309 &self,
310 user_key: &str,
311 session_name: Option<String>,
312 identity: Option<String>,
313 activate: bool,
314 ) -> Result<String, SessionError> {
315 let user_sessions = self
317 .sessions
318 .entry(user_key.to_string())
319 .or_insert_with(|| UserSessions {
320 sessions: DashMap::new(),
321 token_store: Arc::new(pep::MemoryTokenStore::new()),
322 active_session_id: Mutex::new(String::new()),
323 });
324
325 if user_sessions.sessions.len() >= self.max_sessions_per_user {
327 return Err(SessionError::MaxSessionsReached(self.max_sessions_per_user));
328 }
329
330 let (mut session, workflow_rx) = Session::new();
332
333 if let Some(ref config_toml) = self.config_toml {
335 session.config_toml = Some(config_toml.clone());
336 session.parse_auto_handoff_config();
337 if let Ok(table) = config_toml.parse::<toml::Value>() {
338 if let Some(name) = table
339 .get("agent")
340 .and_then(|a| a.get("name"))
341 .and_then(|n| n.as_str())
342 {
343 session.agent_name = name.to_string();
344 }
345 }
346 }
347
348 session.secrets = self.secrets.clone();
349 session.build_info = self.build_info.clone();
350
351 self.apply_user_isolation(&mut session, user_key);
353
354 session.session_name = session_name;
356
357 session.identity = identity;
359
360 let (ws_tx_entry, _) = broadcast::channel::<String>(256);
362
363 let session_id = format!(
365 "session_{}_{}",
366 chrono::Utc::now().format("%Y_%m_%d_%H_%M"),
367 &uuid::Uuid::new_v4().to_string()[..8]
368 );
369
370 let now = chrono::Utc::now();
371
372 user_sessions.sessions.insert(
374 session_id.clone(),
375 UserSessionEntry {
376 session: Arc::new(Mutex::new(session)),
377 ws_tx: ws_tx_entry.clone(),
378 created_at: now,
379 last_active: Arc::new(Mutex::new(now)),
380 },
381 );
382
383 if activate {
392 *user_sessions.active_session_id.lock().await = session_id.clone();
393 }
394
395 let session_arc = user_sessions
397 .sessions
398 .get(&session_id)
399 .map(|e| e.session.clone());
400 if let Some(session_arc) = session_arc {
401 self.spawn_user_drain_task(
402 session_id.clone(),
403 session_arc,
404 ws_tx_entry,
405 workflow_rx,
406 );
407 }
408
409 Ok(session_id)
410 }
411
412 pub async fn get_session(
415 &self,
416 user_key: &str,
417 session_id: &str,
418 ) -> Option<(Arc<Mutex<Session>>, broadcast::Sender<String>)> {
419 let user_sessions = self.sessions.get(user_key)?;
420 let entry = user_sessions.sessions.get(session_id)?;
421
422 let now = chrono::Utc::now();
424 *entry.last_active.lock().await = now;
425
426 Some((entry.session.clone(), entry.ws_tx.clone()))
427 }
428
429 pub async fn get_session_by_any_id(
447 &self,
448 user_key: &str,
449 id: &str,
450 ) -> Option<(String, Arc<Mutex<Session>>, broadcast::Sender<String>)> {
451 let user_sessions = self.sessions.get(user_key)?;
453 if let Some(entry) = user_sessions.sessions.get(id) {
454 let now = chrono::Utc::now();
456 *entry.last_active.lock().await = now;
457 return Some((id.to_string(), entry.session.clone(), entry.ws_tx.clone()));
458 }
459
460 for entry in user_sessions.sessions.iter() {
462 let session = entry.session.lock().await;
463 if session.session_id.as_deref() == Some(id) {
464 let key = entry.key().clone();
465 let ws_tx = entry.ws_tx.clone();
466 drop(session);
467 let now = chrono::Utc::now();
469 *entry.last_active.lock().await = now;
470 return Some((key, entry.session.clone(), ws_tx));
471 }
472 }
473
474 None
475 }
476
477 pub async fn list_sessions(&self, user_key: &str) -> Vec<SessionListItem> {
479 let Some(user_sessions) = self.sessions.get(user_key) else {
480 return Vec::new();
481 };
482
483 let mut items = Vec::new();
484 for entry in user_sessions.sessions.iter() {
485 let session = entry.session.lock().await;
486 let workflow_state = match session.workflow_state {
487 trustee_core::types::WorkflowState::Idle => "Idle",
488 trustee_core::types::WorkflowState::Running => "Running",
489 trustee_core::types::WorkflowState::Cancelling => "Cancelling",
490 };
491 let last_active = entry.last_active.lock().await;
492 items.push(SessionListItem {
493 session_id: entry.key().clone(),
494 session_name: session.session_name.clone(),
495 workflow_state: workflow_state.to_string(),
496 created_at: entry.created_at.to_rfc3339(),
497 last_active: last_active.to_rfc3339(),
498 handoff_count: session.handoff_count,
499 });
500 }
501 drop(user_sessions);
502
503 items.sort_by(|a, b| b.last_active.cmp(&a.last_active));
505 items
506 }
507
508 pub async fn destroy_session(
510 &self,
511 user_key: &str,
512 session_id: &str,
513 ) -> Result<(), SessionError> {
514 let user_sessions = self
515 .sessions
516 .get(user_key)
517 .ok_or_else(|| SessionError::NotFound(session_id.to_string()))?;
518
519 {
521 let entry = user_sessions
522 .sessions
523 .get(session_id)
524 .ok_or_else(|| SessionError::NotFound(session_id.to_string()))?;
525 let session = entry.session.lock().await;
526 if session.workflow_state != trustee_core::types::WorkflowState::Idle {
527 let state_str = match session.workflow_state {
528 trustee_core::types::WorkflowState::Running => "Running",
529 trustee_core::types::WorkflowState::Cancelling => "Cancelling",
530 _ => "Unknown",
531 };
532 return Err(SessionError::NotIdle(state_str.to_string()));
533 }
534 }
535
536 user_sessions.sessions.remove(session_id);
538
539 let mut active_id = user_sessions.active_session_id.lock().await;
541 if &*active_id == session_id {
542 let mut newest: Option<(String, chrono::DateTime<chrono::Utc>)> = None;
544 for entry in user_sessions.sessions.iter() {
545 let la = entry.last_active.lock().await;
546 if newest.as_ref().map_or(true, |(_, t)| *la > *t) {
547 newest = Some((entry.key().clone(), *la));
548 }
549 }
550 *active_id = newest.map(|(id, _)| id).unwrap_or_default();
551 }
552
553 Ok(())
554 }
555
556 pub async fn ensure_active_session(
565 &self,
566 user_key: &str,
567 ) -> (
568 String,
569 Arc<Mutex<Session>>,
570 broadcast::Sender<String>,
571 Arc<pep::MemoryTokenStore>,
572 ) {
573 let token_store = {
575 let user_sessions = self
576 .sessions
577 .entry(user_key.to_string())
578 .or_insert_with(|| UserSessions {
579 sessions: DashMap::new(),
580 token_store: Arc::new(pep::MemoryTokenStore::new()),
581 active_session_id: Mutex::new(String::new()),
582 });
583 user_sessions.token_store.clone()
584 };
585
586 let active_id = {
588 let user_sessions = self.sessions.get(user_key).unwrap();
589 let guard = user_sessions.active_session_id.lock().await;
590 guard.clone()
591 };
592
593 if !active_id.is_empty() {
594 if let Some((session, ws_tx)) = self.get_session(user_key, &active_id).await {
595 return (active_id, session, ws_tx, token_store);
596 }
597 }
599
600 let existing_session: Option<(String, Arc<Mutex<Session>>, broadcast::Sender<String>)> = {
604 let user_sessions = self.sessions.get(user_key).unwrap();
605 let result = user_sessions.sessions.iter().next().map(|first| {
606 (
607 first.key().clone(),
608 first.session.clone(),
609 first.ws_tx.clone(),
610 )
611 });
612 result
613 };
614 if let Some((id, session, ws_tx)) = existing_session {
615 let now = chrono::Utc::now();
616 if let Some(entry) = self.sessions.get(user_key) {
617 if let Some(e) = entry.sessions.get(&id) {
618 *e.last_active.lock().await = now;
619 }
620 *entry.active_session_id.lock().await = id.clone();
621 }
622
623 return (id, session, ws_tx, token_store);
624 }
625
626 let session_id = self
628 .create_session(user_key, None, None, true)
629 .await
630 .unwrap_or_else(|_| "default".to_string());
631
632 let (session, ws_tx) = self
633 .get_session(user_key, &session_id)
634 .await
635 .expect("just-created session must exist");
636
637 (session_id, session, ws_tx, token_store)
638 }
639
640 pub async fn ensure_user_session(
643 &self,
644 user_key: &str,
645 ) -> (Arc<Mutex<Session>>, broadcast::Sender<String>, Arc<pep::MemoryTokenStore>) {
646 let (_id, session, ws_tx, token_store) = self.ensure_active_session(user_key).await;
647 (session, ws_tx, token_store)
648 }
649
650 pub async fn set_active_session(&self, user_key: &str, session_id: &str) {
652 if let Some(user_sessions) = self.sessions.get(user_key) {
653 if user_sessions.sessions.contains_key(session_id) {
654 *user_sessions.active_session_id.lock().await = session_id.to_string();
655 }
656 }
657 }
658
659 pub fn get_user_home_dir(&self, user_key: &str) -> Option<std::path::PathBuf> {
670 let hash = trustee_core::user_hash(user_key);
671 dirs::home_dir().map(|home| home.join(".trustee").join("users").join(&hash))
672 }
673
674 pub fn get_user_config_and_home(&self, user_key: &str) -> (Option<String>, Option<std::path::PathBuf>) {
679 (self.config_toml.clone(), self.get_user_home_dir(user_key))
680 }
681
682 fn apply_user_isolation(&self, session: &mut Session, user_key: &str) {
691 let user_hash = trustee_core::user_hash(user_key);
692
693 let user_home = if let Some(home) = dirs::home_dir() {
695 let user_home = home.join(".trustee").join("users").join(&user_hash);
696 session.home_dir = Some(user_home.clone());
697 Some(user_home)
698 } else {
699 None
700 };
701
702 session.project_id = Some(format!("web{}", &user_hash[..16]));
703
704 let shared_secrets = session.secrets.clone().unwrap_or_default();
711 let mut merged_secrets = shared_secrets.clone();
712
713 if let Some(ref user_home) = user_home {
714 if let Ok(merged) = self.load_user_secrets(user_home, &merged_secrets) {
715 merged_secrets = merged;
716 }
717 }
718
719 if let Some(ref user_home) = user_home {
721 if let Some(merged) = self.merge_user_config(user_home, session.config_toml.as_deref()) {
722 session.config_toml = Some(merged);
723 tracing::debug!("Merged per-user config into session");
724 }
725 }
726
727 if let Some(ref mut config_toml) = session.config_toml {
729 substitute_env_vars(config_toml, &merged_secrets);
730 }
731
732 session.secrets = Some(shared_secrets);
734 }
735
736 fn load_user_secrets(
739 &self,
740 user_home: &std::path::Path,
741 base: &std::collections::HashMap<String, String>,
742 ) -> std::io::Result<std::collections::HashMap<String, String>> {
743 let user_env_path = user_home.join(".env");
744 if !user_env_path.exists() {
745 return Ok(base.clone());
746 }
747 let content = std::fs::read_to_string(&user_env_path)?;
748 let mut merged = base.clone();
749 for line in content.lines() {
750 let line = line.trim();
751 if line.is_empty() || line.starts_with('#') {
752 continue;
753 }
754 if let Some((key, value)) = line.split_once('=') {
755 let key = key.trim().to_string();
756 let value = value
757 .trim()
758 .trim_matches('"')
759 .trim_matches('\'')
760 .to_string();
761 merged.insert(key, value);
762 }
763 }
764 tracing::debug!("Loaded per-user secrets from {}", user_env_path.display());
765 Ok(merged)
766 }
767
768 fn merge_user_config(
783 &self,
784 user_home: &std::path::Path,
785 shared_config: Option<&str>,
786 ) -> Option<String> {
787 let user_config_path = user_home.join("config").join("trustee.toml");
788 if !user_config_path.exists() {
789 return None;
790 }
791 let user_config_toml = std::fs::read_to_string(&user_config_path).ok()?;
792 let shared = shared_config
793 .unwrap_or("")
794 .parse::<toml::Value>()
795 .ok()?;
796 let overlay = user_config_toml.parse::<toml::Value>().ok()?;
797
798 let allowed = |section: &str| {
799 section == "mcp"
800 || section == "thq"
804 || (self.allow_llm_overlay && section == "llm")
805 };
806 let dir_name = user_home
809 .file_name()
810 .and_then(|n| n.to_str())
811 .unwrap_or("<unknown>");
812 let masked_user = dir_name.get(..8).unwrap_or(dir_name);
813
814 let mut filtered_overlay = toml::map::Map::new();
815 if let Some(table) = overlay.as_table() {
816 for (section, value) in table {
817 if allowed(section) {
818 filtered_overlay.insert(section.clone(), value.clone());
819 } else {
820 tracing::warn!(
821 "user config overlay: dropping non-allowlisted section [{}] for user {}",
822 section,
823 masked_user
824 );
825 }
826 }
827 }
828
829 if filtered_overlay.is_empty() {
830 return None;
832 }
833 let overlay = toml::Value::Table(filtered_overlay);
834
835 let mut shared = shared;
836 deep_merge_toml(&mut shared, &overlay);
837 let merged = toml::to_string(&shared).ok()?;
838 tracing::debug!("Merged per-user config from {}", user_config_path.display());
839 Some(merged)
840 }
841
842 pub fn resolve_user_config(&self, user_key: &str) -> Option<String> {
851 let config_toml = self.config_toml.clone()?;
852
853 let user_home = self.get_user_home_dir(user_key)?;
855
856 let mut merged_secrets = self.secrets.clone().unwrap_or_default();
858 if let Ok(merged) = self.load_user_secrets(&user_home, &merged_secrets) {
859 merged_secrets = merged;
860 }
861
862 let mut resolved = config_toml;
864 if let Some(merged) = self.merge_user_config(&user_home, Some(&resolved)) {
865 resolved = merged;
866 }
867
868 substitute_env_vars(&mut resolved, &merged_secrets);
870
871 Some(resolved)
872 }
873
874 pub async fn get_or_build_mcp_loader(
890 &self,
891 user_key: &str,
892 token_store: &Arc<pep::MemoryTokenStore>,
893 ) -> Result<Option<std::sync::Arc<abk::agent::McpToolLoader>>, String> {
894 let user_hash = trustee_core::user_hash(user_key);
895
896 let resolved = self.resolve_user_config(user_key);
898 let fingerprint = fingerprint_mcp_section(resolved.as_deref());
899
900 if let Some(entry) = self.mcp_loaders.get(&user_hash) {
902 if entry.degraded.is_none() {
903 if entry.fingerprint == fingerprint {
904 return Ok(entry.loader.clone());
905 }
906 } else if let (Some(err), Some(failed_at)) = (&entry.degraded, entry.failed_at) {
907 let backoff = chrono::Duration::from_std(MCP_BUILD_RETRY_BACKOFF)
908 .unwrap_or_else(|_| chrono::Duration::seconds(30));
909 if chrono::Utc::now() < failed_at + backoff {
910 return Err(err.clone());
911 }
912 }
913 }
914
915 let lock = self
917 .mcp_build_locks
918 .entry(user_hash.clone())
919 .or_insert_with(|| Arc::new(tokio::sync::Mutex::new(())))
920 .clone();
921 let _guard = lock.lock().await;
922
923 if let Some(entry) = self.mcp_loaders.get(&user_hash) {
925 if entry.degraded.is_none() && entry.fingerprint == fingerprint {
926 return Ok(entry.loader.clone());
927 }
928 }
929
930 match self
931 .build_mcp_loader(&user_hash, resolved.as_deref(), fingerprint, token_store)
932 .await
933 {
934 Ok(entry) => {
935 self.mcp_loaders.insert(user_hash.clone(), entry);
936 Ok(self.mcp_loaders.get(&user_hash).unwrap().loader.clone())
937 }
938 Err(err) => {
939 tracing::warn!(
940 "MCP loader build FAILED for user {}; dispatch fails loud, retry after {:?}",
941 &user_hash[..8.min(user_hash.len())],
942 MCP_BUILD_RETRY_BACKOFF
943 );
944 self.mcp_loaders.insert(
945 user_hash,
946 McpLoaderEntry {
947 loader: None,
948 fingerprint,
949 built_at: chrono::Utc::now(),
950 degraded: Some(err.clone()),
951 failed_at: Some(chrono::Utc::now()),
952 },
953 );
954 Err(err)
955 }
956 }
957 }
958
959 async fn build_mcp_loader(
962 &self,
963 user_hash: &str,
964 resolved: Option<&str>,
965 fingerprint: u64,
966 token_store: &Arc<pep::MemoryTokenStore>,
967 ) -> Result<McpLoaderEntry, String> {
968 let mcp_config: Option<abk::config::McpConfig> = match resolved {
969 Some(toml_str) => {
970 let value = toml_str
971 .parse::<toml::Value>()
972 .map_err(|e| format!("config parse failed: {}", e))?;
973 match value.get("mcp") {
974 Some(section) => {
975 use serde::Deserialize as _;
976 Some(
977 abk::config::McpConfig::deserialize(section.clone())
978 .map_err(|e| format!("invalid [mcp] config: {}", e))?,
979 )
980 }
981 None => None,
982 }
983 }
984 None => None,
985 };
986
987 let loader = match mcp_config {
988 Some(cfg) if cfg.enabled => {
989 let built = abk::agent::McpToolLoader::with_token_store(
990 &cfg,
991 Some(token_store.clone() as std::sync::Arc<dyn pep::token_store::TokenStore>),
992 )
993 .await
994 .map_err(|e| format!("MCP loader build failed: {}", e))?;
995
996 let servers: Vec<String> = built
1000 .server_statuses
1001 .iter()
1002 .map(|s| {
1003 if s.connected {
1004 format!("{}(up,{}tools)", s.name, s.tool_count)
1005 } else {
1006 format!("{}(DOWN)", s.name)
1007 }
1008 })
1009 .collect();
1010 tracing::info!(
1011 "MCP loader built for user {}: servers=[{}] total_tools={}",
1012 &user_hash[..8.min(user_hash.len())],
1013 servers.join(", "),
1014 built.tool_count
1015 );
1016 Some(std::sync::Arc::new(built))
1017 }
1018 _ => None,
1019 };
1020
1021 Ok(McpLoaderEntry {
1022 loader,
1023 fingerprint,
1024 built_at: chrono::Utc::now(),
1025 degraded: None,
1026 failed_at: None,
1027 })
1028 }
1029
1030 fn spawn_user_drain_task(
1032 &self,
1033 session_id: String,
1034 session: Arc<Mutex<Session>>,
1035 ws_tx: broadcast::Sender<String>,
1036 mut workflow_rx: mpsc::UnboundedReceiver<TuiMessage>,
1037 ) {
1038 let mut last_broadcast_state: Option<String> = None;
1044 tokio::spawn(async move {
1045 while let Some(msg) = workflow_rx.recv().await {
1046 {
1047 let mut session = session.lock().await;
1048 session.handle_workflow_message(msg.clone());
1049
1050 let state_str = match session.workflow_state {
1051 trustee_core::types::WorkflowState::Idle => "Idle",
1052 trustee_core::types::WorkflowState::Running => "Running",
1053 trustee_core::types::WorkflowState::Cancelling => "Cancelling",
1054 };
1055 if last_broadcast_state.as_deref() != Some(state_str) {
1056 last_broadcast_state = Some(state_str.to_string());
1057 let state_msg = serde_json::json!({
1058 "type": "StateChanged",
1059 "state": state_str
1060 });
1061 let _ = ws_tx.send(state_msg.to_string());
1062 }
1063 }
1064
1065 let json =
1066 serde_json::to_string(&SerializableMessage(&msg)).unwrap_or_default();
1067 let _ = ws_tx.send(json);
1068 }
1069 tracing::debug!("Drain task ended for session: {}", session_id);
1070 });
1071 }
1072
1073 pub fn spawn_drain_task(self, mut workflow_rx: mpsc::UnboundedReceiver<TuiMessage>) {
1076 let default_user = self
1078 .sessions
1079 .get("default")
1080 .expect("default user must exist");
1081 let first_entry = default_user
1082 .sessions
1083 .iter()
1084 .next()
1085 .expect("default user must have at least one session");
1086 let session = first_entry.session.clone();
1087 let ws_tx = first_entry.ws_tx.clone();
1088 let session_id = first_entry.key().clone();
1089 drop(first_entry);
1090 drop(default_user);
1091
1092 let mut last_broadcast_state = Some("Running".to_string());
1096 tokio::spawn(async move {
1097 while let Some(msg) = workflow_rx.recv().await {
1098 {
1099 let mut session = session.lock().await;
1100 session.handle_workflow_message(msg.clone());
1101
1102 let state_str = match session.workflow_state {
1103 trustee_core::types::WorkflowState::Idle => "Idle",
1104 trustee_core::types::WorkflowState::Running => "Running",
1105 trustee_core::types::WorkflowState::Cancelling => "Cancelling",
1106 };
1107 if last_broadcast_state.as_deref() != Some(state_str) {
1108 last_broadcast_state = Some(state_str.to_string());
1109 let state_msg = serde_json::json!({
1110 "type": "StateChanged",
1111 "state": state_str
1112 });
1113 let _ = ws_tx.send(state_msg.to_string());
1114 }
1115 }
1116
1117 let json =
1118 serde_json::to_string(&SerializableMessage(&msg)).unwrap_or_default();
1119 let _ = ws_tx.send(json);
1120 }
1121 tracing::debug!("Drain task ended for session: {}", session_id);
1122 });
1123 }
1124
1125 pub async fn resolve_user_key(&self, headers: &axum::http::HeaderMap) -> String {
1127 let Some(ref auth) = self.auth else {
1128 return "default".to_string();
1129 };
1130
1131 if let Some(token) = headers
1133 .get(axum::http::header::AUTHORIZATION)
1134 .and_then(|v| v.to_str().ok())
1135 .and_then(|v| v.strip_prefix("Bearer "))
1136 .map(|s| s.to_string())
1137 {
1138 if token.starts_with("dev:") {
1139 let parts: Vec<&str> = token.splitn(4, ':').collect();
1140 if parts.len() >= 4 {
1141 return format!("dev:{}", parts[1]);
1142 }
1143 }
1144 if let Ok(claims) = auth.validate_token(&token).await {
1145 return claims.sub;
1146 }
1147 }
1148
1149 let cookie_session_id = headers
1151 .get(axum::http::header::COOKIE)
1152 .and_then(|v| v.to_str().ok())
1153 .and_then(|cookies| {
1154 cookies
1155 .split(';')
1156 .map(|c| c.trim())
1157 .find_map(|c| {
1158 c.strip_prefix(&format!("{}=", auth.config.cookie_name))
1159 .map(|s| s.to_string())
1160 })
1161 });
1162
1163 if let Some(session_id) = cookie_session_id {
1164 if session_id.starts_with("dev:") {
1165 let parts: Vec<&str> = session_id.splitn(4, ':').collect();
1166 if parts.len() >= 4 {
1167 return format!("dev:{}", parts[1]);
1168 }
1169 }
1170
1171 if let Ok(access_token) = auth.session_manager.get_token(&session_id).await {
1172 if let Ok(claims) = auth.validate_token(&access_token).await {
1173 return claims.sub;
1174 }
1175 }
1176 }
1177
1178 "default".to_string()
1179 }
1180}
1181
1182struct SerializableMessage<'a>(&'a TuiMessage);
1188
1189impl<'a> serde::Serialize for SerializableMessage<'a> {
1190 fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
1191 where
1192 S: serde::Serializer,
1193 {
1194 use serde::ser::SerializeStruct;
1195
1196 match self.0 {
1197 TuiMessage::OutputLine(line) => {
1198 let mut s = serializer.serialize_struct("msg", 2)?;
1199 s.serialize_field("type", "OutputLine")?;
1200 s.serialize_field("line", line)?;
1201 s.end()
1202 }
1203 TuiMessage::StreamDelta(delta) => {
1204 let mut s = serializer.serialize_struct("msg", 2)?;
1205 s.serialize_field("type", "StreamDelta")?;
1206 s.serialize_field("delta", delta)?;
1207 s.end()
1208 }
1209 TuiMessage::ReasoningDelta(delta) => {
1210 let mut s = serializer.serialize_struct("msg", 2)?;
1211 s.serialize_field("type", "ReasoningDelta")?;
1212 s.serialize_field("delta", delta)?;
1213 s.end()
1214 }
1215 TuiMessage::WorkflowCompleted => {
1216 let mut s = serializer.serialize_struct("msg", 2)?;
1217 s.serialize_field("type", "WorkflowCompleted")?;
1218 s.serialize_field("state", "Idle")?;
1219 s.end()
1220 }
1221 TuiMessage::WorkflowError(err) => {
1222 let mut s = serializer.serialize_struct("msg", 2)?;
1223 s.serialize_field("type", "WorkflowError")?;
1224 s.serialize_field("error", err)?;
1225 s.end()
1226 }
1227 TuiMessage::ResumeInfo(info) => match info {
1228 Some(ri) => {
1229 let mut s = serializer.serialize_struct("msg", 5)?;
1230 s.serialize_field("type", "ResumeInfo")?;
1231 s.serialize_field("state", "Idle")?;
1232 s.serialize_field("session_id", &ri.session_id)?;
1233 s.serialize_field("checkpoint_id", &ri.checkpoint_id)?;
1234 s.serialize_field("iteration", &ri.iteration)?;
1235 s.end()
1236 }
1237 None => {
1238 let mut s = serializer.serialize_struct("msg", 2)?;
1239 s.serialize_field("type", "ResumeInfo")?;
1240 s.serialize_field("state", "Idle")?;
1241 s.end()
1242 }
1243 },
1244 TuiMessage::TodoUpdate(content) => {
1245 let mut s = serializer.serialize_struct("msg", 2)?;
1246 s.serialize_field("type", "TodoUpdate")?;
1247 s.serialize_field("content", content)?;
1248 s.end()
1249 }
1250 TuiMessage::WorkflowCancelled => {
1251 let mut s = serializer.serialize_struct("msg", 2)?;
1252 s.serialize_field("type", "WorkflowCancelled")?;
1253 s.serialize_field("state", "Idle")?;
1254 s.end()
1255 }
1256 TuiMessage::HandoffReady(briefing) => {
1257 let mut s = serializer.serialize_struct("msg", 3)?;
1258 s.serialize_field("type", "HandoffReady")?;
1259 s.serialize_field("state", "Idle")?;
1260 s.serialize_field("briefing", briefing)?;
1261 s.end()
1262 }
1263 TuiMessage::SessionRotated { old, new } => {
1264 let mut s = serializer.serialize_struct("msg", 3)?;
1265 s.serialize_field("type", "SessionRotated")?;
1266 s.serialize_field("old", old)?;
1267 s.serialize_field("new", new)?;
1268 s.end()
1269 }
1270 TuiMessage::HandoffFailed => {
1271 let mut s = serializer.serialize_struct("msg", 2)?;
1272 s.serialize_field("type", "HandoffFailed")?;
1273 s.serialize_field("state", "Idle")?;
1274 s.end()
1275 }
1276 TuiMessage::ToolPending {
1277 tool_name,
1278 hint,
1279 } => {
1280 let mut s = serializer.serialize_struct("msg", 3)?;
1281 s.serialize_field("type", "ToolPending")?;
1282 s.serialize_field("tool_name", tool_name)?;
1283 s.serialize_field("hint", hint)?;
1284 s.end()
1285 }
1286 TuiMessage::ToolDone {
1287 tool_name,
1288 success,
1289 hint,
1290 } => {
1291 let mut s = serializer.serialize_struct("msg", 4)?;
1292 s.serialize_field("type", "ToolDone")?;
1293 s.serialize_field("tool_name", tool_name)?;
1294 s.serialize_field("success", success)?;
1295 s.serialize_field("hint", hint)?;
1296 s.end()
1297 }
1298 TuiMessage::ContextTokensUpdated(count) => {
1299 let mut s = serializer.serialize_struct("msg", 2)?;
1300 s.serialize_field("type", "ContextTokensUpdated")?;
1301 s.serialize_field("count", count)?;
1302 s.end()
1303 }
1304 TuiMessage::McpServerStatus {
1305 name,
1306 connected,
1307 tool_count,
1308 error,
1309 } => {
1310 let mut s = serializer.serialize_struct("msg", 5)?;
1311 s.serialize_field("type", "McpServerStatus")?;
1312 s.serialize_field("name", name)?;
1313 s.serialize_field("connected", connected)?;
1314 s.serialize_field("tool_count", tool_count)?;
1315 s.serialize_field("error", error)?;
1316 s.end()
1317 }
1318 TuiMessage::SessionTitleUpdated(title) => {
1319 let mut s = serializer.serialize_struct("msg", 2)?;
1320 s.serialize_field("type", "SessionTitleUpdated")?;
1321 s.serialize_field("title", title)?;
1322 s.end()
1323 }
1324 }
1325 }
1326}
1327
1328fn deep_merge_toml(base: &mut toml::Value, overlay: &toml::Value) {
1339 match (base, overlay) {
1340 (toml::Value::Table(base_table), toml::Value::Table(overlay_table)) => {
1341 for (key, overlay_val) in overlay_table {
1342 match base_table.get_mut(key) {
1343 Some(base_val) => {
1344 deep_merge_toml(base_val, overlay_val);
1346 }
1347 None => {
1348 base_table.insert(key.clone(), overlay_val.clone());
1350 }
1351 }
1352 }
1353 }
1354 (base, overlay) => {
1356 *base = overlay.clone();
1357 }
1358 }
1359}
1360
1361fn substitute_env_vars(s: &mut String, secrets: &std::collections::HashMap<String, String>) {
1366 let mut result = String::with_capacity(s.len());
1368 let bytes = s.as_bytes();
1369 let mut i = 0;
1370
1371 while i < bytes.len() {
1372 if i + 1 < bytes.len() && bytes[i] == b'$' && bytes[i + 1] == b'{' {
1373 if let Some(end) = s[i + 2..].find('}') {
1375 let var_name = &s[i + 2..i + 2 + end];
1376 if let Some(value) = secrets.get(var_name) {
1378 result.push_str(value);
1379 } else if let Ok(value) = std::env::var(var_name) {
1380 result.push_str(&value);
1381 } else {
1382 result.push_str(&s[i..i + 2 + end + 1]);
1384 }
1385 i = i + 2 + end + 1;
1386 } else {
1387 result.push('$');
1389 i += 1;
1390 }
1391 } else {
1392 result.push(bytes[i] as char);
1393 i += 1;
1394 }
1395 }
1396
1397 *s = result;
1398}
1399
1400fn fingerprint_mcp_section(resolved: Option<&str>) -> u64 {
1410 use sha2::{Digest, Sha256};
1411 let section = resolved
1412 .and_then(|s| s.parse::<toml::Value>().ok())
1413 .and_then(|v| v.get("mcp").cloned());
1414 let bytes = match section {
1415 Some(v) => v.to_string().into_bytes(),
1416 None => b"<no-mcp>".to_vec(),
1417 };
1418 let digest = Sha256::digest(&bytes);
1419 u64::from_be_bytes(digest[..8].try_into().expect("sha256 digest >= 8 bytes"))
1420}
1421
1422#[cfg(test)]
1423mod tests {
1424 use super::*;
1425
1426 fn test_state() -> ServerState {
1428 let (session, _rx) = Session::new();
1429 let (ws_tx, _ws_rx) = tokio::sync::broadcast::channel::<String>(16);
1430 ServerState::new(session, ws_tx, None)
1431 }
1432
1433 fn temp_user_home(tag: &str) -> std::path::PathBuf {
1435 static COUNTER: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0);
1436 let n = COUNTER.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
1437 let dir = std::env::temp_dir().join(format!(
1438 "trustee-state-test-{}-{}-{}",
1439 tag,
1440 std::process::id(),
1441 n
1442 ));
1443 std::fs::create_dir_all(dir.join("config")).expect("create temp user home");
1444 dir
1445 }
1446
1447 fn parse(toml_str: &str) -> toml::Value {
1448 toml_str.parse::<toml::Value>().expect("valid test TOML")
1449 }
1450
1451 #[test]
1454 fn overlay_allowlist_drops_non_allowlisted_sections() {
1455 let state = test_state();
1456 let home = temp_user_home("allowlist");
1457 std::fs::write(
1458 home.join("config").join("trustee.toml"),
1459 "[server]\nport = 1\n\n[auth]\nmode = \"kanidm\"\n\n[mcp]\nmode = \"user\"\n",
1460 )
1461 .expect("write overlay");
1462
1463 let shared = "[server]\nport = 8080\n\n[mcp]\nmode = \"shared\"\n";
1464 let merged = state
1465 .merge_user_config(&home, Some(shared))
1466 .expect("overlay has allowlisted content");
1467
1468 let merged_val = parse(&merged);
1469 let expected = parse("[server]\nport = 8080\n\n[mcp]\nmode = \"user\"\n");
1470 assert_eq!(merged_val, expected, "merged config must be shared + [mcp] overlay only");
1471 assert!(merged_val.get("auth").is_none(), "overlay [auth] must be dropped");
1472 }
1473
1474 #[test]
1476 fn no_overlay_file_returns_none() {
1477 let state = test_state();
1478 let home = temp_user_home("empty");
1479 assert!(state
1480 .merge_user_config(&home, Some("[server]\nport = 8080\n"))
1481 .is_none());
1482 }
1483
1484 #[test]
1486 fn overlay_with_no_allowlisted_sections_is_noop() {
1487 let state = test_state();
1488 let home = temp_user_home("all-dropped");
1489 std::fs::write(
1490 home.join("config").join("trustee.toml"),
1491 "[server]\nport = 1\n\n[storage]\npath = \"/tmp/x\"\n",
1492 )
1493 .expect("write overlay");
1494 assert!(state
1495 .merge_user_config(&home, Some("[server]\nport = 8080\n"))
1496 .is_none());
1497 }
1498
1499 #[test]
1502 fn llm_overlay_dropped_by_default_and_kept_when_enabled() {
1503 let shared = "[llm]\nprovider = \"openai\"\n\n[mcp]\nmode = \"shared\"\n";
1504 let overlay = "[llm]\nprovider = \"anthropic\"\n";
1505
1506 let state = test_state();
1508 assert!(!state.allow_llm_overlay);
1509 let home = temp_user_home("llm-off");
1510 std::fs::write(home.join("config").join("trustee.toml"), overlay).expect("write overlay");
1511 assert!(state.merge_user_config(&home, Some(shared)).is_none());
1512
1513 let state = test_state().with_allow_llm_overlay(true);
1515 assert!(state.allow_llm_overlay);
1516 let home = temp_user_home("llm-on");
1517 std::fs::write(home.join("config").join("trustee.toml"), overlay).expect("write overlay");
1518 let merged = state
1519 .merge_user_config(&home, Some(shared))
1520 .expect("[llm] overlay applies when opted in");
1521 let merged_val = parse(&merged);
1522 assert_eq!(merged_val["llm"]["provider"].as_str(), Some("anthropic"));
1523 assert_eq!(merged_val["mcp"]["mode"].as_str(), Some("shared"));
1524 }
1525
1526 #[test]
1530 fn user_home_dir_uses_consolidated_hash() {
1531 let state = test_state();
1532 if let Some(home) = state.get_user_home_dir("farzan@example.com") {
1533 assert_eq!(
1534 home.file_name().and_then(|n| n.to_str()),
1535 Some(trustee_core::user_hash("farzan@example.com")).as_deref()
1536 );
1537 let users_root = dirs::home_dir().unwrap().join(".trustee").join("users");
1538 assert_eq!(home.parent(), Some(&users_root).map(|p| p.as_path()));
1539 }
1540 }
1542
1543 #[tokio::test]
1548 async fn agent_principals_get_isolated_session_buckets() {
1549 let state = test_state();
1550 let key_a = "agent-farzan";
1551 let key_b = "agent-paydar";
1552
1553 let (sid_a, session_a, _tx_a, _ts_a) = state.ensure_active_session(key_a).await;
1554 let (_sid_b, _session_b, _tx_b, _ts_b) = state.ensure_active_session(key_b).await;
1555
1556 assert!(
1558 state.get_session_by_any_id(key_a, &sid_a).await.is_some(),
1559 "owner bucket resolves its own session"
1560 );
1561 assert!(
1562 state.get_session_by_any_id(key_b, &sid_a).await.is_none(),
1563 "cross-agent session access must be 404/None"
1564 );
1565 assert!(Arc::strong_count(&session_a) >= 1);
1566
1567 if let Some(home) = state.get_user_home_dir(key_a) {
1569 assert_eq!(
1570 home.file_name().and_then(|n| n.to_str()),
1571 Some(trustee_core::user_hash(key_a)).as_deref()
1572 );
1573 }
1574 }
1575}
1576
1577#[cfg(test)]
1582mod mcp_loader_cache_tests {
1583 use super::*;
1584
1585 fn state_with_shared(shared: &str) -> ServerState {
1586 let (session, _rx) = Session::new();
1587 let (ws_tx, _ws_rx) = tokio::sync::broadcast::channel::<String>(16);
1588 let mut state = ServerState::new(session, ws_tx, None);
1589 state.config_toml = Some(shared.to_string());
1590 state
1591 }
1592
1593 struct TempUser {
1597 key: String,
1598 home: std::path::PathBuf,
1599 }
1600
1601 impl TempUser {
1602 fn new(tag: &str) -> Self {
1603 let key = format!("16c-{tag}-{}@test.invalid", std::process::id());
1604 let home = dirs::home_dir()
1605 .expect("HOME available in test env")
1606 .join(".trustee")
1607 .join("users")
1608 .join(trustee_core::user_hash(&key));
1609 std::fs::create_dir_all(home.join("config")).expect("create user home");
1610 Self { key, home }
1611 }
1612
1613 fn write_overlay(&self, toml_str: &str) {
1614 std::fs::write(self.home.join("config").join("trustee.toml"), toml_str)
1615 .expect("write overlay");
1616 }
1617 }
1618
1619 impl Drop for TempUser {
1620 fn drop(&mut self) {
1621 let _ = std::fs::remove_dir_all(&self.home);
1622 }
1623 }
1624
1625 #[tokio::test]
1626 async fn no_mcp_config_caches_disabled_marker() {
1627 let state = state_with_shared("[server]\nport = 8080\n");
1628 let user = TempUser::new("nomcp");
1629 let ts = Arc::new(pep::MemoryTokenStore::new());
1630
1631 let first = state.get_or_build_mcp_loader(&user.key, &ts).await.unwrap();
1632 assert!(first.is_none(), "no [mcp] anywhere → disabled marker");
1633 assert_eq!(state.mcp_loaders.len(), 1, "exactly one cache entry");
1634
1635 let second = state.get_or_build_mcp_loader(&user.key, &ts).await.unwrap();
1637 assert!(second.is_none());
1638 assert_eq!(state.mcp_loaders.len(), 1);
1639 }
1640
1641 #[tokio::test]
1642 async fn concurrent_cold_builds_single_flight() {
1643 let state = Arc::new(state_with_shared("[server]\nport = 8080\n"));
1644 let user = TempUser::new("singleflight");
1645 let ts = Arc::new(pep::MemoryTokenStore::new());
1646
1647 let mut handles = Vec::new();
1648 for _ in 0..5 {
1649 let state = state.clone();
1650 let key = user.key.clone();
1651 let ts = ts.clone();
1652 handles.push(tokio::spawn(async move {
1653 state.get_or_build_mcp_loader(&key, &ts).await
1654 }));
1655 }
1656 for h in handles {
1657 h.await.unwrap().expect("all five succeed");
1658 }
1659 assert_eq!(state.mcp_loaders.len(), 1, "single-flight → one entry");
1660 }
1661
1662 #[tokio::test]
1663 async fn fingerprint_change_triggers_rebuild() {
1664 let state = state_with_shared("[server]\nport = 8080\n");
1665 let user = TempUser::new("fpchange");
1666 let ts = Arc::new(pep::MemoryTokenStore::new());
1667
1668 user.write_overlay("[mcp]\nenabled = true\n\n[[mcp.servers]]\nname = \"v1\"\nurl = \"http://127.0.0.1:9/sse\"\n");
1671 let v1 = state.get_or_build_mcp_loader(&user.key, &ts).await.unwrap();
1672 assert!(v1.is_some(), "enabled [mcp] → real loader");
1673 let fp1 = state
1674 .mcp_loaders
1675 .get(&trustee_core::user_hash(&user.key))
1676 .unwrap()
1677 .fingerprint;
1678
1679 std::thread::sleep(std::time::Duration::from_millis(5));
1681 user.write_overlay("[mcp]\nenabled = true\n\n[[mcp.servers]]\nname = \"v2\"\nurl = \"http://127.0.0.1:9/other\"\n");
1682 let v2 = state.get_or_build_mcp_loader(&user.key, &ts).await.unwrap();
1683 assert!(v2.is_some());
1684 let entry = state
1685 .mcp_loaders
1686 .get(&trustee_core::user_hash(&user.key))
1687 .unwrap();
1688 assert_ne!(
1689 entry.fingerprint, fp1,
1690 "fingerprint must change with content"
1691 );
1692 assert!(entry.degraded.is_none());
1693
1694 let _still_usable = v1.as_ref().unwrap().tool_count;
1696 }
1697
1698 #[tokio::test]
1699 async fn degraded_entry_fails_loud_within_backoff_and_isolates_users() {
1700 let state = state_with_shared("[server]\nport = 8080\n");
1701 let bad = TempUser::new("degraded-bad");
1702 let good = TempUser::new("degraded-good");
1703 let ts = Arc::new(pep::MemoryTokenStore::new());
1704
1705 bad.write_overlay("[mcp]\nenabled = \"not-a-bool\"\n");
1707 let err = match state.get_or_build_mcp_loader(&bad.key, &ts).await {
1708 Ok(_) => panic!("invalid [mcp] must fail loud"),
1709 Err(e) => e,
1710 };
1711 assert!(
1712 err.contains("invalid [mcp]"),
1713 "surfaces the parse error: {err}"
1714 );
1715
1716 let entry = state
1717 .mcp_loaders
1718 .get(&trustee_core::user_hash(&bad.key))
1719 .unwrap();
1720 assert!(entry.degraded.is_some(), "poison entry recorded");
1721
1722 let err2 = match state.get_or_build_mcp_loader(&bad.key, &ts).await {
1724 Ok(_) => panic!("still within backoff"),
1725 Err(e) => e,
1726 };
1727 assert_eq!(err, err2, "same cached error");
1728
1729 let good_loader = state
1731 .get_or_build_mcp_loader(&good.key, &ts)
1732 .await
1733 .expect("other user unaffected");
1734 assert!(
1735 good_loader.is_none(),
1736 "good user has no [mcp] → disabled marker"
1737 );
1738 }
1739}