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
190    /// Encrypted credential authority exposed only through metadata/replace/clear APIs.
191    pub credential_store: Arc<bamboo_config::CredentialStore>,
192
193    /// Shared Remote Cluster Fabric deploy engine (one worker registry across the
194    /// HTTP operator handlers and the `cluster` agent tool).
195    pub fabric_deployer: Arc<bamboo_server_tools::FabricDeployer>,
196
197    /// In-process mailbox bus (broker), when not externally configured. Held so
198    /// it lives for the server's lifetime (dropping it aborts the bus). `None`
199    /// when an external broker is configured or the bus couldn't bind. Never read
200    /// — its only job is to keep the bus task alive until AppState drops.
201    #[allow(dead_code)]
202    embedded_broker: Option<builder::EmbeddedBroker>,
203
204    /// The cluster health monitor sweep. Lives for the server's lifetime (dropping
205    /// it aborts the sweep). `None` when the monitor is disabled
206    /// (`health_interval_secs = 0`). Never read — held only to keep the task alive.
207    #[allow(dead_code)]
208    health_monitor: Option<builder::HealthMonitor>,
209
210    /// Hot-reloadable LLM provider with direct access
211    ///
212    /// This eliminates the proxy pattern where we created an AgentAppState
213    /// that called back to web_service via HTTP. Now we have direct provider access.
214    pub provider: Arc<RwLock<Arc<dyn LLMProvider>>>,
215
216    /// Stable handle that always delegates to the latest provider in `provider`.
217    ///
218    /// This avoids stale provider snapshots after runtime config updates.
219    provider_handle: Arc<dyn LLMProvider>,
220
221    /// Active conversation sessions (in-memory cache)
222    ///
223    /// Maps session IDs to Session objects. Persisted to storage
224    /// via the `storage` field.
225    pub sessions: bamboo_engine::SessionCache,
226
227    /// Persistent storage backend for sessions (V2).
228    ///
229    /// Implemented as folder-per-session with a global `sessions.json` index.
230    pub storage: Arc<dyn Storage>,
231
232    /// Concrete session store implementation (for index/list/cleanup APIs).
233    pub session_store: Arc<SessionStoreV2>,
234
235    /// Per-session write serialisation + metadata-merge persistence layer.
236    ///
237    /// Wraps the same [`Storage`] as `self.storage`, adding per-session
238    /// `Mutex` guards and authoritative-metadata-group merge semantics.
239    /// Use `self.persistence.merge_save_runtime(...)` for any write that
240    /// may race with a UI metadata update.
241    pub persistence: Arc<LockedSessionStore>,
242
243    /// Framework-owned session coordinator (cache + storage + persistence).
244    /// The canonical load/save coordination lives here in `bamboo-engine`, not
245    /// on `AppState`; the inherent `AppState::load_session`/`save_and_cache_session`
246    /// methods now delegate to it. Holds clones of the same `Arc`s as the
247    /// `sessions`/`storage`/`persistence` fields above.
248    pub session_repo: bamboo_engine::SessionRepository,
249
250    /// Background scheduler for async sub-session spawning.
251    pub spawn_scheduler: Arc<SpawnScheduler>,
252
253    /// Coordinates child completion notifications into parent resume.
254    pub child_completion_coordinator: Arc<bamboo_engine::ChildCompletionCoordinator>,
255
256    /// Spawner for the guardian adversarial-review child, injected into each run
257    /// so the terminal gate can create a read-only reviewer (the engine runner
258    /// cannot construct a child directly — see [`bamboo_engine::GuardianSpawner`]).
259    /// Backed by a dedicated [`crate::tools::ChildSessionAdapter`].
260    pub guardian_spawner: Arc<dyn bamboo_engine::GuardianSpawner>,
261
262    /// Bash self-resume hook (issue #84 Phase 2b). Backed by the same
263    /// [`ChildCompletionCoordinator`] that handles child-completion resumes —
264    /// it polls the live shell registry and resumes a session once all its
265    /// background bash shells finish.
266    pub bash_resume_hook: Arc<dyn bamboo_engine::BashResumeHook>,
267
268    /// Schedule store (timed tasks).
269    pub schedule_store: Arc<ScheduleStore>,
270
271    /// Background schedule manager that triggers scheduled runs.
272    pub schedule_manager: Arc<ScheduleManager>,
273
274    /// bamboo-connect manager (#452 / epic #447): owns every configured IM
275    /// platform's long-poll/dispatch background task. Fully inert (zero
276    /// tasks) when `config.connect.platforms` is empty. Held so its tasks
277    /// live for the server's lifetime (`ConnectManager::drop` aborts them).
278    pub connect_manager: Arc<crate::connect::ConnectManager>,
279
280    /// Tool surface factory providing pre-built tool executors for each session type.
281    ///
282    /// Use `state.tools_for(ToolSurface::Root)` for root sessions,
283    /// `state.tools_for(ToolSurface::Child)` for child sessions, etc.
284    pub tool_factory: crate::tools::ToolSurfaceFactory,
285
286    /// Shared tool-execution permission checker — the same `Arc` the tool
287    /// executors use. Retained so request handlers can record session grants
288    /// when the user approves a permission prompt (see the respond handler).
289    pub permission_checker: Arc<dyn bamboo_tools::permission::PermissionChecker>,
290
291    /// Durable, revisioned authority for permission policy. The checker is
292    /// updated only after a successful commit to this section.
293    pub permission_section: Arc<bamboo_tools::permission::PermissionSection>,
294
295    /// Serializes the complete permission commit + live-checker publication.
296    pub permission_io_lock: Arc<tokio::sync::Mutex<()>>,
297    pub approval_registry:
298        bamboo_engine::external_agents::approval_registry::SharedApprovalRegistry,
299
300    /// Backend notification policy service (preferences + dedup + per-session
301    /// relays). Classifies agent events into `AgentEvent::Notification` for
302    /// clients to render; preferences are persisted server-side.
303    pub notification_service: Arc<bamboo_notification::NotificationService>,
304
305    /// Live SSE/WS client-subscriber counts per session (see
306    /// [`watchers::SessionWatchers`]). Used to suppress a redundant desktop
307    /// popup for categories the UI already surfaces while a client is
308    /// actively watching a session.
309    pub session_watchers: Arc<watchers::SessionWatchers>,
310
311    /// Cancellation tokens for in-flight requests
312    ///
313    /// Maps request/session IDs to their cancellation tokens,
314    /// allowing graceful shutdown of long-running operations.
315    pub cancel_tokens: Arc<RwLock<HashMap<String, CancellationToken>>>,
316
317    /// Cancels the supervised MCP proxy service (issue #47) on shutdown so the
318    /// reconnect/backoff supervisor stops cleanly instead of looping forever
319    /// after an intended stop. Unused when no broker is configured.
320    pub mcp_proxy_shutdown: CancellationToken,
321
322    /// Skill manager for prompt-based skill execution
323    ///
324    /// Manages the skill registry and handles skill lookup,
325    /// validation, and execution.
326    pub skill_manager: Arc<SkillManager>,
327
328    /// Durable, recovered workflow-run boundary. It owns the production engine
329    /// plus server-derived session/catalog trust adapters.
330    pub workflow_runs: crate::workflow::WorkflowRunAccess,
331
332    /// MCP server manager for external tool servers
333    ///
334    /// Handles lifecycle of Model Context Protocol servers,
335    /// including initialization, tool discovery, and shutdown.
336    pub mcp_manager: Arc<McpServerManager>,
337
338    /// Supervises long-running "service" plugins (issue #479, prereq for
339    /// epic #477). Always constructed, fully inert until a plugin install
340    /// (or the boot-time reconcile) calls `start_service`. See
341    /// `crate::service_manager`'s module docs.
342    pub service_manager: Arc<crate::service_manager::ServiceManager>,
343
344    /// Handle to the background boot-time service reconcile pass
345    /// (`plugin_installer::boot_reconcile_services`, spawned fire-and-forget
346    /// by `app_state::builder` — see its comment). It is deliberately NOT
347    /// synchronized against `plugin_installer::PLUGIN_OP_LOCK`, so it can, in
348    /// principle, race a `ServerPluginInstaller::install`/
349    /// `stop_services_for_upgrade` call that lands on the SAME data dir
350    /// while it is still in flight (e.g. immediately after construction).
351    /// Production code never touches this field; it exists purely as a
352    /// test-only synchronization point (see
353    /// [`AppState::wait_for_boot_reconcile_services`]) so
354    /// `plugin_installer::tests` can deterministically drain that one-shot
355    /// pass before exercising service install/stop/upgrade, instead of
356    /// racing it under CI scheduling jitter (issue #486).
357    #[doc(hidden)]
358    pub boot_reconcile_services_handle: tokio::sync::Mutex<Option<tokio::task::JoinHandle<()>>>,
359
360    /// Metrics collection and persistence service
361    ///
362    /// Tracks token usage, costs, and performance metrics
363    /// across all sessions.
364    pub metrics_service: Arc<MetricsService>,
365
366    /// Active agent runners indexed by session ID
367    ///
368    /// Each runner manages event broadcasting and cancellation
369    /// for an active agent execution.
370    pub agent_runners: Arc<RwLock<HashMap<String, AgentRunner>>>,
371
372    /// Reference-counted execute handlers still preparing a runner, keyed by
373    /// session. This server-scoped registry closes the durable pending-turn
374    /// expiry race without leaking state across AppState instances/tests.
375    pub(crate) execute_startups: Arc<std::sync::Mutex<HashMap<String, usize>>>,
376
377    /// Session-scoped event streams (long-lived).
378    ///
379    /// Unlike `agent_runners`, these senders exist even when no agent execution is running.
380    /// They are used for:
381    /// - UI subscriptions to `/api/v1/events/{session_id}` (background tasks, etc.)
382    /// - sub-session forwarding (child -> parent)
383    pub session_event_senders: Arc<RwLock<HashMap<String, broadcast::Sender<AgentEvent>>>>,
384
385    /// Account-scoped durable change feed (powers `GET /api/v1/stream`).
386    ///
387    /// Unlike `session_event_senders`, this is a single account-wide sink: all
388    /// durable change events (message appended, session metadata, task updates,
389    /// terminal status) across every session are sequenced, journaled to disk,
390    /// and broadcast here for resumable multi-client sync.
391    pub account_sink: Arc<bamboo_engine::events::AccountEventSink>,
392
393    /// Registry for tracking external processes.
394    pub process_registry: Arc<ProcessRegistry>,
395
396    /// Optional metrics bus for event streaming
397    ///
398    /// When enabled, allows subscribing to metrics events
399    /// in real-time.
400    pub metrics_bus: Option<bamboo_metrics::bus::MetricsBus>,
401
402    /// Unified agent execution runtime holding shared resources.
403    pub agent: Arc<bamboo_engine::Agent>,
404
405    /// Multi-provider registry (used when features.provider_model_ref is enabled).
406    pub provider_registry: Arc<bamboo_llm::ProviderRegistry>,
407
408    /// Provider/model router (used when features.provider_model_ref is enabled).
409    pub provider_router: Arc<bamboo_llm::ProviderModelRouter>,
410
411    /// Unified model catalog service (used when features.provider_model_ref is enabled).
412    pub model_catalog: Arc<bamboo_llm::ModelCatalogService>,
413
414    /// Tracks session ids whose auto-title generation is currently in flight.
415    ///
416    /// Used by [`crate::title_gen`] to dedupe concurrent invocations
417    /// (e.g. execute handler firing while a regenerate-title request is running).
418    pub title_gen_in_flight: Arc<dashmap::DashSet<String>>,
419
420    /// v2-P2 (#181, slice 2): in-memory one-time pairing codes. A 6-digit numeric
421    /// code (keyed by the code string) maps to an entry holding its expiry. Codes
422    /// are PROCESS-EPHEMERAL — never persisted to `config.json`; a restart drops
423    /// all outstanding codes by design. Keyed by `Instant`-based expiry; expired
424    /// entries are purged opportunistically on insert/lookup.
425    pub pairing_codes: Arc<dashmap::DashMap<String, crate::handlers::settings::PairingCodeEntry>>,
426
427    /// v2-P2 (#181, slice 2): per-process brute-force guard for the public
428    /// code-redemption path (`POST /v2/pair { code }`). A 6-digit code is only
429    /// ~1M space, so a public redeem endpoint is brute-forceable without a guard.
430    /// Tracks recent FAILED code-redemption attempts and a cooldown deadline.
431    pub pairing_code_guard: Arc<crate::handlers::settings::PairingCodeGuard>,
432
433    /// #190: per-client-IP brute-force guard for the public root-password
434    /// endpoints (`POST /v1/bamboo/access/verify` and the root-password path of
435    /// `POST /v2/pair`). Tracks recent FAILED root-password attempts per IP and a
436    /// per-key cooldown; loopback/desktop requests are exempted by the handlers
437    /// so the desktop can never lock itself out. PROCESS-EPHEMERAL — never
438    /// persisted; a restart clears all counters.
439    pub root_password_guard: Arc<crate::handlers::settings::RootPasswordGuard>,
440
441    /// Process-ephemeral credentials for Codex children that route model calls
442    /// through this server. Tokens are hashed in memory and revoked at the end
443    /// of their owning actor activation.
444    pub(crate) codex_run_tokens: Arc<crate::codex_run_tokens::CodexRunTokenRegistry>,
445}
446
447impl AppState {
448    /// Try to claim the title-generation slot for `session_id`.
449    /// Returns `true` on success, `false` if generation is already in flight.
450    pub fn title_gen_acquire(&self, session_id: &str) -> bool {
451        self.title_gen_in_flight.insert(session_id.to_string())
452    }
453
454    /// Release the title-generation slot for `session_id`. Idempotent.
455    pub fn title_gen_release(&self, session_id: &str) {
456        self.title_gen_in_flight.remove(session_id);
457    }
458
459    /// Test-only synchronization point (see
460    /// [`boot_reconcile_services_handle`](Self::boot_reconcile_services_handle)'s
461    /// doc comment): wait for the background boot-time service reconcile
462    /// pass to finish. Idempotent — a second call (or a call after
463    /// production code never having populated the handle) is a no-op.
464    #[doc(hidden)]
465    pub async fn wait_for_boot_reconcile_services(&self) {
466        let handle = self.boot_reconcile_services_handle.lock().await.take();
467        if let Some(handle) = handle {
468            let _ = handle.await;
469        }
470    }
471}
472
473mod agent_session_context;
474mod builder;
475mod config_runtime;
476pub(crate) use config_runtime::ConfigLiveHealth;
477pub(crate) use config_runtime::ConfigSectionMutationError;
478pub mod init;
479pub mod parent_approval_reviewer;
480mod persistence;
481mod provider_api;
482pub mod resume_adapter;
483pub mod runner_lifecycle;
484// `pub` (not `pub(crate)`): `ScheduleContext::notification_relay` (a public
485// field of the public `schedule_app::ScheduleContext`) is typed
486// `session_events::NotificationRelayDeps`, so external callers that build a
487// `ScheduleContext` by hand (e.g. integration tests) need to name it.
488pub mod session_events;
489mod session_loader;
490mod tools;
491pub mod watchers;
492
493#[cfg(test)]
494mod tests;
495
496#[derive(Debug, Clone, Copy, Default)]
497pub struct ConfigUpdateEffects {
498    pub reload_provider: bool,
499    pub reconcile_mcp: bool,
500}