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 use khive_db::{SQLITE_WAL_CAPACITY_REFUSED_STAGE, SQLITE_WAL_CAPACITY_UNAVAILABLE_STAGE};
173
174pub const WRITER_QUEUE_SATURATED_STAGE: &str = "writer_queue_saturated";
177
178pub const WRITER_TASK_BEGIN_BUSY_STAGE: &str = "writer_task_begin_busy";
181
182pub const WRITER_TASK_REQUEST_FAILED_STAGE: &str = "writer_task_request_failed";
185
186pub const WRITER_TASK_TERMINATED_STAGE: &str = "writer_task_terminated";
189
190pub const STORAGE_ADMISSION_TIMEOUT_STAGE: &str = "storage_admission_timeout";
195
196pub const READ_TX_AGE_EVICTED_STAGE: &str = "read_tx_age_evicted";
207
208pub const WRITER_ADMISSION_SCOPE: &str = "writer_admission";
215
216#[derive(Debug, Clone, PartialEq, Eq)]
223pub struct AdmissionFailureContext {
224 pub stage: &'static str,
228 pub timeout: Duration,
230 pub capability: Option<khive_storage::StorageCapability>,
234 pub operation: Option<String>,
236 pub pool_identity: Option<String>,
238 pub scope: Option<&'static str>,
243 pub retry_after_ms: Option<u64>,
249}
250
251#[derive(Debug, Clone, PartialEq, Eq)]
258pub struct RetryableFailureContext {
259 pub stage: &'static str,
261 pub timeout: Duration,
263 pub capability: Option<khive_storage::StorageCapability>,
265 pub operation: Option<String>,
267 pub pool_identity: Option<String>,
269 pub scope: Option<&'static str>,
271 pub retry_after_ms: Option<u64>,
273}
274
275#[derive(Debug, Clone, Copy, PartialEq, Eq)]
282pub struct WriterTaskFailureContext {
283 pub stage: &'static str,
285 pub request_state: khive_storage::WriterTaskRequestState,
287 pub task_terminated: bool,
289 pub retryable: bool,
291}
292
293impl From<AdmissionFailureContext> for RetryableFailureContext {
294 fn from(context: AdmissionFailureContext) -> Self {
295 Self {
296 stage: context.stage,
297 timeout: context.timeout,
298 capability: context.capability,
299 operation: context.operation,
300 pool_identity: context.pool_identity,
301 scope: context.scope,
302 retry_after_ms: context.retry_after_ms,
303 }
304 }
305}
306
307#[derive(Debug, Clone, Copy, PartialEq, Eq)]
313pub enum ChannelIngestFailureClass {
314 Retryable { reason: &'static str },
316 Permanent { reason: &'static str },
318 Unknown { reason: &'static str },
320}
321
322impl ChannelIngestFailureClass {
323 pub const fn name(self) -> &'static str {
325 match self {
326 Self::Retryable { .. } => "retryable",
327 Self::Permanent { .. } => "permanent",
328 Self::Unknown { .. } => "unknown",
329 }
330 }
331
332 pub const fn reason(self) -> &'static str {
334 match self {
335 Self::Retryable { reason } | Self::Permanent { reason } | Self::Unknown { reason } => {
336 reason
337 }
338 }
339 }
340}
341
342#[derive(Debug, Clone, PartialEq, Eq)]
345pub struct WriterPoolCheckoutTimeoutContext {
346 pub timeout: Duration,
349 pub capability: Option<khive_storage::StorageCapability>,
351 pub operation: Option<String>,
353}
354
355#[derive(Debug, Clone, PartialEq, Eq)]
360pub struct GuardedWriteFailure {
361 pub entry_index: Option<usize>,
364 pub missing_source: Option<Uuid>,
367 pub missing_target: Option<Uuid>,
370}
371
372impl fmt::Display for GuardedWriteFailure {
373 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
374 let mut missing = Vec::new();
375 if let Some(source) = self.missing_source {
376 missing.push(format!("source {source}"));
377 }
378 if let Some(target) = self.missing_target {
379 missing.push(format!("target {target}"));
380 }
381 let missing = if missing.is_empty() {
382 "endpoint(s)".to_string()
383 } else {
384 missing.join(" and ")
385 };
386 match self.entry_index {
387 Some(index) => write!(
388 f,
389 "batch entry {index}: {missing} no longer exist at write time"
390 ),
391 None => write!(f, "{missing} no longer exist at write time"),
392 }
393 }
394}
395
396impl std::error::Error for GuardedWriteFailure {}
397
398#[derive(Debug, Clone, PartialEq, Eq)]
400pub struct MissingPackDependency {
401 pub from: String,
402 pub requires: String,
403}
404
405impl fmt::Display for MissingPackDependency {
406 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
407 write!(
408 f,
409 "pack '{}' requires '{}', but '{}' is not in the loaded pack set",
410 self.from, self.requires, self.requires
411 )
412 }
413}
414
415impl std::error::Error for MissingPackDependency {}
416
417#[derive(Debug, Clone, PartialEq, Eq)]
419pub struct MissingPackDependencies {
420 pub missing: Vec<MissingPackDependency>,
421}
422
423impl fmt::Display for MissingPackDependencies {
424 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
425 let parts: Vec<String> = self.missing.iter().map(ToString::to_string).collect();
426 write!(f, "{}", parts.join("; "))
427 }
428}
429
430impl std::error::Error for MissingPackDependencies {}
431
432#[derive(Debug, Clone, PartialEq, Eq)]
434pub struct CircularPackDependency {
435 pub cycle: Vec<String>,
436}
437
438impl fmt::Display for CircularPackDependency {
439 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
440 write!(
441 f,
442 "circular dependency detected among packs: {}",
443 self.cycle.join(" -> ")
444 )
445 }
446}
447
448impl std::error::Error for CircularPackDependency {}
449
450#[derive(Debug, Clone, Copy, PartialEq, Eq)]
457pub enum DenialAuditOutcome {
458 Committed,
461 NotCommitted(&'static str),
464 NoStore,
466 NotAudited,
469}
470
471impl DenialAuditOutcome {
472 pub fn wire_code(&self) -> String {
475 match self {
476 Self::Committed => "committed".to_string(),
477 Self::NotCommitted(code) => format!("not_committed:{code}"),
478 Self::NoStore => "no_store".to_string(),
479 Self::NotAudited => "not_audited".to_string(),
480 }
481 }
482}
483
484#[derive(Debug, Clone, PartialEq, Eq)]
488pub struct DenialReceipt {
489 pub audit_event_id: Option<uuid::Uuid>,
491 pub audit_outcome: DenialAuditOutcome,
493}
494
495impl DenialReceipt {
496 pub fn not_audited() -> Self {
498 Self {
499 audit_event_id: None,
500 audit_outcome: DenialAuditOutcome::NotAudited,
501 }
502 }
503
504 pub fn no_store() -> Self {
506 Self {
507 audit_event_id: None,
508 audit_outcome: DenialAuditOutcome::NoStore,
509 }
510 }
511}
512
513#[derive(Debug, Clone)]
515pub struct ReceiptRefusal {
516 pub code: &'static str,
518 pub message: String,
520 pub receipt_id: String,
522 pub reason: String,
524 pub detail: serde_json::Value,
526}
527
528#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize)]
530#[serde(rename_all = "snake_case")]
531pub enum RefusalRecordingErrorClass {
532 EventStoreUnavailable,
533 EventAppendFailed,
534}
535
536impl RefusalRecordingErrorClass {
537 pub const fn as_str(self) -> &'static str {
538 match self {
539 Self::EventStoreUnavailable => "event_store_unavailable",
540 Self::EventAppendFailed => "event_append_failed",
541 }
542 }
543}
544
545#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize)]
548#[serde(untagged)]
549pub enum RefusalEventRecording {
550 Recorded {
551 item_index: usize,
552 subject: Uuid,
553 event_id: Uuid,
554 },
555 Failed {
556 item_index: usize,
557 subject: Uuid,
558 error_class: RefusalRecordingErrorClass,
559 },
560}
561
562#[derive(Debug)]
568pub struct RefusalEventContext {
569 pub source: Box<RuntimeError>,
570 pub recordings: Vec<RefusalEventRecording>,
571}
572
573impl std::fmt::Display for RefusalEventContext {
574 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
575 std::fmt::Display::fmt(self.source.as_ref(), formatter)
576 }
577}
578
579impl std::error::Error for RefusalEventContext {
580 fn source(&self) -> Option<&(dyn std::error::Error + 'static)> {
581 Some(self.source.as_ref())
582 }
583}
584
585#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize)]
587pub struct ResolutionFacts {
588 pub project_id: Uuid,
589 pub duplicate_anchor_ids: Vec<Uuid>,
590 pub slug_backfilled: bool,
591 pub project_created: bool,
592 pub orphaned_project_id: Option<Uuid>,
593 pub orphaned_note_count: u64,
594}
595
596#[derive(Debug)]
598pub struct ResolutionFailureContext {
599 pub source: Box<RuntimeError>,
600 pub resolution: ResolutionFacts,
601}
602
603impl std::fmt::Display for ResolutionFailureContext {
604 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
605 std::fmt::Display::fmt(self.source.as_ref(), formatter)
606 }
607}
608
609impl std::error::Error for ResolutionFailureContext {
610 fn source(&self) -> Option<&(dyn std::error::Error + 'static)> {
611 Some(self.source.as_ref())
612 }
613}
614
615#[derive(Debug, Error)]
619pub enum RuntimeError {
620 #[error("{failure}")]
622 AuditObligation {
623 #[source]
625 failure: Box<AuditObligationFailure>,
626 domain_result: serde_json::Value,
628 },
629
630 #[error("storage: {0}")]
631 Storage(#[from] khive_storage::StorageError),
632
633 #[error("sqlite: {0}")]
634 Sqlite(khive_db::SqliteError),
635
636 #[error("query: {0}")]
637 Query(#[from] khive_query::QueryError),
638
639 #[error("not found: {0}")]
640 NotFound(String),
641
642 #[error("invalid input: {0}")]
643 InvalidInput(String),
644
645 #[error("invalid input: {0}")]
646 UnknownVerb(String),
647
648 #[error("unconfigured: {0} is not set")]
649 Unconfigured(String),
650
651 #[error("unknown embedding model: {0}")]
652 UnknownModel(String),
653
654 #[error("embedding: {0}")]
655 Embedding(#[from] lattice_embed::EmbedError),
656
657 #[error("ambiguous: {0}")]
658 Ambiguous(String),
659
660 #[error("fusion: {0}")]
661 Fusion(#[from] khive_fusion::FuseError),
662
663 #[error("unknown fusion strategy: {0}")]
667 UnknownFusionStrategy(String),
668
669 #[error("internal: {0}")]
670 Internal(String),
671
672 #[error("audit batch incompatible event store: {0}")]
678 IncompatibleEventStore(String),
679
680 #[error("guarded edge write refused: {0}")]
681 GuardedWriteFailed(GuardedWriteFailure),
682
683 #[error("missing pack dependency: {0}")]
684 MissingPackDependency(MissingPackDependency),
685
686 #[error("missing pack dependencies: {0}")]
687 MissingPackDependencies(MissingPackDependencies),
688
689 #[error("{0}")]
690 CircularPackDependency(CircularPackDependency),
691
692 #[error("pack '{name}' registered twice (indices {first_idx} and {second_idx})")]
693 PackRedeclared {
694 name: String,
695 first_idx: usize,
696 second_idx: usize,
697 },
698
699 #[error(
703 "verb collision: verb {verb:?} declared by both pack {first_pack:?} and pack \
704 {second_pack:?}; rename one handler or use Visibility::Subhandler for internal verbs"
705 )]
706 VerbCollision {
707 verb: String,
708 first_pack: String,
709 second_pack: String,
710 },
711
712 #[error(
714 "pack {pack:?} handler {verb:?} declares request-envelope parameter {param:?}; rename the verb argument"
715 )]
716 ReservedEnvelopeParam {
717 pack: String,
718 verb: String,
719 param: String,
720 },
721
722 #[error("{}", .0.message)]
740 RefusedWithReceipt(Box<ReceiptRefusal>),
741
742 #[error(transparent)]
746 RefusedWithEvents { context: RefusalEventContext },
747
748 #[error(transparent)]
751 WithResolution { context: ResolutionFailureContext },
752
753 #[error("permission denied for verb {verb:?}: {reason}")]
763 PermissionDenied {
764 verb: String,
765 reason: String,
766 receipt: Box<DenialReceipt>,
767 },
768
769 #[error("gate unavailable for verb {verb:?}: {reason}")]
775 GateUnavailable { verb: String, reason: String },
776
777 #[error("{0}")]
781 Khive(khive_types::KhiveError),
782
783 #[error("not found in this namespace")]
788 NamespaceMismatch { id: uuid::Uuid },
789
790 #[error("ambiguous prefix {prefix:?}: matches {}", format_uuid_list(matches))]
796 AmbiguousPrefix {
797 prefix: String,
798 matches: Vec<uuid::Uuid>,
799 },
800
801 #[error(
806 "cross-backend merge is not supported: \
807 into_id {into_id} is on backend '{into_backend}', \
808 from_id {from_id} is on backend '{from_backend}'. \
809 Both entities must be on the same backend to merge."
810 )]
811 CrossBackendMergeUnsupported {
812 into_id: uuid::Uuid,
813 from_id: uuid::Uuid,
814 into_backend: String,
815 from_backend: String,
816 },
817
818 #[error("unknown remote: {name:?}")]
821 UnknownRemote { name: String },
822
823 #[error("remote cache missing for remote={remote:?} namespace={namespace:?}")]
825 RemoteCacheMissing { remote: String, namespace: String },
826
827 #[error("ambiguous id {id:?}: matched {count} records")]
829 AmbiguousId { id: String, count: usize },
830
831 #[error("cross-namespace write denied: cannot write to remote namespace {namespace:?}")]
833 CrossNamespaceWrite { namespace: String },
834
835 #[error("remote fetch error for remote={remote:?}: {message}")]
838 RemoteFetchError { remote: String, message: String },
839
840 #[error(
846 "write budget exceeded: max_new_entries={max_new_entries}, \
847 attempted_new_entries={attempted_new_entries}"
848 )]
849 WriteBudgetExceeded {
850 max_new_entries: u64,
851 attempted_new_entries: u64,
852 },
853
854 #[error("write blocked: {0}")]
860 SecretDetected(crate::secret_gate::SecretMatch),
861
862 #[error("{operation} exceeded its {budget_ms}ms deadline (elapsed {elapsed_ms}ms)")]
876 DeadlineExceeded {
877 operation: String,
878 budget_ms: u64,
879 elapsed_ms: u64,
880 },
881}
882
883impl From<khive_db::SqliteError> for RuntimeError {
884 fn from(error: khive_db::SqliteError) -> Self {
885 match error {
886 khive_db::SqliteError::RequestReadStopped(error) => Self::Storage(error),
887 khive_db::SqliteError::InheritedWriterTransaction
888 | khive_db::SqliteError::WriterSettlementUnknown => {
889 Self::Storage(khive_storage::StorageError::WriterTaskTerminated {
890 request_state: khive_storage::WriterTaskRequestState::SideEffectsUnknown,
891 })
892 }
893 khive_db::SqliteError::WriterPoisoned => {
894 Self::Storage(khive_storage::StorageError::WriterTaskTerminated {
895 request_state: khive_storage::WriterTaskRequestState::NotStarted,
896 })
897 }
898 error => Self::Sqlite(error),
899 }
900 }
901}
902
903impl RuntimeError {
904 pub fn with_resolution(self, resolution: ResolutionFacts) -> Self {
906 Self::WithResolution {
907 context: ResolutionFailureContext {
908 source: Box::new(self),
909 resolution,
910 },
911 }
912 }
913
914 pub fn with_refusal_events(self, mut recordings: Vec<RefusalEventRecording>) -> Self {
917 if recordings.is_empty() {
918 return self;
919 }
920 match self {
921 Self::RefusedWithEvents { mut context } => {
922 context.recordings.append(&mut recordings);
923 Self::RefusedWithEvents { context }
924 }
925 source => Self::RefusedWithEvents {
926 context: RefusalEventContext {
927 source: Box::new(source),
928 recordings,
929 },
930 },
931 }
932 }
933
934 pub fn refusal_source(&self) -> &Self {
937 let mut error = self;
938 loop {
939 error = match error {
940 Self::RefusedWithEvents { context } => &context.source,
941 Self::WithResolution { context } => context.source.as_ref(),
942 _ => return error,
943 };
944 }
945 }
946
947 pub fn is_stream_policy_refusal(&self) -> bool {
954 matches!(self.refusal_source(), Self::Khive(error)
955 if error.kind() == khive_types::ErrorKind::Conflict
956 && error.details().and_then(|details| details.get("reason"))
957 == Some("stream_member"))
958 }
959
960 pub fn permission_denied(verb: impl Into<String>, reason: impl Into<String>) -> Self {
962 Self::PermissionDenied {
963 verb: verb.into(),
964 reason: reason.into(),
965 receipt: Box::new(DenialReceipt::not_audited()),
966 }
967 }
968
969 pub fn internal_with_context(context: impl fmt::Display, source: impl fmt::Display) -> Self {
974 Self::Internal(format!("{context}: {source}"))
975 }
976
977 pub fn channel_ingest_failure_class(&self) -> ChannelIngestFailureClass {
985 let source = self.refusal_source();
986 let reason = source.variant_name();
987 if source.retryable_failure_context().is_some() {
988 ChannelIngestFailureClass::Retryable { reason }
989 } else if matches!(source, Self::SecretDetected(_)) {
990 ChannelIngestFailureClass::Permanent { reason }
991 } else {
992 ChannelIngestFailureClass::Unknown { reason }
993 }
994 }
995
996 fn variant_name(&self) -> &'static str {
998 match self {
999 Self::AuditObligation { .. } => "AuditObligation",
1000 Self::Storage(_) => "Storage",
1001 Self::Sqlite(_) => "Sqlite",
1002 Self::Query(_) => "Query",
1003 Self::NotFound(_) => "NotFound",
1004 Self::InvalidInput(_) => "InvalidInput",
1005 Self::UnknownVerb(_) => "UnknownVerb",
1006 Self::Unconfigured(_) => "Unconfigured",
1007 Self::UnknownModel(_) => "UnknownModel",
1008 Self::Embedding(_) => "Embedding",
1009 Self::Ambiguous(_) => "Ambiguous",
1010 Self::Fusion(_) => "Fusion",
1011 Self::UnknownFusionStrategy(_) => "UnknownFusionStrategy",
1012 Self::Internal(_) => "Internal",
1013 Self::GuardedWriteFailed(_) => "GuardedWriteFailed",
1014 Self::MissingPackDependency(_) => "MissingPackDependency",
1015 Self::MissingPackDependencies(_) => "MissingPackDependencies",
1016 Self::CircularPackDependency(_) => "CircularPackDependency",
1017 Self::PackRedeclared { .. } => "PackRedeclared",
1018 Self::VerbCollision { .. } => "VerbCollision",
1019 Self::ReservedEnvelopeParam { .. } => "ReservedEnvelopeParam",
1020 Self::PermissionDenied { .. } => "PermissionDenied",
1021 Self::GateUnavailable { .. } => "GateUnavailable",
1022 Self::Khive(_) => "Khive",
1023 Self::NamespaceMismatch { .. } => "NamespaceMismatch",
1024 Self::AmbiguousPrefix { .. } => "AmbiguousPrefix",
1025 Self::CrossBackendMergeUnsupported { .. } => "CrossBackendMergeUnsupported",
1026 Self::UnknownRemote { .. } => "UnknownRemote",
1027 Self::RemoteCacheMissing { .. } => "RemoteCacheMissing",
1028 Self::AmbiguousId { .. } => "AmbiguousId",
1029 Self::CrossNamespaceWrite { .. } => "CrossNamespaceWrite",
1030 Self::RemoteFetchError { .. } => "RemoteFetchError",
1031 Self::WriteBudgetExceeded { .. } => "WriteBudgetExceeded",
1032 Self::SecretDetected(_) => "SecretDetected",
1033 Self::DeadlineExceeded { .. } => "DeadlineExceeded",
1034 Self::IncompatibleEventStore(_) => "IncompatibleEventStore",
1035 Self::RefusedWithReceipt { .. } => "RefusedWithReceipt",
1036 Self::RefusedWithEvents { context } => context.source.variant_name(),
1037 Self::WithResolution { context } => context.source.variant_name(),
1038 }
1039 }
1040
1041 pub fn writer_pool_checkout_timeout_context(&self) -> Option<WriterPoolCheckoutTimeoutContext> {
1049 let (sqlite_error, capability, operation) = match self.refusal_source() {
1050 Self::Sqlite(error) => (error, None, None),
1051 Self::Storage(khive_storage::StorageError::Driver {
1052 capability,
1053 operation,
1054 source,
1055 }) => (
1056 source.downcast_ref::<khive_db::SqliteError>()?,
1057 Some(*capability),
1058 Some(operation.to_string()),
1059 ),
1060 _ => return None,
1061 };
1062
1063 let khive_db::SqliteError::WriterPoolCheckoutTimeout { timeout } = sqlite_error else {
1064 return None;
1065 };
1066 Some(WriterPoolCheckoutTimeoutContext {
1067 timeout: *timeout,
1068 capability,
1069 operation,
1070 })
1071 }
1072
1073 pub fn admission_failure_context(&self) -> Option<AdmissionFailureContext> {
1078 let source = self.refusal_source();
1079 if let Some(context) = self.writer_pool_checkout_timeout_context() {
1080 return Some(AdmissionFailureContext {
1081 stage: WRITER_POOL_CHECKOUT_TIMEOUT_STAGE,
1082 timeout: context.timeout,
1083 capability: context.capability,
1084 operation: context.operation,
1085 pool_identity: None,
1086 scope: None,
1087 retry_after_ms: None,
1088 });
1089 }
1090 if let Self::Storage(khive_storage::StorageError::WriteQueueFull { timeout_ms }) = source {
1091 return Some(AdmissionFailureContext {
1092 stage: WRITER_QUEUE_SATURATED_STAGE,
1093 timeout: Duration::from_millis(*timeout_ms),
1094 capability: None,
1095 operation: None,
1096 pool_identity: None,
1097 scope: Some(WRITER_ADMISSION_SCOPE),
1098 retry_after_ms: Some(*timeout_ms),
1099 });
1100 }
1101 if let Self::Storage(khive_storage::StorageError::AdmissionTimeout {
1102 operation,
1103 timeout_ms,
1104 pool_identity,
1105 }) = source
1106 {
1107 return Some(AdmissionFailureContext {
1108 stage: STORAGE_ADMISSION_TIMEOUT_STAGE,
1109 timeout: Duration::from_millis(*timeout_ms),
1110 capability: None,
1111 operation: Some(operation.to_string()),
1112 pool_identity: pool_identity.clone(),
1113 scope: None,
1114 retry_after_ms: None,
1115 });
1116 }
1117 None
1118 }
1119
1120 pub fn writer_task_failure_context(&self) -> Option<WriterTaskFailureContext> {
1124 let Self::Storage(error) = self.refusal_source() else {
1125 return None;
1126 };
1127 match error {
1128 khive_storage::StorageError::WriterTaskRequestFailed {
1129 request_state,
1130 source,
1131 } => Some(WriterTaskFailureContext {
1132 stage: WRITER_TASK_REQUEST_FAILED_STAGE,
1133 request_state: *request_state,
1134 task_terminated: false,
1135 retryable: source.is_retryable(),
1136 }),
1137 khive_storage::StorageError::WriterTaskTerminated { request_state } => {
1138 Some(WriterTaskFailureContext {
1139 stage: WRITER_TASK_TERMINATED_STAGE,
1140 request_state: *request_state,
1141 task_terminated: true,
1142 retryable: false,
1143 })
1144 }
1145 _ => None,
1146 }
1147 }
1148
1149 pub fn retryable_failure_context(&self) -> Option<RetryableFailureContext> {
1152 let source = self.refusal_source();
1153 if let Some(context) = self.admission_failure_context() {
1154 return Some(context.into());
1155 }
1156 if let Self::Storage(khive_storage::StorageError::WriterTaskBusy { timeout_ms }) = source {
1157 return Some(RetryableFailureContext {
1158 stage: WRITER_TASK_BEGIN_BUSY_STAGE,
1159 timeout: Duration::from_millis(*timeout_ms),
1160 capability: None,
1161 operation: Some("writer_task_begin".to_string()),
1162 pool_identity: None,
1163 scope: None,
1164 retry_after_ms: None,
1165 });
1166 }
1167 if let Self::Storage(khive_storage::StorageError::ReadTransactionAgeEvicted {
1168 operation,
1169 max_age_secs,
1170 }) = source
1171 {
1172 return Some(RetryableFailureContext {
1173 stage: READ_TX_AGE_EVICTED_STAGE,
1174 timeout: Duration::from_secs(*max_age_secs),
1175 capability: Some(khive_storage::StorageCapability::Sql),
1176 operation: Some(operation.to_string()),
1177 pool_identity: None,
1178 scope: None,
1179 retry_after_ms: None,
1180 });
1181 }
1182 if let Self::Storage(
1183 khive_storage::StorageError::ReadTransactionAgeEvictionCleanupFailed {
1184 operation,
1185 max_age_secs,
1186 ..
1187 },
1188 ) = source
1189 {
1190 return Some(RetryableFailureContext {
1191 stage: READ_TX_AGE_EVICTED_STAGE,
1192 timeout: Duration::from_secs(*max_age_secs),
1193 capability: Some(khive_storage::StorageCapability::Sql),
1194 operation: Some(operation.to_string()),
1195 pool_identity: None,
1196 scope: None,
1197 retry_after_ms: None,
1198 });
1199 }
1200 None
1201 }
1202}
1203
1204pub fn fts_text_leg_or_err<T>(
1211 result: Result<Vec<T>, RuntimeError>,
1212 context: &'static str,
1213 query: &str,
1214) -> RuntimeResult<Vec<T>> {
1215 match result {
1216 Ok(hits) => Ok(hits),
1217 Err(RuntimeError::Storage(se)) if se.is_fts5_syntax_error() => {
1218 tracing::warn!(
1219 error = %se,
1220 query = %query,
1221 context,
1222 "FTS text leg failed on a parser syntax error; failing loud (#569)"
1223 );
1224 Err(RuntimeError::InvalidInput(format!(
1225 "{context}: FTS query could not be parsed: {se}"
1226 )))
1227 }
1228 Err(e) => Err(e),
1229 }
1230}
1231
1232fn format_uuid_list(uuids: &[uuid::Uuid]) -> String {
1233 uuids
1234 .iter()
1235 .map(uuid::Uuid::to_string)
1236 .collect::<Vec<_>>()
1237 .join(", ")
1238}
1239
1240impl From<khive_types::EntityTypeError> for RuntimeError {
1244 fn from(e: khive_types::EntityTypeError) -> Self {
1245 Self::InvalidInput(e.to_string())
1246 }
1247}
1248
1249impl From<khive_types::KhiveError> for RuntimeError {
1250 fn from(e: khive_types::KhiveError) -> Self {
1251 Self::Khive(e)
1252 }
1253}
1254
1255#[cfg(test)]
1256mod refusal_event_context_tests {
1257 use super::*;
1258
1259 fn failed() -> Vec<RefusalEventRecording> {
1260 vec![RefusalEventRecording::Failed {
1261 item_index: 2,
1262 subject: Uuid::from_u128(11),
1263 error_class: RefusalRecordingErrorClass::EventAppendFailed,
1264 }]
1265 }
1266
1267 #[test]
1268 fn refusal_events_preserve_the_source_type_display_and_policy_class() {
1269 let original = RuntimeError::SecretDetected(crate::secret_gate::SecretMatch {
1270 detector: "fixture",
1271 trigger: None,
1272 masked: "never-in-display".into(),
1273 location: Some("atoms[2].content".into()),
1274 });
1275 let message = original.to_string();
1276 let wrapped = original.with_refusal_events(failed());
1277 assert_eq!(wrapped.to_string(), message);
1278 assert!(matches!(
1279 wrapped.refusal_source(),
1280 RuntimeError::SecretDetected(_)
1281 ));
1282 assert!(std::error::Error::source(&wrapped)
1283 .unwrap()
1284 .downcast_ref::<RuntimeError>()
1285 .is_some_and(|error| matches!(error, RuntimeError::SecretDetected(_))));
1286 assert_eq!(
1287 wrapped.channel_ingest_failure_class(),
1288 ChannelIngestFailureClass::Permanent {
1289 reason: "SecretDetected"
1290 }
1291 );
1292 assert!(wrapped.retryable_failure_context().is_none());
1293
1294 let lookalike = RuntimeError::InvalidInput(message).with_refusal_events(failed());
1295 assert_eq!(
1296 lookalike.channel_ingest_failure_class(),
1297 ChannelIngestFailureClass::Unknown {
1298 reason: "InvalidInput"
1299 }
1300 );
1301 let policy = RuntimeError::Khive(
1302 khive_types::KhiveError::conflict("immutable")
1303 .with_details(khive_types::Details::new([("reason", "stream_member")])),
1304 )
1305 .with_refusal_events(failed());
1306 assert!(policy.is_stream_policy_refusal());
1307 }
1308
1309 #[test]
1310 fn refusal_events_preserve_retry_and_writer_finality_contexts() {
1311 let checkout = RuntimeError::Sqlite(khive_db::SqliteError::WriterPoolCheckoutTimeout {
1312 timeout: Duration::from_millis(17),
1313 })
1314 .with_refusal_events(failed());
1315 assert_eq!(
1316 checkout
1317 .writer_pool_checkout_timeout_context()
1318 .unwrap()
1319 .timeout,
1320 Duration::from_millis(17)
1321 );
1322 assert_eq!(
1323 checkout.admission_failure_context().unwrap().stage,
1324 WRITER_POOL_CHECKOUT_TIMEOUT_STAGE
1325 );
1326 assert_eq!(
1327 checkout.retryable_failure_context().unwrap().stage,
1328 WRITER_POOL_CHECKOUT_TIMEOUT_STAGE
1329 );
1330
1331 let queued =
1332 RuntimeError::Storage(khive_storage::StorageError::WriteQueueFull { timeout_ms: 23 })
1333 .with_refusal_events(failed());
1334 assert_eq!(
1335 queued.retryable_failure_context().unwrap().stage,
1336 WRITER_QUEUE_SATURATED_STAGE
1337 );
1338 assert_eq!(
1339 queued.channel_ingest_failure_class(),
1340 ChannelIngestFailureClass::Retryable { reason: "Storage" }
1341 );
1342
1343 let busy =
1344 RuntimeError::Storage(khive_storage::StorageError::WriterTaskBusy { timeout_ms: 31 })
1345 .with_refusal_events(failed());
1346 assert_eq!(
1347 busy.retryable_failure_context().unwrap().stage,
1348 WRITER_TASK_BEGIN_BUSY_STAGE
1349 );
1350
1351 let stopped = RuntimeError::Storage(khive_storage::StorageError::WriterTaskTerminated {
1352 request_state: khive_storage::WriterTaskRequestState::SideEffectsUnknown,
1353 })
1354 .with_refusal_events(failed());
1355 let context = stopped.writer_task_failure_context().unwrap();
1356 assert_eq!(
1357 context.request_state,
1358 khive_storage::WriterTaskRequestState::SideEffectsUnknown
1359 );
1360 assert!(context.task_terminated);
1361 assert!(!context.retryable);
1362 assert!(stopped.retryable_failure_context().is_none());
1363 }
1364
1365 #[test]
1366 fn refusal_events_preserve_the_original_error_and_its_underlying_source_chain() {
1367 let original = RuntimeError::Storage(khive_storage::StorageError::driver(
1368 khive_storage::StorageCapability::Notes,
1369 "refusal-context-fixture",
1370 std::io::Error::new(std::io::ErrorKind::PermissionDenied, "fixture refusal"),
1371 ));
1372 let message = original.to_string();
1373 let wrapped = original.with_refusal_events(failed());
1374 assert_eq!(wrapped.to_string(), message);
1375
1376 let runtime = std::error::Error::source(&wrapped)
1377 .unwrap()
1378 .downcast_ref::<RuntimeError>()
1379 .expect("the first source must retain the runtime classification");
1380 assert!(matches!(runtime, RuntimeError::Storage(_)));
1381 let storage = std::error::Error::source(runtime)
1382 .unwrap()
1383 .downcast_ref::<khive_storage::StorageError>()
1384 .expect("the original runtime error must retain its storage source");
1385 let driver = std::error::Error::source(storage)
1386 .unwrap()
1387 .downcast_ref::<std::io::Error>()
1388 .expect("the storage source chain must remain traversable");
1389 assert_eq!(driver.kind(), std::io::ErrorKind::PermissionDenied);
1390 }
1391
1392 #[test]
1393 fn refusal_events_with_no_eligible_targets_do_not_wrap_or_grow_metadata() {
1394 let error = RuntimeError::InvalidInput("unchanged".into()).with_refusal_events(vec![]);
1395 assert!(matches!(error, RuntimeError::InvalidInput(ref message) if message == "unchanged"));
1396 let error = error
1397 .with_refusal_events(failed())
1398 .with_refusal_events(vec![RefusalEventRecording::Recorded {
1399 item_index: 3,
1400 subject: Uuid::from_u128(12),
1401 event_id: Uuid::from_u128(13),
1402 }]);
1403 let source = std::error::Error::source(&error)
1404 .unwrap()
1405 .downcast_ref::<RuntimeError>()
1406 .expect("appending recordings must not add another error wrapper");
1407 assert!(matches!(source, RuntimeError::InvalidInput(_)));
1408 assert!(std::error::Error::source(source).is_none());
1409 let RuntimeError::RefusedWithEvents { context } = error else {
1410 panic!("recordings must accompany the original typed source");
1411 };
1412 let RefusalEventContext { source, recordings } = context;
1413 assert!(matches!(*source, RuntimeError::InvalidInput(_)));
1414 assert_eq!(recordings.len(), 2);
1415 assert_eq!(recordings[0], failed()[0]);
1416 }
1417
1418 #[test]
1419 fn refusal_recording_classes_serialize_only_the_closed_safe_spellings() {
1420 for (class, name) in [
1421 (
1422 RefusalRecordingErrorClass::EventStoreUnavailable,
1423 "event_store_unavailable",
1424 ),
1425 (
1426 RefusalRecordingErrorClass::EventAppendFailed,
1427 "event_append_failed",
1428 ),
1429 ] {
1430 assert_eq!(class.as_str(), name);
1431 assert_eq!(
1432 serde_json::to_value(class).unwrap(),
1433 serde_json::json!(name)
1434 );
1435 }
1436 }
1437}
1438
1439#[cfg(test)]
1440mod stream_policy_refusal_tests {
1441 use super::RuntimeError;
1442 use khive_types::{Details, KhiveError};
1443
1444 #[test]
1445 fn stream_policy_refusal_uses_structured_kind_and_marker() {
1446 let policy = RuntimeError::Khive(
1447 KhiveError::conflict("changed diagnostic wording")
1448 .with_details(Details::new([("reason", "stream_member")])),
1449 );
1450 assert!(policy.is_stream_policy_refusal());
1451
1452 let wrong_kind = RuntimeError::Khive(
1453 KhiveError::invalid_input("stream entries are immutable")
1454 .with_details(Details::new([("reason", "stream_member")])),
1455 );
1456 assert!(!wrong_kind.is_stream_policy_refusal());
1457 }
1458
1459 #[test]
1460 fn stream_policy_refusal_excludes_other_conflicts_and_message_lookalikes() {
1461 for marker in [
1462 None,
1463 Some("seq_conflict"),
1464 Some("key_conflict"),
1465 Some("version_conflict"),
1466 Some("fence_conflict"),
1467 Some("identity_conflict"),
1468 Some("expired"),
1469 Some("live_until_unreadable"),
1470 Some("key_ambiguous"),
1471 Some("unknown_op"),
1472 Some("stream_member_extra"),
1473 ] {
1474 let mut error = KhiveError::conflict("stream entries are immutable");
1475 if let Some(marker) = marker {
1476 error = error.with_details(Details::new([("reason", marker)]));
1477 }
1478 assert!(
1479 !RuntimeError::Khive(error).is_stream_policy_refusal(),
1480 "unrelated conflict marker {marker:?}"
1481 );
1482 }
1483
1484 let internal = RuntimeError::Internal("conflict: stream entries are immutable".into());
1485 assert!(!internal.is_stream_policy_refusal());
1486 let driver = RuntimeError::Storage(khive_storage::StorageError::driver(
1487 khive_storage::StorageCapability::Sql,
1488 "inline.execute",
1489 std::io::Error::other("stream_member"),
1490 ));
1491 assert!(!driver.is_stream_policy_refusal());
1492 }
1493}
1494
1495#[cfg(test)]
1496mod channel_ingest_failure_class_tests {
1497 use super::{ChannelIngestFailureClass, ResolutionFacts, RuntimeError};
1498 use crate::secret_gate::SecretMatch;
1499 use std::time::Duration;
1500 use uuid::Uuid;
1501
1502 #[test]
1503 fn secret_detected_is_permanent_by_typed_variant_not_display_text() {
1504 let first = RuntimeError::SecretDetected(SecretMatch {
1505 detector: "fixture",
1506 trigger: None,
1507 masked: "first-rendering".to_string(),
1508 location: None,
1509 });
1510 let second = RuntimeError::SecretDetected(SecretMatch {
1511 detector: "fixture",
1512 trigger: Some("token"),
1513 masked: "completely-different-rendering".to_string(),
1514 location: None,
1515 });
1516
1517 assert_ne!(first.to_string(), second.to_string());
1518 assert_eq!(
1519 first.channel_ingest_failure_class(),
1520 ChannelIngestFailureClass::Permanent {
1521 reason: "SecretDetected"
1522 }
1523 );
1524 assert_eq!(
1525 second.channel_ingest_failure_class(),
1526 ChannelIngestFailureClass::Permanent {
1527 reason: "SecretDetected"
1528 },
1529 "rendered error details must not participate in ingest classification"
1530 );
1531 }
1532
1533 #[test]
1534 fn admission_failures_are_retryable_and_unclassified_errors_are_unknown() {
1535 let retryable =
1536 RuntimeError::Storage(khive_storage::StorageError::WriteQueueFull { timeout_ms: 25 });
1537 assert_eq!(
1538 retryable.channel_ingest_failure_class(),
1539 ChannelIngestFailureClass::Retryable { reason: "Storage" }
1540 );
1541
1542 let begin_busy =
1543 RuntimeError::Storage(khive_storage::StorageError::WriterTaskBusy { timeout_ms: 175 });
1544 let context = begin_busy
1545 .retryable_failure_context()
1546 .expect("contended BEGIN must remain typed and retryable");
1547 assert_eq!(context.stage, "writer_task_begin_busy");
1548 assert_eq!(context.timeout, Duration::from_millis(175));
1549 assert_eq!(context.operation.as_deref(), Some("writer_task_begin"));
1550 assert_eq!(context.scope, None);
1551 assert_eq!(context.retry_after_ms, None);
1552 assert_eq!(
1553 begin_busy.channel_ingest_failure_class(),
1554 ChannelIngestFailureClass::Retryable { reason: "Storage" }
1555 );
1556
1557 let unknown = RuntimeError::InvalidInput(
1558 "write blocked: SecretDetected text must not affect classification".to_string(),
1559 );
1560 assert_eq!(
1561 unknown.channel_ingest_failure_class(),
1562 ChannelIngestFailureClass::Unknown {
1563 reason: "InvalidInput"
1564 },
1565 "a rendered message resembling SecretDetected must remain Unknown unless its typed variant is SecretDetected"
1566 );
1567 }
1568
1569 #[test]
1570 fn resolution_wrapper_keeps_remote_fetch_retry_classification() {
1571 let wrapped = RuntimeError::RemoteFetchError {
1572 remote: "https://example.com/repo".into(),
1573 message: "cache repair failed".into(),
1574 }
1575 .with_resolution(ResolutionFacts {
1576 project_id: Uuid::from_u128(1),
1577 duplicate_anchor_ids: vec![],
1578 slug_backfilled: false,
1579 project_created: false,
1580 orphaned_project_id: None,
1581 orphaned_note_count: 0,
1582 });
1583 let RuntimeError::WithResolution { context } = &wrapped else {
1584 panic!("expected resolution context");
1585 };
1586 assert!(matches!(
1587 std::error::Error::source(context)
1588 .and_then(|source| source.downcast_ref::<RuntimeError>()),
1589 Some(RuntimeError::RemoteFetchError { .. })
1590 ));
1591 assert!(matches!(
1592 wrapped.refusal_source(),
1593 RuntimeError::RemoteFetchError { .. }
1594 ));
1595 assert!(wrapped.retryable_failure_context().is_none());
1596 assert_eq!(
1597 wrapped.channel_ingest_failure_class(),
1598 ChannelIngestFailureClass::Unknown {
1599 reason: "RemoteFetchError"
1600 }
1601 );
1602 }
1603
1604 #[test]
1605 fn writer_task_failure_context_separates_request_finality_from_task_liveness() {
1606 use khive_storage::{StorageError, WriterTaskRequestState};
1607
1608 let rolled_back = RuntimeError::Storage(StorageError::WriterTaskRequestFailed {
1609 request_state: WriterTaskRequestState::TransactionRolledBack,
1610 source: Box::new(StorageError::Pool {
1611 operation: "writer_task_commit".into(),
1612 message: "commit refused".into(),
1613 }),
1614 });
1615 let rolled_back_context = rolled_back
1616 .writer_task_failure_context()
1617 .expect("proven rollback must remain typed through RuntimeError");
1618 assert_eq!(
1619 rolled_back_context.request_state,
1620 WriterTaskRequestState::TransactionRolledBack
1621 );
1622 assert_eq!(
1623 rolled_back_context.stage,
1624 super::WRITER_TASK_REQUEST_FAILED_STAGE
1625 );
1626 assert!(!rolled_back_context.task_terminated);
1627 assert!(rolled_back_context.retryable);
1628
1629 let unknown = RuntimeError::Storage(StorageError::WriterTaskTerminated {
1630 request_state: WriterTaskRequestState::SideEffectsUnknown,
1631 });
1632 let unknown_context = unknown
1633 .writer_task_failure_context()
1634 .expect("ambiguous finality must remain typed through RuntimeError");
1635 assert_eq!(
1636 unknown_context.request_state,
1637 WriterTaskRequestState::SideEffectsUnknown
1638 );
1639 assert_eq!(unknown_context.stage, super::WRITER_TASK_TERMINATED_STAGE);
1640 assert!(unknown_context.task_terminated);
1641 assert!(!unknown_context.retryable);
1642 }
1643
1644 #[test]
1645 fn storage_admission_timeout_is_a_retryable_admission_failure() {
1646 let admission = RuntimeError::Storage(khive_storage::StorageError::AdmissionTimeout {
1647 operation: "sql_bridge.writer_handle".into(),
1648 timeout_ms: 30_000,
1649 pool_identity: None,
1650 });
1651 let context = admission
1652 .admission_failure_context()
1653 .expect("a storage admission timeout happens before the operation starts");
1654 assert_eq!(context.stage, super::STORAGE_ADMISSION_TIMEOUT_STAGE);
1655 assert_eq!(context.timeout, Duration::from_millis(30_000));
1656 assert_eq!(
1657 context.operation.as_deref(),
1658 Some("sql_bridge.writer_handle")
1659 );
1660 assert_eq!(context.scope, None);
1661 assert_eq!(context.retry_after_ms, None);
1662 assert_eq!(
1663 admission.channel_ingest_failure_class(),
1664 ChannelIngestFailureClass::Retryable { reason: "Storage" }
1665 );
1666
1667 let mid_flight = RuntimeError::Storage(khive_storage::StorageError::Timeout {
1670 operation: "sql_bridge.reader_open".into(),
1671 });
1672 assert!(mid_flight.admission_failure_context().is_none());
1673 assert!(mid_flight.retryable_failure_context().is_none());
1674 }
1675
1676 #[test]
1683 fn both_read_tx_age_eviction_outcomes_map_to_the_same_retryable_stage() {
1684 let clean = RuntimeError::Storage(khive_storage::StorageError::ReadTransactionAgeEvicted {
1685 operation: "query_all".into(),
1686 max_age_secs: 120,
1687 });
1688 let clean_context = clean
1689 .retryable_failure_context()
1690 .expect("a clean age eviction must be typed-retryable");
1691 assert_eq!(clean_context.stage, super::READ_TX_AGE_EVICTED_STAGE);
1692 assert_eq!(clean_context.timeout, Duration::from_secs(120));
1693 assert_eq!(clean_context.operation.as_deref(), Some("query_all"));
1694 assert_eq!(
1695 clean_context.capability,
1696 Some(khive_storage::StorageCapability::Sql)
1697 );
1698
1699 let cleanup_failed = RuntimeError::Storage(
1700 khive_storage::StorageError::ReadTransactionAgeEvictionCleanupFailed {
1701 operation: "query_all".into(),
1702 max_age_secs: 120,
1703 message: "rollback failed: disk I/O error".into(),
1704 },
1705 );
1706 let cleanup_failed_context = cleanup_failed
1707 .retryable_failure_context()
1708 .expect("a failed cleanup rollback must remain typed-retryable");
1709 assert_eq!(
1710 cleanup_failed_context.stage,
1711 super::READ_TX_AGE_EVICTED_STAGE
1712 );
1713 assert_eq!(cleanup_failed_context.timeout, Duration::from_secs(120));
1714 assert_eq!(
1715 cleanup_failed_context.operation.as_deref(),
1716 Some("query_all")
1717 );
1718 assert_eq!(
1719 cleanup_failed_context.capability,
1720 Some(khive_storage::StorageCapability::Sql)
1721 );
1722
1723 assert_ne!(
1724 clean.to_string(),
1725 cleanup_failed.to_string(),
1726 "the rendered message must distinguish a clean eviction from a failed cleanup even \
1727 though both map to the same wire stage"
1728 );
1729 assert!(cleanup_failed.to_string().contains("rollback failed"));
1730 }
1731}