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/// 16F: the service-account issuer declared by the shared config's first
148/// `[mcp.credentials.*]` entry with `type = "service-account"` (see
149/// `ServerState::service_issuer` for why this exists).
150pub fn service_issuers_from_config(config_toml: &str) -> Vec<String> {
151    let Ok(v) = config_toml.parse::<toml::Value>() else {
152        return Vec::new();
153    };
154    let Some(creds) = v.get("mcp").and_then(|m| m.get("credentials")).and_then(|c| c.as_table())
155    else {
156        return Vec::new();
157    };
158    let mut out = Vec::new();
159    for (_name, cred) in creds {
160        if cred.get("type").and_then(|t| t.as_str()) == Some("service-account") {
161            if let Some(issuer) = cred.get("issuer_url").and_then(|i| i.as_str()) {
162                if !out.iter().any(|e: &String| e == issuer) {
163                    out.push(issuer.to_string());
164                }
165            }
166        }
167    }
168    out
169}
170
171/// Shared state accessible by all axum handlers.
172#[derive(Clone)]
173pub struct ServerState {
174    /// Per-user multi-session registry (MSU).
175    pub sessions: SessionRegistry,
176    /// Broadcast sender for backward compat — delegates to the default user's channel.
177    pub ws_tx: broadcast::Sender<String>,
178    /// Auth state (None = auth disabled, all endpoints open).
179    pub auth: Option<Arc<AuthState>>,
180    /// Shared config TOML (all users share the same agent config).
181    pub config_toml: Option<String>,
182    /// Shared secrets (injected into every per-user session).
183    pub secrets: Option<std::collections::HashMap<String, String>>,
184    /// Shared build info (injected into every per-user session).
185    pub build_info: Option<trustee_core::types::BuildInfo>,
186    /// Global concurrency limiter — limits the number of simultaneous workflows
187    /// across all users. Default: 8 concurrent workflows.
188    pub workflow_semaphore: Arc<tokio::sync::Semaphore>,
189    /// Maximum number of concurrent sessions per user. Default: 4.
190    pub max_sessions_per_user: usize,
191    /// Whether per-user config overlays may override the `[llm]` section.
192    /// Default: `false` — overlays are limited to `[mcp]`.
193    pub allow_llm_overlay: bool,
194    /// Per-user McpToolLoader cache (16C), keyed by user hash (never the raw key).
195    pub mcp_loaders: Arc<DashMap<String, McpLoaderEntry>>,
196    /// Single-flight build locks per user hash (16C).
197    mcp_build_locks: Arc<DashMap<String, Arc<tokio::sync::Mutex<()>>>>,
198    /// 16F: THQ per-agent dispatch table (agent_name → target), boot-populated.
199    pub thq_dispatch: Arc<DashMap<String, ThqDispatchEntry>>,
200    /// 16F: cached impersonation Bearers per agent user_key
201    /// `(token, expires_at)` — expiry-buffered, re-minted on 401.
202    pub agent_dispatch_tokens: Arc<DashMap<String, (String, std::time::Instant)>>,
203}
204
205impl ServerState {
206    /// Create new shared state from a default session, broadcast sender, and optional auth.
207    pub fn new(
208        session: Session,
209        ws_tx: broadcast::Sender<String>,
210        auth: Option<Arc<AuthState>>,
211    ) -> Self {
212        let sessions = Arc::new(DashMap::new());
213
214        // Store the default user's UserSessions with an initial session
215        let token_store = Arc::new(pep::MemoryTokenStore::new());
216        let (ws_tx_entry, _) = broadcast::channel::<String>(256);
217
218        let now = chrono::Utc::now();
219        let initial_entry = UserSessionEntry {
220            session: Arc::new(Mutex::new(session)),
221            ws_tx: ws_tx_entry,
222            created_at: now,
223            last_active: Arc::new(Mutex::new(now)),
224        };
225
226        let user_sessions = UserSessions {
227            sessions: DashMap::new(),
228            token_store,
229            active_session_id: Mutex::new(String::new()),
230        };
231        user_sessions.sessions.insert("default".to_string(), initial_entry);
232
233        sessions.insert("default".to_string(), user_sessions);
234
235        Self {
236            sessions,
237            ws_tx,
238            auth,
239            config_toml: None,
240            secrets: None,
241            build_info: None,
242            workflow_semaphore: Arc::new(tokio::sync::Semaphore::new(8)),
243            max_sessions_per_user: 4,
244            allow_llm_overlay: false,
245            mcp_loaders: Arc::new(DashMap::new()),
246            mcp_build_locks: Arc::new(DashMap::new()),
247            thq_dispatch: Arc::new(DashMap::new()),
248            agent_dispatch_tokens: Arc::new(DashMap::new()),
249        }
250    }
251
252    pub fn with_config_toml(mut self, config_toml: String) -> Self {
253        self.config_toml = Some(config_toml);
254        self
255    }
256
257    /// 16F: the issuer URL Kanidm accepts SERVICE-ACCOUNT token exchange on.
258    ///
259    /// Kanidm binds exchange to the OAuth2 client's configured origin: the
260    /// `[oidc]` issuer host serves logins/validation fine but REJECTS token
261    /// exchange with `invalid_request`, while the `idp` vhost — the issuer
262    /// every working `[mcp.credentials.*]` service-account entry already
263    /// uses — accepts it (verified live 2026-08-30: same token, 200 vs 400).
264    /// The exchange therefore takes its issuer from the shared config's
265    /// first `service-account` credential; callers fall back to the auth
266    /// issuer when absent.
267    pub fn service_issuers(&self) -> Vec<String> {
268        service_issuers_from_config(self.config_toml.as_deref().unwrap_or(""))
269    }
270
271    pub fn with_secrets(mut self, secrets: std::collections::HashMap<String, String>) -> Self {
272        self.secrets = Some(secrets);
273        self
274    }
275
276    pub fn with_build_info(mut self, build_info: trustee_core::types::BuildInfo) -> Self {
277        self.build_info = Some(build_info);
278        self
279    }
280
281    pub fn with_max_concurrent_workflows(mut self, max: usize) -> Self {
282        self.workflow_semaphore = Arc::new(tokio::sync::Semaphore::new(max));
283        self
284    }
285
286    /// Set the max sessions per user.
287    pub fn with_max_sessions_per_user(mut self, max: usize) -> Self {
288        self.max_sessions_per_user = max;
289        self
290    }
291
292    /// Allow per-user config overlays to override the `[llm]` section.
293    /// Default: `false` (overlays limited to `[mcp]`).
294    pub fn with_allow_llm_overlay(mut self, allow: bool) -> Self {
295        self.allow_llm_overlay = allow;
296        self
297    }
298
299    // -----------------------------------------------------------------------
300    // MSU: Multi-session methods
301    // -----------------------------------------------------------------------
302
303    /// Create a new session for a user. Returns the session_id.
304    ///
305    /// Creates a fresh `Session::new()`, copies shared config, sets per-user
306    /// isolation, creates a broadcast channel, spawns a drain task, and inserts
307    /// into the user's session DashMap. The new session becomes the "active" one.
308    pub async fn create_session(
309        &self,
310        user_key: &str,
311        session_name: Option<String>,
312        identity: Option<String>,
313        activate: bool,
314    ) -> Result<String, SessionError> {
315        // Get or create the user's UserSessions entry
316        let user_sessions = self
317            .sessions
318            .entry(user_key.to_string())
319            .or_insert_with(|| UserSessions {
320                sessions: DashMap::new(),
321                token_store: Arc::new(pep::MemoryTokenStore::new()),
322                active_session_id: Mutex::new(String::new()),
323            });
324
325        // Check session limit
326        if user_sessions.sessions.len() >= self.max_sessions_per_user {
327            return Err(SessionError::MaxSessionsReached(self.max_sessions_per_user));
328        }
329
330        // Create new Session
331        let (mut session, workflow_rx) = Session::new();
332
333        // Copy shared config
334        if let Some(ref config_toml) = self.config_toml {
335            session.config_toml = Some(config_toml.clone());
336            session.parse_auto_handoff_config();
337            if let Ok(table) = config_toml.parse::<toml::Value>() {
338                if let Some(name) = table
339                    .get("agent")
340                    .and_then(|a| a.get("name"))
341                    .and_then(|n| n.as_str())
342                {
343                    session.agent_name = name.to_string();
344                }
345            }
346        }
347
348        session.secrets = self.secrets.clone();
349        session.build_info = self.build_info.clone();
350
351        // Per-user isolation
352        self.apply_user_isolation(&mut session, user_key);
353
354        // Apply session_name if provided
355        session.session_name = session_name;
356
357        // Apply agent identity if provided
358        session.identity = identity;
359
360        // Create broadcast channel
361        let (ws_tx_entry, _) = broadcast::channel::<String>(256);
362
363        // Generate session_id
364        let session_id = format!(
365            "session_{}_{}",
366            chrono::Utc::now().format("%Y_%m_%d_%H_%M"),
367            &uuid::Uuid::new_v4().to_string()[..8]
368        );
369
370        let now = chrono::Utc::now();
371
372        // Insert into user's sessions DashMap
373        user_sessions.sessions.insert(
374            session_id.clone(),
375            UserSessionEntry {
376                session: Arc::new(Mutex::new(session)),
377                ws_tx: ws_tx_entry.clone(),
378                created_at: now,
379                last_active: Arc::new(Mutex::new(now)),
380            },
381        );
382
383        // Set as active session only when the caller opts in.
384        //
385        // `resume_session` passes `activate: false` so resuming an arbitrary
386        // checkpoint never hijacks the caller's current active live session
387        // pointer. `new_session` / `create_session` (fresh-start paths) pass
388        // `activate: true` to keep the existing "new session becomes active"
389        // behavior. This keeps the per-user active pointer under UI control
390        // rather than letting any authenticated client overwrite it.
391        if activate {
392            *user_sessions.active_session_id.lock().await = session_id.clone();
393        }
394
395        // Spawn drain task
396        let session_arc = user_sessions
397            .sessions
398            .get(&session_id)
399            .map(|e| e.session.clone());
400        if let Some(session_arc) = session_arc {
401            self.spawn_user_drain_task(
402                session_id.clone(),
403                session_arc,
404                ws_tx_entry,
405                workflow_rx,
406            );
407        }
408
409        Ok(session_id)
410    }
411
412    /// Get a specific session by user_key + session_id.
413    /// Updates last_active on the session entry.
414    pub async fn get_session(
415        &self,
416        user_key: &str,
417        session_id: &str,
418    ) -> Option<(Arc<Mutex<Session>>, broadcast::Sender<String>)> {
419        let user_sessions = self.sessions.get(user_key)?;
420        let entry = user_sessions.sessions.get(session_id)?;
421
422        // Update last_active
423        let now = chrono::Utc::now();
424        *entry.last_active.lock().await = now;
425
426        Some((entry.session.clone(), entry.ws_tx.clone()))
427    }
428
429    /// Get a session by EITHER its live MSU registry key OR its
430    /// checkpoint/session identity (`session.session_id`).
431    ///
432    /// The web frontend tracks `currentSessionId` from the `ResumeInfo` WS
433    /// message, which carries the auto-derived checkpoint id
434    /// (`session_YYYY_MM_DD_HH_MM_uuid8`) — NOT the live MSU registry key
435    /// (`"default"` or the key from `create_session()`). External clients
436    /// like Torpi/THQ pass the live registry key. This resolver accepts both:
437    ///
438    /// 1. Try registry-key lookup first (precise, used by Torpi/THQ).
439    /// 2. Fall back to scanning the user's live sessions for one whose
440    ///    `session.session_id` matches the requested id (used by the
441    ///    embedded web UI after a command or resume).
442    ///
443    /// Returns `(live_registry_key, session_arc, ws_tx)`, or `None` if not
444    /// found. The live key is returned so callers that need to set it as
445    /// active (or otherwise reference the registry) use the real key.
446    pub async fn get_session_by_any_id(
447        &self,
448        user_key: &str,
449        id: &str,
450    ) -> Option<(String, Arc<Mutex<Session>>, broadcast::Sender<String>)> {
451        // Fast path: registry key match.
452        let user_sessions = self.sessions.get(user_key)?;
453        if let Some(entry) = user_sessions.sessions.get(id) {
454            // Update last_active
455            let now = chrono::Utc::now();
456            *entry.last_active.lock().await = now;
457            return Some((id.to_string(), entry.session.clone(), entry.ws_tx.clone()));
458        }
459
460        // Slow path: scan live sessions for a matching session.session_id.
461        for entry in user_sessions.sessions.iter() {
462            let session = entry.session.lock().await;
463            if session.session_id.as_deref() == Some(id) {
464                let key = entry.key().clone();
465                let ws_tx = entry.ws_tx.clone();
466                drop(session);
467                // Update last_active
468                let now = chrono::Utc::now();
469                *entry.last_active.lock().await = now;
470                return Some((key, entry.session.clone(), ws_tx));
471            }
472        }
473
474        None
475    }
476
477    /// List all active sessions for a user, sorted by last_active desc.
478    pub async fn list_sessions(&self, user_key: &str) -> Vec<SessionListItem> {
479        let Some(user_sessions) = self.sessions.get(user_key) else {
480            return Vec::new();
481        };
482
483        let mut items = Vec::new();
484        for entry in user_sessions.sessions.iter() {
485            let session = entry.session.lock().await;
486            let workflow_state = match session.workflow_state {
487                trustee_core::types::WorkflowState::Idle => "Idle",
488                trustee_core::types::WorkflowState::Running => "Running",
489                trustee_core::types::WorkflowState::Cancelling => "Cancelling",
490            };
491            let last_active = entry.last_active.lock().await;
492            items.push(SessionListItem {
493                session_id: entry.key().clone(),
494                session_name: session.session_name.clone(),
495                workflow_state: workflow_state.to_string(),
496                created_at: entry.created_at.to_rfc3339(),
497                last_active: last_active.to_rfc3339(),
498                handoff_count: session.handoff_count,
499            });
500        }
501        drop(user_sessions);
502
503        // Sort by last_active descending
504        items.sort_by(|a, b| b.last_active.cmp(&a.last_active));
505        items
506    }
507
508    /// Destroy a session. The session must be Idle.
509    pub async fn destroy_session(
510        &self,
511        user_key: &str,
512        session_id: &str,
513    ) -> Result<(), SessionError> {
514        let user_sessions = self
515            .sessions
516            .get(user_key)
517            .ok_or_else(|| SessionError::NotFound(session_id.to_string()))?;
518
519        // Check workflow state before removing
520        {
521            let entry = user_sessions
522                .sessions
523                .get(session_id)
524                .ok_or_else(|| SessionError::NotFound(session_id.to_string()))?;
525            let session = entry.session.lock().await;
526            if session.workflow_state != trustee_core::types::WorkflowState::Idle {
527                let state_str = match session.workflow_state {
528                    trustee_core::types::WorkflowState::Running => "Running",
529                    trustee_core::types::WorkflowState::Cancelling => "Cancelling",
530                    _ => "Unknown",
531                };
532                return Err(SessionError::NotIdle(state_str.to_string()));
533            }
534        }
535
536        // Remove from DashMap
537        user_sessions.sessions.remove(session_id);
538
539        // If this was the active session, pick a new active
540        let mut active_id = user_sessions.active_session_id.lock().await;
541        if &*active_id == session_id {
542            // Pick the most recently active remaining session
543            let mut newest: Option<(String, chrono::DateTime<chrono::Utc>)> = None;
544            for entry in user_sessions.sessions.iter() {
545                let la = entry.last_active.lock().await;
546                if newest.as_ref().map_or(true, |(_, t)| *la > *t) {
547                    newest = Some((entry.key().clone(), *la));
548                }
549            }
550            *active_id = newest.map(|(id, _)| id).unwrap_or_default();
551        }
552
553        Ok(())
554    }
555
556    /// Get or create the user's "active" session for legacy routes.
557    ///
558    /// Behavior:
559    /// 1. If user has no sessions → create one
560    /// 2. If active session exists → return it
561    /// 3. If active session was destroyed → create a new one
562    ///
563    /// Returns: (session_id, session_arc, ws_tx, token_store)
564    pub async fn ensure_active_session(
565        &self,
566        user_key: &str,
567    ) -> (
568        String,
569        Arc<Mutex<Session>>,
570        broadcast::Sender<String>,
571        Arc<pep::MemoryTokenStore>,
572    ) {
573        // Get or create user's UserSessions
574        let token_store = {
575            let user_sessions = self
576                .sessions
577                .entry(user_key.to_string())
578                .or_insert_with(|| UserSessions {
579                    sessions: DashMap::new(),
580                    token_store: Arc::new(pep::MemoryTokenStore::new()),
581                    active_session_id: Mutex::new(String::new()),
582                });
583            user_sessions.token_store.clone()
584        };
585
586        // Check if active session exists
587        let active_id = {
588            let user_sessions = self.sessions.get(user_key).unwrap();
589            let guard = user_sessions.active_session_id.lock().await;
590            guard.clone()
591        };
592
593        if !active_id.is_empty() {
594            if let Some((session, ws_tx)) = self.get_session(user_key, &active_id).await {
595                return (active_id, session, ws_tx, token_store);
596            }
597            // Active session was destroyed, fall through to create
598        }
599
600        // Need to create a new session
601        // For the "default" user, we may already have a "default" session entry
602        // from ServerState::new() — check for it
603        let existing_session: Option<(String, Arc<Mutex<Session>>, broadcast::Sender<String>)> = {
604            let user_sessions = self.sessions.get(user_key).unwrap();
605            let result = user_sessions.sessions.iter().next().map(|first| {
606                (
607                    first.key().clone(),
608                    first.session.clone(),
609                    first.ws_tx.clone(),
610                )
611            });
612            result
613        };
614        if let Some((id, session, ws_tx)) = existing_session {
615            let now = chrono::Utc::now();
616            if let Some(entry) = self.sessions.get(user_key) {
617                if let Some(e) = entry.sessions.get(&id) {
618                    *e.last_active.lock().await = now;
619                }
620                *entry.active_session_id.lock().await = id.clone();
621            }
622
623            return (id, session, ws_tx, token_store);
624        }
625
626        // Create a brand new session
627        let session_id = self
628            .create_session(user_key, None, None, true)
629            .await
630            .unwrap_or_else(|_| "default".to_string());
631
632        let (session, ws_tx) = self
633            .get_session(user_key, &session_id)
634            .await
635            .expect("just-created session must exist");
636
637        (session_id, session, ws_tx, token_store)
638    }
639
640    /// DEPRECATED: Use ensure_active_session() instead.
641    /// Kept for backward compatibility — same 3-tuple return type.
642    pub async fn ensure_user_session(
643        &self,
644        user_key: &str,
645    ) -> (Arc<Mutex<Session>>, broadcast::Sender<String>, Arc<pep::MemoryTokenStore>) {
646        let (_id, session, ws_tx, token_store) = self.ensure_active_session(user_key).await;
647        (session, ws_tx, token_store)
648    }
649
650    /// Set a session as the user's active session.
651    pub async fn set_active_session(&self, user_key: &str, session_id: &str) {
652        if let Some(user_sessions) = self.sessions.get(user_key) {
653            if user_sessions.sessions.contains_key(session_id) {
654                *user_sessions.active_session_id.lock().await = session_id.to_string();
655            }
656        }
657    }
658
659    // -----------------------------------------------------------------------
660    // Read-only helpers (no session creation side effects)
661    // -----------------------------------------------------------------------
662
663    /// Resolve a user's home_dir without creating an in-memory session.
664    ///
665    /// This is the read-only equivalent of the isolation logic in
666    /// `apply_user_isolation`. Used by endpoints that only need to read
667    /// checkpoint data from disk (history, session list, session detail)
668    /// and must NOT create ghost sessions as a side effect.
669    pub fn get_user_home_dir(&self, user_key: &str) -> Option<std::path::PathBuf> {
670        let hash = trustee_core::user_hash(user_key);
671        dirs::home_dir().map(|home| home.join(".trustee").join("users").join(&hash))
672    }
673
674    /// Resolve config_toml and home_dir without creating an in-memory session.
675    ///
676    /// Returns `(config_toml, home_dir)`. If config is not loaded,
677    /// config_toml will be None.
678    pub fn get_user_config_and_home(&self, user_key: &str) -> (Option<String>, Option<std::path::PathBuf>) {
679        (self.config_toml.clone(), self.get_user_home_dir(user_key))
680    }
681
682    // -----------------------------------------------------------------------
683    // Private helpers
684    // -----------------------------------------------------------------------
685
686    /// Apply per-user isolation: SHA-256 hash → home_dir + project_id.
687    ///
688    /// Hashing goes through the single consolidated [`trustee_core::user_hash`]
689    /// so the web path and the CLI path can never drift apart.
690    fn apply_user_isolation(&self, session: &mut Session, user_key: &str) {
691        let user_hash = trustee_core::user_hash(user_key);
692
693        // Set per-user home directory for checkpoint isolation
694        let user_home = if let Some(home) = dirs::home_dir() {
695            let user_home = home.join(".trustee").join("users").join(&user_hash);
696            session.home_dir = Some(user_home.clone());
697            Some(user_home)
698        } else {
699            None
700        };
701
702        session.project_id = Some(format!("web{}", &user_hash[..16]));
703
704        // ── Per-user .env (Task 2) ──────────────────────────────────────
705        //
706        // Load per-user secrets from ~/.trustee/users/{hash}/.env
707        // These are merged on top of shared secrets (per-user wins).
708        // They are NEVER set as process env vars — used only for ${VAR}
709        // substitution in the config TOML below.
710        let shared_secrets = session.secrets.clone().unwrap_or_default();
711        let mut merged_secrets = shared_secrets.clone();
712
713        if let Some(ref user_home) = user_home {
714            if let Ok(merged) = self.load_user_secrets(user_home, &merged_secrets) {
715                merged_secrets = merged;
716            }
717        }
718
719        // ── Per-user config overlay (Task 3) ────────────────────────────
720        if let Some(ref user_home) = user_home {
721            if let Some(merged) = self.merge_user_config(user_home, session.config_toml.as_deref()) {
722                session.config_toml = Some(merged);
723                tracing::debug!("Merged per-user config into session");
724            }
725        }
726
727        // ── ${VAR} substitution (Task 4) ────────────────────────────────
728        if let Some(ref mut config_toml) = session.config_toml {
729            substitute_env_vars(config_toml, &merged_secrets);
730        }
731
732        // ── Strip per-user secrets (Task 5) ────────────────────────────
733        session.secrets = Some(shared_secrets);
734    }
735
736    /// Load secrets from a user's ~/.trustee/users/{hash}/.env and merge
737    /// on top of `base` (user wins). Returns the merged map.
738    fn load_user_secrets(
739        &self,
740        user_home: &std::path::Path,
741        base: &std::collections::HashMap<String, String>,
742    ) -> std::io::Result<std::collections::HashMap<String, String>> {
743        let user_env_path = user_home.join(".env");
744        if !user_env_path.exists() {
745            return Ok(base.clone());
746        }
747        let content = std::fs::read_to_string(&user_env_path)?;
748        let mut merged = base.clone();
749        for line in content.lines() {
750            let line = line.trim();
751            if line.is_empty() || line.starts_with('#') {
752                continue;
753            }
754            if let Some((key, value)) = line.split_once('=') {
755                let key = key.trim().to_string();
756                let value = value
757                    .trim()
758                    .trim_matches('"')
759                    .trim_matches('\'')
760                    .to_string();
761                merged.insert(key, value);
762            }
763        }
764        tracing::debug!("Loaded per-user secrets from {}", user_env_path.display());
765        Ok(merged)
766    }
767
768    /// Deep-merge a user's per-user config overlay on top of the shared
769    /// config. Returns the merged TOML string, or None if no overlay exists.
770    ///
771    /// Overlay allowlist (task 16B): only allowlisted top-level sections of
772    /// the user overlay are merged; anything else is dropped loudly.
773    /// - `[mcp]` is always allowed — per-user MCP tool sets are the point of
774    ///   the overlay convention.
775    /// - `[llm]` is allowed only when the instance opted in via the
776    ///   `[users].allow_llm_overlay` knob (default **false**).
777    ///
778    /// Everything else (`[server]`, `[auth]`, `[storage]`, `[web]`, …) is
779    /// boot-time, instance-level config and must not be rewritable through a
780    /// per-user file. This is predictability hardening, not a security
781    /// boundary: those sections were never re-read from session config.
782    fn merge_user_config(
783        &self,
784        user_home: &std::path::Path,
785        shared_config: Option<&str>,
786    ) -> Option<String> {
787        let user_config_path = user_home.join("config").join("trustee.toml");
788        if !user_config_path.exists() {
789            return None;
790        }
791        let user_config_toml = std::fs::read_to_string(&user_config_path).ok()?;
792        let shared = shared_config
793            .unwrap_or("")
794            .parse::<toml::Value>()
795            .ok()?;
796        let overlay = user_config_toml.parse::<toml::Value>().ok()?;
797
798        let allowed = |section: &str| {
799            section == "mcp"
800                // 16E: [thq] is per-user identity config (read from the overlay
801                // file at boot); keeping it in the merged session config is
802                // inert but avoids a false "dropping" warn on every dispatch.
803                || section == "thq"
804                // xagent personas: an agent's [lifecycle].system_template IS
805                // her identity — user-owned by design. The WASM extension
806                // switch ([lifecycle].enabled) is stripped below: persona is
807                // user-owned, extension control is not.
808                || section == "lifecycle"
809                || (self.allow_llm_overlay && section == "llm")
810        };
811        // Mask the user in logs: user_home's dir name IS the hash — never
812        // log the raw key (keys may be emails). Log a short hash prefix.
813        let dir_name = user_home
814            .file_name()
815            .and_then(|n| n.to_str())
816            .unwrap_or("<unknown>");
817        let masked_user = dir_name.get(..8).unwrap_or(dir_name);
818
819        let mut filtered_overlay = toml::map::Map::new();
820        if let Some(table) = overlay.as_table() {
821            for (section, value) in table {
822                if section == "lifecycle" {
823                    let mut v = value.clone();
824                    if let Some(t) = v.as_table_mut() {
825                        if t.remove("enabled").is_some() {
826                            tracing::info!(
827                                "user config overlay: [lifecycle].enabled is not user-settable — stripped for user {}",
828                                masked_user
829                            );
830                        }
831                    }
832                    filtered_overlay.insert(section.clone(), v);
833                } else if allowed(section) {
834                    filtered_overlay.insert(section.clone(), value.clone());
835                } else {
836                    tracing::warn!(
837                        "user config overlay: dropping non-allowlisted section [{}] for user {}",
838                        section,
839                        masked_user
840                    );
841                }
842            }
843        }
844
845        if filtered_overlay.is_empty() {
846            // Nothing survived the allowlist — no-op, keep shared config as-is.
847            return None;
848        }
849        let overlay = toml::Value::Table(filtered_overlay);
850
851        let mut shared = shared;
852        deep_merge_toml(&mut shared, &overlay);
853        let merged = toml::to_string(&shared).ok()?;
854        tracing::debug!("Merged per-user config from {}", user_config_path.display());
855        Some(merged)
856    }
857
858    /// Resolve the fully-merged config TOML for a user WITHOUT creating a session.
859    ///
860    /// This is the read-only equivalent of the config resolution in
861    /// `apply_user_isolation`: shared config + per-user overlay + ${VAR}
862    /// substitution. Used by endpoints that need to inspect config (e.g.
863    /// listing available LLM models) without spawning a ghost session.
864    ///
865    /// Returns None if no shared config is loaded.
866    pub fn resolve_user_config(&self, user_key: &str) -> Option<String> {
867        let config_toml = self.config_toml.clone()?;
868
869        // Resolve user home dir (same hash scheme as apply_user_isolation)
870        let user_home = self.get_user_home_dir(user_key)?;
871
872        // Start from shared secrets; merge per-user .env on top
873        let mut merged_secrets = self.secrets.clone().unwrap_or_default();
874        if let Ok(merged) = self.load_user_secrets(&user_home, &merged_secrets) {
875            merged_secrets = merged;
876        }
877
878        // Merge per-user config overlay
879        let mut resolved = config_toml;
880        if let Some(merged) = self.merge_user_config(&user_home, Some(&resolved)) {
881            resolved = merged;
882        }
883
884        // Substitute ${VAR} from merged secrets
885        substitute_env_vars(&mut resolved, &merged_secrets);
886
887        Some(resolved)
888    }
889
890    /// Get or build the per-user MCP tool loader (16C).
891    ///
892    /// Cache semantics:
893    /// - fingerprint = content hash of the effective `[mcp]` section
894    ///   (shared + allowlist-filtered overlay + ${VAR} substitution, i.e.
895    ///   exactly what `resolve_user_config` produces) — stale on ANY change.
896    /// - hit + match → `Arc` clone, zero network I/O.
897    /// - miss/stale → single-flight build (one build per user at a time;
898    ///   late arrivals re-check and reuse the winner's entry).
899    /// - `Ok(None)` = MCP disabled for this user → agent runs MCP-less via
900    ///   abk `McpSource::Prebuilt(None)` (same semantics as
901    ///   `[mcp] enabled = false`, no per-task re-evaluation).
902    /// - build failure → cached degraded entry with
903    ///   [`MCP_BUILD_RETRY_BACKOFF`]; the error is returned so THIS user's
904    ///   dispatch fails loud while other users are unaffected.
905    pub async fn get_or_build_mcp_loader(
906        &self,
907        user_key: &str,
908        token_store: &Arc<pep::MemoryTokenStore>,
909    ) -> Result<Option<std::sync::Arc<abk::agent::McpToolLoader>>, String> {
910        let user_hash = trustee_core::user_hash(user_key);
911
912        // Effective per-user config — the fingerprint source of truth.
913        let resolved = self.resolve_user_config(user_key);
914        let fingerprint = fingerprint_mcp_section(resolved.as_deref());
915
916        // Fast path: fresh, non-degraded entry.
917        if let Some(entry) = self.mcp_loaders.get(&user_hash) {
918            if entry.degraded.is_none() {
919                if entry.fingerprint == fingerprint {
920                    return Ok(entry.loader.clone());
921                }
922            } else if let (Some(err), Some(failed_at)) = (&entry.degraded, entry.failed_at) {
923                let backoff = chrono::Duration::from_std(MCP_BUILD_RETRY_BACKOFF)
924                    .unwrap_or_else(|_| chrono::Duration::seconds(30));
925                if chrono::Utc::now() < failed_at + backoff {
926                    return Err(err.clone());
927                }
928            }
929        }
930
931        // Single-flight: one build per user at a time.
932        let lock = self
933            .mcp_build_locks
934            .entry(user_hash.clone())
935            .or_insert_with(|| Arc::new(tokio::sync::Mutex::new(())))
936            .clone();
937        let _guard = lock.lock().await;
938
939        // Double-check: another task may have built while we waited.
940        if let Some(entry) = self.mcp_loaders.get(&user_hash) {
941            if entry.degraded.is_none() && entry.fingerprint == fingerprint {
942                return Ok(entry.loader.clone());
943            }
944        }
945
946        match self
947            .build_mcp_loader(&user_hash, resolved.as_deref(), fingerprint, token_store)
948            .await
949        {
950            Ok(entry) => {
951                self.mcp_loaders.insert(user_hash.clone(), entry);
952                Ok(self.mcp_loaders.get(&user_hash).unwrap().loader.clone())
953            }
954            Err(err) => {
955                tracing::warn!(
956                    "MCP loader build FAILED for user {}; dispatch fails loud, retry after {:?}",
957                    &user_hash[..8.min(user_hash.len())],
958                    MCP_BUILD_RETRY_BACKOFF
959                );
960                self.mcp_loaders.insert(
961                    user_hash,
962                    McpLoaderEntry {
963                        loader: None,
964                        fingerprint,
965                        built_at: chrono::Utc::now(),
966                        degraded: Some(err.clone()),
967                        failed_at: Some(chrono::Utc::now()),
968                    },
969                );
970                Err(err)
971            }
972        }
973    }
974
975    /// Build a loader entry from the effective config. No caching here —
976    /// the caller owns insertion and degraded handling.
977    async fn build_mcp_loader(
978        &self,
979        user_hash: &str,
980        resolved: Option<&str>,
981        fingerprint: u64,
982        token_store: &Arc<pep::MemoryTokenStore>,
983    ) -> Result<McpLoaderEntry, String> {
984        let mcp_config: Option<abk::config::McpConfig> = match resolved {
985            Some(toml_str) => {
986                let value = toml_str
987                    .parse::<toml::Value>()
988                    .map_err(|e| format!("config parse failed: {}", e))?;
989                match value.get("mcp") {
990                    Some(section) => {
991                        use serde::Deserialize as _;
992                        Some(
993                            abk::config::McpConfig::deserialize(section.clone())
994                                .map_err(|e| format!("invalid [mcp] config: {}", e))?,
995                        )
996                    }
997                    None => None,
998                }
999            }
1000            None => None,
1001        };
1002
1003        let loader = match mcp_config {
1004            Some(cfg) if cfg.enabled => {
1005                let built = abk::agent::McpToolLoader::with_token_store(
1006                    &cfg,
1007                    Some(token_store.clone() as std::sync::Arc<dyn pep::token_store::TokenStore>),
1008                )
1009                .await
1010                .map_err(|e| format!("MCP loader build failed: {}", e))?;
1011
1012                // THE parity-evidence log line (16C acceptance + migration task):
1013                // one INFO line per build; a task loop that reuses the cache
1014                // shows exactly one line per user per fingerprint.
1015                let servers: Vec<String> = built
1016                    .server_statuses
1017                    .iter()
1018                    .map(|s| {
1019                        if s.connected {
1020                            format!("{}(up,{}tools)", s.name, s.tool_count)
1021                        } else {
1022                            format!("{}(DOWN)", s.name)
1023                        }
1024                    })
1025                    .collect();
1026                tracing::info!(
1027                    "MCP loader built for user {}: servers=[{}] total_tools={}",
1028                    &user_hash[..8.min(user_hash.len())],
1029                    servers.join(", "),
1030                    built.tool_count
1031                );
1032                Some(std::sync::Arc::new(built))
1033            }
1034            _ => None,
1035        };
1036
1037        Ok(McpLoaderEntry {
1038            loader,
1039            fingerprint,
1040            built_at: chrono::Utc::now(),
1041            degraded: None,
1042            failed_at: None,
1043        })
1044    }
1045
1046    /// Spawn a background drain task for a specific session's workflow receiver.
1047    fn spawn_user_drain_task(
1048        &self,
1049        session_id: String,
1050        session: Arc<Mutex<Session>>,
1051        ws_tx: broadcast::Sender<String>,
1052        mut workflow_rx: mpsc::UnboundedReceiver<TuiMessage>,
1053    ) {
1054        // Bug 6: broadcast StateChanged only on actual transitions. Without
1055        // this, every WS message was accompanied by a duplicate StateChanged,
1056        // doubling traffic and re-triggering frontend updateState on hot
1057        // StreamDelta/ReasoningDelta bursts. New subscribers learn the current
1058        // state from the WS snapshot, so the initial None is safe.
1059        let mut last_broadcast_state: Option<String> = None;
1060        tokio::spawn(async move {
1061            while let Some(msg) = workflow_rx.recv().await {
1062                {
1063                    let mut session = session.lock().await;
1064                    session.handle_workflow_message(msg.clone());
1065
1066                    let state_str = match session.workflow_state {
1067                        trustee_core::types::WorkflowState::Idle => "Idle",
1068                        trustee_core::types::WorkflowState::Running => "Running",
1069                        trustee_core::types::WorkflowState::Cancelling => "Cancelling",
1070                    };
1071                    if last_broadcast_state.as_deref() != Some(state_str) {
1072                        last_broadcast_state = Some(state_str.to_string());
1073                        let state_msg = serde_json::json!({
1074                            "type": "StateChanged",
1075                            "state": state_str
1076                        });
1077                        let _ = ws_tx.send(state_msg.to_string());
1078                    }
1079                }
1080
1081                let json =
1082                    serde_json::to_string(&SerializableMessage(&msg)).unwrap_or_default();
1083                let _ = ws_tx.send(json);
1084            }
1085            tracing::debug!("Drain task ended for session: {}", session_id);
1086        });
1087    }
1088
1089    /// Spawn the default user's drain task (backward compatibility).
1090    /// Called during server startup for the initial session.
1091    pub fn spawn_drain_task(self, mut workflow_rx: mpsc::UnboundedReceiver<TuiMessage>) {
1092        // Get the default user's first session
1093        let default_user = self
1094            .sessions
1095            .get("default")
1096            .expect("default user must exist");
1097        let first_entry = default_user
1098            .sessions
1099            .iter()
1100            .next()
1101            .expect("default user must have at least one session");
1102        let session = first_entry.session.clone();
1103        let ws_tx = first_entry.ws_tx.clone();
1104        let session_id = first_entry.key().clone();
1105        drop(first_entry);
1106        drop(default_user);
1107
1108        // Bug 6: transition-only StateChanged broadcasts (see
1109        // spawn_user_drain_task). The task starts in Running state because
1110        // spawn_drain_task is only called after a command began executing.
1111        let mut last_broadcast_state = Some("Running".to_string());
1112        tokio::spawn(async move {
1113            while let Some(msg) = workflow_rx.recv().await {
1114                {
1115                    let mut session = session.lock().await;
1116                    session.handle_workflow_message(msg.clone());
1117
1118                    let state_str = match session.workflow_state {
1119                        trustee_core::types::WorkflowState::Idle => "Idle",
1120                        trustee_core::types::WorkflowState::Running => "Running",
1121                        trustee_core::types::WorkflowState::Cancelling => "Cancelling",
1122                    };
1123                    if last_broadcast_state.as_deref() != Some(state_str) {
1124                        last_broadcast_state = Some(state_str.to_string());
1125                        let state_msg = serde_json::json!({
1126                            "type": "StateChanged",
1127                            "state": state_str
1128                        });
1129                        let _ = ws_tx.send(state_msg.to_string());
1130                    }
1131                }
1132
1133                let json =
1134                    serde_json::to_string(&SerializableMessage(&msg)).unwrap_or_default();
1135                let _ = ws_tx.send(json);
1136            }
1137            tracing::debug!("Drain task ended for session: {}", session_id);
1138        });
1139    }
1140
1141    /// Resolve the user key from request headers.
1142    pub async fn resolve_user_key(&self, headers: &axum::http::HeaderMap) -> String {
1143        let Some(ref auth) = self.auth else {
1144            return "default".to_string();
1145        };
1146
1147        // Try Bearer header first
1148        if let Some(token) = headers
1149            .get(axum::http::header::AUTHORIZATION)
1150            .and_then(|v| v.to_str().ok())
1151            .and_then(|v| v.strip_prefix("Bearer "))
1152            .map(|s| s.to_string())
1153        {
1154            if token.starts_with("dev:") {
1155                let parts: Vec<&str> = token.splitn(4, ':').collect();
1156                if parts.len() >= 4 {
1157                    return format!("dev:{}", parts[1]);
1158                }
1159            }
1160            if let Ok(claims) = auth.validate_token(&token).await {
1161                return claims.sub;
1162            }
1163        }
1164
1165        // Try cookie
1166        let cookie_session_id = headers
1167            .get(axum::http::header::COOKIE)
1168            .and_then(|v| v.to_str().ok())
1169            .and_then(|cookies| {
1170                cookies
1171                    .split(';')
1172                    .map(|c| c.trim())
1173                    .find_map(|c| {
1174                        c.strip_prefix(&format!("{}=", auth.config.cookie_name))
1175                            .map(|s| s.to_string())
1176                    })
1177            });
1178
1179        if let Some(session_id) = cookie_session_id {
1180            if session_id.starts_with("dev:") {
1181                let parts: Vec<&str> = session_id.splitn(4, ':').collect();
1182                if parts.len() >= 4 {
1183                    return format!("dev:{}", parts[1]);
1184                }
1185            }
1186
1187            if let Ok(access_token) = auth.session_manager.get_token(&session_id).await {
1188                if let Ok(claims) = auth.validate_token(&access_token).await {
1189                    return claims.sub;
1190                }
1191            }
1192        }
1193
1194        "default".to_string()
1195    }
1196}
1197
1198// ---------------------------------------------------------------------------
1199// SerializableMessage (unchanged)
1200// ---------------------------------------------------------------------------
1201
1202/// Wrapper to serialize `TuiMessage` as JSON with a `type` discriminator.
1203struct SerializableMessage<'a>(&'a TuiMessage);
1204
1205impl<'a> serde::Serialize for SerializableMessage<'a> {
1206    fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
1207    where
1208        S: serde::Serializer,
1209    {
1210        use serde::ser::SerializeStruct;
1211
1212        match self.0 {
1213            TuiMessage::OutputLine(line) => {
1214                let mut s = serializer.serialize_struct("msg", 2)?;
1215                s.serialize_field("type", "OutputLine")?;
1216                s.serialize_field("line", line)?;
1217                s.end()
1218            }
1219            TuiMessage::StreamDelta(delta) => {
1220                let mut s = serializer.serialize_struct("msg", 2)?;
1221                s.serialize_field("type", "StreamDelta")?;
1222                s.serialize_field("delta", delta)?;
1223                s.end()
1224            }
1225            TuiMessage::ReasoningDelta(delta) => {
1226                let mut s = serializer.serialize_struct("msg", 2)?;
1227                s.serialize_field("type", "ReasoningDelta")?;
1228                s.serialize_field("delta", delta)?;
1229                s.end()
1230            }
1231            TuiMessage::WorkflowCompleted => {
1232                let mut s = serializer.serialize_struct("msg", 2)?;
1233                s.serialize_field("type", "WorkflowCompleted")?;
1234                s.serialize_field("state", "Idle")?;
1235                s.end()
1236            }
1237            TuiMessage::WorkflowError(err) => {
1238                let mut s = serializer.serialize_struct("msg", 2)?;
1239                s.serialize_field("type", "WorkflowError")?;
1240                s.serialize_field("error", err)?;
1241                s.end()
1242            }
1243            TuiMessage::ResumeInfo(info) => match info {
1244                Some(ri) => {
1245                    let mut s = serializer.serialize_struct("msg", 5)?;
1246                    s.serialize_field("type", "ResumeInfo")?;
1247                    s.serialize_field("state", "Idle")?;
1248                    s.serialize_field("session_id", &ri.session_id)?;
1249                    s.serialize_field("checkpoint_id", &ri.checkpoint_id)?;
1250                    s.serialize_field("iteration", &ri.iteration)?;
1251                    s.end()
1252                }
1253                None => {
1254                    let mut s = serializer.serialize_struct("msg", 2)?;
1255                    s.serialize_field("type", "ResumeInfo")?;
1256                    s.serialize_field("state", "Idle")?;
1257                    s.end()
1258                }
1259            },
1260            TuiMessage::TodoUpdate(content) => {
1261                let mut s = serializer.serialize_struct("msg", 2)?;
1262                s.serialize_field("type", "TodoUpdate")?;
1263                s.serialize_field("content", content)?;
1264                s.end()
1265            }
1266            TuiMessage::WorkflowCancelled => {
1267                let mut s = serializer.serialize_struct("msg", 2)?;
1268                s.serialize_field("type", "WorkflowCancelled")?;
1269                s.serialize_field("state", "Idle")?;
1270                s.end()
1271            }
1272            TuiMessage::HandoffReady(briefing) => {
1273                let mut s = serializer.serialize_struct("msg", 3)?;
1274                s.serialize_field("type", "HandoffReady")?;
1275                s.serialize_field("state", "Idle")?;
1276                s.serialize_field("briefing", briefing)?;
1277                s.end()
1278            }
1279            TuiMessage::SessionRotated { old, new } => {
1280                let mut s = serializer.serialize_struct("msg", 3)?;
1281                s.serialize_field("type", "SessionRotated")?;
1282                s.serialize_field("old", old)?;
1283                s.serialize_field("new", new)?;
1284                s.end()
1285            }
1286            TuiMessage::HandoffFailed => {
1287                let mut s = serializer.serialize_struct("msg", 2)?;
1288                s.serialize_field("type", "HandoffFailed")?;
1289                s.serialize_field("state", "Idle")?;
1290                s.end()
1291            }
1292            TuiMessage::ToolPending {
1293                tool_name,
1294                hint,
1295            } => {
1296                let mut s = serializer.serialize_struct("msg", 3)?;
1297                s.serialize_field("type", "ToolPending")?;
1298                s.serialize_field("tool_name", tool_name)?;
1299                s.serialize_field("hint", hint)?;
1300                s.end()
1301            }
1302            TuiMessage::ToolDone {
1303                tool_name,
1304                success,
1305                hint,
1306            } => {
1307                let mut s = serializer.serialize_struct("msg", 4)?;
1308                s.serialize_field("type", "ToolDone")?;
1309                s.serialize_field("tool_name", tool_name)?;
1310                s.serialize_field("success", success)?;
1311                s.serialize_field("hint", hint)?;
1312                s.end()
1313            }
1314            TuiMessage::ContextTokensUpdated(count) => {
1315                let mut s = serializer.serialize_struct("msg", 2)?;
1316                s.serialize_field("type", "ContextTokensUpdated")?;
1317                s.serialize_field("count", count)?;
1318                s.end()
1319            }
1320            TuiMessage::McpServerStatus {
1321                name,
1322                connected,
1323                tool_count,
1324                error,
1325            } => {
1326                let mut s = serializer.serialize_struct("msg", 5)?;
1327                s.serialize_field("type", "McpServerStatus")?;
1328                s.serialize_field("name", name)?;
1329                s.serialize_field("connected", connected)?;
1330                s.serialize_field("tool_count", tool_count)?;
1331                s.serialize_field("error", error)?;
1332                s.end()
1333            }
1334            TuiMessage::SessionTitleUpdated(title) => {
1335                let mut s = serializer.serialize_struct("msg", 2)?;
1336                s.serialize_field("type", "SessionTitleUpdated")?;
1337                s.serialize_field("title", title)?;
1338                s.end()
1339            }
1340        }
1341    }
1342}
1343
1344// ---------------------------------------------------------------------------
1345// Per-user config helpers
1346// ---------------------------------------------------------------------------
1347
1348/// Deep-merge a TOML overlay on top of a base value (in-place).
1349///
1350/// - Tables: recursively merge key-by-key (overlay wins on conflict).
1351/// - Arrays: overlay replaces base entirely (no merging).
1352/// - Scalars: overlay replaces base.
1353/// - If a key exists in overlay but not base, it's added.
1354fn deep_merge_toml(base: &mut toml::Value, overlay: &toml::Value) {
1355    match (base, overlay) {
1356        (toml::Value::Table(base_table), toml::Value::Table(overlay_table)) => {
1357            for (key, overlay_val) in overlay_table {
1358                match base_table.get_mut(key) {
1359                    Some(base_val) => {
1360                        // Both exist — recurse if both are tables, else replace
1361                        deep_merge_toml(base_val, overlay_val);
1362                    }
1363                    None => {
1364                        // Key only in overlay — insert
1365                        base_table.insert(key.clone(), overlay_val.clone());
1366                    }
1367                }
1368            }
1369        }
1370        // Non-table: overlay replaces base
1371        (base, overlay) => {
1372            *base = overlay.clone();
1373        }
1374    }
1375}
1376
1377/// Replace `${VAR_NAME}` references in a string with values from a secrets map.
1378///
1379/// Falls back to process environment if the variable is not in the map.
1380/// Variables not found in either are left as-is.
1381fn substitute_env_vars(s: &mut String, secrets: &std::collections::HashMap<String, String>) {
1382    // Simple state machine: scan for ${, read until }, replace.
1383    let mut result = String::with_capacity(s.len());
1384    let bytes = s.as_bytes();
1385    let mut i = 0;
1386
1387    while i < bytes.len() {
1388        if i + 1 < bytes.len() && bytes[i] == b'$' && bytes[i + 1] == b'{' {
1389            // Find closing }
1390            if let Some(end) = s[i + 2..].find('}') {
1391                let var_name = &s[i + 2..i + 2 + end];
1392                // Look up in per-user secrets first, then process env
1393                if let Some(value) = secrets.get(var_name) {
1394                    result.push_str(value);
1395                } else if let Ok(value) = std::env::var(var_name) {
1396                    result.push_str(&value);
1397                } else {
1398                    // Not found — leave as-is
1399                    result.push_str(&s[i..i + 2 + end + 1]);
1400                }
1401                i = i + 2 + end + 1;
1402            } else {
1403                // No closing } — copy as-is
1404                result.push('$');
1405                i += 1;
1406            }
1407        } else {
1408            result.push(bytes[i] as char);
1409            i += 1;
1410        }
1411    }
1412
1413    *s = result;
1414}
1415
1416// ---------------------------------------------------------------------------
1417// Tests (task 16B: overlay allowlist + consolidated user hash)
1418// ---------------------------------------------------------------------------
1419
1420/// Fingerprint the effective `[mcp]` section of a resolved config (16C).
1421///
1422/// Content hash (SHA-256, first 8 bytes as u64) — never mtime: overlays are
1423/// rewritten in place. Configs without `[mcp]` fingerprint to a stable
1424/// constant, so the disabled state caches too.
1425fn fingerprint_mcp_section(resolved: Option<&str>) -> u64 {
1426    use sha2::{Digest, Sha256};
1427    let section = resolved
1428        .and_then(|s| s.parse::<toml::Value>().ok())
1429        .and_then(|v| v.get("mcp").cloned());
1430    let bytes = match section {
1431        Some(v) => v.to_string().into_bytes(),
1432        None => b"<no-mcp>".to_vec(),
1433    };
1434    let digest = Sha256::digest(&bytes);
1435    u64::from_be_bytes(digest[..8].try_into().expect("sha256 digest >= 8 bytes"))
1436}
1437
1438#[cfg(test)]
1439mod tests {
1440    use super::*;
1441
1442    /// Default-knob ServerState for overlay tests.
1443    fn test_state() -> ServerState {
1444        let (session, _rx) = Session::new();
1445        let (ws_tx, _ws_rx) = tokio::sync::broadcast::channel::<String>(16);
1446        ServerState::new(session, ws_tx, None)
1447    }
1448
1449    /// Unique temp dir with a `config/` subdir (std-only; no tempfile dep).
1450    fn temp_user_home(tag: &str) -> std::path::PathBuf {
1451        static COUNTER: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0);
1452        let n = COUNTER.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
1453        let dir = std::env::temp_dir().join(format!(
1454            "trustee-state-test-{}-{}-{}",
1455            tag,
1456            std::process::id(),
1457            n
1458        ));
1459        std::fs::create_dir_all(dir.join("config")).expect("create temp user home");
1460        dir
1461    }
1462
1463    fn parse(toml_str: &str) -> toml::Value {
1464        toml_str.parse::<toml::Value>().expect("valid test TOML")
1465    }
1466
1467    /// (a) Overlay with [mcp] + [server] + [auth]: [mcp] applied (user wins),
1468    /// shared [server] preserved untouched, overlay [auth] dropped entirely.
1469    #[test]
1470    fn overlay_allowlist_drops_non_allowlisted_sections() {
1471        let state = test_state();
1472        let home = temp_user_home("allowlist");
1473        std::fs::write(
1474            home.join("config").join("trustee.toml"),
1475            "[server]\nport = 1\n\n[auth]\nmode = \"kanidm\"\n\n[mcp]\nmode = \"user\"\n",
1476        )
1477        .expect("write overlay");
1478
1479        let shared = "[server]\nport = 8080\n\n[mcp]\nmode = \"shared\"\n";
1480        let merged = state
1481            .merge_user_config(&home, Some(shared))
1482            .expect("overlay has allowlisted content");
1483
1484        let merged_val = parse(&merged);
1485        let expected = parse("[server]\nport = 8080\n\n[mcp]\nmode = \"user\"\n");
1486        assert_eq!(merged_val, expected, "merged config must be shared + [mcp] overlay only");
1487        assert!(merged_val.get("auth").is_none(), "overlay [auth] must be dropped");
1488    }
1489
1490    /// (b) No overlay file → merge is a no-op (None), shared config untouched.
1491    #[test]
1492    fn no_overlay_file_returns_none() {
1493        let state = test_state();
1494        let home = temp_user_home("empty");
1495        assert!(state
1496            .merge_user_config(&home, Some("[server]\nport = 8080\n"))
1497            .is_none());
1498    }
1499
1500    /// Overlay whose sections are ALL non-allowlisted → no-op (None).
1501    #[test]
1502    fn overlay_with_no_allowlisted_sections_is_noop() {
1503        let state = test_state();
1504        let home = temp_user_home("all-dropped");
1505        std::fs::write(
1506            home.join("config").join("trustee.toml"),
1507            "[server]\nport = 1\n\n[storage]\npath = \"/tmp/x\"\n",
1508        )
1509        .expect("write overlay");
1510        assert!(state
1511            .merge_user_config(&home, Some("[server]\nport = 8080\n"))
1512            .is_none());
1513    }
1514
1515    /// (c) [llm] dropped when allow_llm_overlay=false (default), applied when
1516    /// the instance opted in via with_allow_llm_overlay(true).
1517    #[test]
1518    fn llm_overlay_dropped_by_default_and_kept_when_enabled() {
1519        let shared = "[llm]\nprovider = \"openai\"\n\n[mcp]\nmode = \"shared\"\n";
1520        let overlay = "[llm]\nprovider = \"anthropic\"\n";
1521
1522        // Default: knob false → [llm]-only overlay is a no-op.
1523        let state = test_state();
1524        assert!(!state.allow_llm_overlay);
1525        let home = temp_user_home("llm-off");
1526        std::fs::write(home.join("config").join("trustee.toml"), overlay).expect("write overlay");
1527        assert!(state.merge_user_config(&home, Some(shared)).is_none());
1528
1529        // Opted in: [llm] kept, user value wins; shared [mcp] untouched.
1530        let state = test_state().with_allow_llm_overlay(true);
1531        assert!(state.allow_llm_overlay);
1532        let home = temp_user_home("llm-on");
1533        std::fs::write(home.join("config").join("trustee.toml"), overlay).expect("write overlay");
1534        let merged = state
1535            .merge_user_config(&home, Some(shared))
1536            .expect("[llm] overlay applies when opted in");
1537        let merged_val = parse(&merged);
1538        assert_eq!(merged_val["llm"]["provider"].as_str(), Some("anthropic"));
1539        assert_eq!(merged_val["mcp"]["mode"].as_str(), Some("shared"));
1540    }
1541
1542    /// (d) Web-path home dir resolves through the consolidated
1543    /// trustee_core::user_hash (determinism/known-vectors covered in
1544    /// trustee-core's own tests).
1545    #[test]
1546    fn user_home_dir_uses_consolidated_hash() {
1547        let state = test_state();
1548        if let Some(home) = state.get_user_home_dir("farzan@example.com") {
1549            assert_eq!(
1550                home.file_name().and_then(|n| n.to_str()),
1551                Some(trustee_core::user_hash("farzan@example.com")).as_deref()
1552            );
1553            let users_root = dirs::home_dir().unwrap().join(".trustee").join("users");
1554            assert_eq!(home.parent(), Some(&users_root).map(|p| p.as_path()));
1555        }
1556        // dirs::home_dir() unavailable in sandbox → covered by core tests.
1557    }
1558
1559    /// xagent personas: [lifecycle] survives the overlay allowlist (persona
1560    /// is user-owned) but the WASM extension switch does not.
1561    #[test]
1562    fn lifecycle_overlay_survives_allowlist_enabled_stripped() {
1563        let state = test_state();
1564        let home = temp_user_home("lifecycle-persona");
1565        std::fs::write(
1566            home.join("config").join("trustee.toml"),
1567            "[lifecycle]\nsystem_template = \"You are Saman.\"\nenabled = true\n",
1568        )
1569        .expect("write overlay");
1570
1571        let merged = state
1572            .merge_user_config(
1573                &home,
1574                Some("[lifecycle]\nsystem_template = \"shared default\"\n"),
1575            )
1576            .expect("lifecycle overlay is allowlisted");
1577
1578        let v = merged.parse::<toml::Value>().expect("valid TOML");
1579        assert_eq!(
1580            v["lifecycle"]["system_template"].as_str(),
1581            Some("You are Saman."),
1582            "agent persona wins over shared default"
1583        );
1584        assert_eq!(
1585            v["lifecycle"].get("enabled"),
1586            None,
1587            "WASM extension switch is not user-settable"
1588        );
1589    }
1590
1591    /// 16D integration guard: dev-agent principals (user_key `agent-<name>`)
1592    /// get their OWN session bucket and home-dir namespace — one agent can
1593    /// never see another's sessions, and both go through the same
1594    /// user_hash isolation path as humans.
1595    #[tokio::test]
1596    async fn agent_principals_get_isolated_session_buckets() {
1597        let state = test_state();
1598        let key_a = "agent-farzan";
1599        let key_b = "agent-paydar";
1600
1601        let (sid_a, session_a, _tx_a, _ts_a) = state.ensure_active_session(key_a).await;
1602        let (_sid_b, _session_b, _tx_b, _ts_b) = state.ensure_active_session(key_b).await;
1603
1604        // Each agent resolves its own session; B cannot see A's by id.
1605        assert!(
1606            state.get_session_by_any_id(key_a, &sid_a).await.is_some(),
1607            "owner bucket resolves its own session"
1608        );
1609        assert!(
1610            state.get_session_by_any_id(key_b, &sid_a).await.is_none(),
1611            "cross-agent session access must be 404/None"
1612        );
1613        assert!(Arc::strong_count(&session_a) >= 1);
1614
1615        // Same isolation path as humans: home dir = users_root/user_hash(key).
1616        if let Some(home) = state.get_user_home_dir(key_a) {
1617            assert_eq!(
1618                home.file_name().and_then(|n| n.to_str()),
1619                Some(trustee_core::user_hash(key_a)).as_deref()
1620            );
1621        }
1622    }
1623}
1624
1625// ---------------------------------------------------------------------------
1626// 16C — per-user McpToolLoader cache
1627// ---------------------------------------------------------------------------
1628
1629#[cfg(test)]
1630mod mcp_loader_cache_tests {
1631    use super::*;
1632
1633    fn state_with_shared(shared: &str) -> ServerState {
1634        let (session, _rx) = Session::new();
1635        let (ws_tx, _ws_rx) = tokio::sync::broadcast::channel::<String>(16);
1636        let mut state = ServerState::new(session, ws_tx, None);
1637        state.config_toml = Some(shared.to_string());
1638        state
1639    }
1640
1641    /// Deterministic throwaway user: real user-home path under
1642    /// ~/.trustee/users/{hash} (that IS the resolution path under test),
1643    /// with config/ subdir; caller cleans up.
1644    struct TempUser {
1645        key: String,
1646        home: std::path::PathBuf,
1647    }
1648
1649    impl TempUser {
1650        fn new(tag: &str) -> Self {
1651            let key = format!("16c-{tag}-{}@test.invalid", std::process::id());
1652            let home = dirs::home_dir()
1653                .expect("HOME available in test env")
1654                .join(".trustee")
1655                .join("users")
1656                .join(trustee_core::user_hash(&key));
1657            std::fs::create_dir_all(home.join("config")).expect("create user home");
1658            Self { key, home }
1659        }
1660
1661        fn write_overlay(&self, toml_str: &str) {
1662            std::fs::write(self.home.join("config").join("trustee.toml"), toml_str)
1663                .expect("write overlay");
1664        }
1665    }
1666
1667    impl Drop for TempUser {
1668        fn drop(&mut self) {
1669            let _ = std::fs::remove_dir_all(&self.home);
1670        }
1671    }
1672
1673    #[tokio::test]
1674    async fn no_mcp_config_caches_disabled_marker() {
1675        let state = state_with_shared("[server]\nport = 8080\n");
1676        let user = TempUser::new("nomcp");
1677        let ts = Arc::new(pep::MemoryTokenStore::new());
1678
1679        let first = state.get_or_build_mcp_loader(&user.key, &ts).await.unwrap();
1680        assert!(first.is_none(), "no [mcp] anywhere → disabled marker");
1681        assert_eq!(state.mcp_loaders.len(), 1, "exactly one cache entry");
1682
1683        // Second call is a fingerprint hit — same disabled marker, still one entry.
1684        let second = state.get_or_build_mcp_loader(&user.key, &ts).await.unwrap();
1685        assert!(second.is_none());
1686        assert_eq!(state.mcp_loaders.len(), 1);
1687    }
1688
1689    #[tokio::test]
1690    async fn concurrent_cold_builds_single_flight() {
1691        let state = Arc::new(state_with_shared("[server]\nport = 8080\n"));
1692        let user = TempUser::new("singleflight");
1693        let ts = Arc::new(pep::MemoryTokenStore::new());
1694
1695        let mut handles = Vec::new();
1696        for _ in 0..5 {
1697            let state = state.clone();
1698            let key = user.key.clone();
1699            let ts = ts.clone();
1700            handles.push(tokio::spawn(async move {
1701                state.get_or_build_mcp_loader(&key, &ts).await
1702            }));
1703        }
1704        for h in handles {
1705            h.await.unwrap().expect("all five succeed");
1706        }
1707        assert_eq!(state.mcp_loaders.len(), 1, "single-flight → one entry");
1708    }
1709
1710    #[tokio::test]
1711    async fn fingerprint_change_triggers_rebuild() {
1712        let state = state_with_shared("[server]\nport = 8080\n");
1713        let user = TempUser::new("fpchange");
1714        let ts = Arc::new(pep::MemoryTokenStore::new());
1715
1716        // v1: an enabled MCP server pointing at a blackhole port — abk keeps
1717        // the loader with a DOWN status (connect refused is fast on loopback).
1718        user.write_overlay("[mcp]\nenabled = true\n\n[[mcp.servers]]\nname = \"v1\"\nurl = \"http://127.0.0.1:9/sse\"\n");
1719        let v1 = state.get_or_build_mcp_loader(&user.key, &ts).await.unwrap();
1720        assert!(v1.is_some(), "enabled [mcp] → real loader");
1721        let fp1 = state
1722            .mcp_loaders
1723            .get(&trustee_core::user_hash(&user.key))
1724            .unwrap()
1725            .fingerprint;
1726
1727        // v2: different server URL → different content fingerprint → rebuild.
1728        std::thread::sleep(std::time::Duration::from_millis(5));
1729        user.write_overlay("[mcp]\nenabled = true\n\n[[mcp.servers]]\nname = \"v2\"\nurl = \"http://127.0.0.1:9/other\"\n");
1730        let v2 = state.get_or_build_mcp_loader(&user.key, &ts).await.unwrap();
1731        assert!(v2.is_some());
1732        let entry = state
1733            .mcp_loaders
1734            .get(&trustee_core::user_hash(&user.key))
1735            .unwrap();
1736        assert_ne!(
1737            entry.fingerprint, fp1,
1738            "fingerprint must change with content"
1739        );
1740        assert!(entry.degraded.is_none());
1741
1742        // The old Arc stays valid for in-flight sessions (no panic, no revoke).
1743        let _still_usable = v1.as_ref().unwrap().tool_count;
1744    }
1745
1746    #[tokio::test]
1747    async fn degraded_entry_fails_loud_within_backoff_and_isolates_users() {
1748        let state = state_with_shared("[server]\nport = 8080\n");
1749        let bad = TempUser::new("degraded-bad");
1750        let good = TempUser::new("degraded-good");
1751        let ts = Arc::new(pep::MemoryTokenStore::new());
1752
1753        // [mcp] present but not a table → McpConfig deserialize fails → Err.
1754        bad.write_overlay("[mcp]\nenabled = \"not-a-bool\"\n");
1755        let err = match state.get_or_build_mcp_loader(&bad.key, &ts).await {
1756            Ok(_) => panic!("invalid [mcp] must fail loud"),
1757            Err(e) => e,
1758        };
1759        assert!(
1760            err.contains("invalid [mcp]"),
1761            "surfaces the parse error: {err}"
1762        );
1763
1764        let entry = state
1765            .mcp_loaders
1766            .get(&trustee_core::user_hash(&bad.key))
1767            .unwrap();
1768        assert!(entry.degraded.is_some(), "poison entry recorded");
1769
1770        // Within backoff: fast-fail again (cached error).
1771        let err2 = match state.get_or_build_mcp_loader(&bad.key, &ts).await {
1772            Ok(_) => panic!("still within backoff"),
1773            Err(e) => e,
1774        };
1775        assert_eq!(err, err2, "same cached error");
1776
1777        // Other users are completely unaffected.
1778        let good_loader = state
1779            .get_or_build_mcp_loader(&good.key, &ts)
1780            .await
1781            .expect("other user unaffected");
1782        assert!(
1783            good_loader.is_none(),
1784            "good user has no [mcp] → disabled marker"
1785        );
1786    }
1787}