khive_runtime/pack/registry_access.rs
1use std::sync::Arc;
2
3use khive_gate::{GateDecision, GateRequest};
4use khive_storage::EventStore;
5use khive_types::{EventOutcome, Namespace};
6use serde_json::Value;
7
8use crate::error::RuntimeError;
9use crate::runtime::NamespaceToken;
10use crate::KhiveRuntime;
11
12use super::{
13 append_audit_event_best_effort, audit_admission_refused_obligation_count,
14 audit_admission_refused_obligation_last_at_ms, audit_admission_unresolved_obligation_count,
15 audit_admission_unresolved_obligation_last_at_ms, build_audit_storage_event,
16 fold_audit_obligation, masked_audit_event, VerbRegistry, AUDIT_PERSISTENCE_SKIPPED_READ_ONLY,
17};
18#[cfg(doc)]
19use super::{PackRegistry, RequestIdentity, VerbCategory, VerbRegistryBuilder};
20
21impl VerbRegistry {
22 /// Select the owning pack's backend for a note-kind KG read. The caller
23 /// keeps its already-authorized token; this only selects storage.
24 pub fn kg_note_read_runtime_for_kind<'a>(
25 &'a self,
26 runtime: &'a KhiveRuntime,
27 kind: &str,
28 ) -> &'a KhiveRuntime {
29 let Some(resolver) = &self.kg_read_resolver else {
30 return runtime;
31 };
32 let Some(owner) = self
33 .packs
34 .iter()
35 .find(|pack| pack.note_kinds().contains(&kind))
36 else {
37 return runtime;
38 };
39 resolver.runtime_for_pack(owner.name())
40 }
41
42 /// Resolve a KG entity/note handle across the configured backend inventory.
43 ///
44 /// The caller must supply its dispatch-authorized token. By-ID reads do not
45 /// filter the stored namespace (ADR-007); no new token is minted here. With
46 /// ordinary single-runtime registration, retain the supplied runtime's
47 /// existing behavior. This does not route mutations or pack-private records.
48 pub async fn resolve_kg_read_by_id(
49 &self,
50 runtime: &KhiveRuntime,
51 token: &NamespaceToken,
52 id: uuid::Uuid,
53 include_deleted: bool,
54 ) -> Result<Option<crate::Resolved>, RuntimeError> {
55 match &self.kg_read_resolver {
56 Some(resolver) => resolver.by_id(token, id, include_deleted).await,
57 None if include_deleted => runtime.resolve_by_id_including_deleted(token, id).await,
58 None => runtime.resolve_by_id(token, id).await,
59 }
60 }
61
62 /// Find the unique configured backend holding an entity for deletion.
63 /// Includes tombstones so soft deletion cannot hide a duplicate owner.
64 /// The dispatch-authorized token is preserved; lookup is namespace-agnostic.
65 pub async fn resolve_entity_delete_runtime(
66 &self,
67 runtime: &KhiveRuntime,
68 token: &NamespaceToken,
69 id: uuid::Uuid,
70 ) -> Result<Option<KhiveRuntime>, RuntimeError> {
71 match &self.kg_read_resolver {
72 Some(resolver) => resolver.entity_runtime(token, id).await,
73 None => {
74 let store = runtime.entities(token)?;
75 let entity = store.get_entity_including_deleted(id).await?;
76 Ok(entity.map(|_| runtime.clone()))
77 }
78 }
79 }
80
81 /// Clean main-backend attachments after no live or tombstoned owner remains.
82 /// A live or tombstoned entity on any configured backend keeps its roots.
83 pub async fn cleanup_deleted_entity_attachments(
84 &self,
85 runtime: &KhiveRuntime,
86 token: &NamespaceToken,
87 id: uuid::Uuid,
88 ) -> Result<bool, RuntimeError> {
89 if self
90 .resolve_entity_delete_runtime(runtime, token, id)
91 .await?
92 .is_some()
93 {
94 return Ok(false);
95 }
96 runtime.delete_entity_attachments_on_core(id).await
97 }
98
99 /// Recheck a merged-entity read against the kept id before returning it.
100 /// The submitted argument shape is the verb's ordinary shape with the
101 /// effective id substituted. The dispatch's original check remains its
102 /// own audit row; this consultation records the effective target as a
103 /// second row without changing the public GateRequest schema.
104 pub async fn authorize_effective_kg_read(
105 &self,
106 token: &NamespaceToken,
107 verb: &str,
108 mut effective_args: Value,
109 effective_id: uuid::Uuid,
110 ) -> Result<(), RuntimeError> {
111 if let Some(namespace) = token.gate_explicit_namespace() {
112 effective_args["namespace"] = Value::String(namespace.to_owned());
113 }
114 let gate_req = GateRequest::new(
115 token.actor().clone(),
116 token.gate_namespace().clone(),
117 verb,
118 effective_args,
119 );
120 let decision = khive_gate::check_with_mailbox_policy(self.gate.as_ref(), &gate_req);
121 match decision {
122 Ok(decision) => {
123 let audit = masked_audit_event(&gate_req, &decision, self.gate.impl_name());
124 tracing::info!(
125 audit_event = %serde_json::to_string(&audit)
126 .unwrap_or_else(|_| "{\"error\":\"serialize\"}".into()),
127 effective_target_id = %effective_id,
128 "gate.check"
129 );
130 let denied = matches!(&decision, GateDecision::Deny { .. });
131 let receipt = if let Some(store) = &self.event_store {
132 let event = build_audit_storage_event(
133 &gate_req,
134 &audit,
135 if denied {
136 EventOutcome::Denied
137 } else {
138 EventOutcome::Success
139 },
140 Some(crate::cost_unit::base_resource_payload(token.request_id())),
141 )
142 .with_target(effective_id);
143 if denied {
144 self.append_gate_denied_row(store, event, verb).await
145 } else {
146 let outcome = append_audit_event_best_effort(
147 self.audit_batch.as_ref(),
148 store,
149 event,
150 verb,
151 crate::audit_batch::AuditProducer::EffectiveTargetCheck,
152 false,
153 )
154 .await;
155 fold_audit_obligation(Ok(()), outcome, |_| Value::Null)?;
156 crate::error::DenialReceipt::no_store()
157 }
158 } else {
159 crate::error::DenialReceipt::no_store()
160 };
161 match decision {
162 GateDecision::Allow { .. } => Ok(()),
163 GateDecision::Deny { reason } => Err(RuntimeError::PermissionDenied {
164 verb: verb.to_string(),
165 reason,
166 receipt: Box::new(receipt),
167 }),
168 }
169 }
170 Err(error) => Err(self
171 .gate_unavailable_error(&gate_req, &error, token.request_id(), Some(effective_id))
172 .await),
173 }
174 }
175
176 /// Resolve a prefix across the same inventory, rejecting distinct UUIDs.
177 ///
178 /// Retains the local prefix scanner's entity/note/event/edge collision domain,
179 /// including sidecar events. The returned UUID is not a substrate assertion:
180 /// consumers must still fetch/type-check it. All backend failures propagate.
181 pub async fn resolve_kg_read_prefix(
182 &self,
183 runtime: &KhiveRuntime,
184 _token: &NamespaceToken,
185 prefix: &str,
186 include_deleted: bool,
187 ) -> Result<Option<uuid::Uuid>, RuntimeError> {
188 match &self.kg_read_resolver {
189 Some(resolver) => resolver.prefix(prefix, include_deleted).await,
190 None if include_deleted => {
191 runtime
192 .resolve_prefix_unfiltered_including_deleted(prefix)
193 .await
194 }
195 None => runtime.resolve_prefix_unfiltered(prefix).await,
196 }
197 }
198
199 /// This registry's construction-baked default namespace.
200 ///
201 /// Used as the fallback when a request carries no [`RequestIdentity`]
202 /// override (ADR-096 Fork 1) and by transports that need to advertise
203 /// their own resolved identity when forwarding to a warm daemon.
204 pub fn default_namespace(&self) -> &str {
205 &self.default_namespace
206 }
207
208 /// This registry's construction-baked actor identity label, if configured
209 /// (ADR-057). `None` means dispatch mints `ActorRef::anonymous()` absent a
210 /// per-request [`RequestIdentity`] override (ADR-096 Fork 1).
211 pub fn actor_id(&self) -> Option<&str> {
212 self.actor_id.as_deref()
213 }
214
215 /// This registry's construction-baked extra read-visibility namespaces
216 /// (ADR-007 Rev 4 Rule 3b), used absent a per-request [`RequestIdentity`]
217 /// override (ADR-096 Fork 1).
218 pub fn visible_namespaces(&self) -> &[Namespace] {
219 &self.visible_namespaces
220 }
221
222 /// This registry's configured audit `EventStore`, if any (ADR-094).
223 ///
224 /// Lets background tasks that hold a `VerbRegistry` but do not go through
225 /// `dispatch` (e.g. the email channel poll loop) append best-effort
226 /// lifecycle events to the same sink gate-check audit rows use, without
227 /// threading a second `Option<Arc<dyn EventStore>>` field through every
228 /// caller. `None` means either the historical tracing-only default or an
229 /// intentionally read-only audit backend; callers that need to distinguish
230 /// those cases use [`Self::audit_persistence_advisory`].
231 pub fn event_store(&self) -> Option<Arc<dyn EventStore>> {
232 self.event_store.clone()
233 }
234
235 /// Process-lifetime audit-batch health counters for this registry's
236 /// ADR-133 seam, if one is configured. `None` exactly when
237 /// [`Self::event_store`] is `None` — the same condition under which no
238 /// batch exists to report on. The `db_diagnostics` verb feeds this into
239 /// `KhiveRuntime::db_diagnostics_with_audit_metrics` so an operator can
240 /// see flush failures and pure-observability degradation instead of the
241 /// permanently-unavailable placeholder a bare `KhiveRuntime` reports.
242 ///
243 /// `admission_refused_obligations` and `admission_unresolved_obligations`
244 /// are sourced separately from [`audit_admission_refused_obligation_count`]
245 /// and [`audit_admission_unresolved_obligation_count`] rather than from
246 /// `batch.health_metrics()`: they count a decision made in
247 /// `append_audit_event_best_effort` (ADR-103 Amendment 3 / ADR-133
248 /// Amendment 1), not a property of the batch itself, so they are
249 /// process-wide like the rest of this struct's fields rather than
250 /// per-`AuditBatch`.
251 pub fn audit_batch_metrics(&self) -> Option<khive_db::diagnostics::RuntimeAuditBatchMetrics> {
252 self.audit_batch.as_ref().map(|batch| {
253 let m = batch.health_metrics();
254 khive_db::diagnostics::RuntimeAuditBatchMetrics {
255 flush_failures: m.flush_failures,
256 degraded_rows: m.degraded_rows,
257 degraded: m.degraded,
258 admission_refused_obligations: audit_admission_refused_obligation_count(),
259 admission_refused_obligations_last_at_ms:
260 audit_admission_refused_obligation_last_at_ms(),
261 admission_unresolved_obligations: audit_admission_unresolved_obligation_count(),
262 admission_unresolved_obligations_last_at_ms:
263 audit_admission_unresolved_obligation_last_at_ms(),
264 }
265 })
266 }
267
268 /// Test/diagnostic-only accessor for the underlying ADR-133 audit-batch
269 /// seam. `None` when no `EventStore` was configured (the batch is lazily
270 /// constructed from one). Exposed so admission-pressure mechanism tests
271 /// can saturate and drain the SAME instance a real dispatch uses
272 /// (#2117, #2147, #2208, #2217) instead of testing a
273 /// look-alike.
274 pub fn audit_batch_handle(&self) -> Option<Arc<crate::audit_batch::AuditBatch>> {
275 self.audit_batch.clone()
276 }
277
278 /// Stop admitting new audit rows and wait for every already-accepted row
279 /// to reach a terminal state (ADR-133).
280 ///
281 /// A no-op returning `Ok(())` when no `EventStore` — and therefore no
282 /// audit-batch seam — is configured. Callers that own this registry's
283 /// shutdown sequence should call this before tearing down the writer or
284 /// database so no accepted audit row is silently dropped mid-flight.
285 pub async fn shutdown_audit_batch(
286 &self,
287 ) -> Result<(), crate::audit_batch::AuditTerminalReason> {
288 use crate::audit_batch::AuditBatchControl;
289 match &self.audit_batch {
290 Some(audit_batch) => audit_batch.close_and_drain().await,
291 None => Ok(()),
292 }
293 }
294
295 /// Advisory for a dispatch whose configured audit sink is read-only.
296 ///
297 /// The MCP transport places this beside successful per-operation results;
298 /// `None` means audit persistence is configured normally or was never
299 /// configured at all.
300 pub fn audit_persistence_advisory(&self) -> Option<Value> {
301 self.audit_store_read_only.then(|| {
302 serde_json::json!({
303 "code": AUDIT_PERSISTENCE_SKIPPED_READ_ONLY,
304 "severity": "warning",
305 "component": "audit_event_store",
306 "reason": "read_only_backend",
307 "message": "operation completed, but its dispatch audit event was not persisted because the audit backend is read-only",
308 })
309 })
310 }
311
312 /// Explicit, fail-closed opt-in for admission-pressure audit degradation
313 /// (#2147/#2217). `VerbCategory::Assertive` alone is NOT a
314 /// sound proxy for "safe to drop this dispatch's own audit row under
315 /// audit-lane admission pressure": several Assertive handlers have
316 /// their own durable or accounting-bearing side effects. The reviewed
317 /// exclusions are:
318 /// - `memory.recall` dispatches `brain.record_serve` as a background
319 /// write; degrading `memory.recall`'s row raises the risk that a
320 /// serve goes unaccounted for if the ledger dispatch itself later
321 /// also races admission pressure.
322 /// - `db_diagnostics` may backfill WAL frames via a PASSIVE checkpoint
323 /// probe — physical I/O, not a pure in-memory read.
324 /// - `knowledge.search`, `knowledge.suggest`, and auto
325 /// `knowledge.compose` may start persistent ANN consumer/checkpoint
326 /// maintenance from their nominal read path.
327 /// - `git.checkout`, `git.diff` and `git.reconcile` persist a durable
328 /// receipt on every dispatch (checkout and diff also write a manifest
329 /// or diff blob), so their accounting row is not droppable.
330 /// - `tool.check` persists a policy decision receipt. `git.receipts`,
331 /// `git.gates`, `git.status` and `git.log` dispatch that same check,
332 /// so their read results also carry a required decision write.
333 ///
334 /// What membership here means, precisely: the verb performs no domain
335 /// mutation, so its OWN per-dispatch audit/accounting row may be dropped
336 /// under transient admission pressure without the caller losing a
337 /// meaningful result (ADR-103 Amendment 3, ADR-133 Amendment 1). It does
338 /// NOT mean the handler is free of every event-plane write: `search`
339 /// still fires its own best-effort `SearchExecuted` telemetry, and
340 /// `context` still records a one-time `ConfigLocked` event, both on
341 /// independent code paths this mechanism never touches — those events
342 /// commit or fail on their own terms, unaffected by whether this
343 /// dispatch's own audit row degrades.
344 ///
345 /// Every entry here MUST be declared `VerbCategory::Assertive` in its
346 /// named pack's live vocabulary. The
347 /// `admission_degrade_safe_assertive_census_matches_live_pack_sources`
348 /// test below scans every pack that currently declares public Assertive
349 /// handlers and requires every such handler to be classified exactly
350 /// once as safe or as a known incidental writer. A new Assertive verb
351 /// therefore fails closed both at runtime and in the source census until
352 /// it receives an explicit side-effect review.
353 ///
354 /// Entries are `(owning pack name, verb)` pairs, not bare verb names:
355 /// [`Self::admission_degrade_safe`] requires the handler actually
356 /// resolved for `verb` to belong to the exact pack named here. A verb
357 /// name alone is not a sound key — any pack registered through the same
358 /// [`PackRegistry`]/[`VerbRegistryBuilder`] path can declare a handler
359 /// under any name it likes, including one that collides with a name on
360 /// this list, and unique-verb-name validation only rejects that
361 /// collision when the real owning pack is *also* loaded. A deployment
362 /// that omits the real pack (or loads a third-party pack instead) would
363 /// let a same-named write-performing handler inherit degrade-safety it
364 /// never earned. Binding to the pack closes that gap.
365 pub(super) const ADMISSION_DEGRADE_SAFE_VERBS: &'static [(&'static str, &'static str)] = &[
366 // agent
367 ("agent", "agent.observe"),
368 // exec (reads of the blob store, the run receipt and event tables, or
369 // the resolved configuration; the writers are exec.tree and
370 // exec.tree_put, Declarations, and exec.run, a Directive)
371 ("exec", "exec.tree_get"),
372 ("exec", "exec.tree_diff"),
373 ("exec", "exec.receipt"),
374 ("exec", "exec.runs"),
375 ("exec", "exec.events"),
376 ("exec", "exec.identity"),
377 // git
378 // Canonical get project check plus bounded cursor SELECT; no domain writes.
379 ("git", "git.ingest_cursor"),
380 // blob
381 ("blob", "blob.get"),
382 ("blob", "blob.stat"),
383 // brain
384 ("brain", "brain.event_counts"),
385 ("brain", "brain.event_page"),
386 ("brain", "brain.profiles"),
387 ("brain", "brain.profile"),
388 ("brain", "brain.resolve"),
389 ("brain", "brain.bindings"),
390 // comm
391 ("comm", "comm.delivered"),
392 ("comm", "comm.transport_status"),
393 ("comm", "comm.inbox"),
394 ("comm", "comm.unread"),
395 ("comm", "comm.thread"),
396 ("comm", "comm.health"),
397 ("comm", "comm.probe"),
398 // gtd
399 ("gtd", "gtd.census"),
400 ("gtd", "gtd.next"),
401 ("gtd", "gtd.tasks"),
402 // kg
403 ("kg", "get"),
404 ("kg", "list"),
405 ("kg", "stats"),
406 ("kg", "search"),
407 ("kg", "neighbors"),
408 ("kg", "traverse"),
409 ("kg", "context"),
410 ("kg", "query"),
411 ("kg", "resolve"),
412 ("kg", "whoami"),
413 // scan runs the secret gate over caller-supplied text in process: no
414 // store read, no store write, no event.
415 ("kg", "scan"),
416 ("kg", "verbs"),
417 ("kg", "stream.read"),
418 ("kg", "stream.stat"),
419 // knowledge (ANN-maintaining search/suggest/compose are excluded)
420 ("knowledge", "knowledge.get"),
421 ("knowledge", "knowledge.list"),
422 ("knowledge", "knowledge.stats"),
423 ("knowledge", "knowledge.fold"),
424 ("knowledge", "knowledge.topic"),
425 // moodboard
426 ("moodboard", "moodboard.model"),
427 ("moodboard", "moodboard.search"),
428 ("moodboard", "moodboard.preference"),
429 // schedule
430 ("schedule", "schedule.agenda"),
431 // session
432 ("session", "session.list"),
433 ("session", "session.resume"),
434 ("session", "session.export"),
435 ("session", "session.search"),
436 // Fixed SQL reads and file metadata only; no domain or maintenance write.
437 ("session", "session.stats"),
438 // tool (registry, grant and policy reads; tool.suggest runs the same
439 // hybrid search as the kg search and context verbs above)
440 ("tool", "tool.suggest"),
441 ("tool", "tool.describe"),
442 ("tool", "tool.list"),
443 ("tool", "tool.requests"),
444 ("tool", "tool.policies"),
445 ];
446
447 /// Sorted copy of [`Self::ADMISSION_DEGRADE_SAFE_VERBS`], built once, so
448 /// [`VerbRegistryBuilder::build`] can decide each trusted handler's
449 /// eligibility with a binary search instead of a linear scan over every
450 /// entry. Consulted exactly once per registry, at `build()` time — see
451 /// [`Self::admission_degrade_safe`] for why no per-dispatch scan exists
452 /// anymore. The source list above stays grouped by pack (with a `//
453 /// <pack>` comment per group) for human review; this is a derived,
454 /// lookup-shaped view of the same data, not a second source of truth.
455 pub(super) fn admission_degrade_safe_sorted() -> &'static [(&'static str, &'static str)] {
456 static SORTED: std::sync::LazyLock<Vec<(&'static str, &'static str)>> =
457 std::sync::LazyLock::new(|| {
458 let mut pairs = VerbRegistry::ADMISSION_DEGRADE_SAFE_VERBS.to_vec();
459 pairs.sort_unstable();
460 pairs
461 });
462 &SORTED
463 }
464
465 /// Whether `verb` is both declared [`VerbCategory::Assertive`] (the
466 /// speech-act tag for handlers that "retrieve and present facts" rather
467 /// than committing a domain change) AND explicitly opted in to
468 /// admission-pressure audit degradation via
469 /// [`Self::ADMISSION_DEGRADE_SAFE_VERBS`] under the exact pack that
470 /// registered it, AND declared by a pack the composition root actually
471 /// vouches for (see [`VerbRegistryBuilder::register_boxed`]'s doc).
472 /// Unknown, non-opted-in, wrong-pack, or untrusted-pack verbs are
473 /// conservatively `false` — fail-closed, so a new Assertive handler (or
474 /// one registered by a pack other than the one the allowlist names, or
475 /// one registered through [`VerbRegistryBuilder::register`] rather than
476 /// the trusted path) hard-fails its audit obligation like any write
477 /// until someone deliberately reviews it and adds it to the allowlist.
478 ///
479 /// `pack.name()` is a value the `PackRuntime` trait object reports about
480 /// itself — any pack registered through the public
481 /// [`VerbRegistryBuilder::register`] path can claim any name, including
482 /// one on the allowlist, whether or not the pack that name actually
483 /// belongs to is also loaded (verb names are unique per registry, so an
484 /// impostor's same-named handler is only reachable when the real pack
485 /// is absent). Binding eligibility to registration-time trust — decided
486 /// by the *caller*, never by the pack instance — is why this checks
487 /// `degrade_safe_verbs` rather than resolving `pack.name()` at query
488 /// time; [`VerbRegistryBuilder::build`] already excluded every untrusted
489 /// pack's handlers from that set.
490 ///
491 /// The whole decision is precomputed once in `VerbRegistryBuilder::build`
492 /// into [`VerbRegistry::degrade_safe_verbs`] — a verb name is unique
493 /// across `Visibility::Verb` handlers within one registry
494 /// (`validate_unique_verb_names`), so this is a single hash-set lookup,
495 /// not a per-dispatch scan over every registered pack's handler list.
496 ///
497 /// Used only to decide whether a dispatch's own audit-obligation row may
498 /// degrade to best-effort on transient audit-lane admission pressure
499 /// (`append_audit_event_best_effort`) — a read that performed no domain
500 /// write must not fail the caller just because the audit lane is
501 /// momentarily saturated. Never used for permission checking, transport
502 /// routing, or return-shape selection.
503 pub(super) fn admission_degrade_safe(&self, verb: &str) -> bool {
504 self.degrade_safe_verbs.contains(verb)
505 }
506
507 /// Transport replay eligibility from the shared operation-effects table,
508 /// restricted to trusted canonical public handlers. A read that persists a
509 /// fresh serve or telemetry row is excluded because the request id is
510 /// correlation, not deduplication.
511 /// Custom and mounted handlers cannot inherit safety from a name/category.
512 pub fn is_read_replay_safe(&self, verb: &str) -> bool {
513 self.read_replay_safe_verbs.contains(verb)
514 }
515
516 /// White-box accessor for [`Self::admission_degrade_safe`], needed
517 /// because the admission-pressure regression tests in
518 /// `tests/read_verb_admission_exhaustion.rs` compile as a separate
519 /// external binary and cannot reach a crate-private method directly —
520 /// the same reason [`audit_admission_refused_obligation_count`] and
521 /// `AuditBatch::test_snapshot` are `pub` rather than `pub(crate)`.
522 #[cfg(any(test, feature = "test-internals"))]
523 pub fn admission_degrade_safe_probe(&self, verb: &str) -> bool {
524 self.admission_degrade_safe(verb)
525 }
526}