1use std::fmt;
4use std::time::Duration;
5
6use thiserror::Error;
7use uuid::Uuid;
8
9pub type RuntimeResult<T> = Result<T, RuntimeError>;
11
12#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize)]
14#[serde(rename_all = "snake_case")]
15pub enum DomainDisposition {
16 Committed,
18 NotCommitted,
20 Unknown,
22}
23
24impl DomainDisposition {
25 pub const fn as_str(self) -> &'static str {
27 match self {
28 Self::Committed => "committed",
29 Self::NotCommitted => "not_committed",
30 Self::Unknown => "unknown",
31 }
32 }
33}
34
35#[derive(Debug, Error)]
37#[error("{source}")]
38pub struct DispatchError {
39 #[source]
40 source: RuntimeError,
41 disposition: DomainDisposition,
42}
43
44impl DispatchError {
45 pub(crate) fn new(source: RuntimeError, disposition: DomainDisposition) -> Self {
46 Self {
47 source,
48 disposition,
49 }
50 }
51
52 pub(crate) fn before_dispatch(source: RuntimeError) -> Self {
53 Self::new(source, DomainDisposition::NotCommitted)
54 }
55
56 pub(crate) fn after_handler(source: RuntimeError, domain_succeeded: bool) -> Self {
57 Self::new(
58 source,
59 if domain_succeeded {
60 DomainDisposition::Committed
61 } else {
62 DomainDisposition::Unknown
63 },
64 )
65 }
66
67 pub fn source(&self) -> &RuntimeError {
69 &self.source
70 }
71
72 pub fn disposition(&self) -> DomainDisposition {
74 self.disposition
75 }
76
77 pub fn into_parts(self) -> (RuntimeError, DomainDisposition) {
79 (self.source, self.disposition)
80 }
81
82 pub fn into_source(self) -> RuntimeError {
84 self.source
85 }
86}
87
88#[derive(Debug, Clone, Copy, PartialEq, Eq)]
90pub enum AuditObligationReason {
91 Terminal(crate::audit_batch::AuditTerminalReason),
93 GitDigestReceiptFailure,
95}
96
97#[derive(Debug, Error)]
99#[error("{message}")]
100pub struct AuditObligationFailure {
101 pub reason: AuditObligationReason,
103 pub verb: String,
105 pub message: String,
107 #[source]
108 source: Option<khive_storage::StorageError>,
109}
110
111impl AuditObligationFailure {
112 pub fn new(verb: impl Into<String>, reason: crate::audit_batch::AuditTerminalReason) -> Self {
114 let verb = verb.into();
115 Self {
116 message: format!("audit obligation commit failed for verb {verb:?}: {reason:?}"),
117 verb,
118 reason: AuditObligationReason::Terminal(reason),
119 source: None,
120 }
121 }
122
123 pub(crate) fn from_store(verb: &str, source: khive_storage::StorageError) -> Self {
124 let mut failure = Self::new(verb, crate::audit_batch::AuditTerminalReason::StoreFailure);
125 failure.message = format!("audit obligation commit failed for verb {verb:?}: {source}");
126 failure.source = Some(source);
127 failure
128 }
129
130 pub(crate) fn git_digest_receipt(branch: &'static str) -> Self {
131 Self {
132 reason: AuditObligationReason::GitDigestReceiptFailure,
133 verb: "git.digest".into(),
134 message: format!(
135 "git_digest_receipt_persist_failed: {branch}; git.digest writes may have committed, but no durable success receipt was confirmed; inspect ingest state before retrying"
136 ),
137 source: None,
138 }
139 }
140
141 pub const fn wire_code(&self) -> &'static str {
143 use crate::audit_batch::AuditTerminalReason;
144 match self.reason {
145 AuditObligationReason::Terminal(reason) => match reason {
146 AuditTerminalReason::PreflightRejected => "preflight_rejected",
147 AuditTerminalReason::AdmissionClosed => "admission_closed",
148 AuditTerminalReason::QueueAdmissionExhausted => "queue_admission_exhausted",
149 AuditTerminalReason::AdmissionDeadlineExpired => "admission_deadline_expired",
150 AuditTerminalReason::ResolutionDeadlineExpired => "resolution_deadline_expired",
151 AuditTerminalReason::IdentityConflict => "identity_conflict",
152 AuditTerminalReason::StoreFailure => "store_failure",
153 AuditTerminalReason::RetryExhausted => "retry_exhausted",
154 AuditTerminalReason::IdempotencyUnsupported => "idempotency_unsupported",
155 AuditTerminalReason::DriverPanicked => "driver_panicked",
156 AuditTerminalReason::DriverCancelled => "driver_cancelled",
157 AuditTerminalReason::DriverJoinLost => "driver_join_lost",
158 AuditTerminalReason::DriverExitedInconsistent => "driver_exited_inconsistent",
159 AuditTerminalReason::DriverAppendAbandoned => "driver_append_abandoned",
160 AuditTerminalReason::StoreWedged => "store_wedged",
161 },
162 AuditObligationReason::GitDigestReceiptFailure => "git_digest_receipt_failure",
163 }
164 }
165}
166
167pub const WRITER_POOL_CHECKOUT_TIMEOUT_STAGE: &str = "writer_pool_checkout_timeout";
170
171pub const WRITER_QUEUE_SATURATED_STAGE: &str = "writer_queue_saturated";
174
175pub const WRITER_TASK_BEGIN_BUSY_STAGE: &str = "writer_task_begin_busy";
178
179pub const WRITER_TASK_REQUEST_FAILED_STAGE: &str = "writer_task_request_failed";
182
183pub const WRITER_TASK_TERMINATED_STAGE: &str = "writer_task_terminated";
186
187pub const STORAGE_ADMISSION_TIMEOUT_STAGE: &str = "storage_admission_timeout";
192
193pub const READ_TX_AGE_EVICTED_STAGE: &str = "read_tx_age_evicted";
204
205pub const WRITER_ADMISSION_SCOPE: &str = "writer_admission";
212
213#[derive(Debug, Clone, PartialEq, Eq)]
220pub struct AdmissionFailureContext {
221 pub stage: &'static str,
225 pub timeout: Duration,
227 pub capability: Option<khive_storage::StorageCapability>,
231 pub operation: Option<String>,
233 pub pool_identity: Option<String>,
235 pub scope: Option<&'static str>,
240 pub retry_after_ms: Option<u64>,
246}
247
248#[derive(Debug, Clone, PartialEq, Eq)]
255pub struct RetryableFailureContext {
256 pub stage: &'static str,
258 pub timeout: Duration,
260 pub capability: Option<khive_storage::StorageCapability>,
262 pub operation: Option<String>,
264 pub pool_identity: Option<String>,
266 pub scope: Option<&'static str>,
268 pub retry_after_ms: Option<u64>,
270}
271
272#[derive(Debug, Clone, Copy, PartialEq, Eq)]
279pub struct WriterTaskFailureContext {
280 pub stage: &'static str,
282 pub request_state: khive_storage::WriterTaskRequestState,
284 pub task_terminated: bool,
286 pub retryable: bool,
288}
289
290impl From<AdmissionFailureContext> for RetryableFailureContext {
291 fn from(context: AdmissionFailureContext) -> Self {
292 Self {
293 stage: context.stage,
294 timeout: context.timeout,
295 capability: context.capability,
296 operation: context.operation,
297 pool_identity: context.pool_identity,
298 scope: context.scope,
299 retry_after_ms: context.retry_after_ms,
300 }
301 }
302}
303
304#[derive(Debug, Clone, Copy, PartialEq, Eq)]
310pub enum ChannelIngestFailureClass {
311 Retryable { reason: &'static str },
313 Permanent { reason: &'static str },
315 Unknown { reason: &'static str },
317}
318
319impl ChannelIngestFailureClass {
320 pub const fn name(self) -> &'static str {
322 match self {
323 Self::Retryable { .. } => "retryable",
324 Self::Permanent { .. } => "permanent",
325 Self::Unknown { .. } => "unknown",
326 }
327 }
328
329 pub const fn reason(self) -> &'static str {
331 match self {
332 Self::Retryable { reason } | Self::Permanent { reason } | Self::Unknown { reason } => {
333 reason
334 }
335 }
336 }
337}
338
339#[derive(Debug, Clone, PartialEq, Eq)]
342pub struct WriterPoolCheckoutTimeoutContext {
343 pub timeout: Duration,
346 pub capability: Option<khive_storage::StorageCapability>,
348 pub operation: Option<String>,
350}
351
352#[derive(Debug, Clone, PartialEq, Eq)]
357pub struct GuardedWriteFailure {
358 pub entry_index: Option<usize>,
361 pub missing_source: Option<Uuid>,
364 pub missing_target: Option<Uuid>,
367}
368
369impl fmt::Display for GuardedWriteFailure {
370 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
371 let mut missing = Vec::new();
372 if let Some(source) = self.missing_source {
373 missing.push(format!("source {source}"));
374 }
375 if let Some(target) = self.missing_target {
376 missing.push(format!("target {target}"));
377 }
378 let missing = if missing.is_empty() {
379 "endpoint(s)".to_string()
380 } else {
381 missing.join(" and ")
382 };
383 match self.entry_index {
384 Some(index) => write!(
385 f,
386 "batch entry {index}: {missing} no longer exist at write time"
387 ),
388 None => write!(f, "{missing} no longer exist at write time"),
389 }
390 }
391}
392
393impl std::error::Error for GuardedWriteFailure {}
394
395#[derive(Debug, Clone, PartialEq, Eq)]
397pub struct MissingPackDependency {
398 pub from: String,
399 pub requires: String,
400}
401
402impl fmt::Display for MissingPackDependency {
403 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
404 write!(
405 f,
406 "pack '{}' requires '{}', but '{}' is not in the loaded pack set",
407 self.from, self.requires, self.requires
408 )
409 }
410}
411
412impl std::error::Error for MissingPackDependency {}
413
414#[derive(Debug, Clone, PartialEq, Eq)]
416pub struct MissingPackDependencies {
417 pub missing: Vec<MissingPackDependency>,
418}
419
420impl fmt::Display for MissingPackDependencies {
421 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
422 let parts: Vec<String> = self.missing.iter().map(ToString::to_string).collect();
423 write!(f, "{}", parts.join("; "))
424 }
425}
426
427impl std::error::Error for MissingPackDependencies {}
428
429#[derive(Debug, Clone, PartialEq, Eq)]
431pub struct CircularPackDependency {
432 pub cycle: Vec<String>,
433}
434
435impl fmt::Display for CircularPackDependency {
436 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
437 write!(
438 f,
439 "circular dependency detected among packs: {}",
440 self.cycle.join(" -> ")
441 )
442 }
443}
444
445impl std::error::Error for CircularPackDependency {}
446
447#[derive(Debug, Clone, Copy, PartialEq, Eq)]
454pub enum DenialAuditOutcome {
455 Committed,
458 NotCommitted(&'static str),
461 NoStore,
463 NotAudited,
466}
467
468impl DenialAuditOutcome {
469 pub fn wire_code(&self) -> String {
472 match self {
473 Self::Committed => "committed".to_string(),
474 Self::NotCommitted(code) => format!("not_committed:{code}"),
475 Self::NoStore => "no_store".to_string(),
476 Self::NotAudited => "not_audited".to_string(),
477 }
478 }
479}
480
481#[derive(Debug, Clone, PartialEq, Eq)]
485pub struct DenialReceipt {
486 pub audit_event_id: Option<uuid::Uuid>,
488 pub audit_outcome: DenialAuditOutcome,
490}
491
492impl DenialReceipt {
493 pub fn not_audited() -> Self {
495 Self {
496 audit_event_id: None,
497 audit_outcome: DenialAuditOutcome::NotAudited,
498 }
499 }
500
501 pub fn no_store() -> Self {
503 Self {
504 audit_event_id: None,
505 audit_outcome: DenialAuditOutcome::NoStore,
506 }
507 }
508}
509
510#[derive(Debug, Clone)]
512pub struct ReceiptRefusal {
513 pub code: &'static str,
515 pub message: String,
517 pub receipt_id: String,
519 pub reason: String,
521 pub detail: serde_json::Value,
523}
524
525#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize)]
527#[serde(rename_all = "snake_case")]
528pub enum RefusalRecordingErrorClass {
529 EventStoreUnavailable,
530 EventAppendFailed,
531}
532
533impl RefusalRecordingErrorClass {
534 pub const fn as_str(self) -> &'static str {
535 match self {
536 Self::EventStoreUnavailable => "event_store_unavailable",
537 Self::EventAppendFailed => "event_append_failed",
538 }
539 }
540}
541
542#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize)]
545#[serde(untagged)]
546pub enum RefusalEventRecording {
547 Recorded {
548 item_index: usize,
549 subject: Uuid,
550 event_id: Uuid,
551 },
552 Failed {
553 item_index: usize,
554 subject: Uuid,
555 error_class: RefusalRecordingErrorClass,
556 },
557}
558
559#[derive(Debug)]
565pub struct RefusalEventContext {
566 pub source: Box<RuntimeError>,
567 pub recordings: Vec<RefusalEventRecording>,
568}
569
570impl std::fmt::Display for RefusalEventContext {
571 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
572 std::fmt::Display::fmt(self.source.as_ref(), formatter)
573 }
574}
575
576impl std::error::Error for RefusalEventContext {
577 fn source(&self) -> Option<&(dyn std::error::Error + 'static)> {
578 Some(self.source.as_ref())
579 }
580}
581
582#[derive(Debug, Error)]
586pub enum RuntimeError {
587 #[error("{failure}")]
589 AuditObligation {
590 #[source]
592 failure: Box<AuditObligationFailure>,
593 domain_result: serde_json::Value,
595 },
596
597 #[error("storage: {0}")]
598 Storage(#[from] khive_storage::StorageError),
599
600 #[error("sqlite: {0}")]
601 Sqlite(khive_db::SqliteError),
602
603 #[error("query: {0}")]
604 Query(#[from] khive_query::QueryError),
605
606 #[error("not found: {0}")]
607 NotFound(String),
608
609 #[error("invalid input: {0}")]
610 InvalidInput(String),
611
612 #[error("invalid input: {0}")]
613 UnknownVerb(String),
614
615 #[error("unconfigured: {0} is not set")]
616 Unconfigured(String),
617
618 #[error("unknown embedding model: {0}")]
619 UnknownModel(String),
620
621 #[error("embedding: {0}")]
622 Embedding(#[from] lattice_embed::EmbedError),
623
624 #[error("ambiguous: {0}")]
625 Ambiguous(String),
626
627 #[error("fusion: {0}")]
628 Fusion(#[from] khive_fusion::FuseError),
629
630 #[error("unknown fusion strategy: {0}")]
634 UnknownFusionStrategy(String),
635
636 #[error("internal: {0}")]
637 Internal(String),
638
639 #[error("audit batch incompatible event store: {0}")]
645 IncompatibleEventStore(String),
646
647 #[error("guarded edge write refused: {0}")]
648 GuardedWriteFailed(GuardedWriteFailure),
649
650 #[error("missing pack dependency: {0}")]
651 MissingPackDependency(MissingPackDependency),
652
653 #[error("missing pack dependencies: {0}")]
654 MissingPackDependencies(MissingPackDependencies),
655
656 #[error("{0}")]
657 CircularPackDependency(CircularPackDependency),
658
659 #[error("pack '{name}' registered twice (indices {first_idx} and {second_idx})")]
660 PackRedeclared {
661 name: String,
662 first_idx: usize,
663 second_idx: usize,
664 },
665
666 #[error(
670 "verb collision: verb {verb:?} declared by both pack {first_pack:?} and pack \
671 {second_pack:?}; rename one handler or use Visibility::Subhandler for internal verbs"
672 )]
673 VerbCollision {
674 verb: String,
675 first_pack: String,
676 second_pack: String,
677 },
678
679 #[error(
681 "pack {pack:?} handler {verb:?} declares request-envelope parameter {param:?}; rename the verb argument"
682 )]
683 ReservedEnvelopeParam {
684 pack: String,
685 verb: String,
686 param: String,
687 },
688
689 #[error("{}", .0.message)]
707 RefusedWithReceipt(Box<ReceiptRefusal>),
708
709 #[error(transparent)]
713 RefusedWithEvents { context: RefusalEventContext },
714
715 #[error("permission denied for verb {verb:?}: {reason}")]
725 PermissionDenied {
726 verb: String,
727 reason: String,
728 receipt: Box<DenialReceipt>,
729 },
730
731 #[error("gate unavailable for verb {verb:?}: {reason}")]
737 GateUnavailable { verb: String, reason: String },
738
739 #[error("{0}")]
743 Khive(khive_types::KhiveError),
744
745 #[error("not found in this namespace")]
750 NamespaceMismatch { id: uuid::Uuid },
751
752 #[error("ambiguous prefix {prefix:?}: matches {}", format_uuid_list(matches))]
758 AmbiguousPrefix {
759 prefix: String,
760 matches: Vec<uuid::Uuid>,
761 },
762
763 #[error(
768 "cross-backend merge is not supported: \
769 into_id {into_id} is on backend '{into_backend}', \
770 from_id {from_id} is on backend '{from_backend}'. \
771 Both entities must be on the same backend to merge."
772 )]
773 CrossBackendMergeUnsupported {
774 into_id: uuid::Uuid,
775 from_id: uuid::Uuid,
776 into_backend: String,
777 from_backend: String,
778 },
779
780 #[error("unknown remote: {name:?}")]
783 UnknownRemote { name: String },
784
785 #[error("remote cache missing for remote={remote:?} namespace={namespace:?}")]
787 RemoteCacheMissing { remote: String, namespace: String },
788
789 #[error("ambiguous id {id:?}: matched {count} records")]
791 AmbiguousId { id: String, count: usize },
792
793 #[error("cross-namespace write denied: cannot write to remote namespace {namespace:?}")]
795 CrossNamespaceWrite { namespace: String },
796
797 #[error("remote fetch error for remote={remote:?}: {message}")]
800 RemoteFetchError { remote: String, message: String },
801
802 #[error(
808 "write budget exceeded: max_new_entries={max_new_entries}, \
809 attempted_new_entries={attempted_new_entries}"
810 )]
811 WriteBudgetExceeded {
812 max_new_entries: u64,
813 attempted_new_entries: u64,
814 },
815
816 #[error("write blocked: {0}")]
822 SecretDetected(crate::secret_gate::SecretMatch),
823
824 #[error("{operation} exceeded its {budget_ms}ms deadline (elapsed {elapsed_ms}ms)")]
838 DeadlineExceeded {
839 operation: String,
840 budget_ms: u64,
841 elapsed_ms: u64,
842 },
843}
844
845impl From<khive_db::SqliteError> for RuntimeError {
846 fn from(error: khive_db::SqliteError) -> Self {
847 match error {
848 khive_db::SqliteError::RequestReadStopped(error) => Self::Storage(error),
849 error => Self::Sqlite(error),
850 }
851 }
852}
853
854impl RuntimeError {
855 pub fn with_refusal_events(self, mut recordings: Vec<RefusalEventRecording>) -> Self {
858 if recordings.is_empty() {
859 return self;
860 }
861 match self {
862 Self::RefusedWithEvents { mut context } => {
863 context.recordings.append(&mut recordings);
864 Self::RefusedWithEvents { context }
865 }
866 source => Self::RefusedWithEvents {
867 context: RefusalEventContext {
868 source: Box::new(source),
869 recordings,
870 },
871 },
872 }
873 }
874
875 pub fn refusal_source(&self) -> &Self {
878 let mut error = self;
879 while let Self::RefusedWithEvents { context } = error {
880 error = &context.source;
881 }
882 error
883 }
884
885 pub fn is_stream_policy_refusal(&self) -> bool {
892 matches!(self.refusal_source(), Self::Khive(error)
893 if error.kind() == khive_types::ErrorKind::Conflict
894 && error.details().and_then(|details| details.get("reason"))
895 == Some("stream_member"))
896 }
897
898 pub fn permission_denied(verb: impl Into<String>, reason: impl Into<String>) -> Self {
900 Self::PermissionDenied {
901 verb: verb.into(),
902 reason: reason.into(),
903 receipt: Box::new(DenialReceipt::not_audited()),
904 }
905 }
906
907 pub fn channel_ingest_failure_class(&self) -> ChannelIngestFailureClass {
915 let source = self.refusal_source();
916 let reason = source.variant_name();
917 if source.retryable_failure_context().is_some() {
918 ChannelIngestFailureClass::Retryable { reason }
919 } else if matches!(source, Self::SecretDetected(_)) {
920 ChannelIngestFailureClass::Permanent { reason }
921 } else {
922 ChannelIngestFailureClass::Unknown { reason }
923 }
924 }
925
926 fn variant_name(&self) -> &'static str {
928 match self {
929 Self::AuditObligation { .. } => "AuditObligation",
930 Self::Storage(_) => "Storage",
931 Self::Sqlite(_) => "Sqlite",
932 Self::Query(_) => "Query",
933 Self::NotFound(_) => "NotFound",
934 Self::InvalidInput(_) => "InvalidInput",
935 Self::UnknownVerb(_) => "UnknownVerb",
936 Self::Unconfigured(_) => "Unconfigured",
937 Self::UnknownModel(_) => "UnknownModel",
938 Self::Embedding(_) => "Embedding",
939 Self::Ambiguous(_) => "Ambiguous",
940 Self::Fusion(_) => "Fusion",
941 Self::UnknownFusionStrategy(_) => "UnknownFusionStrategy",
942 Self::Internal(_) => "Internal",
943 Self::GuardedWriteFailed(_) => "GuardedWriteFailed",
944 Self::MissingPackDependency(_) => "MissingPackDependency",
945 Self::MissingPackDependencies(_) => "MissingPackDependencies",
946 Self::CircularPackDependency(_) => "CircularPackDependency",
947 Self::PackRedeclared { .. } => "PackRedeclared",
948 Self::VerbCollision { .. } => "VerbCollision",
949 Self::ReservedEnvelopeParam { .. } => "ReservedEnvelopeParam",
950 Self::PermissionDenied { .. } => "PermissionDenied",
951 Self::GateUnavailable { .. } => "GateUnavailable",
952 Self::Khive(_) => "Khive",
953 Self::NamespaceMismatch { .. } => "NamespaceMismatch",
954 Self::AmbiguousPrefix { .. } => "AmbiguousPrefix",
955 Self::CrossBackendMergeUnsupported { .. } => "CrossBackendMergeUnsupported",
956 Self::UnknownRemote { .. } => "UnknownRemote",
957 Self::RemoteCacheMissing { .. } => "RemoteCacheMissing",
958 Self::AmbiguousId { .. } => "AmbiguousId",
959 Self::CrossNamespaceWrite { .. } => "CrossNamespaceWrite",
960 Self::RemoteFetchError { .. } => "RemoteFetchError",
961 Self::WriteBudgetExceeded { .. } => "WriteBudgetExceeded",
962 Self::SecretDetected(_) => "SecretDetected",
963 Self::DeadlineExceeded { .. } => "DeadlineExceeded",
964 Self::IncompatibleEventStore(_) => "IncompatibleEventStore",
965 Self::RefusedWithReceipt { .. } => "RefusedWithReceipt",
966 Self::RefusedWithEvents { context } => context.source.variant_name(),
967 }
968 }
969
970 pub fn writer_pool_checkout_timeout_context(&self) -> Option<WriterPoolCheckoutTimeoutContext> {
978 let (sqlite_error, capability, operation) = match self.refusal_source() {
979 Self::Sqlite(error) => (error, None, None),
980 Self::Storage(khive_storage::StorageError::Driver {
981 capability,
982 operation,
983 source,
984 }) => (
985 source.downcast_ref::<khive_db::SqliteError>()?,
986 Some(*capability),
987 Some(operation.to_string()),
988 ),
989 _ => return None,
990 };
991
992 let khive_db::SqliteError::WriterPoolCheckoutTimeout { timeout } = sqlite_error else {
993 return None;
994 };
995 Some(WriterPoolCheckoutTimeoutContext {
996 timeout: *timeout,
997 capability,
998 operation,
999 })
1000 }
1001
1002 pub fn admission_failure_context(&self) -> Option<AdmissionFailureContext> {
1007 let source = self.refusal_source();
1008 if let Some(context) = self.writer_pool_checkout_timeout_context() {
1009 return Some(AdmissionFailureContext {
1010 stage: WRITER_POOL_CHECKOUT_TIMEOUT_STAGE,
1011 timeout: context.timeout,
1012 capability: context.capability,
1013 operation: context.operation,
1014 pool_identity: None,
1015 scope: None,
1016 retry_after_ms: None,
1017 });
1018 }
1019 if let Self::Storage(khive_storage::StorageError::WriteQueueFull { timeout_ms }) = source {
1020 return Some(AdmissionFailureContext {
1021 stage: WRITER_QUEUE_SATURATED_STAGE,
1022 timeout: Duration::from_millis(*timeout_ms),
1023 capability: None,
1024 operation: None,
1025 pool_identity: None,
1026 scope: Some(WRITER_ADMISSION_SCOPE),
1027 retry_after_ms: Some(*timeout_ms),
1028 });
1029 }
1030 if let Self::Storage(khive_storage::StorageError::AdmissionTimeout {
1031 operation,
1032 timeout_ms,
1033 pool_identity,
1034 }) = source
1035 {
1036 return Some(AdmissionFailureContext {
1037 stage: STORAGE_ADMISSION_TIMEOUT_STAGE,
1038 timeout: Duration::from_millis(*timeout_ms),
1039 capability: None,
1040 operation: Some(operation.to_string()),
1041 pool_identity: pool_identity.clone(),
1042 scope: None,
1043 retry_after_ms: None,
1044 });
1045 }
1046 None
1047 }
1048
1049 pub fn writer_task_failure_context(&self) -> Option<WriterTaskFailureContext> {
1053 let Self::Storage(error) = self.refusal_source() else {
1054 return None;
1055 };
1056 match error {
1057 khive_storage::StorageError::WriterTaskRequestFailed {
1058 request_state,
1059 source,
1060 } => Some(WriterTaskFailureContext {
1061 stage: WRITER_TASK_REQUEST_FAILED_STAGE,
1062 request_state: *request_state,
1063 task_terminated: false,
1064 retryable: source.is_retryable(),
1065 }),
1066 khive_storage::StorageError::WriterTaskTerminated { request_state } => {
1067 Some(WriterTaskFailureContext {
1068 stage: WRITER_TASK_TERMINATED_STAGE,
1069 request_state: *request_state,
1070 task_terminated: true,
1071 retryable: false,
1072 })
1073 }
1074 _ => None,
1075 }
1076 }
1077
1078 pub fn retryable_failure_context(&self) -> Option<RetryableFailureContext> {
1081 let source = self.refusal_source();
1082 if let Some(context) = self.admission_failure_context() {
1083 return Some(context.into());
1084 }
1085 if let Self::Storage(khive_storage::StorageError::WriterTaskBusy { timeout_ms }) = source {
1086 return Some(RetryableFailureContext {
1087 stage: WRITER_TASK_BEGIN_BUSY_STAGE,
1088 timeout: Duration::from_millis(*timeout_ms),
1089 capability: None,
1090 operation: Some("writer_task_begin".to_string()),
1091 pool_identity: None,
1092 scope: None,
1093 retry_after_ms: None,
1094 });
1095 }
1096 if let Self::Storage(khive_storage::StorageError::ReadTransactionAgeEvicted {
1097 operation,
1098 max_age_secs,
1099 }) = source
1100 {
1101 return Some(RetryableFailureContext {
1102 stage: READ_TX_AGE_EVICTED_STAGE,
1103 timeout: Duration::from_secs(*max_age_secs),
1104 capability: Some(khive_storage::StorageCapability::Sql),
1105 operation: Some(operation.to_string()),
1106 pool_identity: None,
1107 scope: None,
1108 retry_after_ms: None,
1109 });
1110 }
1111 if let Self::Storage(
1112 khive_storage::StorageError::ReadTransactionAgeEvictionCleanupFailed {
1113 operation,
1114 max_age_secs,
1115 ..
1116 },
1117 ) = source
1118 {
1119 return Some(RetryableFailureContext {
1120 stage: READ_TX_AGE_EVICTED_STAGE,
1121 timeout: Duration::from_secs(*max_age_secs),
1122 capability: Some(khive_storage::StorageCapability::Sql),
1123 operation: Some(operation.to_string()),
1124 pool_identity: None,
1125 scope: None,
1126 retry_after_ms: None,
1127 });
1128 }
1129 None
1130 }
1131}
1132
1133pub fn fts_text_leg_or_err<T>(
1140 result: Result<Vec<T>, RuntimeError>,
1141 context: &'static str,
1142 query: &str,
1143) -> RuntimeResult<Vec<T>> {
1144 match result {
1145 Ok(hits) => Ok(hits),
1146 Err(RuntimeError::Storage(se)) if se.is_fts5_syntax_error() => {
1147 tracing::warn!(
1148 error = %se,
1149 query = %query,
1150 context,
1151 "FTS text leg failed on a parser syntax error; failing loud (#569)"
1152 );
1153 Err(RuntimeError::InvalidInput(format!(
1154 "{context}: FTS query could not be parsed: {se}"
1155 )))
1156 }
1157 Err(e) => Err(e),
1158 }
1159}
1160
1161fn format_uuid_list(uuids: &[uuid::Uuid]) -> String {
1162 uuids
1163 .iter()
1164 .map(uuid::Uuid::to_string)
1165 .collect::<Vec<_>>()
1166 .join(", ")
1167}
1168
1169impl From<khive_types::EntityTypeError> for RuntimeError {
1173 fn from(e: khive_types::EntityTypeError) -> Self {
1174 Self::InvalidInput(e.to_string())
1175 }
1176}
1177
1178impl From<khive_types::KhiveError> for RuntimeError {
1179 fn from(e: khive_types::KhiveError) -> Self {
1180 Self::Khive(e)
1181 }
1182}
1183
1184#[cfg(test)]
1185mod refusal_event_context_tests {
1186 use super::*;
1187
1188 fn failed() -> Vec<RefusalEventRecording> {
1189 vec![RefusalEventRecording::Failed {
1190 item_index: 2,
1191 subject: Uuid::from_u128(11),
1192 error_class: RefusalRecordingErrorClass::EventAppendFailed,
1193 }]
1194 }
1195
1196 #[test]
1197 fn refusal_events_preserve_the_source_type_display_and_policy_class() {
1198 let original = RuntimeError::SecretDetected(crate::secret_gate::SecretMatch {
1199 detector: "fixture",
1200 trigger: None,
1201 masked: "never-in-display".into(),
1202 location: Some("atoms[2].content".into()),
1203 });
1204 let message = original.to_string();
1205 let wrapped = original.with_refusal_events(failed());
1206 assert_eq!(wrapped.to_string(), message);
1207 assert!(matches!(
1208 wrapped.refusal_source(),
1209 RuntimeError::SecretDetected(_)
1210 ));
1211 assert!(std::error::Error::source(&wrapped)
1212 .unwrap()
1213 .downcast_ref::<RuntimeError>()
1214 .is_some_and(|error| matches!(error, RuntimeError::SecretDetected(_))));
1215 assert_eq!(
1216 wrapped.channel_ingest_failure_class(),
1217 ChannelIngestFailureClass::Permanent {
1218 reason: "SecretDetected"
1219 }
1220 );
1221 assert!(wrapped.retryable_failure_context().is_none());
1222
1223 let lookalike = RuntimeError::InvalidInput(message).with_refusal_events(failed());
1224 assert_eq!(
1225 lookalike.channel_ingest_failure_class(),
1226 ChannelIngestFailureClass::Unknown {
1227 reason: "InvalidInput"
1228 }
1229 );
1230 let policy = RuntimeError::Khive(
1231 khive_types::KhiveError::conflict("immutable")
1232 .with_details(khive_types::Details::new([("reason", "stream_member")])),
1233 )
1234 .with_refusal_events(failed());
1235 assert!(policy.is_stream_policy_refusal());
1236 }
1237
1238 #[test]
1239 fn refusal_events_preserve_retry_and_writer_finality_contexts() {
1240 let checkout = RuntimeError::Sqlite(khive_db::SqliteError::WriterPoolCheckoutTimeout {
1241 timeout: Duration::from_millis(17),
1242 })
1243 .with_refusal_events(failed());
1244 assert_eq!(
1245 checkout
1246 .writer_pool_checkout_timeout_context()
1247 .unwrap()
1248 .timeout,
1249 Duration::from_millis(17)
1250 );
1251 assert_eq!(
1252 checkout.admission_failure_context().unwrap().stage,
1253 WRITER_POOL_CHECKOUT_TIMEOUT_STAGE
1254 );
1255 assert_eq!(
1256 checkout.retryable_failure_context().unwrap().stage,
1257 WRITER_POOL_CHECKOUT_TIMEOUT_STAGE
1258 );
1259
1260 let queued =
1261 RuntimeError::Storage(khive_storage::StorageError::WriteQueueFull { timeout_ms: 23 })
1262 .with_refusal_events(failed());
1263 assert_eq!(
1264 queued.retryable_failure_context().unwrap().stage,
1265 WRITER_QUEUE_SATURATED_STAGE
1266 );
1267 assert_eq!(
1268 queued.channel_ingest_failure_class(),
1269 ChannelIngestFailureClass::Retryable { reason: "Storage" }
1270 );
1271
1272 let busy =
1273 RuntimeError::Storage(khive_storage::StorageError::WriterTaskBusy { timeout_ms: 31 })
1274 .with_refusal_events(failed());
1275 assert_eq!(
1276 busy.retryable_failure_context().unwrap().stage,
1277 WRITER_TASK_BEGIN_BUSY_STAGE
1278 );
1279
1280 let stopped = RuntimeError::Storage(khive_storage::StorageError::WriterTaskTerminated {
1281 request_state: khive_storage::WriterTaskRequestState::SideEffectsUnknown,
1282 })
1283 .with_refusal_events(failed());
1284 let context = stopped.writer_task_failure_context().unwrap();
1285 assert_eq!(
1286 context.request_state,
1287 khive_storage::WriterTaskRequestState::SideEffectsUnknown
1288 );
1289 assert!(context.task_terminated);
1290 assert!(!context.retryable);
1291 assert!(stopped.retryable_failure_context().is_none());
1292 }
1293
1294 #[test]
1295 fn refusal_events_preserve_the_original_error_and_its_underlying_source_chain() {
1296 let original = RuntimeError::Storage(khive_storage::StorageError::driver(
1297 khive_storage::StorageCapability::Notes,
1298 "refusal-context-fixture",
1299 std::io::Error::new(std::io::ErrorKind::PermissionDenied, "fixture refusal"),
1300 ));
1301 let message = original.to_string();
1302 let wrapped = original.with_refusal_events(failed());
1303 assert_eq!(wrapped.to_string(), message);
1304
1305 let runtime = std::error::Error::source(&wrapped)
1306 .unwrap()
1307 .downcast_ref::<RuntimeError>()
1308 .expect("the first source must retain the runtime classification");
1309 assert!(matches!(runtime, RuntimeError::Storage(_)));
1310 let storage = std::error::Error::source(runtime)
1311 .unwrap()
1312 .downcast_ref::<khive_storage::StorageError>()
1313 .expect("the original runtime error must retain its storage source");
1314 let driver = std::error::Error::source(storage)
1315 .unwrap()
1316 .downcast_ref::<std::io::Error>()
1317 .expect("the storage source chain must remain traversable");
1318 assert_eq!(driver.kind(), std::io::ErrorKind::PermissionDenied);
1319 }
1320
1321 #[test]
1322 fn refusal_events_with_no_eligible_targets_do_not_wrap_or_grow_metadata() {
1323 let error = RuntimeError::InvalidInput("unchanged".into()).with_refusal_events(vec![]);
1324 assert!(matches!(error, RuntimeError::InvalidInput(ref message) if message == "unchanged"));
1325 let error = error
1326 .with_refusal_events(failed())
1327 .with_refusal_events(vec![RefusalEventRecording::Recorded {
1328 item_index: 3,
1329 subject: Uuid::from_u128(12),
1330 event_id: Uuid::from_u128(13),
1331 }]);
1332 let source = std::error::Error::source(&error)
1333 .unwrap()
1334 .downcast_ref::<RuntimeError>()
1335 .expect("appending recordings must not add another error wrapper");
1336 assert!(matches!(source, RuntimeError::InvalidInput(_)));
1337 assert!(std::error::Error::source(source).is_none());
1338 let RuntimeError::RefusedWithEvents { context } = error else {
1339 panic!("recordings must accompany the original typed source");
1340 };
1341 let RefusalEventContext { source, recordings } = context;
1342 assert!(matches!(*source, RuntimeError::InvalidInput(_)));
1343 assert_eq!(recordings.len(), 2);
1344 assert_eq!(recordings[0], failed()[0]);
1345 }
1346
1347 #[test]
1348 fn refusal_recording_classes_serialize_only_the_closed_safe_spellings() {
1349 for (class, name) in [
1350 (
1351 RefusalRecordingErrorClass::EventStoreUnavailable,
1352 "event_store_unavailable",
1353 ),
1354 (
1355 RefusalRecordingErrorClass::EventAppendFailed,
1356 "event_append_failed",
1357 ),
1358 ] {
1359 assert_eq!(class.as_str(), name);
1360 assert_eq!(
1361 serde_json::to_value(class).unwrap(),
1362 serde_json::json!(name)
1363 );
1364 }
1365 }
1366}
1367
1368#[cfg(test)]
1369mod stream_policy_refusal_tests {
1370 use super::RuntimeError;
1371 use khive_types::{Details, KhiveError};
1372
1373 #[test]
1374 fn stream_policy_refusal_uses_structured_kind_and_marker() {
1375 let policy = RuntimeError::Khive(
1376 KhiveError::conflict("changed diagnostic wording")
1377 .with_details(Details::new([("reason", "stream_member")])),
1378 );
1379 assert!(policy.is_stream_policy_refusal());
1380
1381 let wrong_kind = RuntimeError::Khive(
1382 KhiveError::invalid_input("stream entries are immutable")
1383 .with_details(Details::new([("reason", "stream_member")])),
1384 );
1385 assert!(!wrong_kind.is_stream_policy_refusal());
1386 }
1387
1388 #[test]
1389 fn stream_policy_refusal_excludes_other_conflicts_and_message_lookalikes() {
1390 for marker in [
1391 None,
1392 Some("seq_conflict"),
1393 Some("key_conflict"),
1394 Some("version_conflict"),
1395 Some("fence_conflict"),
1396 Some("identity_conflict"),
1397 Some("expired"),
1398 Some("live_until_unreadable"),
1399 Some("key_ambiguous"),
1400 Some("unknown_op"),
1401 Some("stream_member_extra"),
1402 ] {
1403 let mut error = KhiveError::conflict("stream entries are immutable");
1404 if let Some(marker) = marker {
1405 error = error.with_details(Details::new([("reason", marker)]));
1406 }
1407 assert!(
1408 !RuntimeError::Khive(error).is_stream_policy_refusal(),
1409 "unrelated conflict marker {marker:?}"
1410 );
1411 }
1412
1413 let internal = RuntimeError::Internal("conflict: stream entries are immutable".into());
1414 assert!(!internal.is_stream_policy_refusal());
1415 let driver = RuntimeError::Storage(khive_storage::StorageError::driver(
1416 khive_storage::StorageCapability::Sql,
1417 "inline.execute",
1418 std::io::Error::other("stream_member"),
1419 ));
1420 assert!(!driver.is_stream_policy_refusal());
1421 }
1422}
1423
1424#[cfg(test)]
1425mod channel_ingest_failure_class_tests {
1426 use super::{ChannelIngestFailureClass, RuntimeError};
1427 use crate::secret_gate::SecretMatch;
1428 use std::time::Duration;
1429
1430 #[test]
1431 fn secret_detected_is_permanent_by_typed_variant_not_display_text() {
1432 let first = RuntimeError::SecretDetected(SecretMatch {
1433 detector: "fixture",
1434 trigger: None,
1435 masked: "first-rendering".to_string(),
1436 location: None,
1437 });
1438 let second = RuntimeError::SecretDetected(SecretMatch {
1439 detector: "fixture",
1440 trigger: Some("token"),
1441 masked: "completely-different-rendering".to_string(),
1442 location: None,
1443 });
1444
1445 assert_ne!(first.to_string(), second.to_string());
1446 assert_eq!(
1447 first.channel_ingest_failure_class(),
1448 ChannelIngestFailureClass::Permanent {
1449 reason: "SecretDetected"
1450 }
1451 );
1452 assert_eq!(
1453 second.channel_ingest_failure_class(),
1454 ChannelIngestFailureClass::Permanent {
1455 reason: "SecretDetected"
1456 },
1457 "rendered error details must not participate in ingest classification"
1458 );
1459 }
1460
1461 #[test]
1462 fn admission_failures_are_retryable_and_unclassified_errors_are_unknown() {
1463 let retryable =
1464 RuntimeError::Storage(khive_storage::StorageError::WriteQueueFull { timeout_ms: 25 });
1465 assert_eq!(
1466 retryable.channel_ingest_failure_class(),
1467 ChannelIngestFailureClass::Retryable { reason: "Storage" }
1468 );
1469
1470 let begin_busy =
1471 RuntimeError::Storage(khive_storage::StorageError::WriterTaskBusy { timeout_ms: 175 });
1472 let context = begin_busy
1473 .retryable_failure_context()
1474 .expect("contended BEGIN must remain typed and retryable");
1475 assert_eq!(context.stage, "writer_task_begin_busy");
1476 assert_eq!(context.timeout, Duration::from_millis(175));
1477 assert_eq!(context.operation.as_deref(), Some("writer_task_begin"));
1478 assert_eq!(context.scope, None);
1479 assert_eq!(context.retry_after_ms, None);
1480 assert_eq!(
1481 begin_busy.channel_ingest_failure_class(),
1482 ChannelIngestFailureClass::Retryable { reason: "Storage" }
1483 );
1484
1485 let unknown = RuntimeError::InvalidInput(
1486 "write blocked: SecretDetected text must not affect classification".to_string(),
1487 );
1488 assert_eq!(
1489 unknown.channel_ingest_failure_class(),
1490 ChannelIngestFailureClass::Unknown {
1491 reason: "InvalidInput"
1492 },
1493 "a rendered message resembling SecretDetected must remain Unknown unless its typed variant is SecretDetected"
1494 );
1495 }
1496
1497 #[test]
1498 fn writer_task_failure_context_separates_request_finality_from_task_liveness() {
1499 use khive_storage::{StorageError, WriterTaskRequestState};
1500
1501 let rolled_back = RuntimeError::Storage(StorageError::WriterTaskRequestFailed {
1502 request_state: WriterTaskRequestState::TransactionRolledBack,
1503 source: Box::new(StorageError::Pool {
1504 operation: "writer_task_commit".into(),
1505 message: "commit refused".into(),
1506 }),
1507 });
1508 let rolled_back_context = rolled_back
1509 .writer_task_failure_context()
1510 .expect("proven rollback must remain typed through RuntimeError");
1511 assert_eq!(
1512 rolled_back_context.request_state,
1513 WriterTaskRequestState::TransactionRolledBack
1514 );
1515 assert_eq!(
1516 rolled_back_context.stage,
1517 super::WRITER_TASK_REQUEST_FAILED_STAGE
1518 );
1519 assert!(!rolled_back_context.task_terminated);
1520 assert!(rolled_back_context.retryable);
1521
1522 let unknown = RuntimeError::Storage(StorageError::WriterTaskTerminated {
1523 request_state: WriterTaskRequestState::SideEffectsUnknown,
1524 });
1525 let unknown_context = unknown
1526 .writer_task_failure_context()
1527 .expect("ambiguous finality must remain typed through RuntimeError");
1528 assert_eq!(
1529 unknown_context.request_state,
1530 WriterTaskRequestState::SideEffectsUnknown
1531 );
1532 assert_eq!(unknown_context.stage, super::WRITER_TASK_TERMINATED_STAGE);
1533 assert!(unknown_context.task_terminated);
1534 assert!(!unknown_context.retryable);
1535 }
1536
1537 #[test]
1538 fn storage_admission_timeout_is_a_retryable_admission_failure() {
1539 let admission = RuntimeError::Storage(khive_storage::StorageError::AdmissionTimeout {
1540 operation: "sql_bridge.writer_handle".into(),
1541 timeout_ms: 30_000,
1542 pool_identity: None,
1543 });
1544 let context = admission
1545 .admission_failure_context()
1546 .expect("a storage admission timeout happens before the operation starts");
1547 assert_eq!(context.stage, super::STORAGE_ADMISSION_TIMEOUT_STAGE);
1548 assert_eq!(context.timeout, Duration::from_millis(30_000));
1549 assert_eq!(
1550 context.operation.as_deref(),
1551 Some("sql_bridge.writer_handle")
1552 );
1553 assert_eq!(context.scope, None);
1554 assert_eq!(context.retry_after_ms, None);
1555 assert_eq!(
1556 admission.channel_ingest_failure_class(),
1557 ChannelIngestFailureClass::Retryable { reason: "Storage" }
1558 );
1559
1560 let mid_flight = RuntimeError::Storage(khive_storage::StorageError::Timeout {
1563 operation: "sql_bridge.reader_open".into(),
1564 });
1565 assert!(mid_flight.admission_failure_context().is_none());
1566 assert!(mid_flight.retryable_failure_context().is_none());
1567 }
1568
1569 #[test]
1576 fn both_read_tx_age_eviction_outcomes_map_to_the_same_retryable_stage() {
1577 let clean = RuntimeError::Storage(khive_storage::StorageError::ReadTransactionAgeEvicted {
1578 operation: "query_all".into(),
1579 max_age_secs: 120,
1580 });
1581 let clean_context = clean
1582 .retryable_failure_context()
1583 .expect("a clean age eviction must be typed-retryable");
1584 assert_eq!(clean_context.stage, super::READ_TX_AGE_EVICTED_STAGE);
1585 assert_eq!(clean_context.timeout, Duration::from_secs(120));
1586 assert_eq!(clean_context.operation.as_deref(), Some("query_all"));
1587 assert_eq!(
1588 clean_context.capability,
1589 Some(khive_storage::StorageCapability::Sql)
1590 );
1591
1592 let cleanup_failed = RuntimeError::Storage(
1593 khive_storage::StorageError::ReadTransactionAgeEvictionCleanupFailed {
1594 operation: "query_all".into(),
1595 max_age_secs: 120,
1596 message: "rollback failed: disk I/O error".into(),
1597 },
1598 );
1599 let cleanup_failed_context = cleanup_failed
1600 .retryable_failure_context()
1601 .expect("a failed cleanup rollback must remain typed-retryable");
1602 assert_eq!(
1603 cleanup_failed_context.stage,
1604 super::READ_TX_AGE_EVICTED_STAGE
1605 );
1606 assert_eq!(cleanup_failed_context.timeout, Duration::from_secs(120));
1607 assert_eq!(
1608 cleanup_failed_context.operation.as_deref(),
1609 Some("query_all")
1610 );
1611 assert_eq!(
1612 cleanup_failed_context.capability,
1613 Some(khive_storage::StorageCapability::Sql)
1614 );
1615
1616 assert_ne!(
1617 clean.to_string(),
1618 cleanup_failed.to_string(),
1619 "the rendered message must distinguish a clean eviction from a failed cleanup even \
1620 though both map to the same wire stage"
1621 );
1622 assert!(cleanup_failed.to_string().contains("rollback failed"));
1623 }
1624}