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}