1use std::borrow::Cow;
4use std::error::Error as StdError;
5use std::fmt;
6
7use thiserror::Error;
8
9use crate::blob::ContentRef;
10use crate::capability::StorageCapability;
11
12#[derive(Debug, Clone, Copy, PartialEq, Eq)]
20pub enum WriterTaskRequestState {
21 NotStarted,
23 TransactionRolledBack,
28 SideEffectsUnknown,
32}
33
34impl fmt::Display for WriterTaskRequestState {
35 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
36 f.write_str(match self {
37 Self::NotStarted => "not_started",
38 Self::TransactionRolledBack => "transaction_rolled_back",
39 Self::SideEffectsUnknown => "side_effects_unknown",
40 })
41 }
42}
43
44#[derive(Debug, Clone, Copy, PartialEq, Eq)]
46pub enum CapacityUnavailablePhase {
47 Identity,
48 Lock,
49 Probe,
50}
51
52impl CapacityUnavailablePhase {
53 pub const fn as_str(self) -> &'static str {
54 match self {
55 Self::Identity => "identity",
56 Self::Lock => "lock",
57 Self::Probe => "probe",
58 }
59 }
60}
61
62impl fmt::Display for CapacityUnavailablePhase {
63 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
64 f.write_str(self.as_str())
65 }
66}
67
68#[derive(Debug, Error)]
70pub enum StorageError {
71 #[error("{capability:?} resource not found: {resource} ({key})")]
72 NotFound {
73 capability: StorageCapability,
74 resource: &'static str,
75 key: String,
76 },
77
78 #[error("{capability:?} resource already exists: {resource} ({key})")]
79 AlreadyExists {
80 capability: StorageCapability,
81 resource: &'static str,
82 key: String,
83 },
84
85 #[error("conflict in {capability:?} during {operation}: {message}")]
86 Conflict {
87 capability: StorageCapability,
88 operation: Cow<'static, str>,
89 message: String,
90 },
91
92 #[error("invalid input for {capability:?} during {operation}: {message}")]
93 InvalidInput {
94 capability: StorageCapability,
95 operation: Cow<'static, str>,
96 message: String,
97 },
98
99 #[error("unsupported operation for {capability:?}: {operation} ({message})")]
100 Unsupported {
101 capability: StorageCapability,
102 operation: Cow<'static, str>,
103 message: String,
104 },
105
106 #[error(
111 "blob {content_ref} exceeds the {max_bytes}-byte read limit (observed at least {observed_at_least} bytes)"
112 )]
113 BlobTooLarge {
114 content_ref: ContentRef,
115 max_bytes: u64,
116 observed_at_least: u64,
117 },
118
119 #[error(
122 "blob {content_ref} metadata reports {metadata_bytes} bytes but the complete body contains {actual_bytes} bytes"
123 )]
124 BlobSizeMismatch {
125 content_ref: ContentRef,
126 metadata_bytes: u64,
127 actual_bytes: u64,
128 },
129
130 #[error("blob digest mismatch: expected {expected}, computed {actual}")]
133 BlobDigestMismatch {
134 expected: ContentRef,
135 actual: ContentRef,
136 },
137
138 #[error("pool failure during {operation}: {message}")]
139 Pool {
140 operation: Cow<'static, str>,
141 message: String,
142 },
143
144 #[error("timeout during {operation}")]
145 Timeout { operation: Cow<'static, str> },
146
147 #[error("admission timeout during {operation} after {timeout_ms}ms{pool}", pool = match .pool_identity {
153 Some(identity) => format!(" (pool: {identity})"),
154 None => String::new(),
155 })]
156 AdmissionTimeout {
157 operation: Cow<'static, str>,
158 timeout_ms: u64,
160 pool_identity: Option<String>,
162 },
163
164 #[error("sql transaction failure during {operation}: {message}")]
165 Transaction {
166 operation: Cow<'static, str>,
167 message: String,
168 },
169
170 #[error(
179 "cached read-only transaction exceeded the maximum read-transaction age \
180 ({max_age_secs}s) during {operation} and was rolled back; retry to open a fresh \
181 read snapshot"
182 )]
183 ReadTransactionAgeEvicted {
184 operation: Cow<'static, str>,
185 max_age_secs: u64,
186 },
187
188 #[error(
200 "cached read-only transaction exceeded the maximum read-transaction age \
201 ({max_age_secs}s) during {operation} but could not be cleanly rolled back \
202 ({message}); the connection was discarded, retry to open a fresh read snapshot"
203 )]
204 ReadTransactionAgeEvictionCleanupFailed {
205 operation: Cow<'static, str>,
206 max_age_secs: u64,
207 message: String,
208 },
209
210 #[error("serialization failure in {capability:?}: {message}")]
211 Serialization {
212 capability: StorageCapability,
213 message: String,
214 },
215
216 #[error("index maintenance failure in {capability:?}: {message}")]
217 IndexMaintenance {
218 capability: StorageCapability,
219 message: String,
220 },
221
222 #[error("backend driver error in {capability:?} during {operation}: {source}")]
223 Driver {
224 capability: StorageCapability,
225 operation: Cow<'static, str>,
226 #[source]
227 source: Box<dyn StdError + Send + Sync>,
228 },
229
230 #[error("write queue full: timed out after {timeout_ms}ms waiting for writer task capacity")]
235 WriteQueueFull { timeout_ms: u64 },
236
237 #[error(
242 "writer task could not begin within {timeout_ms}ms because SQLite remained busy; request was not executed"
243 )]
244 WriterTaskBusy { timeout_ms: u64 },
245
246 #[error("writer task request failed (request_state={request_state}): {source}")]
252 WriterTaskRequestFailed {
253 request_state: WriterTaskRequestState,
254 #[source]
255 source: Box<StorageError>,
256 },
257
258 #[error("writer task terminated (request_state={request_state})")]
265 WriterTaskTerminated {
266 request_state: WriterTaskRequestState,
267 sqlite_full_codes: Option<(i32, i32)>,
270 },
271
272 #[error("internal storage error: {0}")]
275 Internal(String),
276
277 #[error(
282 "KHIVE_WRITE_QUEUE=1 but no Tokio runtime context is available to spawn the writer task"
283 )]
284 WriterTaskNoRuntime,
285
286 #[error(
290 "refusing write on {capability:?} at {volume}: {available_bytes} bytes available, \
291 the {floor_bytes}-byte free-space floor plus {required_headroom_bytes} bytes of \
292 operation headroom would be violated"
293 )]
294 CapacityFloor {
295 capability: StorageCapability,
296 volume: String,
297 available_bytes: u64,
298 floor_bytes: u64,
299 required_headroom_bytes: u64,
300 },
301
302 #[error("capacity admission unavailable for {capability:?} in {phase} phase: {message}")]
304 CapacityUnavailable {
305 capability: StorageCapability,
306 phase: CapacityUnavailablePhase,
307 message: String,
308 },
309}
310
311impl StorageError {
312 pub fn writer_task_terminated(request_state: WriterTaskRequestState) -> Self {
316 Self::WriterTaskTerminated {
317 request_state,
318 sqlite_full_codes: None,
319 }
320 }
321
322 pub fn driver(
324 capability: StorageCapability,
325 operation: impl Into<Cow<'static, str>>,
326 source: impl StdError + Send + Sync + 'static,
327 ) -> Self {
328 Self::Driver {
329 capability,
330 operation: operation.into(),
331 source: Box::new(source),
332 }
333 }
334
335 pub fn capability(&self) -> Option<StorageCapability> {
337 match self {
338 Self::NotFound { capability, .. }
339 | Self::AlreadyExists { capability, .. }
340 | Self::Conflict { capability, .. }
341 | Self::InvalidInput { capability, .. }
342 | Self::Unsupported { capability, .. }
343 | Self::Serialization { capability, .. }
344 | Self::IndexMaintenance { capability, .. }
345 | Self::Driver { capability, .. }
346 | Self::CapacityFloor { capability, .. }
347 | Self::CapacityUnavailable { capability, .. } => Some(*capability),
348 Self::BlobTooLarge { .. }
349 | Self::BlobSizeMismatch { .. }
350 | Self::BlobDigestMismatch { .. } => Some(StorageCapability::Blob),
351 Self::WriterTaskRequestFailed { source, .. } => source.capability(),
352 Self::Pool { .. }
353 | Self::Timeout { .. }
354 | Self::AdmissionTimeout { .. }
355 | Self::Transaction { .. }
356 | Self::ReadTransactionAgeEvicted { .. }
357 | Self::ReadTransactionAgeEvictionCleanupFailed { .. }
358 | Self::WriteQueueFull { .. }
359 | Self::WriterTaskBusy { .. }
360 | Self::WriterTaskTerminated { .. }
361 | Self::Internal(..)
362 | Self::WriterTaskNoRuntime => None,
363 }
364 }
365
366 pub fn is_retryable(&self) -> bool {
368 if let Self::WriterTaskRequestFailed { source, .. } = self {
369 return source.is_retryable();
370 }
371 matches!(
372 self,
373 Self::Pool { .. }
374 | Self::Timeout { .. }
375 | Self::AdmissionTimeout { .. }
376 | Self::Transaction { .. }
377 | Self::ReadTransactionAgeEvicted { .. }
378 | Self::ReadTransactionAgeEvictionCleanupFailed { .. }
379 | Self::WriteQueueFull { .. }
380 | Self::WriterTaskBusy { .. }
381 )
382 }
383
384 pub fn is_fts5_syntax_error(&self) -> bool {
400 if let Self::WriterTaskRequestFailed { source, .. } = self {
401 return source.is_fts5_syntax_error();
402 }
403 let Self::Driver {
404 capability,
405 operation,
406 source,
407 } = self
408 else {
409 return false;
410 };
411 if *capability != StorageCapability::Text || operation.as_ref() != "fts_search" {
412 return false;
413 }
414 let msg = source.to_string();
415 msg.contains("fts5: syntax error")
416 || msg.contains("fts5: parser stack overflow")
417 || msg.contains("fts5: column queries are not supported")
418 || msg.contains("fts5: phrase queries are not supported (detail")
419 || msg.contains("fts5: NEAR queries are not supported (detail")
420 }
421
422 pub fn is_unique_constraint_violation(&self) -> bool {
435 if let Self::WriterTaskRequestFailed { source, .. } = self {
436 return source.is_unique_constraint_violation();
437 }
438 let Self::Driver {
439 capability,
440 operation,
441 source,
442 } = self
443 else {
444 return false;
445 };
446 if *capability != StorageCapability::Sql {
447 return false;
448 }
449 if !matches!(
450 operation.as_ref(),
451 "execute" | "pool_writer.execute" | "tx.execute"
452 ) {
453 return false;
454 }
455 source.to_string().contains("UNIQUE constraint failed")
456 }
457}
458
459#[cfg(test)]
460mod tests {
461 use super::*;
462 use std::fmt;
463
464 #[derive(Debug)]
465 struct FakeSource(String);
466
467 impl fmt::Display for FakeSource {
468 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
469 write!(f, "{}", self.0)
470 }
471 }
472
473 impl StdError for FakeSource {}
474
475 fn driver_err(operation: &'static str, message: &str) -> StorageError {
476 StorageError::driver(
477 StorageCapability::Text,
478 operation,
479 FakeSource(message.into()),
480 )
481 }
482
483 #[test]
484 fn writer_task_request_state_display_is_stable() {
485 assert_eq!(
486 WriterTaskRequestState::NotStarted.to_string(),
487 "not_started"
488 );
489 assert_eq!(
490 WriterTaskRequestState::TransactionRolledBack.to_string(),
491 "transaction_rolled_back"
492 );
493 assert_eq!(
494 WriterTaskRequestState::SideEffectsUnknown.to_string(),
495 "side_effects_unknown"
496 );
497 }
498
499 #[test]
500 fn writer_task_busy_is_retryable_without_claiming_queue_rejection() {
501 let error = StorageError::WriterTaskBusy { timeout_ms: 175 };
502 assert!(error.is_retryable());
503 assert_eq!(
504 error.to_string(),
505 "writer task could not begin within 175ms because SQLite remained busy; request was not executed"
506 );
507 assert_eq!(error.capability(), None);
508 }
509
510 #[test]
511 fn writer_task_request_failure_preserves_source_policy_and_rollback_state() {
512 let error = StorageError::WriterTaskRequestFailed {
513 request_state: WriterTaskRequestState::TransactionRolledBack,
514 source: Box::new(StorageError::Pool {
515 operation: "writer_task_commit".into(),
516 message: "commit refused".into(),
517 }),
518 };
519
520 assert_eq!(error.capability(), None);
521 assert!(
522 error.is_retryable(),
523 "rollback finality must not discard the source error's retry policy"
524 );
525 assert_eq!(
526 error.to_string(),
527 "writer task request failed (request_state=transaction_rolled_back): pool failure during writer_task_commit: commit refused"
528 );
529 assert_eq!(
530 StdError::source(&error).map(ToString::to_string),
531 Some("pool failure during writer_task_commit: commit refused".to_string()),
532 "the original typed storage error must remain the public source"
533 );
534 }
535
536 #[test]
537 fn writer_task_request_failure_does_not_invent_retryability() {
538 let error = StorageError::WriterTaskRequestFailed {
539 request_state: WriterTaskRequestState::TransactionRolledBack,
540 source: Box::new(StorageError::InvalidInput {
541 capability: StorageCapability::Notes,
542 operation: "append_note".into(),
543 message: "deterministic refusal".into(),
544 }),
545 };
546
547 assert!(!error.is_retryable());
548 assert_eq!(error.capability(), Some(StorageCapability::Notes));
549 }
550
551 #[test]
552 fn writer_task_terminated_is_uncapability_scoped_and_not_retryable() {
553 for request_state in [
554 WriterTaskRequestState::NotStarted,
555 WriterTaskRequestState::TransactionRolledBack,
556 WriterTaskRequestState::SideEffectsUnknown,
557 ] {
558 let error = StorageError::writer_task_terminated(request_state);
559 assert_eq!(error.capability(), None);
560 assert!(!error.is_retryable());
561 assert_eq!(
562 error.to_string(),
563 format!("writer task terminated (request_state={request_state})")
564 );
565 }
566 }
567
568 #[test]
569 fn blob_integrity_errors_are_blob_scoped_and_not_retryable() {
570 let requested = crate::blob::ContentRef::from_hex("a".repeat(64)).unwrap();
571 let actual = crate::blob::ContentRef::from_hex("b".repeat(64)).unwrap();
572 let errors = [
573 StorageError::BlobTooLarge {
574 content_ref: requested.clone(),
575 max_bytes: 8,
576 observed_at_least: 9,
577 },
578 StorageError::BlobSizeMismatch {
579 content_ref: requested.clone(),
580 metadata_bytes: 7,
581 actual_bytes: 8,
582 },
583 StorageError::BlobDigestMismatch {
584 expected: requested,
585 actual,
586 },
587 ];
588
589 for error in errors {
590 assert_eq!(error.capability(), Some(StorageCapability::Blob));
591 assert!(!error.is_retryable());
592 }
593 }
594
595 #[test]
596 fn fts5_syntax_error_at_fts_search_is_classified_as_syntax_error() {
597 let e = driver_err("fts_search", "fts5: syntax error near \"@\"");
598 assert!(e.is_fts5_syntax_error());
599 }
600
601 #[test]
602 fn fts5_parser_stack_overflow_is_classified_as_syntax_error() {
603 let e = driver_err("fts_search", "fts5: parser stack overflow");
604 assert!(e.is_fts5_syntax_error());
605 }
606
607 #[test]
608 fn fts5_unsupported_column_query_is_classified_as_syntax_error() {
609 let e = driver_err(
610 "fts_search",
611 "fts5: column queries are not supported (detail=none)",
612 );
613 assert!(e.is_fts5_syntax_error());
614 }
615
616 #[test]
617 fn timeout_is_not_classified_as_syntax_error() {
618 let e = StorageError::Timeout {
619 operation: "fts_search".into(),
620 };
621 assert!(!e.is_fts5_syntax_error());
622 }
623
624 #[test]
625 fn pool_failure_is_not_classified_as_syntax_error() {
626 let e = StorageError::Pool {
627 operation: "fts_search".into(),
628 message: "pool exhausted".into(),
629 };
630 assert!(!e.is_fts5_syntax_error());
631 }
632
633 #[test]
634 fn driver_error_at_non_search_operation_is_not_classified_as_syntax_error() {
635 let e = driver_err("open_fts_reader", "fts5: syntax error near \"@\"");
636 assert!(!e.is_fts5_syntax_error());
637 }
638
639 #[test]
640 fn driver_error_with_unrelated_message_is_not_classified_as_syntax_error() {
641 let e = driver_err("fts_search", "disk I/O error");
642 assert!(!e.is_fts5_syntax_error());
643 }
644
645 #[test]
646 fn fts5_phrase_detail_query_is_classified_as_syntax_error() {
647 let e = driver_err(
648 "fts_search",
649 "fts5: phrase queries are not supported (detail!=full)",
650 );
651 assert!(e.is_fts5_syntax_error());
652 }
653
654 #[test]
655 fn fts5_near_detail_query_is_classified_as_syntax_error() {
656 let e = driver_err(
657 "fts_search",
658 "fts5: NEAR queries are not supported (detail!=full)",
659 );
660 assert!(e.is_fts5_syntax_error());
661 }
662
663 #[test]
664 fn unprefixed_detail_message_is_not_classified_as_syntax_error() {
665 let e = driver_err(
666 "fts_search",
667 "phrase queries are not supported (detail!=full)",
668 );
669 assert!(!e.is_fts5_syntax_error());
670 }
671
672 #[test]
673 fn fts5_shadow_table_corruption_is_not_classified_as_syntax_error() {
674 let e = driver_err(
675 "fts_search",
676 "fts5: error creating shadow table notes_content: no such table",
677 );
678 assert!(!e.is_fts5_syntax_error());
679 }
680
681 #[test]
682 fn non_text_capability_is_not_classified_as_syntax_error() {
683 let e = StorageError::Driver {
684 capability: StorageCapability::Vectors,
685 operation: "fts_search".into(),
686 source: Box::new(FakeSource("fts5: syntax error near \"@\"".into())),
687 };
688 assert!(!e.is_fts5_syntax_error());
689 }
690
691 fn driver_err_sql(operation: &'static str, message: &str) -> StorageError {
692 StorageError::driver(
693 StorageCapability::Sql,
694 operation,
695 FakeSource(message.into()),
696 )
697 }
698
699 #[test]
700 fn unique_constraint_failure_at_execute_sql_capability_is_classified() {
701 let e = driver_err_sql(
702 "execute",
703 "UNIQUE constraint failed: brain_serve_ledger.namespace, \
704 brain_serve_ledger.target_id, brain_serve_ledger.query_class, \
705 brain_serve_ledger.served_at",
706 );
707 assert!(e.is_unique_constraint_violation());
708 }
709
710 #[test]
711 fn unique_constraint_failure_at_pool_writer_execute_is_classified() {
712 let e = driver_err_sql("pool_writer.execute", "UNIQUE constraint failed: t.id");
713 assert!(e.is_unique_constraint_violation());
714 }
715
716 #[test]
717 fn unique_constraint_message_at_non_execute_operation_is_not_classified() {
718 let e = driver_err_sql("query_row", "UNIQUE constraint failed: t.id");
719 assert!(!e.is_unique_constraint_violation());
720 }
721
722 #[test]
723 fn non_unique_driver_error_at_execute_is_not_classified() {
724 let e = driver_err_sql("execute", "disk I/O error");
725 assert!(!e.is_unique_constraint_violation());
726 }
727
728 #[test]
729 fn non_sql_capability_is_not_classified_as_unique_violation() {
730 let e = driver_err("execute", "UNIQUE constraint failed: t.id");
731 assert!(!e.is_unique_constraint_violation());
732 }
733
734 #[test]
735 fn timeout_is_not_classified_as_unique_violation() {
736 let e = StorageError::Timeout {
737 operation: "execute".into(),
738 };
739 assert!(!e.is_unique_constraint_violation());
740 }
741}