pub struct AppState {Show 33 fields
pub version: String,
pub machine: 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<Sender<DaemonEvent>>,
pub started_at: Instant,
pub log_buffer: LogBuffer,
pub error_store: Option<ErrorStore>,
pub multi_tenant_mode: bool,
pub disk_bytes: Arc<AtomicU64>,
pub sys_metrics: Arc<Mutex<SysMetrics>>,
pub bound_addr: Arc<OnceLock<SocketAddr>>,
pub prompt_context_cache: Arc<RwLock<PromptFactsCache>>,
pub tier_s_admission_lock: Arc<Mutex<()>>,
pub activity_log: Arc<ActivityLog>,
pub palace_write_locks: Arc<DashMap<String, Arc<Mutex<()>>>>,
pub pending_activity_writes: Arc<AtomicUsize>,
pub worker_liveness: Arc<WorkerLiveness>,
pub wedge_threshold: Duration,
pub lock_stalls: Arc<LockStallTracker>,
pub palace_names: Arc<DashMap<String, String>>,
pub startup_gate: StartupOpenGate,
pub pin_project_map: Arc<DashMap<String, PathBuf>>,
pub bm25_index_tx: Sender<Bm25IndexRequest>,
pub bm25_dirty: DirtyPalaces,
pub update_available: Arc<Mutex<Option<String>>>,
pub daemon_readiness: Arc<AtomicU8>,
pub write_op_budget: Duration,
pub palace_last_used: StampCache,
pub write_pipeline_budget: Duration,
/* private fields */
}Expand description
Shared application state passed to every request handler.
Why: The stdio loop and HTTP server need the same handles to the registry, data root, and embedder so MCP tools can perform real reads/writes against the live trusty-memory core. The embedder is heavy (loads ONNX weights) so it is resolved lazily through the process-wide singleton on first use.
#4836: this struct used to carry its own Arc<OnceCell<Arc<FastEmbedder>>>,
a SECOND cell independent of retrieval::shared_embedder(). Startup warmed
the shared cell and latched daemon_readiness off it, while every recall
consumed the private cell — so the flag reported one embedder’s state while
the request path used another’s. One cell removes the disagreement (and the
duplicate ~90 MB ONNX session) by construction.
What: Clone-able via Arc fields. The registry / data root are eager; the
embedder is reached via AppState::embedder.
Test: app_state_default_constructs confirms construction without panic.
Fields§
§version: String§machine: MachineBudgetWhat machine this daemon is on, and the memory budget that follows (#6820).
Why: trusty-memory had NO tier detection at all — max_open_palaces
stayed at the fixed DEFAULT_MAX_OPEN_PALACES (64) whether the host had
12 GB or 128 GB, so a sub-minimum machine held the same 64 resident
palaces as a workstation. #6820 detects and exposes the tier; #6823 is
what reads this field to size that cap. Deliberately NOT retuning any cap
here — the issue is scoped to detect + expose + document.
What: trusty_common::machine_tier::MachineBudget — the cgroup-clamped
RAM reading, the tier band, and the two proportional soft limits, from
the same shared module trusty-search resolves its own caps through.
Resolved once per AppState, in AppState::new.
Test: machine_budget_is_detected_and_self_consistent in lib_tests.
registry: Arc<PalaceRegistry>§data_root: PathBuf§default_palace: Option<String>Optional default palace applied to MCP tool calls when the caller
omits the palace argument. Set via trusty-memory serve --palace.
chat_provider: Arc<OnceCell<Option<Arc<dyn ChatProvider>>>>Active chat provider selected at startup. None means no upstream is
configured (no Ollama detected and no OpenRouter key) — callers must
degrade gracefully (chat endpoint returns 412).
session_stores: Arc<SessionStoreCache>Per-palace chat-session stores, opened lazily so cold-start cost is paid only when chat-history endpoints are hit.
#4639: was an unbounded DashMap with no remove/TTL/cap, leaking one
chat_sessions.redb fd per palace for the daemon’s lifetime (844
measured live, all pointing at already-unlinked files). Now an
LRU-bounded cache that evicts cold, unused stores.
events: Arc<Sender<DaemonEvent>>Broadcast sender for live DaemonEvent pushes to SSE subscribers.
Why: Lets mutating handlers emit events that any connected dashboard
receives instantly. Cap of 128 buffers transient slow readers; if a
receiver lags it gets RecvError::Lagged and we emit a lag frame.
started_at: InstantInstant the daemon started, used to compute uptime_secs on /health.
Why (issue #35): GET /health reports how long the daemon has been
up. Capturing a monotonic Instant at AppState construction lets the
handler compute the elapsed seconds cheaply and without a clock-skew
hazard.
What: a wall-monotonic Instant; AppState::new stamps it at startup.
Test: health_reports_idle_worker_pool.
log_buffer: LogBufferIn-memory ring buffer of recent tracing log lines (issue #35).
Why: the GET /api/v1/logs/tail endpoint serves the last N log lines
so operators can inspect a running daemon without tailing a file. The
buffer is shared between the tracing LogBufferLayer (writer) and the
HTTP handler (reader).
What: a cheap Arc-backed clone of the buffer the subscriber writes
to. Defaults to an empty buffer for states that never install the
layer (tests, the stdio path).
Test: rpc_logs_tail_answers_a_bounded_page.
error_store: Option<ErrorStore>Bug-capture ERROR store (bug-reporting #478, Phase 1).
Why: Phase 2 MCP / HTTP endpoints need to query captured errors; stashing
the ErrorStore handle here lets any handler reach it cheaply without
a second global or per-request construction.
What: populated by run_serve from the init_tracing_with_buffer_and_capture
result; the layer writes to this store automatically so every
tracing::error! call site contributes without any changes to call
sites. None in states that do not install the layer (tests, the
stdio path).
Test: compile-presence is verified by the trusty-memory build; Phase 2
will add query tests in web.rs.
multi_tenant_mode: boolMinimal multi-tenant authorization seam (issue #1714). false
(single-tenant, the default) preserves today’s behaviour — every
existing caller of palace_create force=true keeps working. true
opts into authz::authorize_force_palace_create failing closed on
every force=true request until a real capability check lands. Set
via with_multi_tenant_mode_from_env (TRUSTY_MEMORY_MULTI_TENANT=1);
see the authz module docs for the full design rationale.
disk_bytes: Arc<AtomicU64>Most recent on-disk footprint of data_root, in bytes (issue #35).
Why: GET /health reports disk_bytes. Walking the data directory on
every health request would make a frequent health poll do unbounded
I/O; a background task recomputes it every 10 s and stores it here so
the handler reads it lock-free.
What: an AtomicU64 updated by the ticker spawned in run_http_on.
0 until the first walk completes.
Test: health_reports_idle_worker_pool.
sys_metrics: Arc<Mutex<SysMetrics>>Per-process RSS + CPU sampler, refreshed on each /health request
(issue #35).
Why: CPU usage is a delta between two sysinfo refreshes, so the
sampler must persist between requests — hence the shared Mutex.
What: a tokio::sync::Mutex<SysMetrics> so the async health handler
can sample without blocking the runtime.
Test: health_reports_idle_worker_pool.
bound_addr: Arc<OnceLock<SocketAddr>>HTTP listener address the daemon bound to, once run_http_on is running.
Why: clients (and /health responses) need to advertise the live
host:port even though port selection happens dynamically (7070–7079
walk + OS fallback). Stashing it on AppState lets request handlers
surface the discovery value without re-querying the listener.
What: a OnceLock<SocketAddr> so run_http_on writes it exactly once
at bind time and every handler reads it lock-free thereafter. Empty
(None from get()) on the stdio path where no listener exists.
Test: health_endpoint_reports_bound_addr (added below).
prompt_context_cache: Arc<RwLock<PromptFactsCache>>Cached prompt-facts surface served by the MCP get_prompt_context
tool (issue #42).
Why: The original session-init prompts/get design loaded context
once per connection; switching to a per-message tool lets the model
pull fresh, query-filtered context on demand. The cache holds both
the raw triples (for filtered lookups) and a pre-formatted Markdown
block (for the unfiltered hot path) so neither code path re-walks
the KG. The cache is rebuilt by
prompt_facts::rebuild_prompt_cache after any write that touches a
hot predicate. #5524: every caller-supplied assert reaches that rebuild
through kg_write::assert_triple rather than each surface remembering
to call it; the retract side (remove_prompt_fact) still calls the
rebuild directly.
What: An Arc<tokio::sync::RwLock<PromptFactsCache>> so the hot
read path takes a brief read lock and clones the cache; rebuilds
take a write lock for the assignment only. The async-aware lock
(issue #229) yields to the tokio runtime instead of blocking a
runtime thread for the rebuild duration. An empty triples vec ↔
“no context stored yet” (the tool handler renders a hint).
Test: get_prompt_context_returns_cached_or_hint,
get_prompt_context_filters_by_query.
tier_s_admission_lock: Arc<Mutex<()>>Serializes the Tier S check-then-write sequence (#4888).
Why: the 20-fact cap is only a cap if it cannot be raced past.
check_tier_s_admission counts active facts across every palace and
then the caller writes — two callers that both observe 19 would both
pass and the surface would land at 21. Nothing else serializes them:
the KG’s single-writer actor orders writes only within one palace,
and the count spans all of them. “Usually 20, occasionally 21” is not
the invariant ADR-0028 D8 asks for.
What: an async mutex whose guard is acquired by
check_tier_s_admission and returned to the caller, which must hold it
until its kg.assert is enqueued. Cold predicates never acquire it, so
ordinary knowledge-graph writes stay fully concurrent; hot writes are
deliberate and rare, so serializing them costs nothing measurable.
Test: tier_s_cap_holds_under_concurrent_writes.
activity_log: Arc<ActivityLog>Persistent activity log (issue #96).
Why: the dashboard activity feed used to be a pure live-stream over
/sse — opening the UI showed an empty feed and any mutation from
the MCP path was invisible. Holding an ActivityLog on AppState
lets emit record an entry on every push so the
GET /api/v1/activity handler can return historical rows on mount
and the live SSE stream can continue prepending events on top of
the loaded history. None on builds that opt out (tests that use
AppState::new get a real log under their tempdir so behaviour
matches production).
What: an Arc<ActivityLog> shared with every emitter.
Test: web::tests::activity_endpoint_lists_recent_emits.
palace_write_locks: Arc<DashMap<String, Arc<Mutex<()>>>>Per-palace write serialisation locks (issue #230).
Why: the dedup gate in tools.rs previously read a snapshot of
existing drawers, checked for near-duplicates via Jaro-Winkler, and
then issued the write — a classic time-of-check/time-of-use race.
Two concurrent memory_remember calls with the same content could
both see the pre-write snapshot, both pass the gate, and both land
duplicate drawers. Serialising the gate-then-write sequence per
palace closes the window: while one task holds the mutex, any
concurrent writer for the same palace blocks until the first write
finishes and is visible to list_drawers. The lock is per
palace (not global) so writes to different palaces continue to
run in parallel.
What: a DashMap keyed by palace id, where each entry is an
Arc<tokio::sync::Mutex<()>>. The mutex is constructed lazily by
palace_write_lock on first access. Arc lets callers hold a
clone of the lock past the lifetime of the DashMap entry so the
map never needs to be held across an .await.
Test: tools::tests::dedup_gate_blocks_concurrent_duplicate_writes.
pending_activity_writes: Arc<AtomicUsize>Counter of in-flight activity-log writes spawned by emit
(issue #232).
Why: emit offloads the synchronous redb append to the tokio blocking
pool via spawn_blocking so the async runtime is never parked waiting
on fsync. The write is fire-and-forget — emit returns immediately
after spawning. Tests that observe the activity log right after a
burst of emit calls need a deterministic synchronization point;
holding an in-flight counter lets flush_activity_writes poll until
every spawned append has settled, which keeps the assertions
race-free without forcing every caller to .await.
What: an Arc<AtomicUsize> incremented before each spawn_blocking
and decremented inside the closure (after the append completes, even
if it errored). The counter is cheap (one atomic add per emit) and
stays at zero in steady-state production traffic.
Test: web::tests::activity_endpoint_lists_recent_emits and
tests::emit_persists_mutations_but_skips_status_changed call
flush_activity_writes to drain the counter before reading the log.
worker_liveness: Arc<WorkerLiveness>Live occupancy gauge for the palace open path (issue #4001).
Why: during the #3992 incident six daemon threads sat parked in
concurrent_open::backoff_sleep_ms with a memory_remember hung
~1800 s, while both doctors reported HEALTHY. Every existing signal —
HTTP liveness, fastembed cache state, lock-file staleness — describes
the process, not the work. This is the one field that can answer
“is anything actually moving?”, and /health surfaces it so an
out-of-process doctor can report what the daemon actually observed
instead of inferring health from a cheap proxy.
What: an worker_liveness::WorkerLiveness slot table; one CAS on
entry and one store on exit per tracked operation, so the gauge cannot
itself become the load problem it exists to detect.
Test: web::tests::health_tests::health_reports_wedged_worker_pool.
wedge_threshold: DurationHow long an operation may run before the pool is called wedged.
Why this is state rather than an env read on the request path: /health
is polled once a second, so re-reading (and re-parsing) an environment
variable per request is needless work; resolving it once at construction
also makes the value a property of the daemon instead of a process-wide
global, which is what lets tests drive the wedge condition
deterministically instead of mutating shared env state and racing each
other.
What: defaults to worker_liveness::wedge_threshold.
Test: web::tests::health_tests::health_reports_wedged_worker_pool.
lock_stalls: Arc<LockStallTracker>Stamps for handle locks found held (#4001).
Why: a dream cycle, forget or import holds PalaceHandle::write_mutex
without registering in Self::worker_liveness, so only observing the
lock itself can see it held past Self::wedge_threshold.
What: a lock_stall::LockStallTracker, swept by memory.health and by
the daemon’s ticker.
Test: tools::tests::write_liveness_tests::a_dream_cycle_holding_the_handle_write_mutex_reads_as_wedged.
palace_names: Arc<DashMap<String, String>>In-memory cache mapping palace id → Palace.name (issue #228).
Why: every memory_remember / memory_note write used to call
PalaceRegistry::list_palaces (a synchronous filesystem walk of the
data root) just to resolve a friendly palace name for the SSE
DrawerAdded event. With N palaces on disk the cost was O(N) opendirs
plus palace.json reads on every write, blocking the async runtime.
Caching the name in-memory turns the lookup into a DashMap::get.
What: DashMap<String, String> populated by create_palace and
load_palaces_from_disk, kept in sync by rename / delete paths.
Missing entries are treated as “name unknown” so callers fall back to
the palace id and the emit path never fails.
Test: palace_name_cache_populated_after_hydration and
palace_name_cache_updates_on_create.
startup_gate: StartupOpenGateShared bound on concurrent palace opens during startup work (#7106).
Why: hydration, the BM25 backfill sweep and the BM25 repair sweep all
run at boot and all walk the whole estate. Three independent limits
multiply; one gate on the shared state is what makes “at most N palaces
open at once” true of the process. It also carries the high-water mark,
so the bound is observable rather than merely intended.
What: a startup_budget::StartupOpenGate sized by
TRUSTY_MEMORY_STARTUP_OPEN_LIMIT (default 4). Cloning shares the
semaphore and counters.
Test: startup_task_tests::hydration_never_exceeds_the_startup_open_limit.
pin_project_map: Arc<DashMap<String, PathBuf>>Single-pass startup pin-file map: palace id → project root path (issue #470).
Why: after daemon startup we have no record of which on-disk project
directories correspond to which palace ids — that information only
existed inside the pin files on disk. Eager-opening every palace on
startup is too expensive. This field captures the scan-only result of
startup_scan::scan_pin_map so handlers that want to locate a project
by its palace id (e.g. future cwd-inference, project-health checks)
can do a single DashMap::get instead of a filesystem walk.
Populated once, shortly after load_palaces_from_disk returns, by
spawn_startup_tasks. Never mutated after population — it is a
snapshot of what the filesystem looked like at startup.
What: DashMap<String (palace_id), PathBuf (project root)>.
The outer Arc lets spawn_startup_tasks (which holds only a clone
of AppState) write to the same backing map that request handlers
read. Population is asynchronous so callers must treat an absent entry
as “not yet scanned” (or “no pin found”), never as “palace unknown”.
Test: startup_scan::tests::scan_pin_map_* validate the underlying
scanner function; the wiring in spawn_startup_tasks is covered by
the integration-test daemon start path.
bm25_index_tx: Sender<Bm25IndexRequest>Bounded sender for the BM25 index worker (issue #231).
Why: the previous fire-and-forget design tokio::spawned one task per
memory_remember / memory_note call, so a write burst against a slow
or unreachable BM25 daemon grew an unbounded in-flight task queue. A
single long-lived worker draining a bounded mpsc channel caps that
back-pressure: writers try_send (never block), full-queue requests
are dropped with a warn!, and the worker exits cleanly when the last
sender is dropped on shutdown.
What: an mpsc::Sender cloned to every AppState clone (cheap). The
matching receiver is consumed by the worker spawned in
AppState::new via tools::spawn_bm25_index_worker. Capacity is
tools::BM25_INDEX_QUEUE_CAPACITY (256).
Test: bm25_index_queue_drops_when_full exercises the full-queue
branch via bm25_index_enqueue.
bm25_dirty: DirtyPalacesPalaces whose BM25 coverage is known to be incomplete (#5048 review).
Why: bm25_index_enqueue drops on a full queue so memory_remember
never waits on daemon RTT. That trade is only defensible if a drop is
actually repaired, and before this field the sole production trigger for
a backfill was daemon startup — so a drop stayed invisible until the
next restart. Every observer of lost coverage marks the palace here and
bm25_repair::spawn_repair_sweep consumes it on an interval.
What: a DashSet of palace ids, shared by every AppState clone.
Idempotent — forty drops for one palace queue one repair.
Test: bm25_repair_tests.rs, bm25_index_queue_drops_when_full.
update_available: Arc<Mutex<Option<String>>>Cached result of the startup update check (issue #537).
Why: /health should report update_available without hitting crates.io
on every probe. A single background check at daemon startup stores the
result here; the health handler reads it lock-free (well, a brief mutex
lock) without a network call.
What: None = up-to-date or check not yet done; Some("x.y.z") = newer
version available. The field is populated by a tokio::spawn in
spawn_startup_tasks (main.rs) after the daemon binds.
Test: indirectly by the /health endpoint tests in web.rs.
daemon_readiness: Arc<AtomicU8>Two-phase readiness state — Warming until the embedder is initialised,
then Ready (issues #910 / #911).
Why: AppState::embedder() used to call FastEmbedder::new() without
any timeout, so the first memory_recall/memory_remember that arrived
before CoreML finished compiling would block for 5–11 hours until the
OnceCell resolved (issue #910). Exposing this state lets the preflight
guards in tools.rs return an explicit fast error immediately —
"trusty-memory is warming up, retry shortly" — instead of queueing
behind an open-ended init.
What: An AtomicU8 starting at DaemonReadiness::Warming (0) and flipped
to DaemonReadiness::Ready (1) by spawn_startup_tasks after the embedder
warm-up succeeds. The transition is one-way and lock-free.
Test: daemon_readiness_transitions_warming_to_ready.
write_op_budget: DurationTotal wall-clock ceiling for one MCP write operation (issue #4002).
Why: a write waits for the per-palace write mutex and then waits again
to enter the per-palace open queue. Each leg was bounded by its own
timeout, so the effective ceiling was their sum (60 s + ~63 s), not
either configured bound. The handlers stamp one
trusty_common::memory_core::timeouts::OpBudget from this value and
clamp every leg through it, so the later leg spends what the earlier
leg left.
What: defaults to
trusty_common::memory_core::timeouts::write_op_budget
(TRUSTY_WRITE_OP_BUDGET_SECS, 60 s). Stored per-instance rather than
re-read from the environment so
AppState::with_write_op_budget can inject a short deadline in tests
without mutating process-wide state, which would race parallel tests —
the same reason PalaceRegistry::with_open_queue_timeout exists.
Test: tools::tests::write_budget_tests.
palace_last_used: StampCacheWhen each palace’s durable last-used stamp was last written (#6424).
Why: the console’s Last Used column needs a stamp that survives a
restart, and a durable write per recall is not worth a column measured
in days. This is the throttle’s memory — see
palace_last_used for the cadence and what lagging it costs.
What: palace id -> unix seconds of the last PERSISTED stamp, not of the
last use. Empty at boot, so the first use of any palace writes.
Test: palace_last_used::tests::stamp_throttles_within_the_window.
write_pipeline_budget: DurationCeiling on the pipeline one write runs while holding the palace write mutex (issue #6366).
Why: AppState::write_op_budget bounds only the waits BEFORE the
mutex is held. The pipeline that runs once it IS held had no ceiling, so
a slow commit held the mutex for as long as it took and every other
writer on that palace queued behind it — three memory_note calls were
aborted client-side after 1800 s while the daemon stayed healthy.
What: defaults to
trusty_common::memory_core::timeouts::write_pipeline_timeout
(TRUSTY_WRITE_PIPELINE_TIMEOUT_SECS, 240 s), and is threaded into
PalaceHandle::remember_with_options_within. Stored per-instance for
the same reason as write_op_budget: a test can inject a short ceiling
without mutating process-wide env.
Test: tools::tests::write_budget_tests::memory_note_surfaces_the_pipeline_ceiling
proves the ceiling reaches the daemon’s own write handler; the pipeline’s
behaviour under that ceiling lives in trusty-common’s retrieval
write-pipeline tests.
Implementations§
Source§impl AppState
impl AppState
Sourcepub fn new(data_root: PathBuf) -> Self
pub fn new(data_root: PathBuf) -> Self
Construct an AppState rooted at the given on-disk data directory.
Why: The CLI (serve) and integration tests need to point the MCP
server at different roots — production at dirs::data_dir, tests at a
tempfile::tempdir().
What: Builds an empty PalaceRegistry, captures the version, and
allocates an empty OnceCell for the embedder. default_palace is
None; use with_default_palace to set it.
Test: tools::tests::dispatch_palace_create_persists constructs an
AppState pointed at a tempdir and round-trips a palace through it.
Sourcepub fn with_write_op_budget(self, budget: Duration) -> Self
pub fn with_write_op_budget(self, budget: Duration) -> Self
Override the per-operation write budget (issue #4002).
Why: tests must prove the two write legs draw on ONE budget without
setting TRUSTY_WRITE_OP_BUDGET_SECS, which is process-wide and would
race any test running in parallel.
What: consuming builder that overwrites AppState::write_op_budget.
Test: tools::tests::write_budget_tests.
Sourcepub fn with_write_pipeline_budget(self, budget: Duration) -> Self
pub fn with_write_pipeline_budget(self, budget: Duration) -> Self
Override the write-pipeline ceiling (issue #6366).
Why: same reason as AppState::with_write_op_budget — a test proving
the ceiling fires must not set TRUSTY_WRITE_PIPELINE_TIMEOUT_SECS,
which is process-wide and would race any parallel test.
What: consuming builder that overwrites
AppState::write_pipeline_budget.
Test: tools::tests::write_budget_tests::memory_note_surfaces_the_pipeline_ceiling.
Sourcepub fn palace_write_lock(&self, palace_id: &str) -> Arc<Mutex<()>> ⓘ
pub fn palace_write_lock(&self, palace_id: &str) -> Arc<Mutex<()>> ⓘ
Acquire (lazily, then clone) the per-palace write mutex.
Why (issue #230): the dedup-check + remember_with_options write
sequence in tools.rs must be atomic per palace to prevent two
concurrent identical writes from both passing the dedup gate.
Callers hold the returned Arc<Mutex<()>>’s guard across the gate
check and the write so the second writer blocks until the first
write is visible to list_drawers. Returning a clone of the Arc
rather than a borrow into the DashMap lets the caller .await
while holding the lock without risking a deadlock against any
future map mutation (DashMap shards are sync mutexes).
What: looks up the palace id in palace_write_locks and returns
a clone of the existing mutex; on the first call for a palace,
inserts a freshly-constructed tokio::sync::Mutex<()> first. The
DashMap::entry().or_insert_with API guarantees the lazy
construction is racy-safe — only one mutex is ever inserted per
palace id.
Test: tools::tests::dedup_gate_blocks_concurrent_duplicate_writes.
Sourcepub fn pinned_project_path(&self, palace_id: &str) -> Option<PathBuf>
pub fn pinned_project_path(&self, palace_id: &str) -> Option<PathBuf>
Look up a project root path by palace id in the startup pin-scan map.
Why: provides a stable, cheap accessor so handlers do not reach directly
into the DashMap field and so the accessor can be mocked in future
tests without touching AppState internals. The map is populated
asynchronously by spawn_startup_tasks — an absent entry means either
the scan has not completed yet or no pin file claimed that id.
What: returns Some(project_path) when the palace id was found during
startup scan; None otherwise.
Test: covered indirectly via the startup-scan integration path; the
underlying map data is validated by startup_scan::tests.
Sourcepub fn with_bm25_lane_from_env(self) -> Self
pub fn with_bm25_lane_from_env(self) -> Self
Builder-style: opt-in to the BM25 lexical lane (issue #156, #5329).
Why: the lane stays gated behind TRUSTY_BM25_DAEMON=1 even though
#5329 removed the daemon that name refers to. Renaming the variable
would break the only enablement path anyone could have configured, in
the one PR whose purpose is not losing that lane — the compatibility is
worth more than the accuracy.
What: when the gate is set, builds a bm25_lane::Bm25Lane over this
state’s data_root — which is where the retired daemon wrote its
snapshots, so an existing corpus is picked up in place — and rebuilds
the bounded indexer channel so its worker holds the lane. Returns self
unchanged when the var is unset or set to anything other than 1.
Test: bm25_lane_disabled_by_default, bm25_lane_enabled_when_env_set.
Sourcepub fn with_bm25_lane(self, lane: Arc<Bm25Lane>) -> Self
pub fn with_bm25_lane(self, lane: Arc<Bm25Lane>) -> Self
Builder-style: install an explicit BM25 lane, bypassing the env gate.
Why: setting bm25 on its own is a footgun. AppState::new spawns the
indexer worker with no lane, so a caller that assigns the field and
nothing else gets a state whose reads use the lane and whose WRITES are
silently discarded by the placeholder worker — which is exactly what
writes_through_the_tool_surface_survive_eviction caught. Every path
that installs a lane goes through here so the two cannot drift apart.
What: rebuilds the bounded indexer channel + worker so the worker holds
the lane, then stores it. The placeholder worker installed by
AppState::new exits cleanly when the replaced sender closes its
receiver. Tests use this to pin explicit limits without mutating
process-global env vars.
Test: tests/bm25_lane_concurrency.rs, bm25_lane_enabled_when_env_set.
Sourcepub fn bm25_lane(&self) -> Option<&Arc<Bm25Lane>>
pub fn bm25_lane(&self) -> Option<&Arc<Bm25Lane>>
The BM25 lane, if one is installed.
Why: the read half of the pub(crate) field above. Callers outside this
crate — the integration tests, and anything that wants to flush or query
the lane directly — need to reach it without being able to swap it for
one the indexer worker has never heard of.
Test: tests/bm25_lane_concurrency.rs, tests/bm25_alias_write.rs.
Sourcepub async fn load_palaces_from_disk(&self) -> Result<usize>
pub async fn load_palaces_from_disk(&self) -> Result<usize>
Scan the palace registry directory and re-register every persisted
palace into the in-memory PalaceRegistry.
Why: AppState::new builds an empty registry, so after a daemon
restart palace_list / the dashboard reported zero palaces even though
dozens existed on disk — palace metadata was persisted by
palace_create but never re-hydrated on startup. This method closes
that gap by walking the on-disk layout (each subdirectory holding a
palace.json is one palace) and rebuilding a live PalaceHandle for
each, so recall paths see the full set immediately after a restart.
What: runs the blocking filesystem walk + per-palace
PalaceHandle::open_with_intent on a spawn_blocking thread (so it
never stalls the async runtime), registers each successfully opened
palace via register_arc, logs every load at debug!, and returns the
count loaded. A palace that fails to open (corrupt index, unreadable
kg.db, etc.) is logged at warn! and skipped — one bad palace must not
abort startup or crash the daemon — but the skip is RECORDED (#4911) and
readable via PalaceRegistry::unopenable, so a palace whose bytes
survive and whose contents cannot be read stays observable instead of
reading as absent. data_root is expected to already be the palace
registry directory — main.rs resolves it via
resolve_palace_registry_dir before constructing the AppState, so
the flat / legacy-palaces/ layout difference is handled exactly once.
Intent (#1487, #4911): every open uses the registry’s own
OpenIntent, NOT the zero-arg PalaceHandle::open default. This is
the path a restarting daemon takes for every palace it already has on
disk, so hardcoding ReadOnlyClient here made the daemon’s
with_writer_intent() guarantee false in the common case.
Test: tests::load_palaces_from_disk_rehydrates_registry writes two
palaces into a tempdir, constructs an AppState, calls this method, and
asserts the returned count and registry contents;
load_palaces_from_disk_honours_registry_open_intent covers the intent
contract in both directions and
load_palaces_from_disk_records_an_unopenable_palace the skip record.
Sourcepub fn with_log_buffer(self, buffer: LogBuffer) -> Self
pub fn with_log_buffer(self, buffer: LogBuffer) -> Self
Builder-style: attach the daemon’s shared LogBuffer so the
GET /api/v1/logs/tail endpoint serves the same lines the tracing
subscriber captures (issue #35).
Why: main builds the buffer (via init_tracing_with_buffer) before
constructing the AppState, then hands a clone here so the HTTP
handler and the tracing layer observe the same ring.
What: replaces the empty default buffer with the supplied one.
Test: rpc_logs_tail_answers_a_bounded_page.
Sourcepub fn with_writer_intent(self) -> Self
pub fn with_writer_intent(self) -> Self
Builder-style: mark this daemon as the sole palace writer so palace
redb files open with OpenIntent::Writer (issue #1487).
Why: The HTTP daemon owns the write lock on every palace’s kg.redb
and index.usearch.redb. Before this fix, when a second daemon
instance opened the same store it silently degraded to a read-only
snapshot and rejected every memory_remember for its lifetime —
effectively silent data loss when an MCP client routed a write to the
rogue instance. Opening as Writer makes the second instance fail
loud (after a short handoff-retry window that absorbs a graceful
launchd bootout→bootstrap overlap) instead of serving broken
reads-only. CLI, stdio-proxy, and test code paths never call this, so
they keep the snapshot read-fallback (issue #59).
What: Replaces self.registry with a fresh PalaceRegistry carrying
OpenIntent::Writer and this data root’s MaintenanceLease (#8733).
Invariant: MUST be called on a fresh, unhydrated, unshared registry —
during startup, before spawn_startup_tasks/load_palaces_from_disk
registers any PalaceHandle and before the AppState (hence its
Arc<PalaceRegistry>) is cloned to a handler. Replacing the registry
discards the prior Arc; doing so after hydration would silently drop
live handles (data loss), and doing so after the state is shared would
leave other clones on the stale read-only registry. The guard is a
debug_assert! on the strongest cheap signals the registry exposes —
is_empty() (no handles hydrated) and Arc::strong_count == 1 (not yet
shared) — so an ordering violation fails fast as the programmer error it
is (the call site is startup-only and fixed). Release builds elide the
assert; the real call site (run_serve) always satisfies it.
Test: with_writer_intent_marks_registry_writer and
with_writer_intent_panics_on_hydrated_registry in lib_tests.
Sourcepub fn with_error_store(self, store: ErrorStore) -> Self
pub fn with_error_store(self, store: ErrorStore) -> Self
Builder-style: attach the bug-capture ErrorStore handle (bug-reporting #478).
Why: Phase 2 MCP / HTTP endpoints need a handle to the in-memory error
ring so they can serve recent_errors / errors_by_fingerprint
without disk I/O on the hot path. Installing it here — rather than
adding it as a separate global — keeps the state graph explicit and
lets tests skip it by never calling this method.
What: stores Some(store) in AppState::error_store; the BugCaptureLayer
that writes to this store is already installed in the tracing
subscriber by init_tracing_with_buffer_and_capture. The store is
Clone (cheap Arc clone internally) so both the layer and this
field share the same underlying ring.
Test: Phase 2 will add error_store_captures_and_queries in web.rs.
Sourcepub fn with_multi_tenant_mode_from_env(self) -> Self
pub fn with_multi_tenant_mode_from_env(self) -> Self
Builder-style: opt into multi-tenant authorization mode (issue #1714).
Why: mirrors with_bm25_client_from_env’s pattern of keeping env-var
gating in one place. Unset (the default) preserves today’s
single-tenant behaviour with zero change for existing callers. Issue
#2522 review: activation is silent otherwise, which makes a
misconfigured (or unexpectedly enabled) deployment hard to diagnose
from logs alone — log once at startup when the mode flips on.
What: sets multi_tenant_mode from TRUSTY_MEMORY_MULTI_TENANT=1; see
the authz module for what the flag then enforces. Logs via
tracing::info! (stderr only) when enabled; stays silent when
disabled (the default).
Test: authorize_force_palace_create_denies_multi_tenant_without_capability.
Sourcepub fn emit(&self, event: DaemonEvent)
pub fn emit(&self, event: DaemonEvent)
Send a DaemonEvent to all connected SSE subscribers and persist
it to the activity log when the variant carries a source.
Why: Mutating handlers call this after a successful write so the
dashboard can update without polling. The send is best-effort —
broadcast::Sender::send returns Err only when there are no live
receivers, which is fine (no listeners == no work to do). Issue
#96 additionally writes the entry to the persistent activity log
so the feed can serve historical rows on page load and so MCP /
HTTP / Hook origins are visible to the operator. Persistence is
also best-effort — a write failure is logged but never blocks the
SSE broadcast.
Issue #232: the activity-log append is a synchronous redb write +
fsync. Calling it directly on the async caller’s task parked a tokio
worker thread on disk I/O for every SSE event. We now offload the
append to the blocking thread pool via spawn_blocking and return
immediately — emit stays synchronous so every existing caller
(including the sync dispatch_hook_fired JSON-RPC handler) keeps
compiling unchanged. The fire-and-forget pattern matches the
pre-fix semantics (best-effort, never blocks the SSE broadcast)
while freeing the async runtime to do real work during the write.
What: serialises the event for the log (skipping StatusChanged
which is a recomputed aggregate, not a mutation), spawns the redb
append on tokio::task::spawn_blocking keyed by a clone of the
Arc<ActivityLog> and the cloned event, then sends the event over
the broadcast channel. A pending_activity_writes counter is bumped
before the spawn and decremented inside the closure so
Self::flush_activity_writes can drain in tests.
Test: web::tests::sse_stream_receives_palace_created confirms a
subscriber observes the emitted event;
activity_endpoint_lists_recent_emits confirms persistence via
flush_activity_writes.
Sourcepub async fn flush_activity_writes(&self)
pub async fn flush_activity_writes(&self)
Block (asynchronously) until every in-flight activity-log write
spawned by Self::emit has settled.
Why: emit offloads its redb append to tokio::task::spawn_blocking
and returns immediately (issue #232). Tests that observe the
activity log right after a burst of emits would otherwise race the
blocking-pool worker; this helper gives them a deterministic
synchronization point. Production code never needs to call this —
the dashboard reads through GET /api/v1/activity, which already
tolerates writes settling asynchronously.
What: spins on pending_activity_writes with a 1 ms yield until the
counter is zero. Cheap: tests typically emit a handful of events
and the loop exits within a single scheduler tick.
Test: covered indirectly by emit_persists_mutations_but_skips_status_changed
and web::tests::activity_endpoint_lists_recent_emits.
Sourcepub fn session_store(&self, palace_id: &str) -> Result<Arc<ChatSessionStore>>
pub fn session_store(&self, palace_id: &str) -> Result<Arc<ChatSessionStore>>
Open (or return cached) the chat-session store for a palace.
Why: Chat session persistence lives in a dedicated redb file under
the palace’s data dir (chat_sessions.redb) so it doesn’t intermingle
with the KG’s transactional load. The store is cheap to clone via
Arc but the underlying connection should be reused, so cache by id.
What: delegates to the LRU-bounded SessionStoreCache, which creates
the palace data dir if missing, opens (or reuses) a ChatSessionStore,
and evicts cold, unused stores once more than the cap are resident.
Callers keep the returned Arc for as long as they need it — eviction
never closes a store someone still holds.
Test: session_store_cache::tests::open_handles_are_bounded_by_cap,
session_store_cache::tests::in_use_store_is_never_evicted; the call
path is covered indirectly by the session HTTP handlers in web::tests.
Sourcepub fn with_default_palace(self, name: Option<String>) -> Self
pub fn with_default_palace(self, name: Option<String>) -> Self
Builder-style setter for the default palace name.
Why: serve --palace <name> wants to bind every tool call to a
project-scoped namespace without forcing every MCP request to repeat
the palace argument.
What: Returns self with default_palace = Some(name).
Test: default_palace_used_when_arg_omitted covers the resolution
path; this setter is exercised there.
Sourcepub async fn chat_provider(&self) -> Option<Arc<dyn ChatProvider>>
pub async fn chat_provider(&self) -> Option<Arc<dyn ChatProvider>>
Resolve (or initialize) the shared embedder.
Why: FastEmbedder load is expensive — we share one instance across all
tool calls; the OnceCell ensures concurrent first-use races collapse
to a single load.
What: Returns Arc<FastEmbedder> on success. Errors propagate from the
underlying ONNX load.
Test: Indirectly via dispatch_remember_then_recall.
Resolve the active chat provider, auto-detecting on first call.
Why: Provider selection depends on filesystem-loaded config plus a
network probe (Ollama liveness), so it must be lazily initialised at
runtime. Caching the choice in a OnceCell keeps it stable across
concurrent requests without re-probing on every chat call.
What: On first use loads ~/.trusty-memory/config.toml, prefers an
auto-detected Ollama instance (when local_model.enabled), and falls
back to OpenRouter when an API key is set. Returns None when neither
is available, so memory.chat can refuse before opening a stream.
Test: rpc_chat_providers_answers_both_upstreams covers the detection
path indirectly through memory.chat_providers.
Sourcepub fn spawn_alias_discovery(&self, palace: String, project_root: PathBuf)
pub fn spawn_alias_discovery(&self, palace: String, project_root: PathBuf)
Spawn a fire-and-forget background task that auto-discovers project
aliases under project_root and asserts new ones into palace.
Why (issue #42): Projects carry implicit shorthand — cargo package
names that differ from their directory, binary names that differ
from packages, first-letter abbreviations — that should be surfaced
without a user ever calling add_alias. Running discovery as a
detached task on palace-open keeps startup latency unchanged: the
daemon binds and starts serving immediately while the discovery scan
completes in the background, and any newly-asserted aliases land in
the prompt cache before the model’s next get_prompt_context call.
What: clones self (cheap; Arc-backed), spawns a tokio task that
invokes the discover_aliases tool handler directly so the
dedup + cache-rebuild logic runs exactly the same path as the MCP
tool call. Errors are logged at warn!; one failed discovery never
destabilises the daemon.
Test: not unit-tested (timing-dependent fire-and-forget); the
underlying discover_aliases dispatch is covered by
dispatch_discover_aliases_inserts_new_and_dedupes in tools::tests.
Sourcepub fn readiness(&self) -> DaemonReadiness
pub fn readiness(&self) -> DaemonReadiness
Return the current readiness state.
Why: tool handlers and the /health endpoint need a cheap, lock-free
way to check whether the embedder has been initialised yet.
What: loads daemon_readiness with Acquire ordering so the caller
sees all writes the startup task made before setting the state.
#4836: the value read here is only as good as the writes that reach it.
It used to be written exactly once, by the startup warm-up task, so a
single failed attempt pinned the daemon at Warming for the rest of its
life. AppState::embedder now also flips it on success, which is the
signal that actually proves a vector search can run.
Test: daemon_readiness_transitions_warming_to_ready,
resolving_the_embedder_marks_a_warming_daemon_ready.
Sourcepub fn set_ready(&self)
pub fn set_ready(&self)
Flip the readiness state from Warming to Ready.
Why: called by spawn_startup_tasks in main.rs once the embedder
warm-up succeeds — this is the single state-transition site.
What: store(Ready, Release) so subsequent Acquire loads in handlers
observe a consistent state. Idempotent: calling it multiple times is
harmless.
Test: daemon_readiness_transitions_warming_to_ready.
Sourcepub async fn embedder(&self) -> Result<Arc<dyn Embedder + Send + Sync>>
pub async fn embedder(&self) -> Result<Arc<dyn Embedder + Send + Sync>>
Obtain the shared FastEmbedder instance, initialising it on first call.
Why: centralises lazy embedder access so every tool handler goes through
one bounded init path (tracks #910 internally).
What: wraps OnceCell::get_or_try_init with a timeout so a slow
CoreML/CUDA first-compile cannot block a handler indefinitely. On
timeout the OnceCell is left unresolved and the next caller retries.
Callers on the request path SHOULD check readiness() before this
method (issue #1970) — every recall handler now checks
readiness() == Ready first and only calls embedder() on that
branch, falling back to a BM25/L0/L1-only path while Warming instead
of paying this method’s cold-init cost. Reaching this method while
still Warming is not a bug (the warm-up task itself calls
embedder() while in Warming state), just unusual on the request
path.
This timeout is a backstop against a pathological init delay (e.g. the
warm-up task’s own call, or a handler that skips the readiness()
check). If this timeout fires the OnceCell is left in the unresolved
state and the next call retries from scratch.
Trait Implementations§
Auto Trait Implementations§
impl !RefUnwindSafe for AppState
impl !UnwindSafe for AppState
impl Freeze for AppState
impl Send for AppState
impl Sync for AppState
impl Unpin for AppState
impl UnsafeUnpin for AppState
Blanket Implementations§
Source§impl<T> BorrowMut<T> for Twhere
T: ?Sized,
impl<T> BorrowMut<T> for Twhere
T: ?Sized,
Source§fn borrow_mut(&mut self) -> &mut T
fn borrow_mut(&mut self) -> &mut T
Source§impl<T> CloneToUninit for Twhere
T: Clone,
impl<T> CloneToUninit for Twhere
T: Clone,
Source§impl<T> Instrument for T
impl<T> Instrument for T
Source§fn instrument(self, span: Span) -> Instrumented<Self> ⓘ
fn instrument(self, span: Span) -> Instrumented<Self> ⓘ
Source§fn in_current_span(self) -> Instrumented<Self> ⓘ
fn in_current_span(self) -> Instrumented<Self> ⓘ
Source§impl<T> IntoEither for T
impl<T> IntoEither for T
Source§fn into_either(self, into_left: bool) -> Either<Self, Self> ⓘ
fn into_either(self, into_left: bool) -> Either<Self, Self> ⓘ
self into a Left variant of Either<Self, Self>
if into_left is true.
Converts self into a Right variant of Either<Self, Self>
otherwise. Read moreSource§fn into_either_with<F>(self, into_left: F) -> Either<Self, Self> ⓘ
fn into_either_with<F>(self, into_left: F) -> Either<Self, Self> ⓘ
self into a Left variant of Either<Self, Self>
if into_left(&self) returns true.
Converts self into a Right variant of Either<Self, Self>
otherwise. Read more