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 },
268
269 #[error("internal storage error: {0}")]
272 Internal(String),
273
274 #[error(
279 "KHIVE_WRITE_QUEUE=1 but no Tokio runtime context is available to spawn the writer task"
280 )]
281 WriterTaskNoRuntime,
282
283 #[error(
287 "refusing write on {capability:?} at {volume}: {available_bytes} bytes available, \
288 the {floor_bytes}-byte free-space floor plus {required_headroom_bytes} bytes of \
289 operation headroom would be violated"
290 )]
291 CapacityFloor {
292 capability: StorageCapability,
293 volume: String,
294 available_bytes: u64,
295 floor_bytes: u64,
296 required_headroom_bytes: u64,
297 },
298
299 #[error("capacity admission unavailable for {capability:?} in {phase} phase: {message}")]
301 CapacityUnavailable {
302 capability: StorageCapability,
303 phase: CapacityUnavailablePhase,
304 message: String,
305 },
306}
307
308impl StorageError {
309 pub fn driver(
311 capability: StorageCapability,
312 operation: impl Into<Cow<'static, str>>,
313 source: impl StdError + Send + Sync + 'static,
314 ) -> Self {
315 Self::Driver {
316 capability,
317 operation: operation.into(),
318 source: Box::new(source),
319 }
320 }
321
322 pub fn capability(&self) -> Option<StorageCapability> {
324 match self {
325 Self::NotFound { capability, .. }
326 | Self::AlreadyExists { capability, .. }
327 | Self::Conflict { capability, .. }
328 | Self::InvalidInput { capability, .. }
329 | Self::Unsupported { capability, .. }
330 | Self::Serialization { capability, .. }
331 | Self::IndexMaintenance { capability, .. }
332 | Self::Driver { capability, .. }
333 | Self::CapacityFloor { capability, .. }
334 | Self::CapacityUnavailable { capability, .. } => Some(*capability),
335 Self::BlobTooLarge { .. }
336 | Self::BlobSizeMismatch { .. }
337 | Self::BlobDigestMismatch { .. } => Some(StorageCapability::Blob),
338 Self::WriterTaskRequestFailed { source, .. } => source.capability(),
339 Self::Pool { .. }
340 | Self::Timeout { .. }
341 | Self::AdmissionTimeout { .. }
342 | Self::Transaction { .. }
343 | Self::ReadTransactionAgeEvicted { .. }
344 | Self::ReadTransactionAgeEvictionCleanupFailed { .. }
345 | Self::WriteQueueFull { .. }
346 | Self::WriterTaskBusy { .. }
347 | Self::WriterTaskTerminated { .. }
348 | Self::Internal(..)
349 | Self::WriterTaskNoRuntime => None,
350 }
351 }
352
353 pub fn is_retryable(&self) -> bool {
355 if let Self::WriterTaskRequestFailed { source, .. } = self {
356 return source.is_retryable();
357 }
358 matches!(
359 self,
360 Self::Pool { .. }
361 | Self::Timeout { .. }
362 | Self::AdmissionTimeout { .. }
363 | Self::Transaction { .. }
364 | Self::ReadTransactionAgeEvicted { .. }
365 | Self::ReadTransactionAgeEvictionCleanupFailed { .. }
366 | Self::WriteQueueFull { .. }
367 | Self::WriterTaskBusy { .. }
368 )
369 }
370
371 pub fn is_fts5_syntax_error(&self) -> bool {
387 if let Self::WriterTaskRequestFailed { source, .. } = self {
388 return source.is_fts5_syntax_error();
389 }
390 let Self::Driver {
391 capability,
392 operation,
393 source,
394 } = self
395 else {
396 return false;
397 };
398 if *capability != StorageCapability::Text || operation.as_ref() != "fts_search" {
399 return false;
400 }
401 let msg = source.to_string();
402 msg.contains("fts5: syntax error")
403 || msg.contains("fts5: parser stack overflow")
404 || msg.contains("fts5: column queries are not supported")
405 || msg.contains("fts5: phrase queries are not supported (detail")
406 || msg.contains("fts5: NEAR queries are not supported (detail")
407 }
408
409 pub fn is_unique_constraint_violation(&self) -> bool {
422 if let Self::WriterTaskRequestFailed { source, .. } = self {
423 return source.is_unique_constraint_violation();
424 }
425 let Self::Driver {
426 capability,
427 operation,
428 source,
429 } = self
430 else {
431 return false;
432 };
433 if *capability != StorageCapability::Sql {
434 return false;
435 }
436 if !matches!(
437 operation.as_ref(),
438 "execute" | "pool_writer.execute" | "tx.execute"
439 ) {
440 return false;
441 }
442 source.to_string().contains("UNIQUE constraint failed")
443 }
444}
445
446#[cfg(test)]
447mod tests {
448 use super::*;
449 use std::fmt;
450
451 #[derive(Debug)]
452 struct FakeSource(String);
453
454 impl fmt::Display for FakeSource {
455 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
456 write!(f, "{}", self.0)
457 }
458 }
459
460 impl StdError for FakeSource {}
461
462 fn driver_err(operation: &'static str, message: &str) -> StorageError {
463 StorageError::driver(
464 StorageCapability::Text,
465 operation,
466 FakeSource(message.into()),
467 )
468 }
469
470 #[test]
471 fn writer_task_request_state_display_is_stable() {
472 assert_eq!(
473 WriterTaskRequestState::NotStarted.to_string(),
474 "not_started"
475 );
476 assert_eq!(
477 WriterTaskRequestState::TransactionRolledBack.to_string(),
478 "transaction_rolled_back"
479 );
480 assert_eq!(
481 WriterTaskRequestState::SideEffectsUnknown.to_string(),
482 "side_effects_unknown"
483 );
484 }
485
486 #[test]
487 fn writer_task_busy_is_retryable_without_claiming_queue_rejection() {
488 let error = StorageError::WriterTaskBusy { timeout_ms: 175 };
489 assert!(error.is_retryable());
490 assert_eq!(
491 error.to_string(),
492 "writer task could not begin within 175ms because SQLite remained busy; request was not executed"
493 );
494 assert_eq!(error.capability(), None);
495 }
496
497 #[test]
498 fn writer_task_request_failure_preserves_source_policy_and_rollback_state() {
499 let error = StorageError::WriterTaskRequestFailed {
500 request_state: WriterTaskRequestState::TransactionRolledBack,
501 source: Box::new(StorageError::Pool {
502 operation: "writer_task_commit".into(),
503 message: "commit refused".into(),
504 }),
505 };
506
507 assert_eq!(error.capability(), None);
508 assert!(
509 error.is_retryable(),
510 "rollback finality must not discard the source error's retry policy"
511 );
512 assert_eq!(
513 error.to_string(),
514 "writer task request failed (request_state=transaction_rolled_back): pool failure during writer_task_commit: commit refused"
515 );
516 assert_eq!(
517 StdError::source(&error).map(ToString::to_string),
518 Some("pool failure during writer_task_commit: commit refused".to_string()),
519 "the original typed storage error must remain the public source"
520 );
521 }
522
523 #[test]
524 fn writer_task_request_failure_does_not_invent_retryability() {
525 let error = StorageError::WriterTaskRequestFailed {
526 request_state: WriterTaskRequestState::TransactionRolledBack,
527 source: Box::new(StorageError::InvalidInput {
528 capability: StorageCapability::Notes,
529 operation: "append_note".into(),
530 message: "deterministic refusal".into(),
531 }),
532 };
533
534 assert!(!error.is_retryable());
535 assert_eq!(error.capability(), Some(StorageCapability::Notes));
536 }
537
538 #[test]
539 fn writer_task_terminated_is_uncapability_scoped_and_not_retryable() {
540 for request_state in [
541 WriterTaskRequestState::NotStarted,
542 WriterTaskRequestState::TransactionRolledBack,
543 WriterTaskRequestState::SideEffectsUnknown,
544 ] {
545 let error = StorageError::WriterTaskTerminated { request_state };
546 assert_eq!(error.capability(), None);
547 assert!(!error.is_retryable());
548 assert_eq!(
549 error.to_string(),
550 format!("writer task terminated (request_state={request_state})")
551 );
552 }
553 }
554
555 #[test]
556 fn blob_integrity_errors_are_blob_scoped_and_not_retryable() {
557 let requested = crate::blob::ContentRef::from_hex("a".repeat(64)).unwrap();
558 let actual = crate::blob::ContentRef::from_hex("b".repeat(64)).unwrap();
559 let errors = [
560 StorageError::BlobTooLarge {
561 content_ref: requested.clone(),
562 max_bytes: 8,
563 observed_at_least: 9,
564 },
565 StorageError::BlobSizeMismatch {
566 content_ref: requested.clone(),
567 metadata_bytes: 7,
568 actual_bytes: 8,
569 },
570 StorageError::BlobDigestMismatch {
571 expected: requested,
572 actual,
573 },
574 ];
575
576 for error in errors {
577 assert_eq!(error.capability(), Some(StorageCapability::Blob));
578 assert!(!error.is_retryable());
579 }
580 }
581
582 #[test]
583 fn fts5_syntax_error_at_fts_search_is_classified_as_syntax_error() {
584 let e = driver_err("fts_search", "fts5: syntax error near \"@\"");
585 assert!(e.is_fts5_syntax_error());
586 }
587
588 #[test]
589 fn fts5_parser_stack_overflow_is_classified_as_syntax_error() {
590 let e = driver_err("fts_search", "fts5: parser stack overflow");
591 assert!(e.is_fts5_syntax_error());
592 }
593
594 #[test]
595 fn fts5_unsupported_column_query_is_classified_as_syntax_error() {
596 let e = driver_err(
597 "fts_search",
598 "fts5: column queries are not supported (detail=none)",
599 );
600 assert!(e.is_fts5_syntax_error());
601 }
602
603 #[test]
604 fn timeout_is_not_classified_as_syntax_error() {
605 let e = StorageError::Timeout {
606 operation: "fts_search".into(),
607 };
608 assert!(!e.is_fts5_syntax_error());
609 }
610
611 #[test]
612 fn pool_failure_is_not_classified_as_syntax_error() {
613 let e = StorageError::Pool {
614 operation: "fts_search".into(),
615 message: "pool exhausted".into(),
616 };
617 assert!(!e.is_fts5_syntax_error());
618 }
619
620 #[test]
621 fn driver_error_at_non_search_operation_is_not_classified_as_syntax_error() {
622 let e = driver_err("open_fts_reader", "fts5: syntax error near \"@\"");
623 assert!(!e.is_fts5_syntax_error());
624 }
625
626 #[test]
627 fn driver_error_with_unrelated_message_is_not_classified_as_syntax_error() {
628 let e = driver_err("fts_search", "disk I/O error");
629 assert!(!e.is_fts5_syntax_error());
630 }
631
632 #[test]
633 fn fts5_phrase_detail_query_is_classified_as_syntax_error() {
634 let e = driver_err(
635 "fts_search",
636 "fts5: phrase queries are not supported (detail!=full)",
637 );
638 assert!(e.is_fts5_syntax_error());
639 }
640
641 #[test]
642 fn fts5_near_detail_query_is_classified_as_syntax_error() {
643 let e = driver_err(
644 "fts_search",
645 "fts5: NEAR queries are not supported (detail!=full)",
646 );
647 assert!(e.is_fts5_syntax_error());
648 }
649
650 #[test]
651 fn unprefixed_detail_message_is_not_classified_as_syntax_error() {
652 let e = driver_err(
653 "fts_search",
654 "phrase queries are not supported (detail!=full)",
655 );
656 assert!(!e.is_fts5_syntax_error());
657 }
658
659 #[test]
660 fn fts5_shadow_table_corruption_is_not_classified_as_syntax_error() {
661 let e = driver_err(
662 "fts_search",
663 "fts5: error creating shadow table notes_content: no such table",
664 );
665 assert!(!e.is_fts5_syntax_error());
666 }
667
668 #[test]
669 fn non_text_capability_is_not_classified_as_syntax_error() {
670 let e = StorageError::Driver {
671 capability: StorageCapability::Vectors,
672 operation: "fts_search".into(),
673 source: Box::new(FakeSource("fts5: syntax error near \"@\"".into())),
674 };
675 assert!(!e.is_fts5_syntax_error());
676 }
677
678 fn driver_err_sql(operation: &'static str, message: &str) -> StorageError {
679 StorageError::driver(
680 StorageCapability::Sql,
681 operation,
682 FakeSource(message.into()),
683 )
684 }
685
686 #[test]
687 fn unique_constraint_failure_at_execute_sql_capability_is_classified() {
688 let e = driver_err_sql(
689 "execute",
690 "UNIQUE constraint failed: brain_serve_ledger.namespace, \
691 brain_serve_ledger.target_id, brain_serve_ledger.query_class, \
692 brain_serve_ledger.served_at",
693 );
694 assert!(e.is_unique_constraint_violation());
695 }
696
697 #[test]
698 fn unique_constraint_failure_at_pool_writer_execute_is_classified() {
699 let e = driver_err_sql("pool_writer.execute", "UNIQUE constraint failed: t.id");
700 assert!(e.is_unique_constraint_violation());
701 }
702
703 #[test]
704 fn unique_constraint_message_at_non_execute_operation_is_not_classified() {
705 let e = driver_err_sql("query_row", "UNIQUE constraint failed: t.id");
706 assert!(!e.is_unique_constraint_violation());
707 }
708
709 #[test]
710 fn non_unique_driver_error_at_execute_is_not_classified() {
711 let e = driver_err_sql("execute", "disk I/O error");
712 assert!(!e.is_unique_constraint_violation());
713 }
714
715 #[test]
716 fn non_sql_capability_is_not_classified_as_unique_violation() {
717 let e = driver_err("execute", "UNIQUE constraint failed: t.id");
718 assert!(!e.is_unique_constraint_violation());
719 }
720
721 #[test]
722 fn timeout_is_not_classified_as_unique_violation() {
723 let e = StorageError::Timeout {
724 operation: "execute".into(),
725 };
726 assert!(!e.is_unique_constraint_violation());
727 }
728}