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
147#[derive(Clone)]
149pub struct ServerState {
150 pub sessions: SessionRegistry,
152 pub ws_tx: broadcast::Sender<String>,
154 pub auth: Option<Arc<AuthState>>,
156 pub config_toml: Option<String>,
158 pub secrets: Option<std::collections::HashMap<String, String>>,
160 pub build_info: Option<trustee_core::types::BuildInfo>,
162 pub workflow_semaphore: Arc<tokio::sync::Semaphore>,
165 pub max_sessions_per_user: usize,
167 pub allow_llm_overlay: bool,
170 pub mcp_loaders: Arc<DashMap<String, McpLoaderEntry>>,
172 mcp_build_locks: Arc<DashMap<String, Arc<tokio::sync::Mutex<()>>>>,
174 pub thq_dispatch: Arc<DashMap<String, ThqDispatchEntry>>,
176 pub agent_dispatch_tokens: Arc<DashMap<String, (String, std::time::Instant)>>,
179}
180
181impl ServerState {
182 pub fn new(
184 session: Session,
185 ws_tx: broadcast::Sender<String>,
186 auth: Option<Arc<AuthState>>,
187 ) -> Self {
188 let sessions = Arc::new(DashMap::new());
189
190 let token_store = Arc::new(pep::MemoryTokenStore::new());
192 let (ws_tx_entry, _) = broadcast::channel::<String>(256);
193
194 let now = chrono::Utc::now();
195 let initial_entry = UserSessionEntry {
196 session: Arc::new(Mutex::new(session)),
197 ws_tx: ws_tx_entry,
198 created_at: now,
199 last_active: Arc::new(Mutex::new(now)),
200 };
201
202 let user_sessions = UserSessions {
203 sessions: DashMap::new(),
204 token_store,
205 active_session_id: Mutex::new(String::new()),
206 };
207 user_sessions.sessions.insert("default".to_string(), initial_entry);
208
209 sessions.insert("default".to_string(), user_sessions);
210
211 Self {
212 sessions,
213 ws_tx,
214 auth,
215 config_toml: None,
216 secrets: None,
217 build_info: None,
218 workflow_semaphore: Arc::new(tokio::sync::Semaphore::new(8)),
219 max_sessions_per_user: 4,
220 allow_llm_overlay: false,
221 mcp_loaders: Arc::new(DashMap::new()),
222 mcp_build_locks: Arc::new(DashMap::new()),
223 thq_dispatch: Arc::new(DashMap::new()),
224 agent_dispatch_tokens: Arc::new(DashMap::new()),
225 }
226 }
227
228 pub fn with_config_toml(mut self, config_toml: String) -> Self {
229 self.config_toml = Some(config_toml);
230 self
231 }
232
233 pub fn service_issuer(&self) -> Option<String> {
244 let config = self.config_toml.as_deref()?;
245 let v: toml::Value = config.parse().ok()?;
246 let creds = v.get("mcp")?.get("credentials")?.as_table()?;
247 for (_name, cred) in creds {
248 if cred.get("type").and_then(|t| t.as_str()) == Some("service-account") {
249 if let Some(issuer) = cred.get("issuer_url").and_then(|i| i.as_str()) {
250 return Some(issuer.to_string());
251 }
252 }
253 }
254 None
255 }
256
257 pub fn with_secrets(mut self, secrets: std::collections::HashMap<String, String>) -> Self {
258 self.secrets = Some(secrets);
259 self
260 }
261
262 pub fn with_build_info(mut self, build_info: trustee_core::types::BuildInfo) -> Self {
263 self.build_info = Some(build_info);
264 self
265 }
266
267 pub fn with_max_concurrent_workflows(mut self, max: usize) -> Self {
268 self.workflow_semaphore = Arc::new(tokio::sync::Semaphore::new(max));
269 self
270 }
271
272 pub fn with_max_sessions_per_user(mut self, max: usize) -> Self {
274 self.max_sessions_per_user = max;
275 self
276 }
277
278 pub fn with_allow_llm_overlay(mut self, allow: bool) -> Self {
281 self.allow_llm_overlay = allow;
282 self
283 }
284
285 pub async fn create_session(
295 &self,
296 user_key: &str,
297 session_name: Option<String>,
298 identity: Option<String>,
299 activate: bool,
300 ) -> Result<String, SessionError> {
301 let user_sessions = self
303 .sessions
304 .entry(user_key.to_string())
305 .or_insert_with(|| UserSessions {
306 sessions: DashMap::new(),
307 token_store: Arc::new(pep::MemoryTokenStore::new()),
308 active_session_id: Mutex::new(String::new()),
309 });
310
311 if user_sessions.sessions.len() >= self.max_sessions_per_user {
313 return Err(SessionError::MaxSessionsReached(self.max_sessions_per_user));
314 }
315
316 let (mut session, workflow_rx) = Session::new();
318
319 if let Some(ref config_toml) = self.config_toml {
321 session.config_toml = Some(config_toml.clone());
322 session.parse_auto_handoff_config();
323 if let Ok(table) = config_toml.parse::<toml::Value>() {
324 if let Some(name) = table
325 .get("agent")
326 .and_then(|a| a.get("name"))
327 .and_then(|n| n.as_str())
328 {
329 session.agent_name = name.to_string();
330 }
331 }
332 }
333
334 session.secrets = self.secrets.clone();
335 session.build_info = self.build_info.clone();
336
337 self.apply_user_isolation(&mut session, user_key);
339
340 session.session_name = session_name;
342
343 session.identity = identity;
345
346 let (ws_tx_entry, _) = broadcast::channel::<String>(256);
348
349 let session_id = format!(
351 "session_{}_{}",
352 chrono::Utc::now().format("%Y_%m_%d_%H_%M"),
353 &uuid::Uuid::new_v4().to_string()[..8]
354 );
355
356 let now = chrono::Utc::now();
357
358 user_sessions.sessions.insert(
360 session_id.clone(),
361 UserSessionEntry {
362 session: Arc::new(Mutex::new(session)),
363 ws_tx: ws_tx_entry.clone(),
364 created_at: now,
365 last_active: Arc::new(Mutex::new(now)),
366 },
367 );
368
369 if activate {
378 *user_sessions.active_session_id.lock().await = session_id.clone();
379 }
380
381 let session_arc = user_sessions
383 .sessions
384 .get(&session_id)
385 .map(|e| e.session.clone());
386 if let Some(session_arc) = session_arc {
387 self.spawn_user_drain_task(
388 session_id.clone(),
389 session_arc,
390 ws_tx_entry,
391 workflow_rx,
392 );
393 }
394
395 Ok(session_id)
396 }
397
398 pub async fn get_session(
401 &self,
402 user_key: &str,
403 session_id: &str,
404 ) -> Option<(Arc<Mutex<Session>>, broadcast::Sender<String>)> {
405 let user_sessions = self.sessions.get(user_key)?;
406 let entry = user_sessions.sessions.get(session_id)?;
407
408 let now = chrono::Utc::now();
410 *entry.last_active.lock().await = now;
411
412 Some((entry.session.clone(), entry.ws_tx.clone()))
413 }
414
415 pub async fn get_session_by_any_id(
433 &self,
434 user_key: &str,
435 id: &str,
436 ) -> Option<(String, Arc<Mutex<Session>>, broadcast::Sender<String>)> {
437 let user_sessions = self.sessions.get(user_key)?;
439 if let Some(entry) = user_sessions.sessions.get(id) {
440 let now = chrono::Utc::now();
442 *entry.last_active.lock().await = now;
443 return Some((id.to_string(), entry.session.clone(), entry.ws_tx.clone()));
444 }
445
446 for entry in user_sessions.sessions.iter() {
448 let session = entry.session.lock().await;
449 if session.session_id.as_deref() == Some(id) {
450 let key = entry.key().clone();
451 let ws_tx = entry.ws_tx.clone();
452 drop(session);
453 let now = chrono::Utc::now();
455 *entry.last_active.lock().await = now;
456 return Some((key, entry.session.clone(), ws_tx));
457 }
458 }
459
460 None
461 }
462
463 pub async fn list_sessions(&self, user_key: &str) -> Vec<SessionListItem> {
465 let Some(user_sessions) = self.sessions.get(user_key) else {
466 return Vec::new();
467 };
468
469 let mut items = Vec::new();
470 for entry in user_sessions.sessions.iter() {
471 let session = entry.session.lock().await;
472 let workflow_state = match session.workflow_state {
473 trustee_core::types::WorkflowState::Idle => "Idle",
474 trustee_core::types::WorkflowState::Running => "Running",
475 trustee_core::types::WorkflowState::Cancelling => "Cancelling",
476 };
477 let last_active = entry.last_active.lock().await;
478 items.push(SessionListItem {
479 session_id: entry.key().clone(),
480 session_name: session.session_name.clone(),
481 workflow_state: workflow_state.to_string(),
482 created_at: entry.created_at.to_rfc3339(),
483 last_active: last_active.to_rfc3339(),
484 handoff_count: session.handoff_count,
485 });
486 }
487 drop(user_sessions);
488
489 items.sort_by(|a, b| b.last_active.cmp(&a.last_active));
491 items
492 }
493
494 pub async fn destroy_session(
496 &self,
497 user_key: &str,
498 session_id: &str,
499 ) -> Result<(), SessionError> {
500 let user_sessions = self
501 .sessions
502 .get(user_key)
503 .ok_or_else(|| SessionError::NotFound(session_id.to_string()))?;
504
505 {
507 let entry = user_sessions
508 .sessions
509 .get(session_id)
510 .ok_or_else(|| SessionError::NotFound(session_id.to_string()))?;
511 let session = entry.session.lock().await;
512 if session.workflow_state != trustee_core::types::WorkflowState::Idle {
513 let state_str = match session.workflow_state {
514 trustee_core::types::WorkflowState::Running => "Running",
515 trustee_core::types::WorkflowState::Cancelling => "Cancelling",
516 _ => "Unknown",
517 };
518 return Err(SessionError::NotIdle(state_str.to_string()));
519 }
520 }
521
522 user_sessions.sessions.remove(session_id);
524
525 let mut active_id = user_sessions.active_session_id.lock().await;
527 if &*active_id == session_id {
528 let mut newest: Option<(String, chrono::DateTime<chrono::Utc>)> = None;
530 for entry in user_sessions.sessions.iter() {
531 let la = entry.last_active.lock().await;
532 if newest.as_ref().map_or(true, |(_, t)| *la > *t) {
533 newest = Some((entry.key().clone(), *la));
534 }
535 }
536 *active_id = newest.map(|(id, _)| id).unwrap_or_default();
537 }
538
539 Ok(())
540 }
541
542 pub async fn ensure_active_session(
551 &self,
552 user_key: &str,
553 ) -> (
554 String,
555 Arc<Mutex<Session>>,
556 broadcast::Sender<String>,
557 Arc<pep::MemoryTokenStore>,
558 ) {
559 let token_store = {
561 let user_sessions = self
562 .sessions
563 .entry(user_key.to_string())
564 .or_insert_with(|| UserSessions {
565 sessions: DashMap::new(),
566 token_store: Arc::new(pep::MemoryTokenStore::new()),
567 active_session_id: Mutex::new(String::new()),
568 });
569 user_sessions.token_store.clone()
570 };
571
572 let active_id = {
574 let user_sessions = self.sessions.get(user_key).unwrap();
575 let guard = user_sessions.active_session_id.lock().await;
576 guard.clone()
577 };
578
579 if !active_id.is_empty() {
580 if let Some((session, ws_tx)) = self.get_session(user_key, &active_id).await {
581 return (active_id, session, ws_tx, token_store);
582 }
583 }
585
586 let existing_session: Option<(String, Arc<Mutex<Session>>, broadcast::Sender<String>)> = {
590 let user_sessions = self.sessions.get(user_key).unwrap();
591 let result = user_sessions.sessions.iter().next().map(|first| {
592 (
593 first.key().clone(),
594 first.session.clone(),
595 first.ws_tx.clone(),
596 )
597 });
598 result
599 };
600 if let Some((id, session, ws_tx)) = existing_session {
601 let now = chrono::Utc::now();
602 if let Some(entry) = self.sessions.get(user_key) {
603 if let Some(e) = entry.sessions.get(&id) {
604 *e.last_active.lock().await = now;
605 }
606 *entry.active_session_id.lock().await = id.clone();
607 }
608
609 return (id, session, ws_tx, token_store);
610 }
611
612 let session_id = self
614 .create_session(user_key, None, None, true)
615 .await
616 .unwrap_or_else(|_| "default".to_string());
617
618 let (session, ws_tx) = self
619 .get_session(user_key, &session_id)
620 .await
621 .expect("just-created session must exist");
622
623 (session_id, session, ws_tx, token_store)
624 }
625
626 pub async fn ensure_user_session(
629 &self,
630 user_key: &str,
631 ) -> (Arc<Mutex<Session>>, broadcast::Sender<String>, Arc<pep::MemoryTokenStore>) {
632 let (_id, session, ws_tx, token_store) = self.ensure_active_session(user_key).await;
633 (session, ws_tx, token_store)
634 }
635
636 pub async fn set_active_session(&self, user_key: &str, session_id: &str) {
638 if let Some(user_sessions) = self.sessions.get(user_key) {
639 if user_sessions.sessions.contains_key(session_id) {
640 *user_sessions.active_session_id.lock().await = session_id.to_string();
641 }
642 }
643 }
644
645 pub fn get_user_home_dir(&self, user_key: &str) -> Option<std::path::PathBuf> {
656 let hash = trustee_core::user_hash(user_key);
657 dirs::home_dir().map(|home| home.join(".trustee").join("users").join(&hash))
658 }
659
660 pub fn get_user_config_and_home(&self, user_key: &str) -> (Option<String>, Option<std::path::PathBuf>) {
665 (self.config_toml.clone(), self.get_user_home_dir(user_key))
666 }
667
668 fn apply_user_isolation(&self, session: &mut Session, user_key: &str) {
677 let user_hash = trustee_core::user_hash(user_key);
678
679 let user_home = if let Some(home) = dirs::home_dir() {
681 let user_home = home.join(".trustee").join("users").join(&user_hash);
682 session.home_dir = Some(user_home.clone());
683 Some(user_home)
684 } else {
685 None
686 };
687
688 session.project_id = Some(format!("web{}", &user_hash[..16]));
689
690 let shared_secrets = session.secrets.clone().unwrap_or_default();
697 let mut merged_secrets = shared_secrets.clone();
698
699 if let Some(ref user_home) = user_home {
700 if let Ok(merged) = self.load_user_secrets(user_home, &merged_secrets) {
701 merged_secrets = merged;
702 }
703 }
704
705 if let Some(ref user_home) = user_home {
707 if let Some(merged) = self.merge_user_config(user_home, session.config_toml.as_deref()) {
708 session.config_toml = Some(merged);
709 tracing::debug!("Merged per-user config into session");
710 }
711 }
712
713 if let Some(ref mut config_toml) = session.config_toml {
715 substitute_env_vars(config_toml, &merged_secrets);
716 }
717
718 session.secrets = Some(shared_secrets);
720 }
721
722 fn load_user_secrets(
725 &self,
726 user_home: &std::path::Path,
727 base: &std::collections::HashMap<String, String>,
728 ) -> std::io::Result<std::collections::HashMap<String, String>> {
729 let user_env_path = user_home.join(".env");
730 if !user_env_path.exists() {
731 return Ok(base.clone());
732 }
733 let content = std::fs::read_to_string(&user_env_path)?;
734 let mut merged = base.clone();
735 for line in content.lines() {
736 let line = line.trim();
737 if line.is_empty() || line.starts_with('#') {
738 continue;
739 }
740 if let Some((key, value)) = line.split_once('=') {
741 let key = key.trim().to_string();
742 let value = value
743 .trim()
744 .trim_matches('"')
745 .trim_matches('\'')
746 .to_string();
747 merged.insert(key, value);
748 }
749 }
750 tracing::debug!("Loaded per-user secrets from {}", user_env_path.display());
751 Ok(merged)
752 }
753
754 fn merge_user_config(
769 &self,
770 user_home: &std::path::Path,
771 shared_config: Option<&str>,
772 ) -> Option<String> {
773 let user_config_path = user_home.join("config").join("trustee.toml");
774 if !user_config_path.exists() {
775 return None;
776 }
777 let user_config_toml = std::fs::read_to_string(&user_config_path).ok()?;
778 let shared = shared_config
779 .unwrap_or("")
780 .parse::<toml::Value>()
781 .ok()?;
782 let overlay = user_config_toml.parse::<toml::Value>().ok()?;
783
784 let allowed = |section: &str| {
785 section == "mcp"
786 || section == "thq"
790 || (self.allow_llm_overlay && section == "llm")
791 };
792 let dir_name = user_home
795 .file_name()
796 .and_then(|n| n.to_str())
797 .unwrap_or("<unknown>");
798 let masked_user = dir_name.get(..8).unwrap_or(dir_name);
799
800 let mut filtered_overlay = toml::map::Map::new();
801 if let Some(table) = overlay.as_table() {
802 for (section, value) in table {
803 if allowed(section) {
804 filtered_overlay.insert(section.clone(), value.clone());
805 } else {
806 tracing::warn!(
807 "user config overlay: dropping non-allowlisted section [{}] for user {}",
808 section,
809 masked_user
810 );
811 }
812 }
813 }
814
815 if filtered_overlay.is_empty() {
816 return None;
818 }
819 let overlay = toml::Value::Table(filtered_overlay);
820
821 let mut shared = shared;
822 deep_merge_toml(&mut shared, &overlay);
823 let merged = toml::to_string(&shared).ok()?;
824 tracing::debug!("Merged per-user config from {}", user_config_path.display());
825 Some(merged)
826 }
827
828 pub fn resolve_user_config(&self, user_key: &str) -> Option<String> {
837 let config_toml = self.config_toml.clone()?;
838
839 let user_home = self.get_user_home_dir(user_key)?;
841
842 let mut merged_secrets = self.secrets.clone().unwrap_or_default();
844 if let Ok(merged) = self.load_user_secrets(&user_home, &merged_secrets) {
845 merged_secrets = merged;
846 }
847
848 let mut resolved = config_toml;
850 if let Some(merged) = self.merge_user_config(&user_home, Some(&resolved)) {
851 resolved = merged;
852 }
853
854 substitute_env_vars(&mut resolved, &merged_secrets);
856
857 Some(resolved)
858 }
859
860 pub async fn get_or_build_mcp_loader(
876 &self,
877 user_key: &str,
878 token_store: &Arc<pep::MemoryTokenStore>,
879 ) -> Result<Option<std::sync::Arc<abk::agent::McpToolLoader>>, String> {
880 let user_hash = trustee_core::user_hash(user_key);
881
882 let resolved = self.resolve_user_config(user_key);
884 let fingerprint = fingerprint_mcp_section(resolved.as_deref());
885
886 if let Some(entry) = self.mcp_loaders.get(&user_hash) {
888 if entry.degraded.is_none() {
889 if entry.fingerprint == fingerprint {
890 return Ok(entry.loader.clone());
891 }
892 } else if let (Some(err), Some(failed_at)) = (&entry.degraded, entry.failed_at) {
893 let backoff = chrono::Duration::from_std(MCP_BUILD_RETRY_BACKOFF)
894 .unwrap_or_else(|_| chrono::Duration::seconds(30));
895 if chrono::Utc::now() < failed_at + backoff {
896 return Err(err.clone());
897 }
898 }
899 }
900
901 let lock = self
903 .mcp_build_locks
904 .entry(user_hash.clone())
905 .or_insert_with(|| Arc::new(tokio::sync::Mutex::new(())))
906 .clone();
907 let _guard = lock.lock().await;
908
909 if let Some(entry) = self.mcp_loaders.get(&user_hash) {
911 if entry.degraded.is_none() && entry.fingerprint == fingerprint {
912 return Ok(entry.loader.clone());
913 }
914 }
915
916 match self
917 .build_mcp_loader(&user_hash, resolved.as_deref(), fingerprint, token_store)
918 .await
919 {
920 Ok(entry) => {
921 self.mcp_loaders.insert(user_hash.clone(), entry);
922 Ok(self.mcp_loaders.get(&user_hash).unwrap().loader.clone())
923 }
924 Err(err) => {
925 tracing::warn!(
926 "MCP loader build FAILED for user {}; dispatch fails loud, retry after {:?}",
927 &user_hash[..8.min(user_hash.len())],
928 MCP_BUILD_RETRY_BACKOFF
929 );
930 self.mcp_loaders.insert(
931 user_hash,
932 McpLoaderEntry {
933 loader: None,
934 fingerprint,
935 built_at: chrono::Utc::now(),
936 degraded: Some(err.clone()),
937 failed_at: Some(chrono::Utc::now()),
938 },
939 );
940 Err(err)
941 }
942 }
943 }
944
945 async fn build_mcp_loader(
948 &self,
949 user_hash: &str,
950 resolved: Option<&str>,
951 fingerprint: u64,
952 token_store: &Arc<pep::MemoryTokenStore>,
953 ) -> Result<McpLoaderEntry, String> {
954 let mcp_config: Option<abk::config::McpConfig> = match resolved {
955 Some(toml_str) => {
956 let value = toml_str
957 .parse::<toml::Value>()
958 .map_err(|e| format!("config parse failed: {}", e))?;
959 match value.get("mcp") {
960 Some(section) => {
961 use serde::Deserialize as _;
962 Some(
963 abk::config::McpConfig::deserialize(section.clone())
964 .map_err(|e| format!("invalid [mcp] config: {}", e))?,
965 )
966 }
967 None => None,
968 }
969 }
970 None => None,
971 };
972
973 let loader = match mcp_config {
974 Some(cfg) if cfg.enabled => {
975 let built = abk::agent::McpToolLoader::with_token_store(
976 &cfg,
977 Some(token_store.clone() as std::sync::Arc<dyn pep::token_store::TokenStore>),
978 )
979 .await
980 .map_err(|e| format!("MCP loader build failed: {}", e))?;
981
982 let servers: Vec<String> = built
986 .server_statuses
987 .iter()
988 .map(|s| {
989 if s.connected {
990 format!("{}(up,{}tools)", s.name, s.tool_count)
991 } else {
992 format!("{}(DOWN)", s.name)
993 }
994 })
995 .collect();
996 tracing::info!(
997 "MCP loader built for user {}: servers=[{}] total_tools={}",
998 &user_hash[..8.min(user_hash.len())],
999 servers.join(", "),
1000 built.tool_count
1001 );
1002 Some(std::sync::Arc::new(built))
1003 }
1004 _ => None,
1005 };
1006
1007 Ok(McpLoaderEntry {
1008 loader,
1009 fingerprint,
1010 built_at: chrono::Utc::now(),
1011 degraded: None,
1012 failed_at: None,
1013 })
1014 }
1015
1016 fn spawn_user_drain_task(
1018 &self,
1019 session_id: String,
1020 session: Arc<Mutex<Session>>,
1021 ws_tx: broadcast::Sender<String>,
1022 mut workflow_rx: mpsc::UnboundedReceiver<TuiMessage>,
1023 ) {
1024 let mut last_broadcast_state: Option<String> = None;
1030 tokio::spawn(async move {
1031 while let Some(msg) = workflow_rx.recv().await {
1032 {
1033 let mut session = session.lock().await;
1034 session.handle_workflow_message(msg.clone());
1035
1036 let state_str = match session.workflow_state {
1037 trustee_core::types::WorkflowState::Idle => "Idle",
1038 trustee_core::types::WorkflowState::Running => "Running",
1039 trustee_core::types::WorkflowState::Cancelling => "Cancelling",
1040 };
1041 if last_broadcast_state.as_deref() != Some(state_str) {
1042 last_broadcast_state = Some(state_str.to_string());
1043 let state_msg = serde_json::json!({
1044 "type": "StateChanged",
1045 "state": state_str
1046 });
1047 let _ = ws_tx.send(state_msg.to_string());
1048 }
1049 }
1050
1051 let json =
1052 serde_json::to_string(&SerializableMessage(&msg)).unwrap_or_default();
1053 let _ = ws_tx.send(json);
1054 }
1055 tracing::debug!("Drain task ended for session: {}", session_id);
1056 });
1057 }
1058
1059 pub fn spawn_drain_task(self, mut workflow_rx: mpsc::UnboundedReceiver<TuiMessage>) {
1062 let default_user = self
1064 .sessions
1065 .get("default")
1066 .expect("default user must exist");
1067 let first_entry = default_user
1068 .sessions
1069 .iter()
1070 .next()
1071 .expect("default user must have at least one session");
1072 let session = first_entry.session.clone();
1073 let ws_tx = first_entry.ws_tx.clone();
1074 let session_id = first_entry.key().clone();
1075 drop(first_entry);
1076 drop(default_user);
1077
1078 let mut last_broadcast_state = Some("Running".to_string());
1082 tokio::spawn(async move {
1083 while let Some(msg) = workflow_rx.recv().await {
1084 {
1085 let mut session = session.lock().await;
1086 session.handle_workflow_message(msg.clone());
1087
1088 let state_str = match session.workflow_state {
1089 trustee_core::types::WorkflowState::Idle => "Idle",
1090 trustee_core::types::WorkflowState::Running => "Running",
1091 trustee_core::types::WorkflowState::Cancelling => "Cancelling",
1092 };
1093 if last_broadcast_state.as_deref() != Some(state_str) {
1094 last_broadcast_state = Some(state_str.to_string());
1095 let state_msg = serde_json::json!({
1096 "type": "StateChanged",
1097 "state": state_str
1098 });
1099 let _ = ws_tx.send(state_msg.to_string());
1100 }
1101 }
1102
1103 let json =
1104 serde_json::to_string(&SerializableMessage(&msg)).unwrap_or_default();
1105 let _ = ws_tx.send(json);
1106 }
1107 tracing::debug!("Drain task ended for session: {}", session_id);
1108 });
1109 }
1110
1111 pub async fn resolve_user_key(&self, headers: &axum::http::HeaderMap) -> String {
1113 let Some(ref auth) = self.auth else {
1114 return "default".to_string();
1115 };
1116
1117 if let Some(token) = headers
1119 .get(axum::http::header::AUTHORIZATION)
1120 .and_then(|v| v.to_str().ok())
1121 .and_then(|v| v.strip_prefix("Bearer "))
1122 .map(|s| s.to_string())
1123 {
1124 if token.starts_with("dev:") {
1125 let parts: Vec<&str> = token.splitn(4, ':').collect();
1126 if parts.len() >= 4 {
1127 return format!("dev:{}", parts[1]);
1128 }
1129 }
1130 if let Ok(claims) = auth.validate_token(&token).await {
1131 return claims.sub;
1132 }
1133 }
1134
1135 let cookie_session_id = headers
1137 .get(axum::http::header::COOKIE)
1138 .and_then(|v| v.to_str().ok())
1139 .and_then(|cookies| {
1140 cookies
1141 .split(';')
1142 .map(|c| c.trim())
1143 .find_map(|c| {
1144 c.strip_prefix(&format!("{}=", auth.config.cookie_name))
1145 .map(|s| s.to_string())
1146 })
1147 });
1148
1149 if let Some(session_id) = cookie_session_id {
1150 if session_id.starts_with("dev:") {
1151 let parts: Vec<&str> = session_id.splitn(4, ':').collect();
1152 if parts.len() >= 4 {
1153 return format!("dev:{}", parts[1]);
1154 }
1155 }
1156
1157 if let Ok(access_token) = auth.session_manager.get_token(&session_id).await {
1158 if let Ok(claims) = auth.validate_token(&access_token).await {
1159 return claims.sub;
1160 }
1161 }
1162 }
1163
1164 "default".to_string()
1165 }
1166}
1167
1168struct SerializableMessage<'a>(&'a TuiMessage);
1174
1175impl<'a> serde::Serialize for SerializableMessage<'a> {
1176 fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
1177 where
1178 S: serde::Serializer,
1179 {
1180 use serde::ser::SerializeStruct;
1181
1182 match self.0 {
1183 TuiMessage::OutputLine(line) => {
1184 let mut s = serializer.serialize_struct("msg", 2)?;
1185 s.serialize_field("type", "OutputLine")?;
1186 s.serialize_field("line", line)?;
1187 s.end()
1188 }
1189 TuiMessage::StreamDelta(delta) => {
1190 let mut s = serializer.serialize_struct("msg", 2)?;
1191 s.serialize_field("type", "StreamDelta")?;
1192 s.serialize_field("delta", delta)?;
1193 s.end()
1194 }
1195 TuiMessage::ReasoningDelta(delta) => {
1196 let mut s = serializer.serialize_struct("msg", 2)?;
1197 s.serialize_field("type", "ReasoningDelta")?;
1198 s.serialize_field("delta", delta)?;
1199 s.end()
1200 }
1201 TuiMessage::WorkflowCompleted => {
1202 let mut s = serializer.serialize_struct("msg", 2)?;
1203 s.serialize_field("type", "WorkflowCompleted")?;
1204 s.serialize_field("state", "Idle")?;
1205 s.end()
1206 }
1207 TuiMessage::WorkflowError(err) => {
1208 let mut s = serializer.serialize_struct("msg", 2)?;
1209 s.serialize_field("type", "WorkflowError")?;
1210 s.serialize_field("error", err)?;
1211 s.end()
1212 }
1213 TuiMessage::ResumeInfo(info) => match info {
1214 Some(ri) => {
1215 let mut s = serializer.serialize_struct("msg", 5)?;
1216 s.serialize_field("type", "ResumeInfo")?;
1217 s.serialize_field("state", "Idle")?;
1218 s.serialize_field("session_id", &ri.session_id)?;
1219 s.serialize_field("checkpoint_id", &ri.checkpoint_id)?;
1220 s.serialize_field("iteration", &ri.iteration)?;
1221 s.end()
1222 }
1223 None => {
1224 let mut s = serializer.serialize_struct("msg", 2)?;
1225 s.serialize_field("type", "ResumeInfo")?;
1226 s.serialize_field("state", "Idle")?;
1227 s.end()
1228 }
1229 },
1230 TuiMessage::TodoUpdate(content) => {
1231 let mut s = serializer.serialize_struct("msg", 2)?;
1232 s.serialize_field("type", "TodoUpdate")?;
1233 s.serialize_field("content", content)?;
1234 s.end()
1235 }
1236 TuiMessage::WorkflowCancelled => {
1237 let mut s = serializer.serialize_struct("msg", 2)?;
1238 s.serialize_field("type", "WorkflowCancelled")?;
1239 s.serialize_field("state", "Idle")?;
1240 s.end()
1241 }
1242 TuiMessage::HandoffReady(briefing) => {
1243 let mut s = serializer.serialize_struct("msg", 3)?;
1244 s.serialize_field("type", "HandoffReady")?;
1245 s.serialize_field("state", "Idle")?;
1246 s.serialize_field("briefing", briefing)?;
1247 s.end()
1248 }
1249 TuiMessage::SessionRotated { old, new } => {
1250 let mut s = serializer.serialize_struct("msg", 3)?;
1251 s.serialize_field("type", "SessionRotated")?;
1252 s.serialize_field("old", old)?;
1253 s.serialize_field("new", new)?;
1254 s.end()
1255 }
1256 TuiMessage::HandoffFailed => {
1257 let mut s = serializer.serialize_struct("msg", 2)?;
1258 s.serialize_field("type", "HandoffFailed")?;
1259 s.serialize_field("state", "Idle")?;
1260 s.end()
1261 }
1262 TuiMessage::ToolPending {
1263 tool_name,
1264 hint,
1265 } => {
1266 let mut s = serializer.serialize_struct("msg", 3)?;
1267 s.serialize_field("type", "ToolPending")?;
1268 s.serialize_field("tool_name", tool_name)?;
1269 s.serialize_field("hint", hint)?;
1270 s.end()
1271 }
1272 TuiMessage::ToolDone {
1273 tool_name,
1274 success,
1275 hint,
1276 } => {
1277 let mut s = serializer.serialize_struct("msg", 4)?;
1278 s.serialize_field("type", "ToolDone")?;
1279 s.serialize_field("tool_name", tool_name)?;
1280 s.serialize_field("success", success)?;
1281 s.serialize_field("hint", hint)?;
1282 s.end()
1283 }
1284 TuiMessage::ContextTokensUpdated(count) => {
1285 let mut s = serializer.serialize_struct("msg", 2)?;
1286 s.serialize_field("type", "ContextTokensUpdated")?;
1287 s.serialize_field("count", count)?;
1288 s.end()
1289 }
1290 TuiMessage::McpServerStatus {
1291 name,
1292 connected,
1293 tool_count,
1294 error,
1295 } => {
1296 let mut s = serializer.serialize_struct("msg", 5)?;
1297 s.serialize_field("type", "McpServerStatus")?;
1298 s.serialize_field("name", name)?;
1299 s.serialize_field("connected", connected)?;
1300 s.serialize_field("tool_count", tool_count)?;
1301 s.serialize_field("error", error)?;
1302 s.end()
1303 }
1304 TuiMessage::SessionTitleUpdated(title) => {
1305 let mut s = serializer.serialize_struct("msg", 2)?;
1306 s.serialize_field("type", "SessionTitleUpdated")?;
1307 s.serialize_field("title", title)?;
1308 s.end()
1309 }
1310 }
1311 }
1312}
1313
1314fn deep_merge_toml(base: &mut toml::Value, overlay: &toml::Value) {
1325 match (base, overlay) {
1326 (toml::Value::Table(base_table), toml::Value::Table(overlay_table)) => {
1327 for (key, overlay_val) in overlay_table {
1328 match base_table.get_mut(key) {
1329 Some(base_val) => {
1330 deep_merge_toml(base_val, overlay_val);
1332 }
1333 None => {
1334 base_table.insert(key.clone(), overlay_val.clone());
1336 }
1337 }
1338 }
1339 }
1340 (base, overlay) => {
1342 *base = overlay.clone();
1343 }
1344 }
1345}
1346
1347fn substitute_env_vars(s: &mut String, secrets: &std::collections::HashMap<String, String>) {
1352 let mut result = String::with_capacity(s.len());
1354 let bytes = s.as_bytes();
1355 let mut i = 0;
1356
1357 while i < bytes.len() {
1358 if i + 1 < bytes.len() && bytes[i] == b'$' && bytes[i + 1] == b'{' {
1359 if let Some(end) = s[i + 2..].find('}') {
1361 let var_name = &s[i + 2..i + 2 + end];
1362 if let Some(value) = secrets.get(var_name) {
1364 result.push_str(value);
1365 } else if let Ok(value) = std::env::var(var_name) {
1366 result.push_str(&value);
1367 } else {
1368 result.push_str(&s[i..i + 2 + end + 1]);
1370 }
1371 i = i + 2 + end + 1;
1372 } else {
1373 result.push('$');
1375 i += 1;
1376 }
1377 } else {
1378 result.push(bytes[i] as char);
1379 i += 1;
1380 }
1381 }
1382
1383 *s = result;
1384}
1385
1386fn fingerprint_mcp_section(resolved: Option<&str>) -> u64 {
1396 use sha2::{Digest, Sha256};
1397 let section = resolved
1398 .and_then(|s| s.parse::<toml::Value>().ok())
1399 .and_then(|v| v.get("mcp").cloned());
1400 let bytes = match section {
1401 Some(v) => v.to_string().into_bytes(),
1402 None => b"<no-mcp>".to_vec(),
1403 };
1404 let digest = Sha256::digest(&bytes);
1405 u64::from_be_bytes(digest[..8].try_into().expect("sha256 digest >= 8 bytes"))
1406}
1407
1408#[cfg(test)]
1409mod tests {
1410 use super::*;
1411
1412 fn test_state() -> ServerState {
1414 let (session, _rx) = Session::new();
1415 let (ws_tx, _ws_rx) = tokio::sync::broadcast::channel::<String>(16);
1416 ServerState::new(session, ws_tx, None)
1417 }
1418
1419 fn temp_user_home(tag: &str) -> std::path::PathBuf {
1421 static COUNTER: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0);
1422 let n = COUNTER.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
1423 let dir = std::env::temp_dir().join(format!(
1424 "trustee-state-test-{}-{}-{}",
1425 tag,
1426 std::process::id(),
1427 n
1428 ));
1429 std::fs::create_dir_all(dir.join("config")).expect("create temp user home");
1430 dir
1431 }
1432
1433 fn parse(toml_str: &str) -> toml::Value {
1434 toml_str.parse::<toml::Value>().expect("valid test TOML")
1435 }
1436
1437 #[test]
1440 fn overlay_allowlist_drops_non_allowlisted_sections() {
1441 let state = test_state();
1442 let home = temp_user_home("allowlist");
1443 std::fs::write(
1444 home.join("config").join("trustee.toml"),
1445 "[server]\nport = 1\n\n[auth]\nmode = \"kanidm\"\n\n[mcp]\nmode = \"user\"\n",
1446 )
1447 .expect("write overlay");
1448
1449 let shared = "[server]\nport = 8080\n\n[mcp]\nmode = \"shared\"\n";
1450 let merged = state
1451 .merge_user_config(&home, Some(shared))
1452 .expect("overlay has allowlisted content");
1453
1454 let merged_val = parse(&merged);
1455 let expected = parse("[server]\nport = 8080\n\n[mcp]\nmode = \"user\"\n");
1456 assert_eq!(merged_val, expected, "merged config must be shared + [mcp] overlay only");
1457 assert!(merged_val.get("auth").is_none(), "overlay [auth] must be dropped");
1458 }
1459
1460 #[test]
1462 fn no_overlay_file_returns_none() {
1463 let state = test_state();
1464 let home = temp_user_home("empty");
1465 assert!(state
1466 .merge_user_config(&home, Some("[server]\nport = 8080\n"))
1467 .is_none());
1468 }
1469
1470 #[test]
1472 fn overlay_with_no_allowlisted_sections_is_noop() {
1473 let state = test_state();
1474 let home = temp_user_home("all-dropped");
1475 std::fs::write(
1476 home.join("config").join("trustee.toml"),
1477 "[server]\nport = 1\n\n[storage]\npath = \"/tmp/x\"\n",
1478 )
1479 .expect("write overlay");
1480 assert!(state
1481 .merge_user_config(&home, Some("[server]\nport = 8080\n"))
1482 .is_none());
1483 }
1484
1485 #[test]
1488 fn llm_overlay_dropped_by_default_and_kept_when_enabled() {
1489 let shared = "[llm]\nprovider = \"openai\"\n\n[mcp]\nmode = \"shared\"\n";
1490 let overlay = "[llm]\nprovider = \"anthropic\"\n";
1491
1492 let state = test_state();
1494 assert!(!state.allow_llm_overlay);
1495 let home = temp_user_home("llm-off");
1496 std::fs::write(home.join("config").join("trustee.toml"), overlay).expect("write overlay");
1497 assert!(state.merge_user_config(&home, Some(shared)).is_none());
1498
1499 let state = test_state().with_allow_llm_overlay(true);
1501 assert!(state.allow_llm_overlay);
1502 let home = temp_user_home("llm-on");
1503 std::fs::write(home.join("config").join("trustee.toml"), overlay).expect("write overlay");
1504 let merged = state
1505 .merge_user_config(&home, Some(shared))
1506 .expect("[llm] overlay applies when opted in");
1507 let merged_val = parse(&merged);
1508 assert_eq!(merged_val["llm"]["provider"].as_str(), Some("anthropic"));
1509 assert_eq!(merged_val["mcp"]["mode"].as_str(), Some("shared"));
1510 }
1511
1512 #[test]
1516 fn user_home_dir_uses_consolidated_hash() {
1517 let state = test_state();
1518 if let Some(home) = state.get_user_home_dir("farzan@example.com") {
1519 assert_eq!(
1520 home.file_name().and_then(|n| n.to_str()),
1521 Some(trustee_core::user_hash("farzan@example.com")).as_deref()
1522 );
1523 let users_root = dirs::home_dir().unwrap().join(".trustee").join("users");
1524 assert_eq!(home.parent(), Some(&users_root).map(|p| p.as_path()));
1525 }
1526 }
1528
1529 #[tokio::test]
1534 async fn agent_principals_get_isolated_session_buckets() {
1535 let state = test_state();
1536 let key_a = "agent-farzan";
1537 let key_b = "agent-paydar";
1538
1539 let (sid_a, session_a, _tx_a, _ts_a) = state.ensure_active_session(key_a).await;
1540 let (_sid_b, _session_b, _tx_b, _ts_b) = state.ensure_active_session(key_b).await;
1541
1542 assert!(
1544 state.get_session_by_any_id(key_a, &sid_a).await.is_some(),
1545 "owner bucket resolves its own session"
1546 );
1547 assert!(
1548 state.get_session_by_any_id(key_b, &sid_a).await.is_none(),
1549 "cross-agent session access must be 404/None"
1550 );
1551 assert!(Arc::strong_count(&session_a) >= 1);
1552
1553 if let Some(home) = state.get_user_home_dir(key_a) {
1555 assert_eq!(
1556 home.file_name().and_then(|n| n.to_str()),
1557 Some(trustee_core::user_hash(key_a)).as_deref()
1558 );
1559 }
1560 }
1561}
1562
1563#[cfg(test)]
1568mod mcp_loader_cache_tests {
1569 use super::*;
1570
1571 fn state_with_shared(shared: &str) -> ServerState {
1572 let (session, _rx) = Session::new();
1573 let (ws_tx, _ws_rx) = tokio::sync::broadcast::channel::<String>(16);
1574 let mut state = ServerState::new(session, ws_tx, None);
1575 state.config_toml = Some(shared.to_string());
1576 state
1577 }
1578
1579 struct TempUser {
1583 key: String,
1584 home: std::path::PathBuf,
1585 }
1586
1587 impl TempUser {
1588 fn new(tag: &str) -> Self {
1589 let key = format!("16c-{tag}-{}@test.invalid", std::process::id());
1590 let home = dirs::home_dir()
1591 .expect("HOME available in test env")
1592 .join(".trustee")
1593 .join("users")
1594 .join(trustee_core::user_hash(&key));
1595 std::fs::create_dir_all(home.join("config")).expect("create user home");
1596 Self { key, home }
1597 }
1598
1599 fn write_overlay(&self, toml_str: &str) {
1600 std::fs::write(self.home.join("config").join("trustee.toml"), toml_str)
1601 .expect("write overlay");
1602 }
1603 }
1604
1605 impl Drop for TempUser {
1606 fn drop(&mut self) {
1607 let _ = std::fs::remove_dir_all(&self.home);
1608 }
1609 }
1610
1611 #[tokio::test]
1612 async fn no_mcp_config_caches_disabled_marker() {
1613 let state = state_with_shared("[server]\nport = 8080\n");
1614 let user = TempUser::new("nomcp");
1615 let ts = Arc::new(pep::MemoryTokenStore::new());
1616
1617 let first = state.get_or_build_mcp_loader(&user.key, &ts).await.unwrap();
1618 assert!(first.is_none(), "no [mcp] anywhere → disabled marker");
1619 assert_eq!(state.mcp_loaders.len(), 1, "exactly one cache entry");
1620
1621 let second = state.get_or_build_mcp_loader(&user.key, &ts).await.unwrap();
1623 assert!(second.is_none());
1624 assert_eq!(state.mcp_loaders.len(), 1);
1625 }
1626
1627 #[tokio::test]
1628 async fn concurrent_cold_builds_single_flight() {
1629 let state = Arc::new(state_with_shared("[server]\nport = 8080\n"));
1630 let user = TempUser::new("singleflight");
1631 let ts = Arc::new(pep::MemoryTokenStore::new());
1632
1633 let mut handles = Vec::new();
1634 for _ in 0..5 {
1635 let state = state.clone();
1636 let key = user.key.clone();
1637 let ts = ts.clone();
1638 handles.push(tokio::spawn(async move {
1639 state.get_or_build_mcp_loader(&key, &ts).await
1640 }));
1641 }
1642 for h in handles {
1643 h.await.unwrap().expect("all five succeed");
1644 }
1645 assert_eq!(state.mcp_loaders.len(), 1, "single-flight → one entry");
1646 }
1647
1648 #[tokio::test]
1649 async fn fingerprint_change_triggers_rebuild() {
1650 let state = state_with_shared("[server]\nport = 8080\n");
1651 let user = TempUser::new("fpchange");
1652 let ts = Arc::new(pep::MemoryTokenStore::new());
1653
1654 user.write_overlay("[mcp]\nenabled = true\n\n[[mcp.servers]]\nname = \"v1\"\nurl = \"http://127.0.0.1:9/sse\"\n");
1657 let v1 = state.get_or_build_mcp_loader(&user.key, &ts).await.unwrap();
1658 assert!(v1.is_some(), "enabled [mcp] → real loader");
1659 let fp1 = state
1660 .mcp_loaders
1661 .get(&trustee_core::user_hash(&user.key))
1662 .unwrap()
1663 .fingerprint;
1664
1665 std::thread::sleep(std::time::Duration::from_millis(5));
1667 user.write_overlay("[mcp]\nenabled = true\n\n[[mcp.servers]]\nname = \"v2\"\nurl = \"http://127.0.0.1:9/other\"\n");
1668 let v2 = state.get_or_build_mcp_loader(&user.key, &ts).await.unwrap();
1669 assert!(v2.is_some());
1670 let entry = state
1671 .mcp_loaders
1672 .get(&trustee_core::user_hash(&user.key))
1673 .unwrap();
1674 assert_ne!(
1675 entry.fingerprint, fp1,
1676 "fingerprint must change with content"
1677 );
1678 assert!(entry.degraded.is_none());
1679
1680 let _still_usable = v1.as_ref().unwrap().tool_count;
1682 }
1683
1684 #[tokio::test]
1685 async fn degraded_entry_fails_loud_within_backoff_and_isolates_users() {
1686 let state = state_with_shared("[server]\nport = 8080\n");
1687 let bad = TempUser::new("degraded-bad");
1688 let good = TempUser::new("degraded-good");
1689 let ts = Arc::new(pep::MemoryTokenStore::new());
1690
1691 bad.write_overlay("[mcp]\nenabled = \"not-a-bool\"\n");
1693 let err = match state.get_or_build_mcp_loader(&bad.key, &ts).await {
1694 Ok(_) => panic!("invalid [mcp] must fail loud"),
1695 Err(e) => e,
1696 };
1697 assert!(
1698 err.contains("invalid [mcp]"),
1699 "surfaces the parse error: {err}"
1700 );
1701
1702 let entry = state
1703 .mcp_loaders
1704 .get(&trustee_core::user_hash(&bad.key))
1705 .unwrap();
1706 assert!(entry.degraded.is_some(), "poison entry recorded");
1707
1708 let err2 = match state.get_or_build_mcp_loader(&bad.key, &ts).await {
1710 Ok(_) => panic!("still within backoff"),
1711 Err(e) => e,
1712 };
1713 assert_eq!(err, err2, "same cached error");
1714
1715 let good_loader = state
1717 .get_or_build_mcp_loader(&good.key, &ts)
1718 .await
1719 .expect("other user unaffected");
1720 assert!(
1721 good_loader.is_none(),
1722 "good user has no [mcp] → disabled marker"
1723 );
1724 }
1725}