use super::*;
struct MinimalIndexStore;
impl IndexStore for MinimalIndexStore {
fn index_statuses(&self) -> StorageFuture<'_, Vec<IndexStatus>> {
Box::pin(async { Ok(Vec::new()) })
}
fn mark_refresh_complete(
&self,
kind: IndexKind,
graph_version: GraphVersion,
) -> StorageFuture<'_, IndexStatus> {
Box::pin(async move {
Ok(IndexStatus {
kind,
index_version: 1,
indexed_graph_version: graph_version,
state: crate::domain::IndexState::Fresh,
last_error: None,
})
})
}
}
#[test]
fn storage_errors_preserve_boundary_messages() {
let io = StorageError::from(std::io::Error::new(
std::io::ErrorKind::PermissionDenied,
"readonly",
));
let sqlite = StorageError::from(rusqlite::Error::InvalidQuery);
assert!(io.to_string().contains("storage I/O failed: readonly"));
assert_eq!(
sqlite.to_string(),
"sqlite operation failed: Query is not read-only"
);
assert_eq!(
StorageError::LockPoisoned.to_string(),
"sqlite connection lock was poisoned"
);
assert_eq!(
StorageError::InvalidInput("missing graph version".to_owned()).to_string(),
"invalid storage input: missing graph version"
);
}
#[tokio::test]
async fn join_errors_map_to_storage_worker_failures() {
let join_error = tokio::spawn(async { panic!("storage worker panic") })
.await
.expect_err("worker should panic");
let error = StorageError::from(join_error);
assert!(error.to_string().contains("storage worker failed"));
}
#[test]
fn index_refresh_task_states_have_stable_storage_values() {
assert_eq!(IndexRefreshTaskState::Queued.as_str(), "queued");
assert_eq!(IndexRefreshTaskState::Running.as_str(), "running");
assert_eq!(IndexRefreshTaskState::Succeeded.as_str(), "succeeded");
assert_eq!(IndexRefreshTaskState::Retrying.as_str(), "retrying");
assert_eq!(IndexRefreshTaskState::Failed.as_str(), "failed");
assert_eq!(IndexRefreshTaskState::DeadLetter.as_str(), "dead_letter");
}
#[tokio::test]
async fn default_index_refresh_queue_methods_report_unavailable_storage() {
let store = MinimalIndexStore;
let cursors = store
.index_cursors()
.await
.expect_err("default cursor storage should be unavailable");
let queued = store
.queue_index_refreshes(IndexRefreshQueueRequest {
kinds: vec![IndexKind::Bm25],
target_graph_version: GraphVersion::new(1),
max_queue_depth: 1,
reset_dead_letter_tasks: false,
now_ms: 10,
})
.await
.expect_err("default task queue should be unavailable");
let claimed = store
.claim_index_refresh_task(IndexRefreshClaimRequest {
lease_owner: "worker".to_owned(),
lease_duration_ms: 100,
max_attempts: 3,
now_ms: 10,
})
.await
.expect_err("default claim should be unavailable");
let completed = store
.complete_index_refresh_task(IndexRefreshCompletion {
task_id: "task".to_owned(),
lease_owner: "worker".to_owned(),
attempt_count: 1,
indexed_graph_version: GraphVersion::new(1),
model_name: None,
model_dimension: None,
now_ms: 20,
})
.await
.expect_err("default completion should be unavailable");
let failed = store
.fail_index_refresh_task(IndexRefreshFailure {
task_id: "task".to_owned(),
lease_owner: "worker".to_owned(),
attempt_count: 1,
error_kind: "indexer".to_owned(),
error_message: "worker failed".to_owned(),
retry_backoff_ms: 100,
max_attempts: 2,
now_ms: 20,
})
.await
.expect_err("default failure handling should be unavailable");
let diagnostics = store
.index_refresh_diagnostics(30)
.await
.expect_err("default diagnostics should be unavailable");
assert!(cursors.to_string().contains("index cursor storage"));
for error in [queued, claimed, completed, failed] {
assert!(
error
.to_string()
.contains("index refresh task storage is unavailable")
);
}
assert!(
diagnostics
.to_string()
.contains("index refresh diagnostics are unavailable")
);
}
#[tokio::test]
async fn default_operational_methods_are_bounded_and_explicit() {
let store = MinimalIndexStore;
let tasks = store
.queue_worker_tasks(vec![WorkerTaskSeed {
kind: WorkerKind::Extractor,
source_scope: "docs".to_owned(),
evidence_id: Some("ev-1".to_owned()),
target_graph_version: GraphVersion::new(1),
input_fingerprint: "extractor:ev-1:1".to_owned(),
payload_json: "{}".to_owned(),
now_ms: 1,
}])
.await
.expect("default queue is a no-op");
let statuses = store
.worker_statuses()
.await
.expect("default status is empty");
let claimed = store
.claim_worker_task(WorkerTaskClaimRequest {
kind: None,
lease_owner: "worker".to_owned(),
lease_duration_ms: 10,
max_attempts: 1,
now_ms: 1,
})
.await
.expect("default claim is empty");
let proposals = store
.list_proposals(ProposalListRequest {
state: None,
limit: 10,
})
.await
.expect("default proposal list is empty");
let conflicts = store
.proposal_conflicts("proposal".to_owned())
.await
.expect("default conflicts are empty");
let audit = store
.query_audit_events(AuditQueryRequest {
operation: None,
limit: 10,
})
.await
.expect("default audit query is empty");
let audit_count = store
.audit_event_count()
.await
.expect("default audit count is zero");
let operator = store
.service_operator_status()
.await
.expect("default operator is disabled");
assert!(tasks.is_empty());
assert!(statuses.is_empty());
assert!(claimed.is_none());
assert!(proposals.is_empty());
assert!(conflicts.is_empty());
assert!(audit.is_empty());
assert_eq!(audit_count, 0);
assert_eq!(operator.state, ServiceOperatorState::Disabled);
for error in [
store
.complete_worker_task(WorkerTaskCompletion {
task_id: "task".to_owned(),
lease_owner: "worker".to_owned(),
attempt_count: 1,
now_ms: 2,
})
.await
.expect_err("completion should require storage"),
store
.fail_worker_task(WorkerTaskFailure {
task_id: "task".to_owned(),
lease_owner: "worker".to_owned(),
attempt_count: 1,
error_kind: "worker".to_owned(),
error_message: "failed".to_owned(),
retry_backoff_ms: 10,
max_attempts: 1,
now_ms: 2,
})
.await
.expect_err("failure should require storage"),
store
.insert_proposal(NewProposal {
proposal_id: "proposal".to_owned(),
source_scope: "docs".to_owned(),
kind: ProposalKind::Evidence,
title: "title".to_owned(),
summary: "summary".to_owned(),
payload_json: "{}".to_owned(),
origin: "test".to_owned(),
provenance: ProposalProvenance::new("test"),
confidence_basis_points: 1,
conflicts: Vec::new(),
now_ms: 1,
})
.await
.expect_err("proposal insert should require storage"),
store
.decide_proposal(ProposalDecision {
proposal_id: "proposal".to_owned(),
next_state: ProposalState::Rejected,
actor: "tester".to_owned(),
reason: None,
now_ms: 2,
})
.await
.expect_err("proposal decision should require storage"),
store
.insert_audit_event(NewAuditEvent {
operation: "test".to_owned(),
interface: "cli".to_owned(),
request_id: "req".to_owned(),
trace_id: "trace".to_owned(),
status: AuditStatus::Completed,
actor: None,
source_scope: None,
graph_version: 0,
detail_json: "{}".to_owned(),
message: None,
now_ms: 1,
})
.await
.expect_err("audit insert should require storage"),
store
.update_service_operator(ServiceOperatorUpdate {
state: ServiceOperatorState::Enabled,
silent_updates_enabled: true,
allowed_scopes: vec!["docs".to_owned()],
last_error: None,
now_ms: 2,
})
.await
.expect_err("operator update should require storage"),
] {
assert!(error.to_string().contains("storage is unavailable"));
}
}