Skip to main content

bamboo_server/app_state/
mod.rs

1//! Unified application state management for the Bamboo server
2//!
3//! This module provides the central AppState struct that consolidates all
4//! server state including sessions, storage, LLM providers, tools, and metrics.
5//!
6//! # Architecture
7//!
8//! The AppState uses a unified design that eliminates the proxy pattern where
9//! web_service created an AgentAppState that called back via HTTP. Instead, it
10//! provides direct access to all components.
11//!
12//! ```text
13//! ┌────────────────────────────────────────────────────┐
14//! │              AppState (Unified)                    │
15//! │                                                    │
16//! │  ┌──────────────┐      ┌──────────────┐          │
17//! │  │   Config     │      │   Provider   │          │
18//! │  │  (Hot-reload)│◄────►│   (LLM)      │          │
19//! │  └──────────────┘      └──────────────┘          │
20//! │                                                    │
21//! │  ┌──────────────┐      ┌──────────────┐          │
22//! │  │   Sessions   │      │   Storage    │          │
23//! │  │  (In-memory) │      │  (Persistent)│          │
24//! │  └──────────────┘      └──────────────┘          │
25//! │                                                    │
26//! │  ┌──────────────┐      ┌──────────────┐          │
27//! │  │    Tools     │      │    Skills    │          │
28//! │  │ (Builtin+MCP)│      │   Manager    │          │
29//! │  └──────────────┘      └──────────────┘          │
30//! │                                                    │
31//! │  ┌──────────────┐      ┌──────────────┐          │
32//! │  │     MCP      │      │   Metrics    │          │
33//! │  │   Manager    │      │   Service    │          │
34//! │  └──────────────┘      └──────────────┘          │
35//! └────────────────────────────────────────────────────┘
36//! ```
37//!
38//! # Key Features
39//!
40//! - **Hot-reloadable configuration**: Config and provider can be reloaded at runtime
41//! - **Direct provider access**: No HTTP proxy overhead
42//! - **Session management**: In-memory session cache with persistent storage
43//! - **Tool composition**: Combines built-in and MCP tools
44//! - **Metrics collection**: Integrated metrics and event tracking
45//!
46//! # Usage Example
47//!
48//! ```rust,no_run
49//! use bamboo_server::app_state::AppState;
50//! use std::path::PathBuf;
51//!
52//! #[tokio::main]
53//! async fn main() {
54//!     // Initialize app state
55//!     let app_data_dir = PathBuf::from("/path/to/.bamboo");
56//!     let state = AppState::new(app_data_dir)
57//!         .await
58//!         .expect("failed to initialize app state");
59//!
60//!     // Access components
61//!     let provider = state.get_provider().await;
62//!     let schemas = state.get_all_tool_schemas();
63//!
64//!     // Hot reload configuration
65//!     state.reload_config().await;
66//!     state.reload_provider().await.ok();
67//! }
68//! ```
69
70use std::collections::HashMap;
71use std::path::PathBuf;
72use std::sync::Arc;
73
74use async_trait::async_trait;
75use tokio::sync::{broadcast, RwLock};
76use tokio_util::sync::CancellationToken;
77
78use crate::error::AppError;
79use crate::schedule_app::{ScheduleManager, ScheduleStore};
80use bamboo_agent_core::storage::Storage;
81use bamboo_agent_core::AgentEvent;
82use bamboo_agent_core::{tools::ToolSchema, Message};
83use bamboo_engine::execution::spawn::SpawnScheduler;
84use bamboo_infrastructure::process::registry::ProcessRegistry;
85use bamboo_llm::Config;
86use bamboo_llm::{LLMError, LLMProvider, LLMStream};
87use bamboo_mcp::manager::McpServerManager;
88use bamboo_metrics::metrics_service::MetricsService;
89use bamboo_skills::SkillManager;
90use bamboo_storage::LockedSessionStore;
91use bamboo_storage::SessionStoreV2;
92
93// Context functions moved to bamboo-agent-runtime::context
94pub use bamboo_engine::context::{
95    build_env_prompt_context, build_workspace_prompt_context, workspace_prompt_guidance,
96    DEFAULT_BASE_PROMPT, ENV_CONTEXT_END_MARKER, ENV_CONTEXT_START_MARKER,
97    WORKSPACE_CONTEXT_END_MARKER, WORKSPACE_CONTEXT_PREFIX, WORKSPACE_CONTEXT_START_MARKER,
98};
99
100/// Placeholder provider used when the configured provider cannot be initialized.
101///
102/// This keeps the server usable for configuration/UX flows while ensuring we fail fast
103/// (instead of silently switching to a different provider or model).
104struct UnconfiguredProvider {
105    message: String,
106}
107
108#[async_trait]
109impl LLMProvider for UnconfiguredProvider {
110    async fn chat_stream(
111        &self,
112        _messages: &[Message],
113        _tools: &[ToolSchema],
114        _max_output_tokens: Option<u32>,
115        _model: &str,
116    ) -> bamboo_llm::provider::Result<LLMStream> {
117        Err(LLMError::Auth(format!(
118            "LLM provider is not configured: {}",
119            self.message
120        )))
121    }
122
123    async fn list_models(&self) -> bamboo_llm::provider::Result<Vec<String>> {
124        Err(LLMError::Auth(format!(
125            "LLM provider is not configured: {}",
126            self.message
127        )))
128    }
129}
130
131// Re-export execution types from the runtime crate.
132pub use bamboo_engine::execution::runner_state::{AgentRunner, AgentStatus};
133
134/// Unified application state consolidating web_service and agent/server state
135///
136/// This struct holds all the state needed to run the Bamboo server, including
137/// configuration, LLM providers, sessions, storage, tools, skills, and metrics.
138///
139/// # Design Goals
140///
141/// - **Direct access**: Components are directly accessible without HTTP proxies
142/// - **Hot reload**: Configuration and providers can be reloaded at runtime
143/// - **Thread safety**: Uses Arc<RwLock> for concurrent access
144/// - **Persistence**: Integrates with JsonlStorage for session persistence
145///
146/// # Component Overview
147///
148/// | Component | Purpose | Thread-Safe |
149/// |-----------|---------|--------------|
150/// | `config` | Application configuration | Yes (RwLock) |
151/// | `provider` | Hot-reloadable LLM provider | Yes (RwLock) |
152/// | `sessions` | Active conversation sessions | Yes (RwLock) |
153/// | `storage` | Persistent session storage | Yes (Arc) |
154/// | `tools` | Tool execution (builtin + MCP) | Yes (Arc) |
155/// | `skill_manager` | Skill registry and execution | Yes (Arc) |
156/// | `mcp_manager` | MCP server lifecycle | Yes (Arc) |
157/// | `metrics_service` | Usage metrics collection | Yes (Arc) |
158/// | `agent_runners` | Active agent executions | Yes (RwLock) |
159pub struct AppState {
160    /// Application data directory (configured via `BAMBOO_DATA_DIR`; default `${HOME}/.bamboo`)
161    pub app_data_dir: PathBuf,
162
163    /// Hot-reloadable application configuration
164    ///
165    /// Can be reloaded from disk at runtime using `reload_config()`.
166    pub config: Arc<RwLock<Config>>,
167
168    /// Process-owned modular configuration authority. Production bootstrap
169    /// always installs one after the recoverable legacy split; injected test
170    /// states may omit it and retain the compatibility-only config path.
171    pub config_facade: Option<Arc<bamboo_config::ConfigFacade>>,
172
173    /// Serializes a config WRITE's whole [in-memory mutation + disk persist] with
174    /// a `reload_config`'s [disk read + in-memory swap], so a reload can never
175    /// observe an in-flight-but-not-yet-persisted update and clobber it with the
176    /// stale disk copy (the residual of #41). It is NOT the `config` RwLock —
177    /// using a separate mutex keeps config READERS (the hot agent-loop path)
178    /// unblocked during a write's disk I/O. #126.
179    pub config_io_lock: Arc<tokio::sync::Mutex<()>>,
180
181    /// Server-owned live configuration watcher and its health envelope.
182    /// The runtime handle keeps the directory watcher tasks alive.
183    pub config_live_health: Arc<std::sync::RwLock<config_runtime::ConfigLiveHealth>>,
184    /// MCP section health is independent from provider health so an invalid or
185    /// degraded MCP candidate cannot make unrelated sections appear unhealthy.
186    pub mcp_config_live_health: Arc<std::sync::RwLock<config_runtime::ConfigLiveHealth>>,
187    #[allow(dead_code)]
188    config_watcher: config_runtime::ConfigWatcherRuntime,
189    /// Project shared-resource watcher. Held for the server lifetime.
190    #[allow(dead_code)]
191    pub(crate) project_resource_watcher: project_watcher::ProjectResourceWatcher,
192
193    /// Encrypted credential authority exposed only through metadata/replace/clear APIs.
194    pub credential_store: Arc<bamboo_config::CredentialStore>,
195
196    /// Shared Remote Cluster Fabric deploy engine (one worker registry across the
197    /// HTTP operator handlers and the `cluster` agent tool).
198    pub fabric_deployer: Arc<bamboo_server_tools::FabricDeployer>,
199
200    /// In-process mailbox bus (broker), when not externally configured. Held so
201    /// it lives for the server's lifetime (dropping it aborts the bus). `None`
202    /// when an external broker is configured or the bus couldn't bind. Never read
203    /// — its only job is to keep the bus task alive until AppState drops.
204    #[allow(dead_code)]
205    embedded_broker: Option<builder::EmbeddedBroker>,
206
207    /// The cluster health monitor sweep. Lives for the server's lifetime (dropping
208    /// it aborts the sweep). `None` when the monitor is disabled
209    /// (`health_interval_secs = 0`). Never read — held only to keep the task alive.
210    #[allow(dead_code)]
211    health_monitor: Option<builder::HealthMonitor>,
212
213    /// Hot-reloadable LLM provider with direct access
214    ///
215    /// This eliminates the proxy pattern where we created an AgentAppState
216    /// that called back to web_service via HTTP. Now we have direct provider access.
217    pub provider: Arc<RwLock<Arc<dyn LLMProvider>>>,
218
219    /// Stable handle that always delegates to the latest provider in `provider`.
220    ///
221    /// This avoids stale provider snapshots after runtime config updates.
222    provider_handle: Arc<dyn LLMProvider>,
223
224    /// Active conversation sessions (in-memory cache)
225    ///
226    /// Maps session IDs to Session objects. Persisted to storage
227    /// via the `storage` field.
228    pub sessions: bamboo_engine::SessionCache,
229
230    /// Persistent storage backend for sessions (V2).
231    ///
232    /// Implemented as folder-per-session with a global `sessions.json` index.
233    pub storage: Arc<dyn Storage>,
234
235    /// Concrete session store implementation (for index/list/cleanup APIs).
236    pub session_store: Arc<SessionStoreV2>,
237
238    /// Authoritative first-class Project registry and shared-resource paths.
239    pub project_store: Arc<bamboo_projects::ProjectStore>,
240
241    /// Redacted adapter used by HTTP creation paths and the agent runtime to
242    /// resolve one authoritative Project/workspace identity.
243    pub project_context_resolver: Arc<bamboo_engine::project_context::ProjectContextResolver>,
244
245    /// Per-session write serialisation + metadata-merge persistence layer.
246    ///
247    /// Wraps the same [`Storage`] as `self.storage`, adding per-session
248    /// `Mutex` guards and authoritative-metadata-group merge semantics.
249    /// Use `self.persistence.merge_save_runtime(...)` for any write that
250    /// may race with a UI metadata update.
251    pub persistence: Arc<LockedSessionStore>,
252
253    /// Durable logical-session delivery plane. These are internal runtime
254    /// capabilities; no public messaging endpoint is registered.
255    pub session_inbox: Arc<dyn bamboo_domain::SessionInboxPort>,
256    pub session_activation_router: Arc<bamboo_engine::SessionActivationRouter>,
257    pub session_messenger: Arc<bamboo_engine::SessionMessenger>,
258
259    /// Framework-owned session coordinator (cache + storage + persistence).
260    /// The canonical load/save coordination lives here in `bamboo-engine`, not
261    /// on `AppState`; the inherent `AppState::load_session`/`save_and_cache_session`
262    /// methods now delegate to it. Holds clones of the same `Arc`s as the
263    /// `sessions`/`storage`/`persistence` fields above.
264    pub session_repo: bamboo_engine::SessionRepository,
265
266    /// Background scheduler for async sub-session spawning.
267    pub spawn_scheduler: Arc<SpawnScheduler>,
268
269    /// Coordinates child completion notifications into parent resume.
270    pub child_completion_coordinator: Arc<bamboo_engine::ChildCompletionCoordinator>,
271
272    /// Spawner for the guardian adversarial-review child, injected into each run
273    /// so the terminal gate can create a read-only reviewer (the engine runner
274    /// cannot construct a child directly — see [`bamboo_engine::GuardianSpawner`]).
275    /// Backed by a dedicated [`crate::tools::ChildSessionAdapter`].
276    pub guardian_spawner: Arc<dyn bamboo_engine::GuardianSpawner>,
277
278    /// Bash self-resume hook (issue #84 Phase 2b). Backed by the same
279    /// [`ChildCompletionCoordinator`] that handles child-completion resumes —
280    /// it polls the live shell registry and resumes a session once all its
281    /// background bash shells finish.
282    pub bash_resume_hook: Arc<dyn bamboo_engine::BashResumeHook>,
283
284    /// Schedule store (timed tasks).
285    pub schedule_store: Arc<ScheduleStore>,
286
287    /// Background schedule manager that triggers scheduled runs.
288    pub schedule_manager: Arc<ScheduleManager>,
289
290    /// bamboo-connect manager (#452 / epic #447): owns every configured IM
291    /// platform's long-poll/dispatch background task. Fully inert (zero
292    /// tasks) when `config.connect.platforms` is empty. Held so its tasks
293    /// live for the server's lifetime (`ConnectManager::drop` aborts them).
294    pub connect_manager: Arc<crate::connect::ConnectManager>,
295
296    /// Tool surface factory providing pre-built tool executors for each session type.
297    ///
298    /// Use `state.tools_for(ToolSurface::Root)` for root sessions,
299    /// `state.tools_for(ToolSurface::Child)` for child sessions, etc.
300    pub tool_factory: crate::tools::ToolSurfaceFactory,
301
302    /// Shared tool-execution permission checker — the same `Arc` the tool
303    /// executors use. Retained so request handlers can record session grants
304    /// when the user approves a permission prompt (see the respond handler).
305    pub permission_checker: Arc<dyn bamboo_tools::permission::PermissionChecker>,
306
307    /// Durable, revisioned authority for permission policy. The checker is
308    /// updated only after a successful commit to this section.
309    pub permission_section: Arc<bamboo_tools::permission::PermissionSection>,
310
311    /// Serializes the complete permission commit + live-checker publication.
312    pub permission_io_lock: Arc<tokio::sync::Mutex<()>>,
313    pub approval_registry:
314        bamboo_engine::external_agents::approval_registry::SharedApprovalRegistry,
315
316    /// Backend notification policy service (preferences + dedup + per-session
317    /// relays). Classifies agent events into `AgentEvent::Notification` for
318    /// clients to render; preferences are persisted server-side.
319    pub notification_service: Arc<bamboo_notification::NotificationService>,
320
321    /// Live SSE/WS client-subscriber counts per session (see
322    /// [`watchers::SessionWatchers`]). Used to suppress a redundant desktop
323    /// popup for categories the UI already surfaces while a client is
324    /// actively watching a session.
325    pub session_watchers: Arc<watchers::SessionWatchers>,
326
327    /// Cancellation tokens for in-flight requests
328    ///
329    /// Maps request/session IDs to their cancellation tokens,
330    /// allowing graceful shutdown of long-running operations.
331    pub cancel_tokens: Arc<RwLock<HashMap<String, CancellationToken>>>,
332
333    /// Cancels the supervised MCP proxy service (issue #47) on shutdown so the
334    /// reconnect/backoff supervisor stops cleanly instead of looping forever
335    /// after an intended stop. Unused when no broker is configured.
336    pub mcp_proxy_shutdown: CancellationToken,
337
338    /// Skill manager for prompt-based skill execution
339    ///
340    /// Manages the skill registry and handles skill lookup,
341    /// validation, and execution.
342    pub skill_manager: Arc<SkillManager>,
343
344    /// Durable, recovered workflow-run boundary. It owns the production engine
345    /// plus server-derived session/catalog trust adapters.
346    pub workflow_runs: crate::workflow::WorkflowRunAccess,
347
348    /// MCP server manager for external tool servers
349    ///
350    /// Handles lifecycle of Model Context Protocol servers,
351    /// including initialization, tool discovery, and shutdown.
352    pub mcp_manager: Arc<McpServerManager>,
353
354    /// Supervises long-running "service" plugins (issue #479, prereq for
355    /// epic #477). Always constructed, fully inert until a plugin install
356    /// (or the boot-time reconcile) calls `start_service`. See
357    /// `crate::service_manager`'s module docs.
358    pub service_manager: Arc<crate::service_manager::ServiceManager>,
359
360    /// Handle to the background boot-time service reconcile pass
361    /// (`plugin_installer::boot_reconcile_services`, spawned fire-and-forget
362    /// by `app_state::builder` — see its comment). It is deliberately NOT
363    /// synchronized against `plugin_installer::PLUGIN_OP_LOCK`, so it can, in
364    /// principle, race a `ServerPluginInstaller::install`/
365    /// `stop_services_for_upgrade` call that lands on the SAME data dir
366    /// while it is still in flight (e.g. immediately after construction).
367    /// Production code never touches this field; it exists purely as a
368    /// test-only synchronization point (see
369    /// [`AppState::wait_for_boot_reconcile_services`]) so
370    /// `plugin_installer::tests` can deterministically drain that one-shot
371    /// pass before exercising service install/stop/upgrade, instead of
372    /// racing it under CI scheduling jitter (issue #486).
373    #[doc(hidden)]
374    pub boot_reconcile_services_handle: tokio::sync::Mutex<Option<tokio::task::JoinHandle<()>>>,
375
376    /// Metrics collection and persistence service
377    ///
378    /// Tracks token usage, costs, and performance metrics
379    /// across all sessions.
380    pub metrics_service: Arc<MetricsService>,
381
382    /// Active agent runners indexed by session ID
383    ///
384    /// Each runner manages event broadcasting and cancellation
385    /// for an active agent execution.
386    pub agent_runners: Arc<RwLock<HashMap<String, AgentRunner>>>,
387
388    /// Reference-counted execute handlers still preparing a runner, keyed by
389    /// session. This server-scoped registry closes the durable pending-turn
390    /// expiry race without leaking state across AppState instances/tests.
391    pub(crate) execute_startups: Arc<std::sync::Mutex<HashMap<String, usize>>>,
392
393    /// Session-scoped event streams (long-lived).
394    ///
395    /// Unlike `agent_runners`, these senders exist even when no agent execution is running.
396    /// They are used for:
397    /// - UI subscriptions to `/api/v1/events/{session_id}` (background tasks, etc.)
398    /// - sub-session forwarding (child -> parent)
399    pub session_event_senders: Arc<RwLock<HashMap<String, broadcast::Sender<AgentEvent>>>>,
400
401    /// Account-scoped durable change feed (powers `GET /api/v1/stream`).
402    ///
403    /// Unlike `session_event_senders`, this is a single account-wide sink: all
404    /// durable change events (message appended, session metadata, task updates,
405    /// terminal status) across every session are sequenced, journaled to disk,
406    /// and broadcast here for resumable multi-client sync.
407    pub account_sink: Arc<bamboo_engine::events::AccountEventSink>,
408
409    /// Registry for tracking external processes.
410    pub process_registry: Arc<ProcessRegistry>,
411
412    /// Optional metrics bus for event streaming
413    ///
414    /// When enabled, allows subscribing to metrics events
415    /// in real-time.
416    pub metrics_bus: Option<bamboo_metrics::bus::MetricsBus>,
417
418    /// Unified agent execution runtime holding shared resources.
419    pub agent: Arc<bamboo_engine::Agent>,
420
421    /// Multi-provider registry (used when features.provider_model_ref is enabled).
422    pub provider_registry: Arc<bamboo_llm::ProviderRegistry>,
423
424    /// Provider/model router (used when features.provider_model_ref is enabled).
425    pub provider_router: Arc<bamboo_llm::ProviderModelRouter>,
426
427    /// Unified model catalog service (used when features.provider_model_ref is enabled).
428    pub model_catalog: Arc<bamboo_llm::ModelCatalogService>,
429
430    /// Tracks session ids whose auto-title generation is currently in flight.
431    ///
432    /// Used by [`crate::title_gen`] to dedupe concurrent invocations
433    /// (e.g. execute handler firing while a regenerate-title request is running).
434    pub title_gen_in_flight: Arc<dashmap::DashSet<String>>,
435
436    /// v2-P2 (#181, slice 2): in-memory one-time pairing codes. A 6-digit numeric
437    /// code (keyed by the code string) maps to an entry holding its expiry. Codes
438    /// are PROCESS-EPHEMERAL — never persisted to `config.json`; a restart drops
439    /// all outstanding codes by design. Keyed by `Instant`-based expiry; expired
440    /// entries are purged opportunistically on insert/lookup.
441    pub pairing_codes: Arc<dashmap::DashMap<String, crate::handlers::settings::PairingCodeEntry>>,
442
443    /// v2-P2 (#181, slice 2): per-process brute-force guard for the public
444    /// code-redemption path (`POST /v2/pair { code }`). A 6-digit code is only
445    /// ~1M space, so a public redeem endpoint is brute-forceable without a guard.
446    /// Tracks recent FAILED code-redemption attempts and a cooldown deadline.
447    pub pairing_code_guard: Arc<crate::handlers::settings::PairingCodeGuard>,
448
449    /// #190: per-client-IP brute-force guard for the public root-password
450    /// endpoints (`POST /v1/bamboo/access/verify` and the root-password path of
451    /// `POST /v2/pair`). Tracks recent FAILED root-password attempts per IP and a
452    /// per-key cooldown; loopback/desktop requests are exempted by the handlers
453    /// so the desktop can never lock itself out. PROCESS-EPHEMERAL — never
454    /// persisted; a restart clears all counters.
455    pub root_password_guard: Arc<crate::handlers::settings::RootPasswordGuard>,
456
457    /// Process-ephemeral credentials for Codex children that route model calls
458    /// through this server. Tokens are hashed in memory and revoked at the end
459    /// of their owning actor activation.
460    pub(crate) codex_run_tokens: Arc<crate::codex_run_tokens::CodexRunTokenRegistry>,
461}
462
463impl AppState {
464    /// Try to claim the title-generation slot for `session_id`.
465    /// Returns `true` on success, `false` if generation is already in flight.
466    pub fn title_gen_acquire(&self, session_id: &str) -> bool {
467        self.title_gen_in_flight.insert(session_id.to_string())
468    }
469
470    /// Release the title-generation slot for `session_id`. Idempotent.
471    pub fn title_gen_release(&self, session_id: &str) {
472        self.title_gen_in_flight.remove(session_id);
473    }
474
475    /// Test-only synchronization point (see
476    /// [`boot_reconcile_services_handle`](Self::boot_reconcile_services_handle)'s
477    /// doc comment): wait for the background boot-time service reconcile
478    /// pass to finish. Idempotent — a second call (or a call after
479    /// production code never having populated the handle) is a no-op.
480    #[doc(hidden)]
481    pub async fn wait_for_boot_reconcile_services(&self) {
482        let handle = self.boot_reconcile_services_handle.lock().await.take();
483        if let Some(handle) = handle {
484            let _ = handle.await;
485        }
486    }
487}
488
489mod agent_session_context;
490mod builder;
491mod config_runtime;
492pub(crate) use config_runtime::ConfigLiveHealth;
493pub(crate) use config_runtime::ConfigSectionMutationError;
494pub mod init;
495pub mod parent_approval_reviewer;
496mod persistence;
497mod project_watcher;
498mod provider_api;
499pub mod resume_adapter;
500pub mod runner_lifecycle;
501// `pub` (not `pub(crate)`): `ScheduleContext::notification_relay` (a public
502// field of the public `schedule_app::ScheduleContext`) is typed
503// `session_events::NotificationRelayDeps`, so external callers that build a
504// `ScheduleContext` by hand (e.g. integration tests) need to name it.
505pub mod session_events;
506mod session_loader;
507mod tools;
508pub mod watchers;
509
510#[cfg(test)]
511mod tests;
512
513#[derive(Debug, Clone, Copy, Default)]
514pub struct ConfigUpdateEffects {
515    pub reload_provider: bool,
516    pub reconcile_mcp: bool,
517}