#![deny(rustdoc::broken_intra_doc_links)]
use crate::session_store_cache::SessionStoreCache;
use anyhow::Result;
use serde_json::{json, Value};
use std::net::SocketAddr;
use std::path::PathBuf;
use std::sync::atomic::{AtomicU8, AtomicUsize, Ordering};
use std::sync::{Arc, OnceLock};
use tokio::sync::{broadcast, OnceCell, RwLock};
use trusty_common::memory_core::embed::Embedder;
use trusty_common::memory_core::{store::ChatSessionStore, PalaceRegistry};
use trusty_common::ChatProvider;
use trusty_mcp::initialize_response;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum DaemonReadiness {
Warming = 0,
Ready = 1,
}
impl DaemonReadiness {
pub fn from_u8(v: u8) -> Self {
if v == 0 {
Self::Warming
} else {
Self::Ready
}
}
}
pub mod activity;
pub mod attribution;
pub mod authz;
pub mod bm25_backfill;
pub mod bm25_index;
pub mod bm25_lane;
pub mod bm25_repair;
pub mod bootstrap;
pub mod dream_scheduler;
pub mod exit_runtime;
pub mod fd_metrics;
pub mod idle_evict;
pub mod worker_liveness;
pub mod lock_stall;
#[cfg(feature = "daemon")]
pub mod chat;
pub mod client;
pub mod commands;
pub mod console_metrics;
pub mod discovery;
mod events;
pub mod hook_emit;
pub mod kg_extract;
pub mod kg_write;
pub mod mcp_service;
pub mod messaging;
pub mod openrpc;
pub mod palace_id_derive;
pub mod palace_last_used;
pub mod project_root;
pub mod prompt_facts;
pub mod prompt_log;
pub mod service;
pub mod session_store_cache;
pub mod startup_budget;
pub mod startup_scan;
#[cfg(all(test, feature = "daemon"))]
pub(crate) mod test_daemon;
pub mod tools;
pub mod transport;
pub mod wordnet_pos;
pub use activity::{ActivityEntry, ActivityFilter, ActivityLog, ActivitySource};
pub use attribution::{CreatorInfo, CreatorSource};
pub(crate) use events::open_activity_log_with_fallback;
#[cfg(test)]
pub(crate) use events::open_activity_log_with_fallback_in;
pub use events::{DaemonEvent, HookType, InjectionKind};
#[cfg(feature = "daemon")]
pub use transport::serve;
pub use transport::socket_path;
#[inline]
pub fn is_data_dir_override_active() -> bool {
matches!(
std::env::var(trusty_common::DATA_DIR_OVERRIDE_ENV),
Ok(v) if !v.trim().is_empty()
)
}
pub const HOOK_PROMPT_EXCERPT_CHARS: usize = 80;
pub fn hook_prompt_excerpt(prompt: &str) -> String {
let normalised: String = prompt.split_whitespace().collect::<Vec<_>>().join(" ");
if normalised.chars().count() <= HOOK_PROMPT_EXCERPT_CHARS {
normalised
} else {
let kept: String = normalised
.chars()
.take(HOOK_PROMPT_EXCERPT_CHARS.saturating_sub(1))
.collect();
format!("{kept}…")
}
}
pub use mcp_service::MemoryMcpService;
pub use tools::MemoryMcpServer;
pub fn resolve_palace_registry_dir(data_dir: PathBuf) -> PathBuf {
trusty_common::palace_alias::palace_registry_dir_from(data_dir)
}
#[derive(Clone)]
pub struct AppState {
pub version: String,
pub machine: trusty_common::machine_tier::MachineBudget,
pub registry: Arc<PalaceRegistry>,
pub data_root: PathBuf,
pub default_palace: Option<String>,
pub chat_provider: Arc<OnceCell<Option<Arc<dyn ChatProvider>>>>,
pub session_stores: Arc<SessionStoreCache>,
pub events: Arc<broadcast::Sender<DaemonEvent>>,
pub started_at: std::time::Instant,
pub log_buffer: trusty_common::log_buffer::LogBuffer,
pub error_store: Option<trusty_common::error_capture::ErrorStore>,
pub multi_tenant_mode: bool,
pub disk_bytes: Arc<std::sync::atomic::AtomicU64>,
pub sys_metrics: Arc<tokio::sync::Mutex<trusty_common::sys_metrics::SysMetrics>>,
pub bound_addr: Arc<OnceLock<SocketAddr>>,
pub prompt_context_cache: Arc<RwLock<prompt_facts::PromptFactsCache>>,
pub tier_s_admission_lock: Arc<tokio::sync::Mutex<()>>,
pub activity_log: Arc<ActivityLog>,
pub(crate) bm25: Option<Arc<bm25_lane::Bm25Lane>>,
pub palace_write_locks: Arc<dashmap::DashMap<String, Arc<tokio::sync::Mutex<()>>>>,
pub pending_activity_writes: Arc<AtomicUsize>,
pub worker_liveness: Arc<worker_liveness::WorkerLiveness>,
pub wedge_threshold: std::time::Duration,
pub lock_stalls: Arc<lock_stall::LockStallTracker>,
pub palace_names: Arc<dashmap::DashMap<String, String>>,
pub startup_gate: crate::startup_budget::StartupOpenGate,
pub pin_project_map: Arc<dashmap::DashMap<String, PathBuf>>,
pub bm25_index_tx: tokio::sync::mpsc::Sender<tools::Bm25IndexRequest>,
pub bm25_dirty: bm25_repair::DirtyPalaces,
pub update_available: Arc<std::sync::Mutex<Option<String>>>,
pub daemon_readiness: Arc<AtomicU8>,
pub write_op_budget: std::time::Duration,
pub palace_last_used: palace_last_used::StampCache,
pub write_pipeline_budget: std::time::Duration,
}
impl AppState {
pub fn new(data_root: PathBuf) -> Self {
let (events_tx, _) = broadcast::channel::<DaemonEvent>(128);
let activity_log = open_activity_log_with_fallback(&data_root);
let (bm25_index_tx, bm25_index_rx) =
tokio::sync::mpsc::channel::<tools::Bm25IndexRequest>(tools::BM25_INDEX_QUEUE_CAPACITY);
let bm25_dirty: bm25_repair::DirtyPalaces = Arc::new(dashmap::DashSet::new());
tools::spawn_bm25_index_worker(bm25_index_rx, None, Arc::clone(&bm25_dirty));
Self {
version: env!("CARGO_PKG_VERSION").to_string(),
machine: trusty_common::machine_tier::MachineBudget::detect(),
registry: Arc::new(PalaceRegistry::from_env()),
data_root,
default_palace: None,
chat_provider: Arc::new(OnceCell::new()),
session_stores: Arc::new(SessionStoreCache::from_env()),
events: Arc::new(events_tx),
started_at: std::time::Instant::now(),
log_buffer: trusty_common::log_buffer::LogBuffer::new(
trusty_common::log_buffer::DEFAULT_LOG_CAPACITY,
),
error_store: None,
multi_tenant_mode: false,
disk_bytes: Arc::new(std::sync::atomic::AtomicU64::new(0)),
sys_metrics: Arc::new(tokio::sync::Mutex::new(
trusty_common::sys_metrics::SysMetrics::new(),
)),
bound_addr: Arc::new(OnceLock::new()),
prompt_context_cache: Arc::new(RwLock::new(prompt_facts::PromptFactsCache::default())),
tier_s_admission_lock: Arc::new(tokio::sync::Mutex::new(())),
activity_log,
bm25: None,
palace_write_locks: Arc::new(dashmap::DashMap::new()),
pending_activity_writes: Arc::new(AtomicUsize::new(0)),
worker_liveness: Arc::new(worker_liveness::WorkerLiveness::new()),
wedge_threshold: worker_liveness::wedge_threshold(),
lock_stalls: Arc::new(lock_stall::LockStallTracker::default()),
palace_names: Arc::new(dashmap::DashMap::new()),
startup_gate: crate::startup_budget::StartupOpenGate::from_env(),
pin_project_map: Arc::new(dashmap::DashMap::new()),
bm25_index_tx,
bm25_dirty,
update_available: Arc::new(std::sync::Mutex::new(None)),
daemon_readiness: Arc::new(AtomicU8::new(DaemonReadiness::Warming as u8)),
write_op_budget: trusty_common::memory_core::timeouts::write_op_budget(),
palace_last_used: Arc::new(dashmap::DashMap::new()),
write_pipeline_budget: trusty_common::memory_core::timeouts::write_pipeline_timeout(),
}
}
#[must_use]
pub fn with_write_op_budget(mut self, budget: std::time::Duration) -> Self {
self.write_op_budget = budget;
self
}
#[must_use]
pub fn with_write_pipeline_budget(mut self, budget: std::time::Duration) -> Self {
self.write_pipeline_budget = budget;
self
}
pub fn palace_write_lock(&self, palace_id: &str) -> Arc<tokio::sync::Mutex<()>> {
if let Some(existing) = self.palace_write_locks.get(palace_id) {
return existing.clone();
}
self.palace_write_locks
.entry(palace_id.to_string())
.or_insert_with(|| Arc::new(tokio::sync::Mutex::new(())))
.clone()
}
pub fn pinned_project_path(&self, palace_id: &str) -> Option<PathBuf> {
self.pin_project_map.get(palace_id).map(|e| e.clone())
}
#[must_use]
pub fn with_bm25_lane_from_env(self) -> Self {
if std::env::var("TRUSTY_BM25_DAEMON").as_deref() != Ok("1") {
return self;
}
let lane = bm25_lane::Bm25Lane::new(self.data_root.clone());
tracing::info!(
max_resident = lane.max_resident(),
text_budget_bytes = ?lane.text_budget_bytes(),
"in-process BM25 lane enabled (TRUSTY_BM25_DAEMON=1)"
);
self.with_bm25_lane(lane)
}
#[must_use]
pub fn with_bm25_lane(mut self, lane: Arc<bm25_lane::Bm25Lane>) -> Self {
let (tx, rx) =
tokio::sync::mpsc::channel::<tools::Bm25IndexRequest>(tools::BM25_INDEX_QUEUE_CAPACITY);
tools::spawn_bm25_index_worker(rx, Some(Arc::clone(&lane)), Arc::clone(&self.bm25_dirty));
self.bm25_index_tx = tx;
self.bm25 = Some(lane);
self
}
pub fn bm25_lane(&self) -> Option<&Arc<bm25_lane::Bm25Lane>> {
self.bm25.as_ref()
}
pub async fn load_palaces_from_disk(&self) -> Result<usize> {
let registry_dir = self.data_root.clone();
let registry = self.registry.clone();
let palace_names = self.palace_names.clone();
let gate = self.startup_gate.clone();
let palaces =
tokio::task::spawn_blocking(move || PalaceRegistry::list_palaces(®istry_dir))
.await
.map_err(|e| anyhow::anyhow!("join list_palaces: {e}"))??;
let total = palaces.len();
let mut loaded = 0usize;
let mut skipped = 0usize;
for palace in palaces {
let permit = gate.acquire().await;
let registry = Arc::clone(®istry);
let palace_names = Arc::clone(&palace_names);
let opened = tokio::task::spawn_blocking(move || -> bool {
let _permit = permit;
match registry.open_handle(&palace) {
Ok(handle) => {
tracing::debug!(
palace = %palace.id,
data_dir = %palace.data_dir.display(),
"loaded palace from disk"
);
palace_names.insert(palace.id.0.clone(), palace.name.clone());
registry.register_arc(handle);
true
}
Err(e) => {
tracing::warn!(
palace = %palace.id,
data_dir = %palace.data_dir.display(),
"skipping palace during startup hydration: {e:#}; \
will retry lazily on first access"
);
registry.record_unopenable(palace.id.clone(), format!("{e:#}"));
false
}
}
})
.await;
match opened {
Ok(true) => loaded += 1,
Ok(false) => skipped += 1,
Err(e) => {
tracing::warn!("palace hydration task failed: {e}");
skipped += 1;
}
}
}
tracing::info!(
open_limit = gate.limit(),
peak_concurrent_opens = gate.peak_concurrent(),
"palace hydration summary: loaded {loaded}/{total} ({skipped} skipped due to errors)"
);
Ok(loaded)
}
#[must_use]
pub fn with_log_buffer(mut self, buffer: trusty_common::log_buffer::LogBuffer) -> Self {
self.log_buffer = buffer;
self
}
#[must_use]
pub fn with_writer_intent(mut self) -> Self {
debug_assert!(self.registry.is_empty() && Arc::strong_count(&self.registry) == 1);
let lease = trusty_common::memory_core::MaintenanceLease::new(&self.data_root);
self.registry = Arc::new(
PalaceRegistry::from_env()
.with_writer_intent()
.with_maintenance_lease(Arc::new(lease)),
);
self
}
#[must_use]
pub fn with_error_store(mut self, store: trusty_common::error_capture::ErrorStore) -> Self {
self.error_store = Some(store);
self
}
#[must_use]
pub fn with_multi_tenant_mode_from_env(mut self) -> Self {
self.multi_tenant_mode = std::env::var("TRUSTY_MEMORY_MULTI_TENANT").as_deref() == Ok("1");
if self.multi_tenant_mode {
tracing::info!(
"multi-tenant mode enabled (TRUSTY_MEMORY_MULTI_TENANT=1): force=true palace_create will be refused"
);
}
self
}
pub fn emit(&self, event: DaemonEvent) {
if let Some(source) = event.source() {
let event_type = event.type_str();
let palace_id = event.palace_id().map(|s| s.to_string());
let log = Arc::clone(&self.activity_log);
let event_for_log = event.clone();
let pending = Arc::clone(&self.pending_activity_writes);
let id = log.alloc_id();
pending.fetch_add(1, Ordering::SeqCst);
tokio::task::spawn_blocking(move || {
let result = log.append_with_id(id, source, palace_id, event_type, &event_for_log);
if let Err(e) = result {
tracing::warn!("activity_log.append failed for {event_type}: {e:#}");
}
pending.fetch_sub(1, Ordering::SeqCst);
});
}
let _ = self.events.send(event);
}
pub async fn flush_activity_writes(&self) {
while self.pending_activity_writes.load(Ordering::SeqCst) > 0 {
tokio::time::sleep(std::time::Duration::from_millis(1)).await;
}
}
pub fn session_store(&self, palace_id: &str) -> Result<Arc<ChatSessionStore>> {
self.session_stores
.get_or_open(palace_id, &self.data_root.join(palace_id))
}
pub fn with_default_palace(mut self, name: Option<String>) -> Self {
self.default_palace = name;
self
}
pub async fn chat_provider(&self) -> Option<Arc<dyn ChatProvider>> {
self.chat_provider
.get_or_init(|| async {
let cfg = crate::service::load_user_config().unwrap_or_default();
if cfg.local_model.enabled {
if let Some(mut p) =
trusty_common::auto_detect_local_provider(&cfg.local_model.base_url).await
{
p.model = cfg.local_model.model.clone();
return Some(Arc::new(p) as Arc<dyn ChatProvider>);
}
}
if !cfg.openrouter_api_key.is_empty() {
return Some(Arc::new(trusty_common::OpenRouterProvider::new(
cfg.openrouter_api_key,
cfg.openrouter_model,
)) as Arc<dyn ChatProvider>);
}
None
})
.await
.clone()
}
pub fn spawn_alias_discovery(&self, palace: String, project_root: PathBuf) {
let state = self.clone();
tokio::spawn(async move {
let args = serde_json::json!({
"palace": palace,
"project_root": project_root.to_string_lossy(),
});
match tools::dispatch_tool(&state, "discover_aliases", args).await {
Ok(result) => tracing::info!(
new = ?result.get("new"),
already_known = ?result.get("already_known"),
"alias discovery complete"
),
Err(e) => tracing::warn!("alias discovery failed: {e:#}"),
}
});
}
pub fn readiness(&self) -> DaemonReadiness {
DaemonReadiness::from_u8(self.daemon_readiness.load(Ordering::Acquire))
}
pub fn set_ready(&self) {
self.daemon_readiness
.store(DaemonReadiness::Ready as u8, Ordering::Release);
}
pub async fn embedder(&self) -> Result<Arc<dyn Embedder + Send + Sync>> {
let embedder = trusty_common::memory_core::retrieval::shared_embedder().await?;
self.set_ready();
Ok(embedder)
}
}
impl std::fmt::Debug for AppState {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("AppState")
.field("version", &self.version)
.field("data_root", &self.data_root)
.field("registry_len", &self.registry.len())
.finish()
}
}
pub async fn handle_message(state: &AppState, msg: Value) -> Value {
let id = msg.get("id").cloned().unwrap_or(Value::Null);
let method = msg.get("method").and_then(|m| m.as_str()).unwrap_or("");
match method {
"initialize" => {
let extra = state
.default_palace
.as_ref()
.map(|dp| json!({ "default_palace": dp }));
let result = initialize_response("trusty-memory", &state.version, extra);
json!({
"jsonrpc": "2.0",
"id": id,
"result": result,
})
}
"notifications/initialized" | "notifications/cancelled" => Value::Null,
"tools/list" => json!({
"jsonrpc": "2.0",
"id": id,
"result": tools::tool_definitions_with(state.default_palace.is_some())
}),
"rpc.discover" => json!({
"jsonrpc": "2.0",
"id": id,
"result": openrpc::build_discover_response(
&state.version,
state.default_palace.is_some(),
),
}),
"tools/call" => {
let params = msg.get("params").cloned().unwrap_or_default();
let tool_name = params
.get("name")
.and_then(|n| n.as_str())
.unwrap_or("")
.to_string();
let args = params.get("arguments").cloned().unwrap_or_default();
match tools::dispatch_tool(state, &tool_name, args).await {
Ok(content) => {
let text = match &content {
Value::String(s) => s.clone(),
other => other.to_string(),
};
json!({
"jsonrpc": "2.0",
"id": id,
"result": {
"content": [{"type": "text", "text": text}]
}
})
}
Err(e) => json!({
"jsonrpc": "2.0",
"id": id,
"error": {"code": -32603, "message": format!("{e:#}")}
}),
}
}
"ping" => json!({"jsonrpc": "2.0", "id": id, "result": {}}),
_ => json!({
"jsonrpc": "2.0",
"id": id,
"error": {
"code": -32601,
"message": format!("Method not found: {method}")
}
}),
}
}
#[cfg(test)]
mod lib_tests;