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 || section == "lifecycle"
809 || (self.allow_llm_overlay && section == "llm")
810 };
811 let dir_name = user_home
814 .file_name()
815 .and_then(|n| n.to_str())
816 .unwrap_or("<unknown>");
817 let masked_user = dir_name.get(..8).unwrap_or(dir_name);
818
819 let mut filtered_overlay = toml::map::Map::new();
820 if let Some(table) = overlay.as_table() {
821 for (section, value) in table {
822 if section == "lifecycle" {
823 let mut v = value.clone();
824 if let Some(t) = v.as_table_mut() {
825 if t.remove("enabled").is_some() {
826 tracing::info!(
827 "user config overlay: [lifecycle].enabled is not user-settable — stripped for user {}",
828 masked_user
829 );
830 }
831 }
832 filtered_overlay.insert(section.clone(), v);
833 } else if allowed(section) {
834 filtered_overlay.insert(section.clone(), value.clone());
835 } else {
836 tracing::warn!(
837 "user config overlay: dropping non-allowlisted section [{}] for user {}",
838 section,
839 masked_user
840 );
841 }
842 }
843 }
844
845 if filtered_overlay.is_empty() {
846 return None;
848 }
849 let overlay = toml::Value::Table(filtered_overlay);
850
851 let mut shared = shared;
852 deep_merge_toml(&mut shared, &overlay);
853 let merged = toml::to_string(&shared).ok()?;
854 tracing::debug!("Merged per-user config from {}", user_config_path.display());
855 Some(merged)
856 }
857
858 pub fn resolve_user_config(&self, user_key: &str) -> Option<String> {
867 let config_toml = self.config_toml.clone()?;
868
869 let user_home = self.get_user_home_dir(user_key)?;
871
872 let mut merged_secrets = self.secrets.clone().unwrap_or_default();
874 if let Ok(merged) = self.load_user_secrets(&user_home, &merged_secrets) {
875 merged_secrets = merged;
876 }
877
878 let mut resolved = config_toml;
880 if let Some(merged) = self.merge_user_config(&user_home, Some(&resolved)) {
881 resolved = merged;
882 }
883
884 substitute_env_vars(&mut resolved, &merged_secrets);
886
887 Some(resolved)
888 }
889
890 pub async fn get_or_build_mcp_loader(
906 &self,
907 user_key: &str,
908 token_store: &Arc<pep::MemoryTokenStore>,
909 ) -> Result<Option<std::sync::Arc<abk::agent::McpToolLoader>>, String> {
910 let user_hash = trustee_core::user_hash(user_key);
911
912 let resolved = self.resolve_user_config(user_key);
914 let fingerprint = fingerprint_mcp_section(resolved.as_deref());
915
916 if let Some(entry) = self.mcp_loaders.get(&user_hash) {
918 if entry.degraded.is_none() {
919 if entry.fingerprint == fingerprint {
920 return Ok(entry.loader.clone());
921 }
922 } else if let (Some(err), Some(failed_at)) = (&entry.degraded, entry.failed_at) {
923 let backoff = chrono::Duration::from_std(MCP_BUILD_RETRY_BACKOFF)
924 .unwrap_or_else(|_| chrono::Duration::seconds(30));
925 if chrono::Utc::now() < failed_at + backoff {
926 return Err(err.clone());
927 }
928 }
929 }
930
931 let lock = self
933 .mcp_build_locks
934 .entry(user_hash.clone())
935 .or_insert_with(|| Arc::new(tokio::sync::Mutex::new(())))
936 .clone();
937 let _guard = lock.lock().await;
938
939 if let Some(entry) = self.mcp_loaders.get(&user_hash) {
941 if entry.degraded.is_none() && entry.fingerprint == fingerprint {
942 return Ok(entry.loader.clone());
943 }
944 }
945
946 match self
947 .build_mcp_loader(&user_hash, resolved.as_deref(), fingerprint, token_store)
948 .await
949 {
950 Ok(entry) => {
951 self.mcp_loaders.insert(user_hash.clone(), entry);
952 Ok(self.mcp_loaders.get(&user_hash).unwrap().loader.clone())
953 }
954 Err(err) => {
955 tracing::warn!(
956 "MCP loader build FAILED for user {}; dispatch fails loud, retry after {:?}",
957 &user_hash[..8.min(user_hash.len())],
958 MCP_BUILD_RETRY_BACKOFF
959 );
960 self.mcp_loaders.insert(
961 user_hash,
962 McpLoaderEntry {
963 loader: None,
964 fingerprint,
965 built_at: chrono::Utc::now(),
966 degraded: Some(err.clone()),
967 failed_at: Some(chrono::Utc::now()),
968 },
969 );
970 Err(err)
971 }
972 }
973 }
974
975 async fn build_mcp_loader(
978 &self,
979 user_hash: &str,
980 resolved: Option<&str>,
981 fingerprint: u64,
982 token_store: &Arc<pep::MemoryTokenStore>,
983 ) -> Result<McpLoaderEntry, String> {
984 let mcp_config: Option<abk::config::McpConfig> = match resolved {
985 Some(toml_str) => {
986 let value = toml_str
987 .parse::<toml::Value>()
988 .map_err(|e| format!("config parse failed: {}", e))?;
989 match value.get("mcp") {
990 Some(section) => {
991 use serde::Deserialize as _;
992 Some(
993 abk::config::McpConfig::deserialize(section.clone())
994 .map_err(|e| format!("invalid [mcp] config: {}", e))?,
995 )
996 }
997 None => None,
998 }
999 }
1000 None => None,
1001 };
1002
1003 let loader = match mcp_config {
1004 Some(cfg) if cfg.enabled => {
1005 let built = abk::agent::McpToolLoader::with_token_store(
1006 &cfg,
1007 Some(token_store.clone() as std::sync::Arc<dyn pep::token_store::TokenStore>),
1008 )
1009 .await
1010 .map_err(|e| format!("MCP loader build failed: {}", e))?;
1011
1012 let servers: Vec<String> = built
1016 .server_statuses
1017 .iter()
1018 .map(|s| {
1019 if s.connected {
1020 format!("{}(up,{}tools)", s.name, s.tool_count)
1021 } else {
1022 format!("{}(DOWN)", s.name)
1023 }
1024 })
1025 .collect();
1026 tracing::info!(
1027 "MCP loader built for user {}: servers=[{}] total_tools={}",
1028 &user_hash[..8.min(user_hash.len())],
1029 servers.join(", "),
1030 built.tool_count
1031 );
1032 Some(std::sync::Arc::new(built))
1033 }
1034 _ => None,
1035 };
1036
1037 Ok(McpLoaderEntry {
1038 loader,
1039 fingerprint,
1040 built_at: chrono::Utc::now(),
1041 degraded: None,
1042 failed_at: None,
1043 })
1044 }
1045
1046 fn spawn_user_drain_task(
1048 &self,
1049 session_id: String,
1050 session: Arc<Mutex<Session>>,
1051 ws_tx: broadcast::Sender<String>,
1052 mut workflow_rx: mpsc::UnboundedReceiver<TuiMessage>,
1053 ) {
1054 let mut last_broadcast_state: Option<String> = None;
1060 tokio::spawn(async move {
1061 while let Some(msg) = workflow_rx.recv().await {
1062 {
1063 let mut session = session.lock().await;
1064 session.handle_workflow_message(msg.clone());
1065
1066 let state_str = match session.workflow_state {
1067 trustee_core::types::WorkflowState::Idle => "Idle",
1068 trustee_core::types::WorkflowState::Running => "Running",
1069 trustee_core::types::WorkflowState::Cancelling => "Cancelling",
1070 };
1071 if last_broadcast_state.as_deref() != Some(state_str) {
1072 last_broadcast_state = Some(state_str.to_string());
1073 let state_msg = serde_json::json!({
1074 "type": "StateChanged",
1075 "state": state_str
1076 });
1077 let _ = ws_tx.send(state_msg.to_string());
1078 }
1079 }
1080
1081 let json =
1082 serde_json::to_string(&SerializableMessage(&msg)).unwrap_or_default();
1083 let _ = ws_tx.send(json);
1084 }
1085 tracing::debug!("Drain task ended for session: {}", session_id);
1086 });
1087 }
1088
1089 pub fn spawn_drain_task(self, mut workflow_rx: mpsc::UnboundedReceiver<TuiMessage>) {
1092 let default_user = self
1094 .sessions
1095 .get("default")
1096 .expect("default user must exist");
1097 let first_entry = default_user
1098 .sessions
1099 .iter()
1100 .next()
1101 .expect("default user must have at least one session");
1102 let session = first_entry.session.clone();
1103 let ws_tx = first_entry.ws_tx.clone();
1104 let session_id = first_entry.key().clone();
1105 drop(first_entry);
1106 drop(default_user);
1107
1108 let mut last_broadcast_state = Some("Running".to_string());
1112 tokio::spawn(async move {
1113 while let Some(msg) = workflow_rx.recv().await {
1114 {
1115 let mut session = session.lock().await;
1116 session.handle_workflow_message(msg.clone());
1117
1118 let state_str = match session.workflow_state {
1119 trustee_core::types::WorkflowState::Idle => "Idle",
1120 trustee_core::types::WorkflowState::Running => "Running",
1121 trustee_core::types::WorkflowState::Cancelling => "Cancelling",
1122 };
1123 if last_broadcast_state.as_deref() != Some(state_str) {
1124 last_broadcast_state = Some(state_str.to_string());
1125 let state_msg = serde_json::json!({
1126 "type": "StateChanged",
1127 "state": state_str
1128 });
1129 let _ = ws_tx.send(state_msg.to_string());
1130 }
1131 }
1132
1133 let json =
1134 serde_json::to_string(&SerializableMessage(&msg)).unwrap_or_default();
1135 let _ = ws_tx.send(json);
1136 }
1137 tracing::debug!("Drain task ended for session: {}", session_id);
1138 });
1139 }
1140
1141 pub async fn resolve_user_key(&self, headers: &axum::http::HeaderMap) -> String {
1143 let Some(ref auth) = self.auth else {
1144 return "default".to_string();
1145 };
1146
1147 if let Some(token) = headers
1149 .get(axum::http::header::AUTHORIZATION)
1150 .and_then(|v| v.to_str().ok())
1151 .and_then(|v| v.strip_prefix("Bearer "))
1152 .map(|s| s.to_string())
1153 {
1154 if token.starts_with("dev:") {
1155 let parts: Vec<&str> = token.splitn(4, ':').collect();
1156 if parts.len() >= 4 {
1157 return format!("dev:{}", parts[1]);
1158 }
1159 }
1160 if let Ok(claims) = auth.validate_token(&token).await {
1161 return claims.sub;
1162 }
1163 }
1164
1165 let cookie_session_id = headers
1167 .get(axum::http::header::COOKIE)
1168 .and_then(|v| v.to_str().ok())
1169 .and_then(|cookies| {
1170 cookies
1171 .split(';')
1172 .map(|c| c.trim())
1173 .find_map(|c| {
1174 c.strip_prefix(&format!("{}=", auth.config.cookie_name))
1175 .map(|s| s.to_string())
1176 })
1177 });
1178
1179 if let Some(session_id) = cookie_session_id {
1180 if session_id.starts_with("dev:") {
1181 let parts: Vec<&str> = session_id.splitn(4, ':').collect();
1182 if parts.len() >= 4 {
1183 return format!("dev:{}", parts[1]);
1184 }
1185 }
1186
1187 if let Ok(access_token) = auth.session_manager.get_token(&session_id).await {
1188 if let Ok(claims) = auth.validate_token(&access_token).await {
1189 return claims.sub;
1190 }
1191 }
1192 }
1193
1194 "default".to_string()
1195 }
1196}
1197
1198struct SerializableMessage<'a>(&'a TuiMessage);
1204
1205impl<'a> serde::Serialize for SerializableMessage<'a> {
1206 fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
1207 where
1208 S: serde::Serializer,
1209 {
1210 use serde::ser::SerializeStruct;
1211
1212 match self.0 {
1213 TuiMessage::OutputLine(line) => {
1214 let mut s = serializer.serialize_struct("msg", 2)?;
1215 s.serialize_field("type", "OutputLine")?;
1216 s.serialize_field("line", line)?;
1217 s.end()
1218 }
1219 TuiMessage::StreamDelta(delta) => {
1220 let mut s = serializer.serialize_struct("msg", 2)?;
1221 s.serialize_field("type", "StreamDelta")?;
1222 s.serialize_field("delta", delta)?;
1223 s.end()
1224 }
1225 TuiMessage::ReasoningDelta(delta) => {
1226 let mut s = serializer.serialize_struct("msg", 2)?;
1227 s.serialize_field("type", "ReasoningDelta")?;
1228 s.serialize_field("delta", delta)?;
1229 s.end()
1230 }
1231 TuiMessage::WorkflowCompleted => {
1232 let mut s = serializer.serialize_struct("msg", 2)?;
1233 s.serialize_field("type", "WorkflowCompleted")?;
1234 s.serialize_field("state", "Idle")?;
1235 s.end()
1236 }
1237 TuiMessage::WorkflowError(err) => {
1238 let mut s = serializer.serialize_struct("msg", 2)?;
1239 s.serialize_field("type", "WorkflowError")?;
1240 s.serialize_field("error", err)?;
1241 s.end()
1242 }
1243 TuiMessage::ResumeInfo(info) => match info {
1244 Some(ri) => {
1245 let mut s = serializer.serialize_struct("msg", 5)?;
1246 s.serialize_field("type", "ResumeInfo")?;
1247 s.serialize_field("state", "Idle")?;
1248 s.serialize_field("session_id", &ri.session_id)?;
1249 s.serialize_field("checkpoint_id", &ri.checkpoint_id)?;
1250 s.serialize_field("iteration", &ri.iteration)?;
1251 s.end()
1252 }
1253 None => {
1254 let mut s = serializer.serialize_struct("msg", 2)?;
1255 s.serialize_field("type", "ResumeInfo")?;
1256 s.serialize_field("state", "Idle")?;
1257 s.end()
1258 }
1259 },
1260 TuiMessage::TodoUpdate(content) => {
1261 let mut s = serializer.serialize_struct("msg", 2)?;
1262 s.serialize_field("type", "TodoUpdate")?;
1263 s.serialize_field("content", content)?;
1264 s.end()
1265 }
1266 TuiMessage::WorkflowCancelled => {
1267 let mut s = serializer.serialize_struct("msg", 2)?;
1268 s.serialize_field("type", "WorkflowCancelled")?;
1269 s.serialize_field("state", "Idle")?;
1270 s.end()
1271 }
1272 TuiMessage::HandoffReady(briefing) => {
1273 let mut s = serializer.serialize_struct("msg", 3)?;
1274 s.serialize_field("type", "HandoffReady")?;
1275 s.serialize_field("state", "Idle")?;
1276 s.serialize_field("briefing", briefing)?;
1277 s.end()
1278 }
1279 TuiMessage::SessionRotated { old, new } => {
1280 let mut s = serializer.serialize_struct("msg", 3)?;
1281 s.serialize_field("type", "SessionRotated")?;
1282 s.serialize_field("old", old)?;
1283 s.serialize_field("new", new)?;
1284 s.end()
1285 }
1286 TuiMessage::HandoffFailed => {
1287 let mut s = serializer.serialize_struct("msg", 2)?;
1288 s.serialize_field("type", "HandoffFailed")?;
1289 s.serialize_field("state", "Idle")?;
1290 s.end()
1291 }
1292 TuiMessage::ToolPending {
1293 tool_name,
1294 hint,
1295 } => {
1296 let mut s = serializer.serialize_struct("msg", 3)?;
1297 s.serialize_field("type", "ToolPending")?;
1298 s.serialize_field("tool_name", tool_name)?;
1299 s.serialize_field("hint", hint)?;
1300 s.end()
1301 }
1302 TuiMessage::ToolDone {
1303 tool_name,
1304 success,
1305 hint,
1306 } => {
1307 let mut s = serializer.serialize_struct("msg", 4)?;
1308 s.serialize_field("type", "ToolDone")?;
1309 s.serialize_field("tool_name", tool_name)?;
1310 s.serialize_field("success", success)?;
1311 s.serialize_field("hint", hint)?;
1312 s.end()
1313 }
1314 TuiMessage::ContextTokensUpdated(count) => {
1315 let mut s = serializer.serialize_struct("msg", 2)?;
1316 s.serialize_field("type", "ContextTokensUpdated")?;
1317 s.serialize_field("count", count)?;
1318 s.end()
1319 }
1320 TuiMessage::McpServerStatus {
1321 name,
1322 connected,
1323 tool_count,
1324 error,
1325 } => {
1326 let mut s = serializer.serialize_struct("msg", 5)?;
1327 s.serialize_field("type", "McpServerStatus")?;
1328 s.serialize_field("name", name)?;
1329 s.serialize_field("connected", connected)?;
1330 s.serialize_field("tool_count", tool_count)?;
1331 s.serialize_field("error", error)?;
1332 s.end()
1333 }
1334 TuiMessage::SessionTitleUpdated(title) => {
1335 let mut s = serializer.serialize_struct("msg", 2)?;
1336 s.serialize_field("type", "SessionTitleUpdated")?;
1337 s.serialize_field("title", title)?;
1338 s.end()
1339 }
1340 }
1341 }
1342}
1343
1344fn deep_merge_toml(base: &mut toml::Value, overlay: &toml::Value) {
1355 match (base, overlay) {
1356 (toml::Value::Table(base_table), toml::Value::Table(overlay_table)) => {
1357 for (key, overlay_val) in overlay_table {
1358 match base_table.get_mut(key) {
1359 Some(base_val) => {
1360 deep_merge_toml(base_val, overlay_val);
1362 }
1363 None => {
1364 base_table.insert(key.clone(), overlay_val.clone());
1366 }
1367 }
1368 }
1369 }
1370 (base, overlay) => {
1372 *base = overlay.clone();
1373 }
1374 }
1375}
1376
1377fn substitute_env_vars(s: &mut String, secrets: &std::collections::HashMap<String, String>) {
1382 let mut result = String::with_capacity(s.len());
1384 let bytes = s.as_bytes();
1385 let mut i = 0;
1386
1387 while i < bytes.len() {
1388 if i + 1 < bytes.len() && bytes[i] == b'$' && bytes[i + 1] == b'{' {
1389 if let Some(end) = s[i + 2..].find('}') {
1391 let var_name = &s[i + 2..i + 2 + end];
1392 if let Some(value) = secrets.get(var_name) {
1394 result.push_str(value);
1395 } else if let Ok(value) = std::env::var(var_name) {
1396 result.push_str(&value);
1397 } else {
1398 result.push_str(&s[i..i + 2 + end + 1]);
1400 }
1401 i = i + 2 + end + 1;
1402 } else {
1403 result.push('$');
1405 i += 1;
1406 }
1407 } else {
1408 result.push(bytes[i] as char);
1409 i += 1;
1410 }
1411 }
1412
1413 *s = result;
1414}
1415
1416fn fingerprint_mcp_section(resolved: Option<&str>) -> u64 {
1426 use sha2::{Digest, Sha256};
1427 let section = resolved
1428 .and_then(|s| s.parse::<toml::Value>().ok())
1429 .and_then(|v| v.get("mcp").cloned());
1430 let bytes = match section {
1431 Some(v) => v.to_string().into_bytes(),
1432 None => b"<no-mcp>".to_vec(),
1433 };
1434 let digest = Sha256::digest(&bytes);
1435 u64::from_be_bytes(digest[..8].try_into().expect("sha256 digest >= 8 bytes"))
1436}
1437
1438#[cfg(test)]
1439mod tests {
1440 use super::*;
1441
1442 fn test_state() -> ServerState {
1444 let (session, _rx) = Session::new();
1445 let (ws_tx, _ws_rx) = tokio::sync::broadcast::channel::<String>(16);
1446 ServerState::new(session, ws_tx, None)
1447 }
1448
1449 fn temp_user_home(tag: &str) -> std::path::PathBuf {
1451 static COUNTER: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0);
1452 let n = COUNTER.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
1453 let dir = std::env::temp_dir().join(format!(
1454 "trustee-state-test-{}-{}-{}",
1455 tag,
1456 std::process::id(),
1457 n
1458 ));
1459 std::fs::create_dir_all(dir.join("config")).expect("create temp user home");
1460 dir
1461 }
1462
1463 fn parse(toml_str: &str) -> toml::Value {
1464 toml_str.parse::<toml::Value>().expect("valid test TOML")
1465 }
1466
1467 #[test]
1470 fn overlay_allowlist_drops_non_allowlisted_sections() {
1471 let state = test_state();
1472 let home = temp_user_home("allowlist");
1473 std::fs::write(
1474 home.join("config").join("trustee.toml"),
1475 "[server]\nport = 1\n\n[auth]\nmode = \"kanidm\"\n\n[mcp]\nmode = \"user\"\n",
1476 )
1477 .expect("write overlay");
1478
1479 let shared = "[server]\nport = 8080\n\n[mcp]\nmode = \"shared\"\n";
1480 let merged = state
1481 .merge_user_config(&home, Some(shared))
1482 .expect("overlay has allowlisted content");
1483
1484 let merged_val = parse(&merged);
1485 let expected = parse("[server]\nport = 8080\n\n[mcp]\nmode = \"user\"\n");
1486 assert_eq!(merged_val, expected, "merged config must be shared + [mcp] overlay only");
1487 assert!(merged_val.get("auth").is_none(), "overlay [auth] must be dropped");
1488 }
1489
1490 #[test]
1492 fn no_overlay_file_returns_none() {
1493 let state = test_state();
1494 let home = temp_user_home("empty");
1495 assert!(state
1496 .merge_user_config(&home, Some("[server]\nport = 8080\n"))
1497 .is_none());
1498 }
1499
1500 #[test]
1502 fn overlay_with_no_allowlisted_sections_is_noop() {
1503 let state = test_state();
1504 let home = temp_user_home("all-dropped");
1505 std::fs::write(
1506 home.join("config").join("trustee.toml"),
1507 "[server]\nport = 1\n\n[storage]\npath = \"/tmp/x\"\n",
1508 )
1509 .expect("write overlay");
1510 assert!(state
1511 .merge_user_config(&home, Some("[server]\nport = 8080\n"))
1512 .is_none());
1513 }
1514
1515 #[test]
1518 fn llm_overlay_dropped_by_default_and_kept_when_enabled() {
1519 let shared = "[llm]\nprovider = \"openai\"\n\n[mcp]\nmode = \"shared\"\n";
1520 let overlay = "[llm]\nprovider = \"anthropic\"\n";
1521
1522 let state = test_state();
1524 assert!(!state.allow_llm_overlay);
1525 let home = temp_user_home("llm-off");
1526 std::fs::write(home.join("config").join("trustee.toml"), overlay).expect("write overlay");
1527 assert!(state.merge_user_config(&home, Some(shared)).is_none());
1528
1529 let state = test_state().with_allow_llm_overlay(true);
1531 assert!(state.allow_llm_overlay);
1532 let home = temp_user_home("llm-on");
1533 std::fs::write(home.join("config").join("trustee.toml"), overlay).expect("write overlay");
1534 let merged = state
1535 .merge_user_config(&home, Some(shared))
1536 .expect("[llm] overlay applies when opted in");
1537 let merged_val = parse(&merged);
1538 assert_eq!(merged_val["llm"]["provider"].as_str(), Some("anthropic"));
1539 assert_eq!(merged_val["mcp"]["mode"].as_str(), Some("shared"));
1540 }
1541
1542 #[test]
1546 fn user_home_dir_uses_consolidated_hash() {
1547 let state = test_state();
1548 if let Some(home) = state.get_user_home_dir("farzan@example.com") {
1549 assert_eq!(
1550 home.file_name().and_then(|n| n.to_str()),
1551 Some(trustee_core::user_hash("farzan@example.com")).as_deref()
1552 );
1553 let users_root = dirs::home_dir().unwrap().join(".trustee").join("users");
1554 assert_eq!(home.parent(), Some(&users_root).map(|p| p.as_path()));
1555 }
1556 }
1558
1559 #[test]
1562 fn lifecycle_overlay_survives_allowlist_enabled_stripped() {
1563 let state = test_state();
1564 let home = temp_user_home("lifecycle-persona");
1565 std::fs::write(
1566 home.join("config").join("trustee.toml"),
1567 "[lifecycle]\nsystem_template = \"You are Saman.\"\nenabled = true\n",
1568 )
1569 .expect("write overlay");
1570
1571 let merged = state
1572 .merge_user_config(
1573 &home,
1574 Some("[lifecycle]\nsystem_template = \"shared default\"\n"),
1575 )
1576 .expect("lifecycle overlay is allowlisted");
1577
1578 let v = merged.parse::<toml::Value>().expect("valid TOML");
1579 assert_eq!(
1580 v["lifecycle"]["system_template"].as_str(),
1581 Some("You are Saman."),
1582 "agent persona wins over shared default"
1583 );
1584 assert_eq!(
1585 v["lifecycle"].get("enabled"),
1586 None,
1587 "WASM extension switch is not user-settable"
1588 );
1589 }
1590
1591 #[tokio::test]
1596 async fn agent_principals_get_isolated_session_buckets() {
1597 let state = test_state();
1598 let key_a = "agent-farzan";
1599 let key_b = "agent-paydar";
1600
1601 let (sid_a, session_a, _tx_a, _ts_a) = state.ensure_active_session(key_a).await;
1602 let (_sid_b, _session_b, _tx_b, _ts_b) = state.ensure_active_session(key_b).await;
1603
1604 assert!(
1606 state.get_session_by_any_id(key_a, &sid_a).await.is_some(),
1607 "owner bucket resolves its own session"
1608 );
1609 assert!(
1610 state.get_session_by_any_id(key_b, &sid_a).await.is_none(),
1611 "cross-agent session access must be 404/None"
1612 );
1613 assert!(Arc::strong_count(&session_a) >= 1);
1614
1615 if let Some(home) = state.get_user_home_dir(key_a) {
1617 assert_eq!(
1618 home.file_name().and_then(|n| n.to_str()),
1619 Some(trustee_core::user_hash(key_a)).as_deref()
1620 );
1621 }
1622 }
1623}
1624
1625#[cfg(test)]
1630mod mcp_loader_cache_tests {
1631 use super::*;
1632
1633 fn state_with_shared(shared: &str) -> ServerState {
1634 let (session, _rx) = Session::new();
1635 let (ws_tx, _ws_rx) = tokio::sync::broadcast::channel::<String>(16);
1636 let mut state = ServerState::new(session, ws_tx, None);
1637 state.config_toml = Some(shared.to_string());
1638 state
1639 }
1640
1641 struct TempUser {
1645 key: String,
1646 home: std::path::PathBuf,
1647 }
1648
1649 impl TempUser {
1650 fn new(tag: &str) -> Self {
1651 let key = format!("16c-{tag}-{}@test.invalid", std::process::id());
1652 let home = dirs::home_dir()
1653 .expect("HOME available in test env")
1654 .join(".trustee")
1655 .join("users")
1656 .join(trustee_core::user_hash(&key));
1657 std::fs::create_dir_all(home.join("config")).expect("create user home");
1658 Self { key, home }
1659 }
1660
1661 fn write_overlay(&self, toml_str: &str) {
1662 std::fs::write(self.home.join("config").join("trustee.toml"), toml_str)
1663 .expect("write overlay");
1664 }
1665 }
1666
1667 impl Drop for TempUser {
1668 fn drop(&mut self) {
1669 let _ = std::fs::remove_dir_all(&self.home);
1670 }
1671 }
1672
1673 #[tokio::test]
1674 async fn no_mcp_config_caches_disabled_marker() {
1675 let state = state_with_shared("[server]\nport = 8080\n");
1676 let user = TempUser::new("nomcp");
1677 let ts = Arc::new(pep::MemoryTokenStore::new());
1678
1679 let first = state.get_or_build_mcp_loader(&user.key, &ts).await.unwrap();
1680 assert!(first.is_none(), "no [mcp] anywhere → disabled marker");
1681 assert_eq!(state.mcp_loaders.len(), 1, "exactly one cache entry");
1682
1683 let second = state.get_or_build_mcp_loader(&user.key, &ts).await.unwrap();
1685 assert!(second.is_none());
1686 assert_eq!(state.mcp_loaders.len(), 1);
1687 }
1688
1689 #[tokio::test]
1690 async fn concurrent_cold_builds_single_flight() {
1691 let state = Arc::new(state_with_shared("[server]\nport = 8080\n"));
1692 let user = TempUser::new("singleflight");
1693 let ts = Arc::new(pep::MemoryTokenStore::new());
1694
1695 let mut handles = Vec::new();
1696 for _ in 0..5 {
1697 let state = state.clone();
1698 let key = user.key.clone();
1699 let ts = ts.clone();
1700 handles.push(tokio::spawn(async move {
1701 state.get_or_build_mcp_loader(&key, &ts).await
1702 }));
1703 }
1704 for h in handles {
1705 h.await.unwrap().expect("all five succeed");
1706 }
1707 assert_eq!(state.mcp_loaders.len(), 1, "single-flight → one entry");
1708 }
1709
1710 #[tokio::test]
1711 async fn fingerprint_change_triggers_rebuild() {
1712 let state = state_with_shared("[server]\nport = 8080\n");
1713 let user = TempUser::new("fpchange");
1714 let ts = Arc::new(pep::MemoryTokenStore::new());
1715
1716 user.write_overlay("[mcp]\nenabled = true\n\n[[mcp.servers]]\nname = \"v1\"\nurl = \"http://127.0.0.1:9/sse\"\n");
1719 let v1 = state.get_or_build_mcp_loader(&user.key, &ts).await.unwrap();
1720 assert!(v1.is_some(), "enabled [mcp] → real loader");
1721 let fp1 = state
1722 .mcp_loaders
1723 .get(&trustee_core::user_hash(&user.key))
1724 .unwrap()
1725 .fingerprint;
1726
1727 std::thread::sleep(std::time::Duration::from_millis(5));
1729 user.write_overlay("[mcp]\nenabled = true\n\n[[mcp.servers]]\nname = \"v2\"\nurl = \"http://127.0.0.1:9/other\"\n");
1730 let v2 = state.get_or_build_mcp_loader(&user.key, &ts).await.unwrap();
1731 assert!(v2.is_some());
1732 let entry = state
1733 .mcp_loaders
1734 .get(&trustee_core::user_hash(&user.key))
1735 .unwrap();
1736 assert_ne!(
1737 entry.fingerprint, fp1,
1738 "fingerprint must change with content"
1739 );
1740 assert!(entry.degraded.is_none());
1741
1742 let _still_usable = v1.as_ref().unwrap().tool_count;
1744 }
1745
1746 #[tokio::test]
1747 async fn degraded_entry_fails_loud_within_backoff_and_isolates_users() {
1748 let state = state_with_shared("[server]\nport = 8080\n");
1749 let bad = TempUser::new("degraded-bad");
1750 let good = TempUser::new("degraded-good");
1751 let ts = Arc::new(pep::MemoryTokenStore::new());
1752
1753 bad.write_overlay("[mcp]\nenabled = \"not-a-bool\"\n");
1755 let err = match state.get_or_build_mcp_loader(&bad.key, &ts).await {
1756 Ok(_) => panic!("invalid [mcp] must fail loud"),
1757 Err(e) => e,
1758 };
1759 assert!(
1760 err.contains("invalid [mcp]"),
1761 "surfaces the parse error: {err}"
1762 );
1763
1764 let entry = state
1765 .mcp_loaders
1766 .get(&trustee_core::user_hash(&bad.key))
1767 .unwrap();
1768 assert!(entry.degraded.is_some(), "poison entry recorded");
1769
1770 let err2 = match state.get_or_build_mcp_loader(&bad.key, &ts).await {
1772 Ok(_) => panic!("still within backoff"),
1773 Err(e) => e,
1774 };
1775 assert_eq!(err, err2, "same cached error");
1776
1777 let good_loader = state
1779 .get_or_build_mcp_loader(&good.key, &ts)
1780 .await
1781 .expect("other user unaffected");
1782 assert!(
1783 good_loader.is_none(),
1784 "good user has no [mcp] → disabled marker"
1785 );
1786 }
1787}