khive_runtime/pack/dispatch.rs
1use std::sync::Arc;
2use std::time::Instant;
3
4use khive_gate::{AuditEvent, GateDecision, GateRequest};
5use khive_storage::{Event, EventStore, EventView, SubstrateKind};
6use khive_types::{EventKind, EventOutcome, Namespace, Visibility};
7use serde_json::Value;
8
9use crate::error::{AuditObligationFailure, DispatchError, RuntimeError};
10use crate::runtime::NamespaceToken;
11
12use super::request_identity::edge_endpoint_table;
13#[cfg(doc)]
14use super::IdResolutionMode;
15use super::{
16 append_audit_event_best_effort, build_audit_storage_event, fold_audit_obligation,
17 identifier_resolution_help, link_audit_success_from_result, masked_audit_event,
18 persist_git_digest_receipt, resolution_mode_contract, resolve_explicit_namespace,
19 GitDigestReceiptOutcome, InterceptedDispatchResult, RequestIdentity, VerbRegistry,
20 VerifiedActor,
21};
22
23impl VerbRegistry {
24 /// Return the help schema envelope for a verb.
25 ///
26 /// Walks registered packs for the first matching `HandlerDef` and returns a
27 /// structured JSON envelope. Subhandlers carry `callable_via_mcp: false`.
28 /// Every envelope carries the shared `identifier_resolution` contract.
29 /// `link`'s envelope additionally carries `endpoint_rules` — the composed
30 /// per-relation source/target allowlist (issue #964) — so batch callers can
31 /// defer to the kernel's own table instead of re-implementing it locally.
32 /// Every `uuid`/`array of uuid` parameter description has its
33 /// declared [`IdResolutionMode`]'s contract text appended — the same
34 /// text rendered under `identifier_resolution.resolution_modes` — so the
35 /// full-UUID-vs-short-prefix rule is stated once per mode (in
36 /// `resolution_mode_contract`) and inherited by every matching param,
37 /// instead of restating it per param across every `HandlerDef` in every
38 /// pack. Parameters whose mode is [`IdResolutionMode::NotApplicable`]
39 /// (every non-identifier parameter) are left unchanged.
40 /// Unknown verbs return `RuntimeError::InvalidInput`. Full shape documented
41 /// in `docs/protocol.md` §Request Schema.
42 pub fn describe_verb(&self, verb: &str) -> Result<Value, RuntimeError> {
43 for pack in self.packs.iter() {
44 for handler in pack.handlers().iter() {
45 if handler.name == verb {
46 let category = format!("{:?}", handler.category);
47 let params_arr: Vec<Value> = handler
48 .params
49 .iter()
50 .map(|p| {
51 let description = match resolution_mode_contract(p.resolution_mode) {
52 Some(contract) => format!("{} {}", p.description, contract),
53 None => p.description.to_string(),
54 };
55 serde_json::json!({
56 "name": p.name,
57 "type": p.param_type,
58 "required": p.required,
59 "description": description,
60 })
61 })
62 .collect();
63 // Subhandlers are not callable via the MCP request surface;
64 // the help payload must match the behaviour the dispatch
65 // path enforces so callers reading `help=true` before
66 // probing see accurate availability.
67 if matches!(handler.visibility, Visibility::Subhandler) {
68 return Ok(serde_json::json!({
69 "verb": verb,
70 "pack": pack.name(),
71 "description": handler.description,
72 "category": category,
73 "params": params_arr,
74 "identifier_resolution": identifier_resolution_help(),
75 "visibility": "internal",
76 "callable_via_mcp": false,
77 "note": "This is an internal subhandler. Calling it via the MCP \
78 request surface returns permission denied. It can only be \
79 invoked by internal runtime callers.",
80 }));
81 }
82 let mut envelope = serde_json::json!({
83 "verb": verb,
84 "pack": pack.name(),
85 "description": handler.description,
86 "category": category,
87 "params": params_arr,
88 "identifier_resolution": identifier_resolution_help(),
89 });
90 // A pack that authored its own schema keeps it; every other
91 // verb gets one derived from the declarations the runtime
92 // already holds, so a bridged model has a schema to read
93 // instead of parsing the prose `params[].type`.
94 if let Some(schema) = pack.input_schema(verb) {
95 envelope["input_schema"] = schema;
96 } else {
97 let described: Vec<(String, String)> = params_arr
98 .iter()
99 .map(|p| {
100 (
101 p["name"].as_str().unwrap_or_default().to_string(),
102 p["description"].as_str().unwrap_or_default().to_string(),
103 )
104 })
105 .collect();
106 if let Some(schema) =
107 crate::input_schema::derive_input_schema(handler.params, &described)
108 {
109 envelope["input_schema"] = schema;
110 }
111 }
112 if verb == "link" {
113 envelope["endpoint_rules"] = Value::Array(edge_endpoint_table(&self.packs));
114 }
115 return Ok(envelope);
116 }
117 }
118 }
119 // Verb-visibility handler names, precomputed at build() time (internal
120 // subhandlers are excluded so they are not advertised in the
121 // unknown-verb error).
122 Err(RuntimeError::UnknownVerb(format!(
123 "unknown verb {verb:?}; available: {}",
124 self.available_verbs.join(", ")
125 )))
126 }
127
128 /// Check whether the gate permits writes into `ns`.
129 ///
130 /// Performs a gate evaluation with verb `"authorize"` before any background
131 /// loop is spawned (ADR-056 §6). Returns `Ok(())` when the gate allows the
132 /// namespace, or `Err(RuntimeError::PermissionDenied{..})` when denied.
133 /// Gate errors (implementation failures) are surfaced as
134 /// `RuntimeError::Internal` carrying the stable classified reason; the
135 /// bounded, masked backend detail goes to the server-side log here, since
136 /// callers log the returned error.
137 pub fn authorize_namespace(&self, ns: Namespace) -> Result<(), RuntimeError> {
138 let actor = crate::actor_identity::resolve_actor(self.actor_id.as_deref());
139 let req = GateRequest::new(actor, ns, "authorize", serde_json::Value::Null);
140 match self.gate.check(&req) {
141 Ok(decision) if decision.is_allow() => Ok(()),
142 Ok(GateDecision::Deny { reason }) => {
143 Err(RuntimeError::permission_denied("authorize", reason))
144 }
145 Ok(_) => Err(RuntimeError::permission_denied("authorize", "gate denied")),
146 Err(e) => {
147 tracing::warn!(
148 error = %crate::secret_gate::bounded_masked_log_text(&e.to_string()),
149 "authorize_namespace: gate check failed (fail-closed)"
150 );
151 Err(RuntimeError::Internal(format!(
152 "gate error: {}",
153 e.wire_reason()
154 )))
155 }
156 }
157 }
158
159 /// Gate and execute an operation handled outside normal pack dispatch.
160 ///
161 /// Multi-backend transports use this to route an operation through a
162 /// coordinator while retaining [`Self::dispatch_with_identity`]'s gate and
163 /// audit lifecycle. Deny is authoritative, gate errors fail closed, and an
164 /// allowed audit is persisted after the intercepted operation resolves so
165 /// its outcome and duration reflect the operation result. Successful
166 /// `git.digest` interception uses the same strict durable-receipt exception
167 /// as normal pack dispatch.
168 pub async fn dispatch_intercepted_with_identity<F, Fut>(
169 &self,
170 verb: &str,
171 params: &Value,
172 identity: Option<&RequestIdentity>,
173 dispatch: F,
174 ) -> Result<Value, RuntimeError>
175 where
176 F: FnOnce(Namespace) -> Fut,
177 Fut: std::future::Future<Output = Result<Value, RuntimeError>>,
178 {
179 self.dispatch_intercepted_with_metadata_with_identity(
180 verb,
181 params,
182 identity,
183 |namespace| async move {
184 dispatch(namespace)
185 .await
186 .map(|result| InterceptedDispatchResult::new(result, ()))
187 },
188 )
189 .await
190 .map(|outcome| outcome.result)
191 }
192
193 /// Gate and execute an intercepted operation whose transport needs typed
194 /// metadata in addition to the canonical verb result.
195 ///
196 /// Audit accounting always receives `outcome.result`; `outcome.metadata`
197 /// crosses the dispatch seam unchanged for the transport to place beside
198 /// that result in its own envelope.
199 pub async fn dispatch_intercepted_with_metadata_with_identity<M, F, Fut>(
200 &self,
201 verb: &str,
202 params: &Value,
203 identity: Option<&RequestIdentity>,
204 dispatch: F,
205 ) -> Result<InterceptedDispatchResult<M>, RuntimeError>
206 where
207 F: FnOnce(Namespace) -> Fut,
208 Fut: std::future::Future<Output = Result<InterceptedDispatchResult<M>, RuntimeError>>,
209 {
210 self.dispatch_intercepted_with_metadata_and_disposition(verb, params, identity, dispatch)
211 .await
212 .map_err(DispatchError::into_source)
213 }
214
215 /// Append the `GateDenied` row of a refused dispatch and report what the
216 /// caller may cite: the row's id when it committed, otherwise why not.
217 pub(super) async fn append_gate_denied_row(
218 &self,
219 store: &Arc<dyn EventStore>,
220 event: Event,
221 verb: &str,
222 ) -> crate::error::DenialReceipt {
223 let audit_event_id = event.id;
224 match append_audit_event_best_effort(
225 self.audit_batch.as_ref(),
226 store,
227 event,
228 verb,
229 crate::audit_batch::AuditProducer::GateDenied,
230 false,
231 )
232 .await
233 {
234 Ok(()) => crate::error::DenialReceipt {
235 audit_event_id: Some(audit_event_id),
236 audit_outcome: crate::error::DenialAuditOutcome::Committed,
237 },
238 Err(failure) => crate::error::DenialReceipt {
239 audit_event_id: None,
240 audit_outcome: crate::error::DenialAuditOutcome::NotCommitted(failure.wire_code()),
241 },
242 }
243 }
244
245 /// Execute an intercepted operation while retaining this boundary's failure provenance.
246 /// Successful canonical results and typed metadata are returned unchanged.
247 pub async fn dispatch_intercepted_with_metadata_and_disposition<M, F, Fut>(
248 &self,
249 verb: &str,
250 params: &Value,
251 identity: Option<&RequestIdentity>,
252 dispatch: F,
253 ) -> Result<InterceptedDispatchResult<M>, DispatchError>
254 where
255 F: FnOnce(Namespace) -> Fut,
256 Fut: std::future::Future<Output = Result<InterceptedDispatchResult<M>, RuntimeError>>,
257 {
258 self.dispatch_intercepted_with_token_and_disposition(verb, params, identity, |token| {
259 dispatch(token.gate_namespace().clone())
260 })
261 .await
262 }
263
264 /// Intercept an operation with the sealed caller token minted after the
265 /// gate decision. Coordinated reads retain the resolved actor and their
266 /// existing namespace selection without reconstructing identity from args.
267 pub async fn dispatch_intercepted_with_token_and_disposition<M, F, Fut>(
268 &self,
269 verb: &str,
270 params: &Value,
271 identity: Option<&RequestIdentity>,
272 dispatch: F,
273 ) -> Result<InterceptedDispatchResult<M>, DispatchError>
274 where
275 F: FnOnce(NamespaceToken) -> Fut,
276 Fut: std::future::Future<Output = Result<InterceptedDispatchResult<M>, RuntimeError>>,
277 {
278 let request_id = identity.and_then(|id| id.request_id);
279 let gate_req = self
280 .gate_request_with_identity(verb, params, identity)
281 .map_err(DispatchError::before_dispatch)?;
282 let gate_decision = khive_gate::check_with_mailbox_policy(self.gate.as_ref(), &gate_req);
283 let mut deferred_audit = match gate_decision {
284 Ok(decision) => {
285 let audit = masked_audit_event(&gate_req, &decision, self.gate.impl_name());
286 tracing::info!(
287 audit_event = %serde_json::to_string(&audit)
288 .unwrap_or_else(|_| "{\"error\":\"serialize\"}".into()),
289 "gate.check"
290 );
291 if let GateDecision::Deny { reason } = decision {
292 let receipt = match &self.event_store {
293 Some(store) => {
294 let event = build_audit_storage_event(
295 &gate_req,
296 &audit,
297 EventOutcome::Denied,
298 Some(crate::cost_unit::base_resource_payload(request_id)),
299 );
300 // The dispatch returns `PermissionDenied` below
301 // whether or not this row commits — a deny never
302 // reports success — so a commit failure has no
303 // caller-visible outcome to fold into; the receipt
304 // on the refusal says whether the row the caller
305 // could cite exists.
306 self.append_gate_denied_row(store, event, verb).await
307 }
308 None => crate::error::DenialReceipt::no_store(),
309 };
310 return Err(DispatchError::before_dispatch(
311 RuntimeError::PermissionDenied {
312 verb: verb.to_string(),
313 reason,
314 receipt: Box::new(receipt),
315 },
316 ));
317 }
318 Some(audit)
319 }
320 Err(err) => {
321 return Err(DispatchError::before_dispatch(
322 self.gate_unavailable_error(&gate_req, &err, request_id, None)
323 .await,
324 ));
325 }
326 };
327
328 let started = Instant::now();
329 let token = self.mint_intercepted_read_token(&gate_req, params, identity);
330 let mut result = dispatch(token).await;
331 let domain_succeeded = result.is_ok();
332 let duration_us = started.elapsed().as_micros() as i64;
333 let receipt_outcome = if verb == "git.digest" && result.is_ok() {
334 let resource = result.as_ref().ok().map(|outcome| {
335 crate::cost_unit::resource_payload(
336 verb,
337 &gate_req.args,
338 &outcome.result,
339 || 0,
340 request_id,
341 )
342 });
343 // The receipt helper operates on the canonical verb result. Move
344 // that value out temporarily so it can turn receipt failures into
345 // the outer dispatch error without discarding successful typed
346 // transport metadata.
347 let mut receipt_result: Result<Value, RuntimeError> = match result.as_mut() {
348 Ok(outcome) => Ok(std::mem::take(&mut outcome.result)),
349 Err(_) => unreachable!("git.digest receipt path is guarded by result.is_ok()"),
350 };
351 let outcome = persist_git_digest_receipt(
352 self.event_store.as_ref(),
353 self.audit_batch.as_ref(),
354 &gate_req,
355 deferred_audit.as_ref(),
356 &mut receipt_result,
357 duration_us,
358 resource,
359 )
360 .await;
361 match receipt_result {
362 Ok(receipted_result) => {
363 if let Ok(intercepted) = &mut result {
364 intercepted.result = receipted_result;
365 }
366 }
367 Err(error) => result = Err(error),
368 }
369 Some(outcome)
370 } else {
371 None
372 };
373 if receipt_outcome.is_none()
374 || receipt_outcome == Some(GitDigestReceiptOutcome::BuildRejected)
375 {
376 if let Some(audit) = deferred_audit.take() {
377 let audit_outcome = self
378 .persist_intercepted_audit(
379 verb,
380 &gate_req,
381 audit,
382 result.as_ref().map(|outcome| &outcome.result),
383 duration_us,
384 request_id,
385 )
386 .await;
387 result = fold_audit_obligation(result, audit_outcome, |outcome| outcome.result);
388 }
389 }
390 result.map_err(|error| DispatchError::after_handler(error, domain_succeeded))
391 }
392
393 async fn persist_intercepted_audit(
394 &self,
395 verb: &str,
396 gate_req: &GateRequest,
397 audit: AuditEvent,
398 result: Result<&Value, &RuntimeError>,
399 duration_us: i64,
400 request_id: Option<u64>,
401 ) -> Result<(), AuditObligationFailure> {
402 let Some(store) = &self.event_store else {
403 return Ok(());
404 };
405 let event = match result {
406 Ok(value) if verb == "link" && gate_req.args.get("links").is_none() => {
407 let resource = crate::cost_unit::resource_payload(
408 verb,
409 &gate_req.args,
410 value,
411 || 0,
412 request_id,
413 );
414 match link_audit_success_from_result(audit.clone(), value) {
415 Some((edge_id, mut payload)) => {
416 if let Value::Object(ref mut map) = payload {
417 map.insert("resource".to_string(), resource);
418 }
419 Event::new(
420 gate_req.namespace.as_str(),
421 gate_req.verb.as_str(),
422 EventKind::Audit,
423 SubstrateKind::Event,
424 format!("{}:{}", gate_req.actor.kind, gate_req.actor.id),
425 )
426 .with_outcome(EventOutcome::Success)
427 .with_target(edge_id)
428 .with_payload(payload)
429 .with_payload_schema_version(2)
430 .with_duration_us(duration_us)
431 }
432 None => build_audit_storage_event(
433 gate_req,
434 &audit,
435 EventOutcome::Success,
436 Some(resource),
437 )
438 .with_duration_us(duration_us),
439 }
440 }
441 Ok(value) => build_audit_storage_event(
442 gate_req,
443 &audit,
444 EventOutcome::Success,
445 Some(crate::cost_unit::resource_payload(
446 verb,
447 &gate_req.args,
448 value,
449 || 0,
450 request_id,
451 )),
452 )
453 .with_duration_us(duration_us),
454 Err(_) => build_audit_storage_event(
455 gate_req,
456 &audit,
457 EventOutcome::Error,
458 Some(crate::cost_unit::base_resource_payload(request_id)),
459 )
460 .with_duration_us(duration_us),
461 };
462 let producer = if result.is_ok() {
463 crate::audit_batch::AuditProducer::DispatchSucceeded
464 } else {
465 crate::audit_batch::AuditProducer::DispatchFailed
466 };
467 append_audit_event_best_effort(
468 self.audit_batch.as_ref(),
469 store,
470 event,
471 verb,
472 producer,
473 self.admission_degrade_safe(verb),
474 )
475 .await
476 }
477
478 /// A create refusal may reveal its key holder only when the same caller can list it.
479 pub fn allows_note_key_disclosure(
480 &self,
481 token: &NamespaceToken,
482 kind: &str,
483 key: &str,
484 ) -> bool {
485 let request = GateRequest::new(
486 token.actor().clone(),
487 token.namespace().clone(),
488 "list",
489 serde_json::json!({"kind":"note", "note_kind":kind, "key_prefix":key}),
490 );
491 self.gate
492 .check(&request)
493 .is_ok_and(|decision| decision.is_allow())
494 }
495
496 fn mint_intercepted_read_token(
497 &self,
498 gate_req: &GateRequest,
499 params: &Value,
500 identity: Option<&RequestIdentity>,
501 ) -> NamespaceToken {
502 // Preserve the coordinator's existing namespace selection: the gate
503 // namespace is primary, and explicit namespace input stays narrow.
504 let visible = if params.get("namespace").is_some() {
505 Vec::new()
506 } else {
507 let mut visible = match identity {
508 Some(identity) => identity
509 .visible_namespaces
510 .iter()
511 .filter_map(|namespace| Namespace::parse(namespace).ok())
512 .collect(),
513 None => self.visible_namespaces.clone(),
514 };
515 visible.push(Namespace::local());
516 visible
517 };
518 NamespaceToken::mint_with_visibility(
519 gate_req.namespace.clone(),
520 visible,
521 gate_req.actor.clone(),
522 )
523 .with_gate_namespace(gate_req.namespace.clone())
524 .with_gate_explicit_namespace(
525 params
526 .get("namespace")
527 .and_then(Value::as_str)
528 .map(str::to_owned),
529 )
530 .with_request_id(identity.and_then(|identity| identity.request_id))
531 .with_process_ref(match identity {
532 Some(identity) => identity.process_ref.clone(),
533 None => crate::config::process_ref_from_env(),
534 })
535 }
536
537 fn gate_request_with_identity(
538 &self,
539 verb: &str,
540 params: &Value,
541 identity: Option<&RequestIdentity>,
542 ) -> Result<GateRequest, RuntimeError> {
543 let default_namespace = identity
544 .map(|id| id.namespace.as_str())
545 .unwrap_or(self.default_namespace.as_str());
546 let namespace = resolve_explicit_namespace(params, default_namespace)?;
547 let actor_id = identity
548 .map(|id| id.actor_id.as_deref())
549 .unwrap_or(self.actor_id.as_deref());
550 let actor = crate::actor_identity::resolve_actor(actor_id);
551 // GateRequest.args deliberately captures submitted dispatch arguments.
552 // The handler's canonicalization and kind hooks have not run; a policy
553 // requiring their effective values belongs after that handler work.
554 let req = GateRequest::new(actor, namespace, verb, params.clone());
555 crate::mailbox_view::validate_mailbox_request(&req)?;
556 Ok(req)
557 }
558
559 pub(super) async fn gate_unavailable_error(
560 &self,
561 gate_req: &GateRequest,
562 error: &khive_gate::GateError,
563 request_id: Option<u64>,
564 effective_target: Option<uuid::Uuid>,
565 ) -> RuntimeError {
566 let audit = AuditEvent::gate_unavailable(gate_req, self.gate.impl_name())
567 .with_operation_attribution(
568 khive_storage::operation_context::current_operation_attribution(),
569 );
570 tracing::info!(
571 audit_event = %serde_json::to_string(&audit)
572 .unwrap_or_else(|_| "{\"error\":\"serialize\"}".into()),
573 "gate.check"
574 );
575 tracing::warn!(
576 verb = %gate_req.verb,
577 error = %crate::secret_gate::bounded_masked_log_text(&error.to_string()),
578 "gate check failed (fail-closed)"
579 );
580 if let Some(store) = &self.event_store {
581 let mut event = build_audit_storage_event(
582 gate_req,
583 &audit,
584 EventOutcome::Error,
585 Some(crate::cost_unit::base_resource_payload(request_id)),
586 );
587 if let Some(target) = effective_target {
588 event = event.with_target(target);
589 }
590 let _ = append_audit_event_best_effort(
591 self.audit_batch.as_ref(),
592 store,
593 event,
594 gate_req.verb.as_str(),
595 crate::audit_batch::AuditProducer::GateUnavailable,
596 false,
597 )
598 .await;
599 }
600 RuntimeError::GateUnavailable {
601 verb: gate_req.verb.clone(),
602 // Caller-visible: a stable, classified reason derived from the
603 // `GateError` variant only. `error`'s `Display` text is logged
604 // above (server-side, via `tracing::warn!`) and must never be
605 // interpolated here — a gate backend's error message can embed
606 // connection details, addresses, or credentials.
607 reason: error.wire_reason().to_string(),
608 }
609 }
610
611 /// Dispatch a verb to the first pack that handles it.
612 ///
613 /// Routes through the gate, then invokes the matching pack handler. When
614 /// `params["help"] == true`, short-circuits to `describe_verb` with no side effects.
615 /// Gate errors fail closed. Full dispatch flow documented in `docs/protocol.md`.
616 ///
617 /// Equivalent to `self.dispatch_with_identity(verb, params, None)` — uses
618 /// this registry's construction-baked `default_namespace` / `actor_id` /
619 /// `visible_namespaces`.
620 pub async fn dispatch(&self, verb: &str, params: Value) -> Result<Value, RuntimeError> {
621 self.dispatch_with_identity(verb, params, None).await
622 }
623
624 /// Dispatch a verb, optionally overriding this registry's baked identity
625 /// scalars for exactly this call (ADR-096 Fork 1).
626 ///
627 /// `identity = None` behaves exactly like [`Self::dispatch`]. `identity =
628 /// Some(id)` uses `id.namespace` / `id.actor_id` / `id.visible_namespaces`
629 /// in place of `self.default_namespace` / `self.actor_id` /
630 /// `self.visible_namespaces` for this call's namespace resolution, gate
631 /// request, and token minting. The registry's own fields are never mutated,
632 /// so concurrent calls with different (or no) identity are independent.
633 /// See `docs/api/pack.md#dispatch_with_identity` for why this enables one warm
634 /// registry to serve many attribution identities over a shared backend.
635 pub async fn dispatch_with_identity(
636 &self,
637 verb: &str,
638 params: Value,
639 identity: Option<RequestIdentity>,
640 ) -> Result<Value, RuntimeError> {
641 self.dispatch_with_disposition(verb, params, identity)
642 .await
643 .map_err(DispatchError::into_source)
644 }
645
646 /// Dispatch with provenance for this operation's own domain result.
647 /// Errors returned by a nested dispatch remain handler errors at this boundary.
648 pub async fn dispatch_with_disposition(
649 &self,
650 verb: &str,
651 params: Value,
652 identity: Option<RequestIdentity>,
653 ) -> Result<Value, DispatchError> {
654 // help=true interception: short-circuit before gate/pack.
655 if params.get("help").and_then(Value::as_bool) == Some(true) {
656 let result = match self.describe_verb(verb) {
657 Ok(value) => Ok(value),
658 Err(error) => match self.mounted_verb_catalog().await {
659 Ok(catalog) => catalog
660 .into_iter()
661 .find(|entry| entry["verb"] == verb)
662 .ok_or(error),
663 Err(error) => Err(error),
664 },
665 };
666 return result.map_err(DispatchError::before_dispatch);
667 }
668 // Resolve namespace before `params` is moved into pack.dispatch, so the
669 // post-dispatch hook can reference it.
670 //
671 // Absent `namespace` and a present-but-malformed `namespace` are
672 // different cases. A present non-string value (null, number, bool,
673 // array, object) is explicit caller input that failed to parse and
674 // must fail closed, not silently coerce to the default namespace.
675 // Only a genuinely absent key defaults. Shared with the multi-backend
676 // coordinator intercept via `resolve_explicit_namespace` so every MCP
677 // ingress path applies the same fail-closed rule.
678 let explicit_namespace = params.get("namespace").is_some_and(Value::is_string);
679 // The caller-supplied correlation id (khive#948), if any. Read once
680 // here so it is in scope for every audit-append site below,
681 // including the ones that run before pack dispatch is attempted.
682 let request_id: Option<u64> = identity.as_ref().and_then(|id| id.request_id);
683 // Thread the configured actor identity into the gate request so the
684 // gate can distinguish human vs agent callers at the dispatch seam.
685 // Resolved once via the shared actor-identity policy and reused for
686 // token minting below, so the gate's notion of "who is the caller"
687 // and the storage token's notion can never drift apart.
688 let gate_req = self
689 .gate_request_with_identity(verb, ¶ms, identity.as_ref())
690 .map_err(DispatchError::before_dispatch)?;
691 let ns = gate_req.namespace.clone();
692 let resolved_actor = gate_req.actor.clone();
693
694 // Consult the gate.
695 //
696 // - Ok(Allow) → proceed to pack dispatch (tracing + optional EventStore).
697 // - Ok(Deny) → emit audit, persist if store configured, return PermissionDenied.
698 // - Err(_) → emit an outage audit and return GateUnavailable.
699 let gate_decision = khive_gate::check_with_mailbox_policy(self.gate.as_ref(), &gate_req);
700 let (gate_blocked, mut deferred_audit) = match gate_decision {
701 Ok(decision) => {
702 let is_deny = matches!(decision, GateDecision::Deny { .. });
703
704 // Emit audit event via tracing.
705 let audit = masked_audit_event(&gate_req, &decision, self.gate.impl_name());
706 tracing::info!(
707 audit_event = %serde_json::to_string(&audit)
708 .unwrap_or_else(|_| "{\"error\":\"serialize\"}".into()),
709 "gate.check"
710 );
711
712 // Drain any process-lifetime `OnceLock` config locks queued
713 // since the last dispatch and persist them as `ConfigLocked`
714 // events, riding this same audit-persistence gate. The
715 // namespace/actor stamped on these rows are whichever
716 // dispatch happens to observe the queue non-empty first:
717 // an accepted provenance quirk, preferred over threading an
718 // `EventStore` handle into every synchronous
719 // `OnceLock::get_or_init` call site. The verb column is NOT
720 // inherited from that bystander dispatch: a config-lock row
721 // wearing an operation verb pollutes verb-filtered queries
722 // (e.g. per-verb receipt counts), so these rows carry their
723 // own `config.lock` pseudo-verb and remain discoverable by
724 // `EventKind::ConfigLocked`.
725 if let Some(store) = &self.event_store {
726 if crate::config_ledger::PENDING
727 .swap(false, std::sync::atomic::Ordering::AcqRel)
728 {
729 for (key, value) in crate::config_ledger::drain_config_locked() {
730 let payload = serde_json::json!({ "key": key, "value": value });
731 let storage_event = Event::new(
732 gate_req.namespace.as_str(),
733 "config.lock",
734 EventKind::ConfigLocked,
735 SubstrateKind::Event,
736 format!("{}:{}", gate_req.actor.kind, gate_req.actor.id),
737 )
738 .with_payload(payload);
739 // ConfigLocked is pure observability: the helper
740 // never returns `Err` for it, so there is
741 // nothing to fold.
742 let _ = append_audit_event_best_effort(
743 self.audit_batch.as_ref(),
744 store,
745 storage_event,
746 "config.lock",
747 crate::audit_batch::AuditProducer::ConfigLocked,
748 false,
749 )
750 .await;
751 }
752 }
753 }
754
755 // Every Allow-outcome audit row defers its append until pack
756 // dispatch returns, so the row can carry the measured
757 // dispatch time in `duration_us` (persisting before dispatch
758 // ran always recorded the `Event::new` default of 0). A
759 // singleton `link` call (no `links` bulk array) additionally
760 // enriches the deferred row with the created/resolved edge
761 // fields (schema v2) once dispatch resolves. Denied calls
762 // have no dispatch to wait for and keep the immediate v1
763 // append below.
764 //
765 // Accepted trade-off for ordinary verbs: a crash between this
766 // Allow decision and the deferred append loses the audit row.
767 // `git.digest` narrows the caller-visible contract below: it
768 // never returns success until the deferred receipt append is
769 // confirmed, though a process crash can still leave committed
770 // ingest writes with no response and no completed receipt.
771 let defer_audit = !is_deny;
772
773 // Persist to EventStore immediately only for denied calls;
774 // the receipt rides on the refusal so the caller can cite
775 // the row.
776 let reason = if is_deny {
777 let reason = match decision {
778 GateDecision::Deny { reason } => reason,
779 _ => String::new(),
780 };
781 let receipt = match &self.event_store {
782 Some(store) => {
783 // ADR-103 Decision (a): the closed `work_class` enum
784 // is stamped on every event, denial included -- only
785 // `resource.cost_unit` is scoped to a successful
786 // dispatch by Amendment 1. `base_resource_payload()`
787 // carries `work_class` alone, no `cost_unit` key.
788 let storage_event = build_audit_storage_event(
789 &gate_req,
790 &audit,
791 EventOutcome::Denied,
792 Some(crate::cost_unit::base_resource_payload(request_id)),
793 );
794 // This path always returns `PermissionDenied`
795 // below, so there is no success outcome to fold a
796 // commit failure into; the receipt says whether
797 // the row exists.
798 self.append_gate_denied_row(store, storage_event, verb)
799 .await
800 }
801 None => crate::error::DenialReceipt::no_store(),
802 };
803 Some((reason, receipt))
804 } else {
805 None
806 };
807 let deferred = if defer_audit { Some(audit) } else { None };
808 (reason, deferred)
809 }
810 Err(err) => {
811 return Err(DispatchError::before_dispatch(
812 self.gate_unavailable_error(&gate_req, &err, request_id, None)
813 .await,
814 ));
815 }
816 };
817
818 // Hard enforcement: Deny is authoritative.
819 if let Some((reason, receipt)) = gate_blocked {
820 return Err(DispatchError::before_dispatch(
821 RuntimeError::PermissionDenied {
822 verb: verb.to_string(),
823 reason,
824 receipt: Box::new(receipt),
825 },
826 ));
827 }
828
829 // Mint the authorized storage token at the dispatch boundary.
830 //
831 // Writes pin to `local` by default. Actor identity and config
832 // `[actor] id` are attribution and gate-context inputs only: they
833 // never route storage. The explicit `namespace=` request param is a
834 // precise single-namespace escape: the caller deliberately
835 // reads/writes exactly that one set; it is NOT widened by `visible_namespaces`.
836 //
837 // When actor_id is configured, mint a token carrying that actor
838 // label so that comm.inbox applies the to_actor filter for directed delivery.
839 // Otherwise, use ActorRef::anonymous() and inbox falls back to party-line.
840 // `actor_id_str` already reflects the per-request identity override
841 // when supplied (resolved above into `resolved_actor`, mirrored into
842 // the gate request). Reusing the same value here guarantees the
843 // gate's actor and the storage token's actor can never diverge.
844 //
845 // On the default (no explicit `namespace=`) path, the read scope
846 // widens to `['local'] ∪ visible_namespaces` (baked, or the
847 // per-request override). `'local'` is always included
848 // (mint_with_visibility deduplicates). Writes remain pinned to
849 // `'local'`. Per-actor distinctions use view-layer tag filters
850 // (assignee, actor_id, from/to), not namespace partitions. `ns`/
851 // `explicit_namespace` were already validated above: reuse them
852 // instead of re-reading `params["namespace"]` with `as_str()`, which
853 // would silently drop malformed non-string values again.
854 let token = if explicit_namespace {
855 // Explicit escape: precise single-namespace scope, read+write. NOT widened.
856 NamespaceToken::mint_with_visibility(ns.clone(), vec![], resolved_actor)
857 } else {
858 // Default path: write namespace = local; read scope = ['local'] ∪ visible_namespaces.
859 let primary = Namespace::local();
860 let mut extra_visible: Vec<Namespace> = match identity.as_ref() {
861 Some(id) => id
862 .visible_namespaces
863 .iter()
864 .filter_map(|s| match Namespace::parse(s) {
865 Ok(parsed) => Some(parsed),
866 Err(e) => {
867 tracing::warn!(
868 namespace = %s,
869 error = %e,
870 "dispatch_with_identity: skipping invalid visible_namespace \
871 entry from per-request identity"
872 );
873 None
874 }
875 })
876 .collect(),
877 None => self.visible_namespaces.clone(),
878 };
879 // ADR-007 Rev 4 Rule 3b, applied once at the seam every identity
880 // path shares: a non-`local` actor reads its own namespace by
881 // default (its episodic memories land there), whether the identity
882 // came from the config loader, a daemon frame, a scheduled replay
883 // or an embedding host. Writes stay pinned to `local` (Rule 0).
884 if let Some(actor_namespace) = resolved_actor
885 .binding_id()
886 .filter(|id| *id != Namespace::LOCAL)
887 .and_then(|id| Namespace::parse(id).ok())
888 {
889 extra_visible.push(actor_namespace);
890 }
891 extra_visible.push(Namespace::local()); // 'local' always readable; mint dedups
892 NamespaceToken::mint_with_visibility(primary, extra_visible, resolved_actor)
893 }
894 .with_gate_namespace(ns.clone())
895 .with_gate_explicit_namespace(
896 params
897 .get("namespace")
898 .and_then(Value::as_str)
899 .map(str::to_owned),
900 )
901 .with_request_id(request_id)
902 .with_process_ref(match identity.as_ref() {
903 Some(id) => id.process_ref.clone(),
904 None => crate::config::process_ref_from_env(),
905 });
906
907 for pack in self.packs.iter() {
908 let handler_def = pack.handlers().iter().find(|v| v.name == verb);
909 let mounted_name = pack.mounted_namespace().and_then(|prefix| {
910 verb.strip_prefix(prefix)
911 .and_then(|suffix| suffix.strip_prefix('.'))
912 });
913 if handler_def.is_some() || mounted_name.is_some() {
914 let definition = if let Some(name) = mounted_name {
915 pack.mounted_catalog().await.and_then(|catalog| {
916 catalog
917 .into_iter()
918 .find(|definition| definition.name == name)
919 .map(Some)
920 .ok_or_else(|| RuntimeError::UnknownVerb(verb.to_owned()))
921 })
922 } else {
923 Ok(None)
924 };
925 // Strip `namespace` from params before forwarding to packs.
926 // The registry has already consumed it to mint the NamespaceToken.
927 //
928 // Exception: if the handler's own `params` schema declares
929 // `"namespace"` as a valid field (e.g. brain.bind, brain.unbind,
930 // brain.bindings, brain.resolve), the field is a *business* argument
931 // — not a transport routing key — and must be passed through
932 // unchanged. Stripping it would silently default the binding to the
933 // "*" wildcard, broadening profile scope across namespaces.
934 let handler_accepts_namespace = handler_def
935 .is_some_and(|h| h.params.iter().any(|p| p.name == "namespace"))
936 || definition
937 .as_ref()
938 .ok()
939 .and_then(|value| value.as_ref())
940 .is_some_and(|definition| {
941 definition
942 .input_schema
943 .get("properties")
944 .is_some_and(|properties| properties.get("namespace").is_some())
945 });
946 let params = if !handler_accepts_namespace {
947 if let Value::Object(mut map) = params {
948 map.remove("namespace");
949 Value::Object(map)
950 } else {
951 params
952 }
953 } else {
954 params
955 };
956 let dispatch_start = Instant::now();
957 let mounted_audit = definition.as_ref().ok().and_then(|v| v.as_ref()).map(|v| {
958 serde_json::json!({"mount": pack.name(), "effect": v.effect, "generation": v.generation})
959 });
960 let mut result = match definition {
961 Ok(Some(definition)) => {
962 pack.dispatch_mounted(&definition, verb, params, self, &token)
963 .await
964 }
965 Ok(None) => pack.dispatch(verb, params, self, &token).await,
966 Err(error) => Err(error),
967 };
968 let domain_succeeded = result.is_ok();
969 let dispatch_us = dispatch_start.elapsed().as_micros() as i64;
970
971 // Unlike ordinary audit rows, a successful `git.digest`
972 // response is returned only after its complete report has
973 // been durably persisted as a schema-v2 audit receipt. The
974 // receipt helper borrows the deferred audit row so malformed
975 // handler output can still fall back to one generic Error
976 // audit. Handler errors use that same ordinary path below.
977 let git_digest_receipt_outcome = if verb == "git.digest" && result.is_ok() {
978 let resource = result.as_ref().ok().map(|value| {
979 crate::cost_unit::resource_payload(
980 verb,
981 &gate_req.args,
982 value,
983 || pack.registered_embedding_model_names().len() as i64,
984 request_id,
985 )
986 });
987 Some(
988 persist_git_digest_receipt(
989 self.event_store.as_ref(),
990 self.audit_batch.as_ref(),
991 &gate_req,
992 deferred_audit.as_ref(),
993 &mut result,
994 dispatch_us,
995 resource,
996 )
997 .await,
998 )
999 } else {
1000 None
1001 };
1002
1003 // Append the deferred Allow-outcome audit row now that
1004 // dispatch has resolved, so `duration_us` carries the
1005 // measured `dispatch_us` instead of the `Event::new` default
1006 // of 0. A successful singleton `link` call enriches the row
1007 // with the created/resolved edge (schema v2); anything that
1008 // cannot be enriched, or is not a singleton `link` call,
1009 // falls back to the generic v1 audit shape so no audit row
1010 // is ever dropped for the deferred path.
1011 let needs_generic_audit = git_digest_receipt_outcome.is_none()
1012 || git_digest_receipt_outcome == Some(GitDigestReceiptOutcome::BuildRejected);
1013 if let (true, Some(audit)) = (needs_generic_audit, deferred_audit.take()) {
1014 if let Some(store) = &self.event_store {
1015 let is_link_singleton =
1016 verb == "link" && gate_req.args.get("links").is_none();
1017 // Read-only pass over `result` first: every arm below
1018 // only needs `audit_outcome` afterward, and folding a
1019 // failure into `result` requires a mutable borrow
1020 // that cannot coexist with the `&result` match below.
1021 let audit_outcome: Result<(), AuditObligationFailure> = match &result {
1022 Ok(ok_val) if is_link_singleton => {
1023 // ADR-103 Amendment 1: `link` (singleton or
1024 // bulk) has no embedding-bearing path — edges
1025 // carry no embedded body — so cost_unit is
1026 // always base_weight("link") alone. The
1027 // registered-model closure is never invoked
1028 // (per_item_weight("link", ..) short-circuits
1029 // to 0 before `model_count` reads it).
1030 let resource = crate::cost_unit::resource_payload(
1031 verb,
1032 &gate_req.args,
1033 ok_val,
1034 || pack.registered_embedding_model_names().len() as i64,
1035 request_id,
1036 );
1037 match link_audit_success_from_result(audit.clone(), ok_val) {
1038 Some((edge_id, mut payload)) => {
1039 if let Value::Object(ref mut map) = payload {
1040 map.insert("resource".to_string(), resource);
1041 }
1042 let storage_event = Event::new(
1043 gate_req.namespace.as_str(),
1044 gate_req.verb.as_str(),
1045 EventKind::Audit,
1046 SubstrateKind::Event,
1047 format!(
1048 "{}:{}",
1049 gate_req.actor.kind, gate_req.actor.id
1050 ),
1051 )
1052 .with_outcome(EventOutcome::Success)
1053 .with_target(edge_id)
1054 .with_payload(payload)
1055 .with_payload_schema_version(2)
1056 .with_duration_us(dispatch_us);
1057 append_audit_event_best_effort(
1058 self.audit_batch.as_ref(),
1059 store,
1060 storage_event,
1061 verb,
1062 crate::audit_batch::AuditProducer::DispatchSucceeded,
1063 self.admission_degrade_safe(verb),
1064 )
1065 .await
1066 }
1067 None => {
1068 tracing::warn!(
1069 verb,
1070 "link audit v2 enrichment parse failed; \
1071 falling back to v1 audit shape"
1072 );
1073 let storage_event = build_audit_storage_event(
1074 &gate_req,
1075 &audit,
1076 EventOutcome::Success,
1077 Some(resource),
1078 )
1079 .with_duration_us(dispatch_us);
1080 append_audit_event_best_effort(
1081 self.audit_batch.as_ref(),
1082 store,
1083 storage_event,
1084 verb,
1085 crate::audit_batch::AuditProducer::DispatchSucceeded,
1086 self.admission_degrade_safe(verb),
1087 )
1088 .await
1089 }
1090 }
1091 }
1092 _ => {
1093 // The persisted audit outcome must reflect
1094 // the dispatch result, not be hardcoded to
1095 // Success — otherwise a failed dispatch is
1096 // recorded as successful work and disappears
1097 // from `outcome=error` queries.
1098 //
1099 // ADR-103 Amendment 1: `resource.cost_unit` is
1100 // computed ONLY on a successful dispatch —
1101 // there is no handler `Value` to read
1102 // `item_count` from on an error, and the
1103 // amendment's "absence has exactly two
1104 // meanings" rule requires the field be
1105 // omitted, never defaulted to 0, on an
1106 // errored dispatch. `work_class` itself is
1107 // NOT one of those two omission cases
1108 // (ADR-103 Decision (a) stamps it on every
1109 // event), so an errored dispatch still gets
1110 // `resource: {"work_class": "interactive"}`,
1111 // just with no `cost_unit` key.
1112 let (outcome, resource) = match &result {
1113 Ok(ok_val) => (
1114 EventOutcome::Success,
1115 Some(crate::cost_unit::resource_payload(
1116 verb,
1117 &gate_req.args,
1118 ok_val,
1119 || pack.registered_embedding_model_names().len() as i64,
1120 request_id,
1121 )),
1122 ),
1123 Err(_) => (
1124 EventOutcome::Error,
1125 Some(crate::cost_unit::base_resource_payload(request_id)),
1126 ),
1127 };
1128 let producer = if result.is_ok() {
1129 crate::audit_batch::AuditProducer::DispatchSucceeded
1130 } else {
1131 crate::audit_batch::AuditProducer::DispatchFailed
1132 };
1133 let mut storage_event =
1134 build_audit_storage_event(&gate_req, &audit, outcome, resource)
1135 .with_duration_us(dispatch_us);
1136 if let Some(metadata) = &mounted_audit {
1137 storage_event.payload["mounted_tool"] = metadata.clone();
1138 }
1139 append_audit_event_best_effort(
1140 self.audit_batch.as_ref(),
1141 store,
1142 storage_event,
1143 verb,
1144 producer,
1145 self.admission_degrade_safe(verb),
1146 )
1147 .await
1148 }
1149 };
1150 // Only a would-be-success dispatch can be flipped by
1151 // an obligation failure (ADR-133 D2/D3/D4): an
1152 // already-erroring dispatch (DispatchFailed producer)
1153 // keeps its original error, matching
1154 // `fold_audit_obligation`'s contract.
1155 result =
1156 fold_audit_obligation(result, audit_outcome, std::convert::identity);
1157 }
1158 }
1159
1160 // Post-dispatch hook: fires on success, opt-in.
1161 if let (Ok(ref ok_val), Some(hook)) = (&result, &self.dispatch_hook) {
1162 let mut dispatch_event = Event::new(
1163 ns.as_str(),
1164 verb,
1165 EventKind::Audit,
1166 SubstrateKind::Event,
1167 pack.name(),
1168 )
1169 .with_outcome(EventOutcome::Success)
1170 .with_duration_us(dispatch_us);
1171
1172 // For recall verbs: extract the first result's id as
1173 // target_id so the brain temporal posterior can observe
1174 // real hit/miss and latency. Copy the serve-attribution
1175 // fields from that same hit so the hook credits the profile
1176 // that actually served instead of always crediting default.
1177 if verb == "memory.recall" {
1178 let first_result =
1179 ok_val.as_array().and_then(|arr| arr.first()).or_else(|| {
1180 ok_val
1181 .get("results")
1182 .and_then(Value::as_array)
1183 .and_then(|arr| arr.first())
1184 });
1185 let first_note_id = first_result
1186 .and_then(|v| v.get("id"))
1187 .and_then(|v| v.as_str())
1188 .and_then(|s| s.parse::<uuid::Uuid>().ok());
1189 if let Some(note_id) = first_note_id {
1190 dispatch_event = dispatch_event.with_target(note_id);
1191 }
1192 let mut payload = serde_json::Map::new();
1193 if let Some(profile_id) = first_result
1194 .and_then(|v| v.get("served_by_profile_id"))
1195 .and_then(Value::as_str)
1196 {
1197 payload.insert(
1198 "served_by_profile_id".to_string(),
1199 Value::String(profile_id.to_string()),
1200 );
1201 }
1202 if let Some(attribution) = first_result
1203 .and_then(|v| v.get("serve_attribution"))
1204 .and_then(Value::as_str)
1205 {
1206 payload.insert(
1207 "serve_attribution".to_string(),
1208 Value::String(attribution.to_string()),
1209 );
1210 }
1211 dispatch_event = dispatch_event.with_payload(Value::Object(payload));
1212 // No first result → target_id stays None (RecallMiss
1213 // in brain's event interpreter).
1214 }
1215
1216 let dispatch_view = EventView {
1217 event: dispatch_event,
1218 observations: Vec::new(),
1219 };
1220 let hook = Arc::clone(hook);
1221 hook.on_dispatch(&dispatch_view).await;
1222 }
1223
1224 // Recently-referenced ring admission: only by-id touches admit
1225 // an id. Runs unconditionally (not gated on `dispatch_hook`,
1226 // which is opt-in) because the ring is a core
1227 // dispatch-boundary capability, not an observer.
1228 //
1229 // Keyed on `token.namespace()`, NOT `ns`: `ns` is the
1230 // gate-resolved namespace, which on the default
1231 // (non-explicit) dispatch path can be a non-local
1232 // `default_namespace` (e.g. "foreign") while the storage
1233 // token that actually created/touched the record is pinned
1234 // to `local`. The ring must be keyed on the namespace the
1235 // record actually lives in: the same namespace
1236 // `resolve_reference`'s ring lookup uses: or admission and
1237 // lookup silently diverge on any non-local `default_namespace`
1238 // config.
1239 if let Ok(ref ok_val) = result {
1240 let admissions = crate::reference_ring::ring_admissions_for(verb, ok_val);
1241 if !admissions.is_empty() {
1242 let actor_key = format!("{}:{}", gate_req.actor.kind, gate_req.actor.id);
1243 for (id, name) in admissions {
1244 self.reference_ring.admit(
1245 token.namespace().as_str(),
1246 &actor_key,
1247 id,
1248 name,
1249 );
1250 }
1251 }
1252 }
1253
1254 return result
1255 .map_err(|error| DispatchError::after_handler(error, domain_succeeded));
1256 }
1257 }
1258
1259 // No pack owns this verb: the gate allowed it, but no dispatch runs.
1260 // Persist the deferred audit row now (duration stays at the
1261 // `Event::new` default of 0 — no dispatch occurred to measure) so an
1262 // allowed-but-unknown verb is never silently dropped from the audit
1263 // trail (matches the "no audit row is ever dropped" contract above).
1264 if let Some(audit) = deferred_audit.take() {
1265 if let Some(store) = &self.event_store {
1266 // Dispatch is about to return `UnknownVerb` below (no pack
1267 // owns this verb), so the persisted outcome must be `Error`,
1268 // not `Success`. `work_class` is still stamped (ADR-103
1269 // Decision (a)); `resource.cost_unit` is omitted, matching
1270 // every other errored-dispatch row.
1271 let storage_event = build_audit_storage_event(
1272 &gate_req,
1273 &audit,
1274 EventOutcome::Error,
1275 Some(crate::cost_unit::base_resource_payload(request_id)),
1276 );
1277 // Dispatch already returns `UnknownVerb` below regardless, so
1278 // — as with the deny paths above — there is no success
1279 // outcome to fold a commit failure into.
1280 let _ = append_audit_event_best_effort(
1281 self.audit_batch.as_ref(),
1282 store,
1283 storage_event,
1284 verb,
1285 crate::audit_batch::AuditProducer::UnknownVerb,
1286 false,
1287 )
1288 .await;
1289 }
1290 }
1291
1292 // Verb-visibility handler names, precomputed at build() time (internal
1293 // subhandlers are excluded so they are not advertised in the
1294 // unknown-verb error).
1295 Err(DispatchError::before_dispatch(RuntimeError::UnknownVerb(
1296 format!(
1297 "unknown verb {verb:?}; available: {}",
1298 self.available_verbs.join(", ")
1299 ),
1300 )))
1301 }
1302
1303 /// Dispatch a verb under an out-of-band verified actor identity.
1304 ///
1305 /// `verified_actor` is a typed [`VerifiedActor`] (constructor rejects blank
1306 /// identifiers) — only code holding a `VerbRegistry` handle can supply it.
1307 /// `dispatch_as` never reads `params["actor"]` to derive the effective actor;
1308 /// individual verbs may still accept an `actor` field for their own documented
1309 /// business semantics, unrelated to the acting principal. Every pack handler
1310 /// that reads "who is calling" resolves it from the `NamespaceToken` the
1311 /// dispatch boundary mints, so `verified_actor` becomes exactly the principal
1312 /// those handlers observe.
1313 ///
1314 /// Equivalent to `dispatch_with_identity(verb, params, Some(identity))` with
1315 /// `identity.actor_id = Some(verified_actor)` and every other identity scalar
1316 /// (namespace, visible namespaces) left at this registry's construction-baked
1317 /// value. [`Self::dispatch`] and [`Self::dispatch_with_identity`] are unaffected.
1318 /// See `docs/api/pack.md#dispatch_as` for the embedding-host use case and the
1319 /// blank-identifier safety rationale.
1320 pub async fn dispatch_as(
1321 &self,
1322 verb: &str,
1323 params: Value,
1324 verified_actor: VerifiedActor,
1325 ) -> Result<Value, RuntimeError> {
1326 let identity = RequestIdentity {
1327 namespace: self.default_namespace.clone(),
1328 actor_id: Some(verified_actor.into_inner()),
1329 visible_namespaces: self
1330 .visible_namespaces
1331 .iter()
1332 .map(|ns| ns.as_str().to_string())
1333 .collect(),
1334 process_ref: crate::config::process_ref_from_env(),
1335 request_id: None,
1336 };
1337 self.dispatch_with_identity(verb, params, Some(identity))
1338 .await
1339 }
1340}