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}