use std::borrow::Cow;
use std::error::Error as StdError;
use std::fmt;
use thiserror::Error;
use crate::blob::ContentRef;
use crate::capability::StorageCapability;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum WriterTaskRequestState {
NotStarted,
TransactionRolledBack,
SideEffectsUnknown,
}
impl fmt::Display for WriterTaskRequestState {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.write_str(match self {
Self::NotStarted => "not_started",
Self::TransactionRolledBack => "transaction_rolled_back",
Self::SideEffectsUnknown => "side_effects_unknown",
})
}
}
#[derive(Debug, Error)]
pub enum StorageError {
#[error("{capability:?} resource not found: {resource} ({key})")]
NotFound {
capability: StorageCapability,
resource: &'static str,
key: String,
},
#[error("{capability:?} resource already exists: {resource} ({key})")]
AlreadyExists {
capability: StorageCapability,
resource: &'static str,
key: String,
},
#[error("conflict in {capability:?} during {operation}: {message}")]
Conflict {
capability: StorageCapability,
operation: Cow<'static, str>,
message: String,
},
#[error("invalid input for {capability:?} during {operation}: {message}")]
InvalidInput {
capability: StorageCapability,
operation: Cow<'static, str>,
message: String,
},
#[error("unsupported operation for {capability:?}: {operation} ({message})")]
Unsupported {
capability: StorageCapability,
operation: Cow<'static, str>,
message: String,
},
#[error(
"blob {content_ref} exceeds the {max_bytes}-byte read limit (observed at least {observed_at_least} bytes)"
)]
BlobTooLarge {
content_ref: ContentRef,
max_bytes: u64,
observed_at_least: u64,
},
#[error(
"blob {content_ref} metadata reports {metadata_bytes} bytes but the complete body contains {actual_bytes} bytes"
)]
BlobSizeMismatch {
content_ref: ContentRef,
metadata_bytes: u64,
actual_bytes: u64,
},
#[error("blob digest mismatch: expected {expected}, computed {actual}")]
BlobDigestMismatch {
expected: ContentRef,
actual: ContentRef,
},
#[error("pool failure during {operation}: {message}")]
Pool {
operation: Cow<'static, str>,
message: String,
},
#[error("timeout during {operation}")]
Timeout { operation: Cow<'static, str> },
#[error("admission timeout during {operation} after {timeout_ms}ms{pool}", pool = match .pool_identity {
Some(identity) => format!(" (pool: {identity})"),
None => String::new(),
})]
AdmissionTimeout {
operation: Cow<'static, str>,
timeout_ms: u64,
pool_identity: Option<String>,
},
#[error("sql transaction failure during {operation}: {message}")]
Transaction {
operation: Cow<'static, str>,
message: String,
},
#[error(
"cached read-only transaction exceeded the maximum read-transaction age \
({max_age_secs}s) during {operation} and was rolled back; retry to open a fresh \
read snapshot"
)]
ReadTransactionAgeEvicted {
operation: Cow<'static, str>,
max_age_secs: u64,
},
#[error(
"cached read-only transaction exceeded the maximum read-transaction age \
({max_age_secs}s) during {operation} but could not be cleanly rolled back \
({message}); the connection was discarded, retry to open a fresh read snapshot"
)]
ReadTransactionAgeEvictionCleanupFailed {
operation: Cow<'static, str>,
max_age_secs: u64,
message: String,
},
#[error("serialization failure in {capability:?}: {message}")]
Serialization {
capability: StorageCapability,
message: String,
},
#[error("index maintenance failure in {capability:?}: {message}")]
IndexMaintenance {
capability: StorageCapability,
message: String,
},
#[error("backend driver error in {capability:?} during {operation}: {source}")]
Driver {
capability: StorageCapability,
operation: Cow<'static, str>,
#[source]
source: Box<dyn StdError + Send + Sync>,
},
#[error("write queue full: timed out after {timeout_ms}ms waiting for writer task capacity")]
WriteQueueFull { timeout_ms: u64 },
#[error(
"writer task could not begin within {timeout_ms}ms because SQLite remained busy; request was not executed"
)]
WriterTaskBusy { timeout_ms: u64 },
#[error("writer task request failed (request_state={request_state}): {source}")]
WriterTaskRequestFailed {
request_state: WriterTaskRequestState,
#[source]
source: Box<StorageError>,
},
#[error("writer task terminated (request_state={request_state})")]
WriterTaskTerminated {
request_state: WriterTaskRequestState,
},
#[error("internal storage error: {0}")]
Internal(String),
#[error(
"KHIVE_WRITE_QUEUE=1 but no Tokio runtime context is available to spawn the writer task"
)]
WriterTaskNoRuntime,
#[error(
"refusing write on {capability:?} at {volume}: {available_bytes} bytes available, \
below the {floor_bytes}-byte floor"
)]
CapacityFloor {
capability: StorageCapability,
volume: String,
available_bytes: u64,
floor_bytes: u64,
},
}
impl StorageError {
pub fn driver(
capability: StorageCapability,
operation: impl Into<Cow<'static, str>>,
source: impl StdError + Send + Sync + 'static,
) -> Self {
Self::Driver {
capability,
operation: operation.into(),
source: Box::new(source),
}
}
pub fn capability(&self) -> Option<StorageCapability> {
match self {
Self::NotFound { capability, .. }
| Self::AlreadyExists { capability, .. }
| Self::Conflict { capability, .. }
| Self::InvalidInput { capability, .. }
| Self::Unsupported { capability, .. }
| Self::Serialization { capability, .. }
| Self::IndexMaintenance { capability, .. }
| Self::Driver { capability, .. }
| Self::CapacityFloor { capability, .. } => Some(*capability),
Self::BlobTooLarge { .. }
| Self::BlobSizeMismatch { .. }
| Self::BlobDigestMismatch { .. } => Some(StorageCapability::Blob),
Self::WriterTaskRequestFailed { source, .. } => source.capability(),
Self::Pool { .. }
| Self::Timeout { .. }
| Self::AdmissionTimeout { .. }
| Self::Transaction { .. }
| Self::ReadTransactionAgeEvicted { .. }
| Self::ReadTransactionAgeEvictionCleanupFailed { .. }
| Self::WriteQueueFull { .. }
| Self::WriterTaskBusy { .. }
| Self::WriterTaskTerminated { .. }
| Self::Internal(..)
| Self::WriterTaskNoRuntime => None,
}
}
pub fn is_retryable(&self) -> bool {
if let Self::WriterTaskRequestFailed { source, .. } = self {
return source.is_retryable();
}
matches!(
self,
Self::Pool { .. }
| Self::Timeout { .. }
| Self::AdmissionTimeout { .. }
| Self::Transaction { .. }
| Self::ReadTransactionAgeEvicted { .. }
| Self::ReadTransactionAgeEvictionCleanupFailed { .. }
| Self::WriteQueueFull { .. }
| Self::WriterTaskBusy { .. }
)
}
pub fn is_fts5_syntax_error(&self) -> bool {
if let Self::WriterTaskRequestFailed { source, .. } = self {
return source.is_fts5_syntax_error();
}
let Self::Driver {
capability,
operation,
source,
} = self
else {
return false;
};
if *capability != StorageCapability::Text || operation.as_ref() != "fts_search" {
return false;
}
let msg = source.to_string();
msg.contains("fts5: syntax error")
|| msg.contains("fts5: parser stack overflow")
|| msg.contains("fts5: column queries are not supported")
|| msg.contains("fts5: phrase queries are not supported (detail")
|| msg.contains("fts5: NEAR queries are not supported (detail")
}
pub fn is_unique_constraint_violation(&self) -> bool {
if let Self::WriterTaskRequestFailed { source, .. } = self {
return source.is_unique_constraint_violation();
}
let Self::Driver {
capability,
operation,
source,
} = self
else {
return false;
};
if *capability != StorageCapability::Sql {
return false;
}
if !matches!(
operation.as_ref(),
"execute" | "pool_writer.execute" | "tx.execute"
) {
return false;
}
source.to_string().contains("UNIQUE constraint failed")
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::fmt;
#[derive(Debug)]
struct FakeSource(String);
impl fmt::Display for FakeSource {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
write!(f, "{}", self.0)
}
}
impl StdError for FakeSource {}
fn driver_err(operation: &'static str, message: &str) -> StorageError {
StorageError::driver(
StorageCapability::Text,
operation,
FakeSource(message.into()),
)
}
#[test]
fn writer_task_request_state_display_is_stable() {
assert_eq!(
WriterTaskRequestState::NotStarted.to_string(),
"not_started"
);
assert_eq!(
WriterTaskRequestState::TransactionRolledBack.to_string(),
"transaction_rolled_back"
);
assert_eq!(
WriterTaskRequestState::SideEffectsUnknown.to_string(),
"side_effects_unknown"
);
}
#[test]
fn writer_task_busy_is_retryable_without_claiming_queue_rejection() {
let error = StorageError::WriterTaskBusy { timeout_ms: 175 };
assert!(error.is_retryable());
assert_eq!(
error.to_string(),
"writer task could not begin within 175ms because SQLite remained busy; request was not executed"
);
assert_eq!(error.capability(), None);
}
#[test]
fn writer_task_request_failure_preserves_source_policy_and_rollback_state() {
let error = StorageError::WriterTaskRequestFailed {
request_state: WriterTaskRequestState::TransactionRolledBack,
source: Box::new(StorageError::Pool {
operation: "writer_task_commit".into(),
message: "commit refused".into(),
}),
};
assert_eq!(error.capability(), None);
assert!(
error.is_retryable(),
"rollback finality must not discard the source error's retry policy"
);
assert_eq!(
error.to_string(),
"writer task request failed (request_state=transaction_rolled_back): pool failure during writer_task_commit: commit refused"
);
assert_eq!(
StdError::source(&error).map(ToString::to_string),
Some("pool failure during writer_task_commit: commit refused".to_string()),
"the original typed storage error must remain the public source"
);
}
#[test]
fn writer_task_request_failure_does_not_invent_retryability() {
let error = StorageError::WriterTaskRequestFailed {
request_state: WriterTaskRequestState::TransactionRolledBack,
source: Box::new(StorageError::InvalidInput {
capability: StorageCapability::Notes,
operation: "append_note".into(),
message: "deterministic refusal".into(),
}),
};
assert!(!error.is_retryable());
assert_eq!(error.capability(), Some(StorageCapability::Notes));
}
#[test]
fn writer_task_terminated_is_uncapability_scoped_and_not_retryable() {
for request_state in [
WriterTaskRequestState::NotStarted,
WriterTaskRequestState::TransactionRolledBack,
WriterTaskRequestState::SideEffectsUnknown,
] {
let error = StorageError::WriterTaskTerminated { request_state };
assert_eq!(error.capability(), None);
assert!(!error.is_retryable());
assert_eq!(
error.to_string(),
format!("writer task terminated (request_state={request_state})")
);
}
}
#[test]
fn blob_integrity_errors_are_blob_scoped_and_not_retryable() {
let requested = crate::blob::ContentRef::from_hex("a".repeat(64)).unwrap();
let actual = crate::blob::ContentRef::from_hex("b".repeat(64)).unwrap();
let errors = [
StorageError::BlobTooLarge {
content_ref: requested.clone(),
max_bytes: 8,
observed_at_least: 9,
},
StorageError::BlobSizeMismatch {
content_ref: requested.clone(),
metadata_bytes: 7,
actual_bytes: 8,
},
StorageError::BlobDigestMismatch {
expected: requested,
actual,
},
];
for error in errors {
assert_eq!(error.capability(), Some(StorageCapability::Blob));
assert!(!error.is_retryable());
}
}
#[test]
fn fts5_syntax_error_at_fts_search_is_classified_as_syntax_error() {
let e = driver_err("fts_search", "fts5: syntax error near \"@\"");
assert!(e.is_fts5_syntax_error());
}
#[test]
fn fts5_parser_stack_overflow_is_classified_as_syntax_error() {
let e = driver_err("fts_search", "fts5: parser stack overflow");
assert!(e.is_fts5_syntax_error());
}
#[test]
fn fts5_unsupported_column_query_is_classified_as_syntax_error() {
let e = driver_err(
"fts_search",
"fts5: column queries are not supported (detail=none)",
);
assert!(e.is_fts5_syntax_error());
}
#[test]
fn timeout_is_not_classified_as_syntax_error() {
let e = StorageError::Timeout {
operation: "fts_search".into(),
};
assert!(!e.is_fts5_syntax_error());
}
#[test]
fn pool_failure_is_not_classified_as_syntax_error() {
let e = StorageError::Pool {
operation: "fts_search".into(),
message: "pool exhausted".into(),
};
assert!(!e.is_fts5_syntax_error());
}
#[test]
fn driver_error_at_non_search_operation_is_not_classified_as_syntax_error() {
let e = driver_err("open_fts_reader", "fts5: syntax error near \"@\"");
assert!(!e.is_fts5_syntax_error());
}
#[test]
fn driver_error_with_unrelated_message_is_not_classified_as_syntax_error() {
let e = driver_err("fts_search", "disk I/O error");
assert!(!e.is_fts5_syntax_error());
}
#[test]
fn fts5_phrase_detail_query_is_classified_as_syntax_error() {
let e = driver_err(
"fts_search",
"fts5: phrase queries are not supported (detail!=full)",
);
assert!(e.is_fts5_syntax_error());
}
#[test]
fn fts5_near_detail_query_is_classified_as_syntax_error() {
let e = driver_err(
"fts_search",
"fts5: NEAR queries are not supported (detail!=full)",
);
assert!(e.is_fts5_syntax_error());
}
#[test]
fn unprefixed_detail_message_is_not_classified_as_syntax_error() {
let e = driver_err(
"fts_search",
"phrase queries are not supported (detail!=full)",
);
assert!(!e.is_fts5_syntax_error());
}
#[test]
fn fts5_shadow_table_corruption_is_not_classified_as_syntax_error() {
let e = driver_err(
"fts_search",
"fts5: error creating shadow table notes_content: no such table",
);
assert!(!e.is_fts5_syntax_error());
}
#[test]
fn non_text_capability_is_not_classified_as_syntax_error() {
let e = StorageError::Driver {
capability: StorageCapability::Vectors,
operation: "fts_search".into(),
source: Box::new(FakeSource("fts5: syntax error near \"@\"".into())),
};
assert!(!e.is_fts5_syntax_error());
}
fn driver_err_sql(operation: &'static str, message: &str) -> StorageError {
StorageError::driver(
StorageCapability::Sql,
operation,
FakeSource(message.into()),
)
}
#[test]
fn unique_constraint_failure_at_execute_sql_capability_is_classified() {
let e = driver_err_sql(
"execute",
"UNIQUE constraint failed: brain_serve_ledger.namespace, \
brain_serve_ledger.target_id, brain_serve_ledger.query_class, \
brain_serve_ledger.served_at",
);
assert!(e.is_unique_constraint_violation());
}
#[test]
fn unique_constraint_failure_at_pool_writer_execute_is_classified() {
let e = driver_err_sql("pool_writer.execute", "UNIQUE constraint failed: t.id");
assert!(e.is_unique_constraint_violation());
}
#[test]
fn unique_constraint_message_at_non_execute_operation_is_not_classified() {
let e = driver_err_sql("query_row", "UNIQUE constraint failed: t.id");
assert!(!e.is_unique_constraint_violation());
}
#[test]
fn non_unique_driver_error_at_execute_is_not_classified() {
let e = driver_err_sql("execute", "disk I/O error");
assert!(!e.is_unique_constraint_violation());
}
#[test]
fn non_sql_capability_is_not_classified_as_unique_violation() {
let e = driver_err("execute", "UNIQUE constraint failed: t.id");
assert!(!e.is_unique_constraint_violation());
}
#[test]
fn timeout_is_not_classified_as_unique_violation() {
let e = StorageError::Timeout {
operation: "execute".into(),
};
assert!(!e.is_unique_constraint_violation());
}
}