Skip to main content

trustee_api/
state.rs

1//! Shared server state: per-user multi-session registry, broadcast channels, and auth state.
2//!
3//! ## Multi-Session Per User (MSU)
4//!
5//! Each authenticated user gets their own [`UserSessions`] containing N independent
6//! [`UserSessionEntry`] instances (default max 4). Each entry has:
7//! - An independent `Session` (workflow state, output, etc.)
8//! - A dedicated broadcast channel for WebSocket fan-out
9//! - Creation and last-active timestamps
10//!
11//! Sessions are keyed by user identity (`sub` claim from JWT, or `dev:email` for
12//! dev mode). Unauthenticated deployments use a single `"default"` key, preserving
13//! backward compatibility with single-user CLI operation.
14
15use 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
24// ---------------------------------------------------------------------------
25// Multi-session types
26// ---------------------------------------------------------------------------
27
28/// A single session with its own broadcast channel.
29pub struct UserSessionEntry {
30    /// The agent session, protected by a mutex.
31    pub session: Arc<Mutex<Session>>,
32    /// Broadcast sender for this session's WebSocket fan-out.
33    pub ws_tx: broadcast::Sender<String>,
34    /// When this session was created.
35    pub created_at: chrono::DateTime<chrono::Utc>,
36    /// Last time a command was submitted or state changed.
37    /// Updated on every /sessions/{id}/command and /sessions/{id}/cancel call.
38    pub last_active: Arc<Mutex<chrono::DateTime<chrono::Utc>>>,
39}
40
41/// All sessions belonging to one authenticated user.
42pub struct UserSessions {
43    /// session_id → session entry
44    pub sessions: DashMap<String, UserSessionEntry>,
45    /// Shared token store for all this user's sessions (MCP credential isolation).
46    pub token_store: Arc<pep::MemoryTokenStore>,
47    /// Which session_id is "active" for legacy /session/* routes.
48    pub active_session_id: Mutex<String>,
49}
50
51/// Summary of an active session for listing (serializable for API responses).
52#[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    /// Number of handoff rotations this session has performed (Bug 5). The
60    /// in-memory registry key never changes across a rotation — only the
61    /// checkpoint-chain name does — so a card whose name changed without this
62    /// hint looks like it silently renamed. 0 = never rotated.
63    pub handoff_count: u32,
64}
65
66/// Errors from multi-session operations.
67#[derive(Debug)]
68pub enum SessionError {
69    /// User has reached max_sessions_per_user limit.
70    MaxSessionsReached(usize),
71    /// Session ID not found for this user.
72    NotFound(String),
73    /// Session is not Idle (cannot destroy/overwrite a running session).
74    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
91/// Top-level registry: user_key → user's session collection.
92pub type SessionRegistry = Arc<DashMap<String, UserSessions>>;
93
94// ---------------------------------------------------------------------------
95// ServerState
96// ---------------------------------------------------------------------------
97
98/// Cached per-user MCP tool loader (16C).
99///
100/// One entry per user hash. `loader: None` means the user's effective
101/// config has MCP disabled — agents run MCP-less via abk's
102/// `McpSource::Prebuilt(None)` (semantically identical to
103/// `[mcp] enabled = false` today, but with zero config re-parsing).
104pub struct McpLoaderEntry {
105    /// Built loader; `None` = MCP disabled for this user.
106    pub loader: Option<std::sync::Arc<abk::agent::McpToolLoader>>,
107    /// Content fingerprint of the effective `[mcp]` config (SHA-256, first
108    /// 8 bytes as u64). Content, never mtime — overlays are rewritten in place.
109    pub fingerprint: u64,
110    pub built_at: chrono::DateTime<chrono::Utc>,
111    /// Set when the last build failed; the message surfaces to the user's
112    /// next dispatch (fail loud). Other users are never affected.
113    pub degraded: Option<String>,
114    /// When `degraded` was set. Rebuilds are held back for
115    /// [`MCP_BUILD_RETRY_BACKOFF`] — a poison entry never sticks forever.
116    pub failed_at: Option<chrono::DateTime<chrono::Utc>>,
117}
118
119/// Minimum delay before retrying a failed MCP loader build (16C).
120pub const MCP_BUILD_RETRY_BACKOFF: std::time::Duration = std::time::Duration::from_secs(30);
121
122/// 16F: per-agent dispatch target — THQ agent_name → the agent-user that
123/// THQ-dispatched sessions must run AS.
124///
125/// Populated once at boot by the 16E discovery scan
126/// (`thq_register::spawn_all`); restart-only lifecycle, exactly like the
127/// THQ registration itself.
128#[derive(Debug, Clone)]
129pub struct ThqDispatchEntry {
130    /// The agent-user's stable key = its Kanidm `sub` (16E sub-pin). This is
131    /// the identity impersonated on the inner dispatch: session bucket,
132    /// per-user home, MCP loader cache, and Cedar principal all resolve
133    /// through it.
134    pub user_key: String,
135    /// The agent's Kanidm service token from its per-user `.env` — exchanged
136    /// for a short-lived `role=agent` Bearer at dispatch time. `None` = the
137    /// agent has no provisioned credential (dispatch fails loud, 502).
138    pub service_token: Option<String>,
139    /// The issuer the token's own credential declares (`[mcp.credentials.*]
140    /// .issuer_url` from the agent's overlay). Kanidm accepts a token
141    /// exchange only on the origin the token was minted for — the issuer
142    /// MUST travel with the credential that owns the token (verified live
143    /// 2026-08-30: same token, 200 on its own vhost, 400 elsewhere).
144    pub issuer_url: Option<String>,
145}
146
147/// Shared state accessible by all axum handlers.
148#[derive(Clone)]
149pub struct ServerState {
150    /// Per-user multi-session registry (MSU).
151    pub sessions: SessionRegistry,
152    /// Broadcast sender for backward compat — delegates to the default user's channel.
153    pub ws_tx: broadcast::Sender<String>,
154    /// Auth state (None = auth disabled, all endpoints open).
155    pub auth: Option<Arc<AuthState>>,
156    /// Shared config TOML (all users share the same agent config).
157    pub config_toml: Option<String>,
158    /// Shared secrets (injected into every per-user session).
159    pub secrets: Option<std::collections::HashMap<String, String>>,
160    /// Shared build info (injected into every per-user session).
161    pub build_info: Option<trustee_core::types::BuildInfo>,
162    /// Global concurrency limiter — limits the number of simultaneous workflows
163    /// across all users. Default: 8 concurrent workflows.
164    pub workflow_semaphore: Arc<tokio::sync::Semaphore>,
165    /// Maximum number of concurrent sessions per user. Default: 4.
166    pub max_sessions_per_user: usize,
167    /// Whether per-user config overlays may override the `[llm]` section.
168    /// Default: `false` — overlays are limited to `[mcp]`.
169    pub allow_llm_overlay: bool,
170    /// Per-user McpToolLoader cache (16C), keyed by user hash (never the raw key).
171    pub mcp_loaders: Arc<DashMap<String, McpLoaderEntry>>,
172    /// Single-flight build locks per user hash (16C).
173    mcp_build_locks: Arc<DashMap<String, Arc<tokio::sync::Mutex<()>>>>,
174    /// 16F: THQ per-agent dispatch table (agent_name → target), boot-populated.
175    pub thq_dispatch: Arc<DashMap<String, ThqDispatchEntry>>,
176    /// 16F: cached impersonation Bearers per agent user_key
177    /// `(token, expires_at)` — expiry-buffered, re-minted on 401.
178    pub agent_dispatch_tokens: Arc<DashMap<String, (String, std::time::Instant)>>,
179}
180
181impl ServerState {
182    /// Create new shared state from a default session, broadcast sender, and optional auth.
183    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        // Store the default user's UserSessions with an initial session
191        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    /// 16F: the issuer URL Kanidm accepts SERVICE-ACCOUNT token exchange on.
234    ///
235    /// Kanidm binds exchange to the OAuth2 client's configured origin: the
236    /// `[oidc]` issuer host serves logins/validation fine but REJECTS token
237    /// exchange with `invalid_request`, while the `idp` vhost — the issuer
238    /// every working `[mcp.credentials.*]` service-account entry already
239    /// uses — accepts it (verified live 2026-08-30: same token, 200 vs 400).
240    /// The exchange therefore takes its issuer from the shared config's
241    /// first `service-account` credential; callers fall back to the auth
242    /// issuer when absent.
243    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    /// Set the max sessions per user.
273    pub fn with_max_sessions_per_user(mut self, max: usize) -> Self {
274        self.max_sessions_per_user = max;
275        self
276    }
277
278    /// Allow per-user config overlays to override the `[llm]` section.
279    /// Default: `false` (overlays limited to `[mcp]`).
280    pub fn with_allow_llm_overlay(mut self, allow: bool) -> Self {
281        self.allow_llm_overlay = allow;
282        self
283    }
284
285    // -----------------------------------------------------------------------
286    // MSU: Multi-session methods
287    // -----------------------------------------------------------------------
288
289    /// Create a new session for a user. Returns the session_id.
290    ///
291    /// Creates a fresh `Session::new()`, copies shared config, sets per-user
292    /// isolation, creates a broadcast channel, spawns a drain task, and inserts
293    /// into the user's session DashMap. The new session becomes the "active" one.
294    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        // Get or create the user's UserSessions entry
302        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        // Check session limit
312        if user_sessions.sessions.len() >= self.max_sessions_per_user {
313            return Err(SessionError::MaxSessionsReached(self.max_sessions_per_user));
314        }
315
316        // Create new Session
317        let (mut session, workflow_rx) = Session::new();
318
319        // Copy shared config
320        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        // Per-user isolation
338        self.apply_user_isolation(&mut session, user_key);
339
340        // Apply session_name if provided
341        session.session_name = session_name;
342
343        // Apply agent identity if provided
344        session.identity = identity;
345
346        // Create broadcast channel
347        let (ws_tx_entry, _) = broadcast::channel::<String>(256);
348
349        // Generate session_id
350        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        // Insert into user's sessions DashMap
359        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        // Set as active session only when the caller opts in.
370        //
371        // `resume_session` passes `activate: false` so resuming an arbitrary
372        // checkpoint never hijacks the caller's current active live session
373        // pointer. `new_session` / `create_session` (fresh-start paths) pass
374        // `activate: true` to keep the existing "new session becomes active"
375        // behavior. This keeps the per-user active pointer under UI control
376        // rather than letting any authenticated client overwrite it.
377        if activate {
378            *user_sessions.active_session_id.lock().await = session_id.clone();
379        }
380
381        // Spawn drain task
382        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    /// Get a specific session by user_key + session_id.
399    /// Updates last_active on the session entry.
400    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        // Update last_active
409        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    /// Get a session by EITHER its live MSU registry key OR its
416    /// checkpoint/session identity (`session.session_id`).
417    ///
418    /// The web frontend tracks `currentSessionId` from the `ResumeInfo` WS
419    /// message, which carries the auto-derived checkpoint id
420    /// (`session_YYYY_MM_DD_HH_MM_uuid8`) — NOT the live MSU registry key
421    /// (`"default"` or the key from `create_session()`). External clients
422    /// like Torpi/THQ pass the live registry key. This resolver accepts both:
423    ///
424    /// 1. Try registry-key lookup first (precise, used by Torpi/THQ).
425    /// 2. Fall back to scanning the user's live sessions for one whose
426    ///    `session.session_id` matches the requested id (used by the
427    ///    embedded web UI after a command or resume).
428    ///
429    /// Returns `(live_registry_key, session_arc, ws_tx)`, or `None` if not
430    /// found. The live key is returned so callers that need to set it as
431    /// active (or otherwise reference the registry) use the real key.
432    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        // Fast path: registry key match.
438        let user_sessions = self.sessions.get(user_key)?;
439        if let Some(entry) = user_sessions.sessions.get(id) {
440            // Update last_active
441            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        // Slow path: scan live sessions for a matching session.session_id.
447        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                // Update last_active
454                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    /// List all active sessions for a user, sorted by last_active desc.
464    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        // Sort by last_active descending
490        items.sort_by(|a, b| b.last_active.cmp(&a.last_active));
491        items
492    }
493
494    /// Destroy a session. The session must be Idle.
495    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        // Check workflow state before removing
506        {
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        // Remove from DashMap
523        user_sessions.sessions.remove(session_id);
524
525        // If this was the active session, pick a new active
526        let mut active_id = user_sessions.active_session_id.lock().await;
527        if &*active_id == session_id {
528            // Pick the most recently active remaining session
529            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    /// Get or create the user's "active" session for legacy routes.
543    ///
544    /// Behavior:
545    /// 1. If user has no sessions → create one
546    /// 2. If active session exists → return it
547    /// 3. If active session was destroyed → create a new one
548    ///
549    /// Returns: (session_id, session_arc, ws_tx, token_store)
550    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        // Get or create user's UserSessions
560        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        // Check if active session exists
573        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            // Active session was destroyed, fall through to create
584        }
585
586        // Need to create a new session
587        // For the "default" user, we may already have a "default" session entry
588        // from ServerState::new() — check for it
589        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        // Create a brand new session
613        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    /// DEPRECATED: Use ensure_active_session() instead.
627    /// Kept for backward compatibility — same 3-tuple return type.
628    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    /// Set a session as the user's active session.
637    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    // -----------------------------------------------------------------------
646    // Read-only helpers (no session creation side effects)
647    // -----------------------------------------------------------------------
648
649    /// Resolve a user's home_dir without creating an in-memory session.
650    ///
651    /// This is the read-only equivalent of the isolation logic in
652    /// `apply_user_isolation`. Used by endpoints that only need to read
653    /// checkpoint data from disk (history, session list, session detail)
654    /// and must NOT create ghost sessions as a side effect.
655    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    /// Resolve config_toml and home_dir without creating an in-memory session.
661    ///
662    /// Returns `(config_toml, home_dir)`. If config is not loaded,
663    /// config_toml will be None.
664    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    // -----------------------------------------------------------------------
669    // Private helpers
670    // -----------------------------------------------------------------------
671
672    /// Apply per-user isolation: SHA-256 hash → home_dir + project_id.
673    ///
674    /// Hashing goes through the single consolidated [`trustee_core::user_hash`]
675    /// so the web path and the CLI path can never drift apart.
676    fn apply_user_isolation(&self, session: &mut Session, user_key: &str) {
677        let user_hash = trustee_core::user_hash(user_key);
678
679        // Set per-user home directory for checkpoint isolation
680        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        // ── Per-user .env (Task 2) ──────────────────────────────────────
691        //
692        // Load per-user secrets from ~/.trustee/users/{hash}/.env
693        // These are merged on top of shared secrets (per-user wins).
694        // They are NEVER set as process env vars — used only for ${VAR}
695        // substitution in the config TOML below.
696        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        // ── Per-user config overlay (Task 3) ────────────────────────────
706        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        // ── ${VAR} substitution (Task 4) ────────────────────────────────
714        if let Some(ref mut config_toml) = session.config_toml {
715            substitute_env_vars(config_toml, &merged_secrets);
716        }
717
718        // ── Strip per-user secrets (Task 5) ────────────────────────────
719        session.secrets = Some(shared_secrets);
720    }
721
722    /// Load secrets from a user's ~/.trustee/users/{hash}/.env and merge
723    /// on top of `base` (user wins). Returns the merged map.
724    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    /// Deep-merge a user's per-user config overlay on top of the shared
755    /// config. Returns the merged TOML string, or None if no overlay exists.
756    ///
757    /// Overlay allowlist (task 16B): only allowlisted top-level sections of
758    /// the user overlay are merged; anything else is dropped loudly.
759    /// - `[mcp]` is always allowed — per-user MCP tool sets are the point of
760    ///   the overlay convention.
761    /// - `[llm]` is allowed only when the instance opted in via the
762    ///   `[users].allow_llm_overlay` knob (default **false**).
763    ///
764    /// Everything else (`[server]`, `[auth]`, `[storage]`, `[web]`, …) is
765    /// boot-time, instance-level config and must not be rewritable through a
766    /// per-user file. This is predictability hardening, not a security
767    /// boundary: those sections were never re-read from session config.
768    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                // 16E: [thq] is per-user identity config (read from the overlay
787                // file at boot); keeping it in the merged session config is
788                // inert but avoids a false "dropping" warn on every dispatch.
789                || section == "thq"
790                || (self.allow_llm_overlay && section == "llm")
791        };
792        // Mask the user in logs: user_home's dir name IS the hash — never
793        // log the raw key (keys may be emails). Log a short hash prefix.
794        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            // Nothing survived the allowlist — no-op, keep shared config as-is.
817            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    /// Resolve the fully-merged config TOML for a user WITHOUT creating a session.
829    ///
830    /// This is the read-only equivalent of the config resolution in
831    /// `apply_user_isolation`: shared config + per-user overlay + ${VAR}
832    /// substitution. Used by endpoints that need to inspect config (e.g.
833    /// listing available LLM models) without spawning a ghost session.
834    ///
835    /// Returns None if no shared config is loaded.
836    pub fn resolve_user_config(&self, user_key: &str) -> Option<String> {
837        let config_toml = self.config_toml.clone()?;
838
839        // Resolve user home dir (same hash scheme as apply_user_isolation)
840        let user_home = self.get_user_home_dir(user_key)?;
841
842        // Start from shared secrets; merge per-user .env on top
843        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        // Merge per-user config overlay
849        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 ${VAR} from merged secrets
855        substitute_env_vars(&mut resolved, &merged_secrets);
856
857        Some(resolved)
858    }
859
860    /// Get or build the per-user MCP tool loader (16C).
861    ///
862    /// Cache semantics:
863    /// - fingerprint = content hash of the effective `[mcp]` section
864    ///   (shared + allowlist-filtered overlay + ${VAR} substitution, i.e.
865    ///   exactly what `resolve_user_config` produces) — stale on ANY change.
866    /// - hit + match → `Arc` clone, zero network I/O.
867    /// - miss/stale → single-flight build (one build per user at a time;
868    ///   late arrivals re-check and reuse the winner's entry).
869    /// - `Ok(None)` = MCP disabled for this user → agent runs MCP-less via
870    ///   abk `McpSource::Prebuilt(None)` (same semantics as
871    ///   `[mcp] enabled = false`, no per-task re-evaluation).
872    /// - build failure → cached degraded entry with
873    ///   [`MCP_BUILD_RETRY_BACKOFF`]; the error is returned so THIS user's
874    ///   dispatch fails loud while other users are unaffected.
875    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        // Effective per-user config — the fingerprint source of truth.
883        let resolved = self.resolve_user_config(user_key);
884        let fingerprint = fingerprint_mcp_section(resolved.as_deref());
885
886        // Fast path: fresh, non-degraded entry.
887        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        // Single-flight: one build per user at a time.
902        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        // Double-check: another task may have built while we waited.
910        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    /// Build a loader entry from the effective config. No caching here —
946    /// the caller owns insertion and degraded handling.
947    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                // THE parity-evidence log line (16C acceptance + migration task):
983                // one INFO line per build; a task loop that reuses the cache
984                // shows exactly one line per user per fingerprint.
985                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    /// Spawn a background drain task for a specific session's workflow receiver.
1017    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        // Bug 6: broadcast StateChanged only on actual transitions. Without
1025        // this, every WS message was accompanied by a duplicate StateChanged,
1026        // doubling traffic and re-triggering frontend updateState on hot
1027        // StreamDelta/ReasoningDelta bursts. New subscribers learn the current
1028        // state from the WS snapshot, so the initial None is safe.
1029        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    /// Spawn the default user's drain task (backward compatibility).
1060    /// Called during server startup for the initial session.
1061    pub fn spawn_drain_task(self, mut workflow_rx: mpsc::UnboundedReceiver<TuiMessage>) {
1062        // Get the default user's first session
1063        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        // Bug 6: transition-only StateChanged broadcasts (see
1079        // spawn_user_drain_task). The task starts in Running state because
1080        // spawn_drain_task is only called after a command began executing.
1081        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    /// Resolve the user key from request headers.
1112    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        // Try Bearer header first
1118        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        // Try cookie
1136        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
1168// ---------------------------------------------------------------------------
1169// SerializableMessage (unchanged)
1170// ---------------------------------------------------------------------------
1171
1172/// Wrapper to serialize `TuiMessage` as JSON with a `type` discriminator.
1173struct 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
1314// ---------------------------------------------------------------------------
1315// Per-user config helpers
1316// ---------------------------------------------------------------------------
1317
1318/// Deep-merge a TOML overlay on top of a base value (in-place).
1319///
1320/// - Tables: recursively merge key-by-key (overlay wins on conflict).
1321/// - Arrays: overlay replaces base entirely (no merging).
1322/// - Scalars: overlay replaces base.
1323/// - If a key exists in overlay but not base, it's added.
1324fn 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                        // Both exist — recurse if both are tables, else replace
1331                        deep_merge_toml(base_val, overlay_val);
1332                    }
1333                    None => {
1334                        // Key only in overlay — insert
1335                        base_table.insert(key.clone(), overlay_val.clone());
1336                    }
1337                }
1338            }
1339        }
1340        // Non-table: overlay replaces base
1341        (base, overlay) => {
1342            *base = overlay.clone();
1343        }
1344    }
1345}
1346
1347/// Replace `${VAR_NAME}` references in a string with values from a secrets map.
1348///
1349/// Falls back to process environment if the variable is not in the map.
1350/// Variables not found in either are left as-is.
1351fn substitute_env_vars(s: &mut String, secrets: &std::collections::HashMap<String, String>) {
1352    // Simple state machine: scan for ${, read until }, replace.
1353    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            // Find closing }
1360            if let Some(end) = s[i + 2..].find('}') {
1361                let var_name = &s[i + 2..i + 2 + end];
1362                // Look up in per-user secrets first, then process env
1363                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                    // Not found — leave as-is
1369                    result.push_str(&s[i..i + 2 + end + 1]);
1370                }
1371                i = i + 2 + end + 1;
1372            } else {
1373                // No closing } — copy as-is
1374                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
1386// ---------------------------------------------------------------------------
1387// Tests (task 16B: overlay allowlist + consolidated user hash)
1388// ---------------------------------------------------------------------------
1389
1390/// Fingerprint the effective `[mcp]` section of a resolved config (16C).
1391///
1392/// Content hash (SHA-256, first 8 bytes as u64) — never mtime: overlays are
1393/// rewritten in place. Configs without `[mcp]` fingerprint to a stable
1394/// constant, so the disabled state caches too.
1395fn 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    /// Default-knob ServerState for overlay tests.
1413    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    /// Unique temp dir with a `config/` subdir (std-only; no tempfile dep).
1420    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    /// (a) Overlay with [mcp] + [server] + [auth]: [mcp] applied (user wins),
1438    /// shared [server] preserved untouched, overlay [auth] dropped entirely.
1439    #[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    /// (b) No overlay file → merge is a no-op (None), shared config untouched.
1461    #[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    /// Overlay whose sections are ALL non-allowlisted → no-op (None).
1471    #[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    /// (c) [llm] dropped when allow_llm_overlay=false (default), applied when
1486    /// the instance opted in via with_allow_llm_overlay(true).
1487    #[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        // Default: knob false → [llm]-only overlay is a no-op.
1493        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        // Opted in: [llm] kept, user value wins; shared [mcp] untouched.
1500        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    /// (d) Web-path home dir resolves through the consolidated
1513    /// trustee_core::user_hash (determinism/known-vectors covered in
1514    /// trustee-core's own tests).
1515    #[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        // dirs::home_dir() unavailable in sandbox → covered by core tests.
1527    }
1528
1529    /// 16D integration guard: dev-agent principals (user_key `agent-<name>`)
1530    /// get their OWN session bucket and home-dir namespace — one agent can
1531    /// never see another's sessions, and both go through the same
1532    /// user_hash isolation path as humans.
1533    #[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        // Each agent resolves its own session; B cannot see A's by id.
1543        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        // Same isolation path as humans: home dir = users_root/user_hash(key).
1554        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// ---------------------------------------------------------------------------
1564// 16C — per-user McpToolLoader cache
1565// ---------------------------------------------------------------------------
1566
1567#[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    /// Deterministic throwaway user: real user-home path under
1580    /// ~/.trustee/users/{hash} (that IS the resolution path under test),
1581    /// with config/ subdir; caller cleans up.
1582    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        // Second call is a fingerprint hit — same disabled marker, still one entry.
1622        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        // v1: an enabled MCP server pointing at a blackhole port — abk keeps
1655        // the loader with a DOWN status (connect refused is fast on loopback).
1656        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        // v2: different server URL → different content fingerprint → rebuild.
1666        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        // The old Arc stays valid for in-flight sessions (no panic, no revoke).
1681        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        // [mcp] present but not a table → McpConfig deserialize fails → Err.
1692        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        // Within backoff: fast-fail again (cached error).
1709        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        // Other users are completely unaffected.
1716        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}