use std::time::Duration;
use async_trait::async_trait;
use chrono::{DateTime, Utc};
use sqlx::Row;
use uuid::Uuid;
use super::dialect::apply_projection_summary_bucket;
use super::sqlite::SqliteCanonicalStore;
use super::system_store::{
DeadLetterGroup, PendingTaskMetric, ProjectionClaimFilter, ProjectionOperation,
ProjectionTaskInsert, ProjectionTaskRow, ProjectionTaskStatus, ProjectionTaskStore,
ProjectionTaskSummary, SystemStoreError, SystemStoreResult,
};
const TABLE: &str = "udb_projection_tasks";
impl SqliteCanonicalStore {
pub(crate) fn pool_ref(&self) -> &sqlx::SqlitePool {
&self.pool
}
}
fn parse_iso(s: &str) -> DateTime<Utc> {
DateTime::parse_from_rfc3339(s)
.map(|dt| dt.with_timezone(&Utc))
.unwrap_or_else(|_| Utc::now())
}
const REQUIRED_SQLITE_VERSION: &str = "3.35";
fn projection_retry_delay_secs(retry_count: i32) -> i64 {
let attempt = retry_count.max(1).min(12) as u32;
(1_i64 << (attempt - 1)).min(3600)
}
fn row_to_projection_task(row: sqlx::sqlite::SqliteRow) -> SystemStoreResult<ProjectionTaskRow> {
let task_id_str: String = row
.try_get("task_id")
.map_err(|e| SystemStoreError::query("sqlite", "SELECT task_id", e))?;
let task_id = Uuid::parse_str(&task_id_str).map_err(|e| {
SystemStoreError::InvalidInput(format!("task_id '{task_id_str}' is not a valid UUID: {e}"))
})?;
let operation_str: String = row
.try_get("operation")
.map_err(|e| SystemStoreError::query("sqlite", "SELECT operation", e))?;
let operation = ProjectionOperation::parse(&operation_str).ok_or_else(|| {
SystemStoreError::InvalidInput(format!(
"unknown projection operation '{operation_str}' in SQLite row"
))
})?;
let status_str: String = row
.try_get("status")
.map_err(|e| SystemStoreError::query("sqlite", "SELECT status", e))?;
let status = ProjectionTaskStatus::parse(&status_str).ok_or_else(|| {
SystemStoreError::InvalidInput(format!(
"unknown projection status '{status_str}' in SQLite row"
))
})?;
let source_row_key_text: String = row.try_get("source_row_key").unwrap_or_default();
let target_options_text: String = row.try_get("target_options").unwrap_or_default();
let source_payload_text: String = row.try_get("source_payload").unwrap_or_default();
let parse_json = |s: &str, field: &'static str| -> SystemStoreResult<serde_json::Value> {
if s.is_empty() {
return Ok(serde_json::Value::Null);
}
serde_json::from_str(s).map_err(|e| {
SystemStoreError::InvalidInput(format!(
"field '{field}' is not valid JSON: {e} (raw: '{s}')"
))
})
};
Ok(ProjectionTaskRow {
task_id,
idempotency_key: row.try_get("idempotency_key").unwrap_or_default(),
project_id: row.try_get("project_id").unwrap_or_default(),
target_backend: row.try_get("target_backend").unwrap_or_default(),
target_instance: row.try_get("target_instance").unwrap_or_default(),
projection_kind: row.try_get("projection_kind").unwrap_or_default(),
resource_name: row.try_get("resource_name").unwrap_or_default(),
operation,
source_row_key: parse_json(&source_row_key_text, "source_row_key")?,
target_options: parse_json(&target_options_text, "target_options")?,
source_payload: parse_json(&source_payload_text, "source_payload")?,
source_checksum: row.try_get("source_checksum").unwrap_or_default(),
status,
retry_count: row.try_get("retry_count").unwrap_or(0),
last_error: row.try_get("last_error").unwrap_or_default(),
created_at: row
.try_get::<String, _>("created_at")
.map(|s| parse_iso(&s))
.unwrap_or_else(|_| Utc::now()),
updated_at: row
.try_get::<String, _>("updated_at")
.map(|s| parse_iso(&s))
.unwrap_or_else(|_| Utc::now()),
next_retry_at: row
.try_get::<Option<String>, _>("next_retry_at")
.ok()
.flatten()
.map(|s| parse_iso(&s)),
completed_at: row
.try_get::<Option<String>, _>("completed_at")
.ok()
.flatten()
.map(|s| parse_iso(&s)),
})
}
#[async_trait]
impl ProjectionTaskStore for SqliteCanonicalStore {
fn backend_label(&self) -> &'static str {
"sqlite"
}
async fn ensure_projection_tables(&self) -> SystemStoreResult<()> {
let stmts = super::sql_schema::sqlite_projection_tasks_ddl(TABLE);
for sql in stmts.iter() {
if let Err(e) = sqlx::query(sql).execute(self.pool_ref()).await {
let msg = e.to_string();
if !msg.contains("duplicate column name") {
return Err(SystemStoreError::query("sqlite", sql.clone(), e));
}
}
}
Ok(())
}
async fn enqueue_projection_task(
&self,
task: &ProjectionTaskInsert,
) -> SystemStoreResult<Uuid> {
let task_id = Uuid::new_v4();
let task_id_text = task_id.to_string();
let source_row_key = serde_json::to_string(&task.source_row_key)
.map_err(|e| SystemStoreError::InvalidInput(format!("source_row_key: {e}")))?;
let target_options = serde_json::to_string(&task.target_options)
.map_err(|e| SystemStoreError::InvalidInput(format!("target_options: {e}")))?;
let source_payload = serde_json::to_string(&task.source_payload)
.map_err(|e| SystemStoreError::InvalidInput(format!("source_payload: {e}")))?;
let sql = format!(
"INSERT INTO {TABLE} (
task_id, idempotency_key, project_id, manifest_checksum, message_type,
source_schema, source_table, source_row_key, operation,
target_backend, target_instance, projection_kind, resource_name,
target_options, source_payload, source_checksum
) VALUES (
?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, ?13, ?14, ?15, ?16
)
ON CONFLICT(idempotency_key) DO NOTHING"
);
sqlx::query(&sql)
.bind(&task_id_text)
.bind(&task.idempotency_key)
.bind(&task.project_id)
.bind(&task.manifest_checksum)
.bind(&task.message_type)
.bind(&task.source_schema)
.bind(&task.source_table)
.bind(&source_row_key)
.bind(task.operation.as_str())
.bind(&task.target_backend)
.bind(&task.target_instance)
.bind(&task.projection_kind)
.bind(&task.resource_name)
.bind(&target_options)
.bind(&source_payload)
.bind(&task.source_checksum)
.execute(self.pool_ref())
.await
.map_err(|e| SystemStoreError::query("sqlite", sql.clone(), e))?;
let lookup_sql = format!("SELECT task_id FROM {TABLE} WHERE idempotency_key = ?1");
let existing: String = sqlx::query_scalar(&lookup_sql)
.bind(&task.idempotency_key)
.fetch_one(self.pool_ref())
.await
.map_err(|e| SystemStoreError::query("sqlite", lookup_sql.clone(), e))?;
Uuid::parse_str(&existing).map_err(|e| {
SystemStoreError::InvalidInput(format!(
"stored task_id '{existing}' is not a valid UUID: {e}"
))
})
}
async fn claim_projection_tasks(
&self,
filter: &ProjectionClaimFilter,
) -> SystemStoreResult<Vec<ProjectionTaskRow>> {
if filter.batch_size <= 0 {
return Ok(Vec::new());
}
let target_clause = match (&filter.target_backend, &filter.target_instance) {
(Some(b), Some(i)) => format!(
"AND target_backend = '{}' AND target_instance = '{}'",
escape_sql_literal(b),
escape_sql_literal(i)
),
(Some(b), None) => {
format!("AND target_backend = '{}'", escape_sql_literal(b))
}
(None, Some(i)) => {
format!("AND target_instance = '{}'", escape_sql_literal(i))
}
(None, None) => String::new(),
};
let sql = format!(
"UPDATE {TABLE}
SET status = 'IN_PROGRESS',
updated_at = strftime('%Y-%m-%dT%H:%M:%fZ', 'now')
WHERE task_id IN (
SELECT task_id FROM {TABLE}
WHERE status IN ('PENDING', 'FAILED')
AND retry_count < ?1
AND (?3 = '' OR project_id = ?3)
AND (next_retry_at IS NULL OR next_retry_at <= strftime('%Y-%m-%dT%H:%M:%fZ', 'now'))
{target_clause}
ORDER BY created_at
LIMIT ?2
)
RETURNING task_id, idempotency_key, project_id,
target_backend, target_instance, projection_kind, resource_name,
operation, source_row_key, target_options, source_payload,
source_checksum, status, retry_count, last_error,
created_at, updated_at, next_retry_at, completed_at"
);
let rows = sqlx::query(&sql)
.bind(filter.max_retries)
.bind(filter.batch_size)
.bind(filter.project_id.as_deref().unwrap_or(""))
.fetch_all(self.pool_ref())
.await
.map_err(|e| SystemStoreError::query("sqlite", sql.clone(), e))?;
let mut out = Vec::with_capacity(rows.len());
for row in rows {
out.push(row_to_projection_task(row)?);
}
Ok(out)
}
async fn mark_projection_task_completed(&self, task_id: Uuid) -> SystemStoreResult<()> {
let sql = format!(
"UPDATE {TABLE}
SET status = 'COMPLETED',
completed_at = strftime('%Y-%m-%dT%H:%M:%fZ', 'now'),
next_retry_at = NULL,
updated_at = strftime('%Y-%m-%dT%H:%M:%fZ', 'now')
WHERE task_id = ?1"
);
let task_id_text = task_id.to_string();
sqlx::query(&sql)
.bind(&task_id_text)
.execute(self.pool_ref())
.await
.map_err(|e| SystemStoreError::query("sqlite", sql.clone(), e))?;
Ok(())
}
async fn mark_projection_task_failed(
&self,
task_id: Uuid,
new_retry_count: i32,
new_status: ProjectionTaskStatus,
error: &str,
) -> SystemStoreResult<()> {
if !matches!(
new_status,
ProjectionTaskStatus::Failed | ProjectionTaskStatus::DeadLetter
) {
return Err(SystemStoreError::InvalidInput(format!(
"mark_projection_task_failed only accepts FAILED or DEAD_LETTER, got {}",
new_status.as_str()
)));
}
let sql = format!(
"UPDATE {TABLE}
SET status = ?,
retry_count = ?,
last_error = ?,
next_retry_at = NULL,
updated_at = strftime('%Y-%m-%dT%H:%M:%fZ', 'now')
WHERE task_id = ?"
);
let task_id_text = task_id.to_string();
sqlx::query(&sql)
.bind(new_status.as_str())
.bind(new_retry_count)
.bind(error)
.bind(&task_id_text)
.execute(self.pool_ref())
.await
.map_err(|e| SystemStoreError::query("sqlite", sql.clone(), e))?;
Ok(())
}
async fn requeue_dead_letter_tasks(
&self,
target_backend: Option<&str>,
) -> SystemStoreResult<i64> {
let where_clause = match target_backend {
Some(b) => format!("AND target_backend = '{}'", escape_sql_literal(b)),
None => String::new(),
};
let sql = format!(
"UPDATE {TABLE}
SET status = 'PENDING',
retry_count = 0,
last_error = '',
next_retry_at = NULL,
updated_at = strftime('%Y-%m-%dT%H:%M:%fZ', 'now')
WHERE status = 'DEAD_LETTER' {where_clause}"
);
let result = sqlx::query(&sql)
.execute(self.pool_ref())
.await
.map_err(|e| SystemStoreError::query("sqlite", sql.clone(), e))?;
Ok(result.rows_affected() as i64)
}
async fn reset_stale_in_progress_tasks(&self, stale_after: Duration) -> SystemStoreResult<i64> {
let cutoff = Utc::now() - chrono::Duration::seconds(stale_after.as_secs() as i64);
let cutoff_str = cutoff.to_rfc3339();
let sql = format!(
"UPDATE {TABLE}
SET status = 'PENDING',
last_error = 'stale in-progress reconciliation',
updated_at = strftime('%Y-%m-%dT%H:%M:%fZ', 'now')
WHERE status = 'IN_PROGRESS' AND updated_at < ?1"
);
let result = sqlx::query(&sql)
.bind(&cutoff_str)
.execute(self.pool_ref())
.await
.map_err(|e| SystemStoreError::query("sqlite", sql.clone(), e))?;
Ok(result.rows_affected() as i64)
}
async fn pending_task_metrics(&self, limit: i64) -> SystemStoreResult<Vec<PendingTaskMetric>> {
let sql = format!(
"SELECT project_id, target_backend, target_instance, projection_kind,
COUNT(*) AS pending,
CAST(
(julianday('now') -
julianday(MIN(created_at))
) * 86400.0
AS REAL) AS oldest_age_seconds
FROM {TABLE}
WHERE status IN ('PENDING', 'FAILED')
GROUP BY project_id, target_backend, target_instance, projection_kind
LIMIT ?"
);
let rows = sqlx::query(&sql)
.bind(limit.max(1))
.fetch_all(self.pool_ref())
.await
.map_err(|e| SystemStoreError::query("sqlite", sql.clone(), e))?;
let mut out = Vec::with_capacity(rows.len());
for row in rows {
out.push(PendingTaskMetric {
project_id: row.try_get("project_id").unwrap_or_default(),
target_backend: row.try_get("target_backend").unwrap_or_default(),
target_instance: row.try_get("target_instance").unwrap_or_default(),
projection_kind: row.try_get("projection_kind").unwrap_or_default(),
pending: row.try_get("pending").unwrap_or(0),
oldest_age_seconds: row.try_get("oldest_age_seconds").unwrap_or(0.0),
});
}
Ok(out)
}
async fn dead_letter_groups(&self, limit: i64) -> SystemStoreResult<Vec<DeadLetterGroup>> {
let sql = format!(
"SELECT source_table, target_backend, target_instance,
COUNT(*) AS dead_count
FROM {TABLE}
WHERE status = 'DEAD_LETTER'
GROUP BY source_table, target_backend, target_instance
LIMIT ?"
);
let rows = sqlx::query(&sql)
.bind(limit.max(1))
.fetch_all(self.pool_ref())
.await
.map_err(|e| SystemStoreError::query("sqlite", sql.clone(), e))?;
let mut out = Vec::with_capacity(rows.len());
for row in rows {
out.push(DeadLetterGroup {
source_table: row.try_get("source_table").unwrap_or_default(),
target_backend: row.try_get("target_backend").unwrap_or_default(),
target_instance: row.try_get("target_instance").unwrap_or_default(),
dead_count: row.try_get("dead_count").unwrap_or(0),
});
}
Ok(out)
}
async fn requeue_dead_letter_by_source(
&self,
source_table: &str,
target_backend: &str,
target_instance: &str,
) -> SystemStoreResult<i64> {
let sql = format!(
"UPDATE {TABLE}
SET status = 'PENDING',
retry_count = 0,
last_error = 'reconciliation repair',
updated_at = strftime('%Y-%m-%dT%H:%M:%fZ', 'now')
WHERE status = 'DEAD_LETTER'
AND source_table = ?
AND target_backend = ?
AND target_instance = ?"
);
let result = sqlx::query(&sql)
.bind(source_table)
.bind(target_backend)
.bind(target_instance)
.execute(self.pool_ref())
.await
.map_err(|e| SystemStoreError::query("sqlite", sql.clone(), e))?;
Ok(result.rows_affected() as i64)
}
async fn pending_projection_task_count(
&self,
idempotency_keys: &[String],
) -> SystemStoreResult<i64> {
if idempotency_keys.is_empty() {
return Ok(0);
}
let placeholders: Vec<&str> = idempotency_keys.iter().map(|_| "?").collect();
let sql = format!(
"SELECT COUNT(*) FROM {TABLE}
WHERE idempotency_key IN ({})
AND status NOT IN ('COMPLETED','DEAD_LETTER','FAILED')",
placeholders.join(",")
);
let mut q = sqlx::query_scalar::<_, i64>(&sql);
for k in idempotency_keys {
q = q.bind(k.clone());
}
let n: i64 = q
.fetch_one(self.pool_ref())
.await
.map_err(|e| SystemStoreError::query("sqlite", sql.clone(), e))?;
Ok(n)
}
async fn projection_task_summary(&self) -> SystemStoreResult<ProjectionTaskSummary> {
let sql = format!("SELECT status, COUNT(*) AS n FROM {TABLE} GROUP BY status");
let rows = sqlx::query(&sql)
.fetch_all(self.pool_ref())
.await
.map_err(|e| SystemStoreError::query("sqlite", sql.clone(), e))?;
let mut s = ProjectionTaskSummary::default();
for row in rows {
let status: String = row.try_get("status").unwrap_or_default();
let n: i64 = row.try_get("n").unwrap_or(0);
apply_projection_summary_bucket(&mut s, "sqlite", &status, n)?;
}
Ok(s)
}
}
fn escape_sql_literal(s: &str) -> String {
s.replace('\'', "''")
}
#[cfg(test)]
mod tests {
use super::*;
use crate::runtime::canonical_store::sqlite::SqliteCanonicalStore;
use sqlx::sqlite::SqlitePoolOptions;
async fn fresh_store() -> SqliteCanonicalStore {
let pool = SqlitePoolOptions::new()
.max_connections(1)
.connect("sqlite::memory:")
.await
.expect("in-memory sqlite");
let store = SqliteCanonicalStore::new(pool, "test", "udb_outbox_events");
ProjectionTaskStore::ensure_projection_tables(&store)
.await
.expect("ensure_projection_tables");
store
}
fn sample_insert(idempotency_key: &str) -> ProjectionTaskInsert {
ProjectionTaskInsert {
idempotency_key: idempotency_key.to_string(),
project_id: "default".to_string(),
manifest_checksum: "sha256:abc".to_string(),
message_type: "udb.test.v1.User".to_string(),
source_schema: "public".to_string(),
source_table: "users".to_string(),
source_row_key: serde_json::json!({"id": "u-1"}),
operation: ProjectionOperation::Upsert,
target_backend: "qdrant".to_string(),
target_instance: "primary".to_string(),
projection_kind: "vector".to_string(),
resource_name: "users_vec".to_string(),
target_options: serde_json::json!([]),
source_payload: serde_json::json!({"name": "Ada"}),
source_checksum: "sha256:row1".to_string(),
}
}
#[tokio::test]
async fn enqueue_summary_claim_roundtrip() {
let store = fresh_store().await;
let id1 = store
.enqueue_projection_task(&sample_insert("k1"))
.await
.expect("enqueue");
assert_ne!(id1, Uuid::nil());
let summary = store
.projection_task_summary()
.await
.expect("summary after one enqueue");
assert_eq!(summary.pending, 1);
assert_eq!(summary.total(), 1);
assert_eq!(summary.claimable(), 1);
let claimed = store
.claim_projection_tasks(&ProjectionClaimFilter::default())
.await
.expect("claim");
assert_eq!(claimed.len(), 1);
assert_eq!(claimed[0].task_id, id1);
assert_eq!(claimed[0].status, ProjectionTaskStatus::InProgress);
assert_eq!(claimed[0].source_row_key, serde_json::json!({"id": "u-1"}));
let summary = store.projection_task_summary().await.unwrap();
assert_eq!(summary.in_progress, 1);
assert_eq!(summary.pending, 0);
}
#[tokio::test]
async fn enqueue_is_idempotent_on_idempotency_key() {
let store = fresh_store().await;
let id_a = store
.enqueue_projection_task(&sample_insert("dup"))
.await
.unwrap();
let id_b = store
.enqueue_projection_task(&sample_insert("dup"))
.await
.unwrap();
assert_eq!(
id_a, id_b,
"second enqueue with same idempotency_key must return existing task_id"
);
let summary = store.projection_task_summary().await.unwrap();
assert_eq!(summary.pending, 1, "no duplicate row created");
}
#[tokio::test]
async fn mark_failed_then_reclaim_increments_retries() {
let store = fresh_store().await;
let id = store
.enqueue_projection_task(&sample_insert("retry-me"))
.await
.unwrap();
let first = store
.claim_projection_tasks(&ProjectionClaimFilter::default())
.await
.unwrap();
assert_eq!(first.len(), 1);
assert_eq!(first[0].retry_count, 0);
store
.mark_projection_task_failed(
id,
1,
ProjectionTaskStatus::Failed,
"synthetic failure on attempt 1",
)
.await
.unwrap();
let summary = store.projection_task_summary().await.unwrap();
assert_eq!(summary.failed, 1);
assert_eq!(summary.in_progress, 0);
let second = store
.claim_projection_tasks(&ProjectionClaimFilter::default())
.await
.unwrap();
assert_eq!(second.len(), 1);
assert_eq!(second[0].task_id, id);
assert_eq!(second[0].retry_count, 1);
assert_eq!(second[0].last_error, "synthetic failure on attempt 1");
assert_eq!(second[0].status, ProjectionTaskStatus::InProgress);
}
#[tokio::test]
async fn claim_respects_max_retries_threshold() {
let store = fresh_store().await;
let id = store
.enqueue_projection_task(&sample_insert("hot"))
.await
.unwrap();
store
.claim_projection_tasks(&ProjectionClaimFilter::default())
.await
.unwrap();
store
.mark_projection_task_failed(id, 3, ProjectionTaskStatus::Failed, "errors keep coming")
.await
.unwrap();
let claimed = store
.claim_projection_tasks(&ProjectionClaimFilter {
batch_size: 50,
max_retries: 3,
..ProjectionClaimFilter::default()
})
.await
.unwrap();
assert!(
claimed.is_empty(),
"row at retry_count = max_retries must not be claimed"
);
let claimed = store
.claim_projection_tasks(&ProjectionClaimFilter {
batch_size: 50,
max_retries: 4,
..ProjectionClaimFilter::default()
})
.await
.unwrap();
assert_eq!(claimed.len(), 1);
}
#[tokio::test]
async fn completed_tasks_are_not_reclaimed() {
let store = fresh_store().await;
let id = store
.enqueue_projection_task(&sample_insert("done"))
.await
.unwrap();
store
.claim_projection_tasks(&ProjectionClaimFilter::default())
.await
.unwrap();
store.mark_projection_task_completed(id).await.unwrap();
let again = store
.claim_projection_tasks(&ProjectionClaimFilter::default())
.await
.unwrap();
assert!(again.is_empty());
let summary = store.projection_task_summary().await.unwrap();
assert_eq!(summary.completed, 1);
assert_eq!(summary.claimable(), 0);
}
#[tokio::test]
async fn claim_filter_isolates_by_target_backend() {
let store = fresh_store().await;
let mut a = sample_insert("a");
a.target_backend = "qdrant".to_string();
let mut b = sample_insert("b");
b.target_backend = "s3".to_string();
store.enqueue_projection_task(&a).await.unwrap();
store.enqueue_projection_task(&b).await.unwrap();
let only_qdrant = store
.claim_projection_tasks(&ProjectionClaimFilter {
target_backend: Some("qdrant".to_string()),
..ProjectionClaimFilter::default()
})
.await
.unwrap();
assert_eq!(only_qdrant.len(), 1);
assert_eq!(only_qdrant[0].idempotency_key, "a");
}
#[tokio::test]
async fn requeue_dead_letter_resets_state() {
let store = fresh_store().await;
let id = store
.enqueue_projection_task(&sample_insert("dlq"))
.await
.unwrap();
store
.claim_projection_tasks(&ProjectionClaimFilter::default())
.await
.unwrap();
store
.mark_projection_task_failed(
id,
5,
ProjectionTaskStatus::DeadLetter,
"fatal upstream error",
)
.await
.unwrap();
let summary = store.projection_task_summary().await.unwrap();
assert_eq!(summary.dead_letter, 1);
let requeued = store.requeue_dead_letter_tasks(None).await.unwrap();
assert_eq!(requeued, 1);
let summary = store.projection_task_summary().await.unwrap();
assert_eq!(summary.dead_letter, 0);
assert_eq!(summary.pending, 1);
let claimed = store
.claim_projection_tasks(&ProjectionClaimFilter::default())
.await
.unwrap();
assert_eq!(claimed.len(), 1);
assert_eq!(claimed[0].retry_count, 0);
assert_eq!(claimed[0].last_error, "");
}
#[tokio::test]
async fn reset_stale_in_progress_recovers_crashed_worker() {
let store = fresh_store().await;
store
.enqueue_projection_task(&sample_insert("crash"))
.await
.unwrap();
store
.claim_projection_tasks(&ProjectionClaimFilter::default())
.await
.unwrap();
let summary = store.projection_task_summary().await.unwrap();
assert_eq!(summary.in_progress, 1);
tokio::time::sleep(Duration::from_millis(50)).await;
let reset = store
.reset_stale_in_progress_tasks(Duration::from_secs(0))
.await
.unwrap();
assert_eq!(reset, 1);
let summary = store.projection_task_summary().await.unwrap();
assert_eq!(summary.pending, 1);
assert_eq!(summary.in_progress, 0);
}
#[tokio::test]
async fn pending_count_excludes_terminal_and_failed() {
let store = fresh_store().await;
for k in ["a", "b", "c", "d"] {
store
.enqueue_projection_task(&sample_insert(k))
.await
.unwrap();
}
let claim = store
.claim_projection_tasks(&ProjectionClaimFilter::default())
.await
.unwrap();
for row in &claim {
match row.idempotency_key.as_str() {
"a" => store
.mark_projection_task_completed(row.task_id)
.await
.unwrap(),
"b" => store
.mark_projection_task_failed(row.task_id, 1, ProjectionTaskStatus::Failed, "x")
.await
.unwrap(),
"c" => store
.mark_projection_task_failed(
row.task_id,
5,
ProjectionTaskStatus::DeadLetter,
"y",
)
.await
.unwrap(),
_ => {}
}
}
let n = store
.pending_projection_task_count(&[
"a".to_string(),
"b".to_string(),
"c".to_string(),
"d".to_string(),
])
.await
.unwrap();
assert_eq!(n, 1, "only `d` (IN_PROGRESS) should count");
let n = store.pending_projection_task_count(&[]).await.unwrap();
assert_eq!(n, 0);
let n = store
.pending_projection_task_count(&["nope".to_string()])
.await
.unwrap();
assert_eq!(n, 0);
}
#[tokio::test]
async fn mark_failed_rejects_invalid_status() {
let store = fresh_store().await;
let id = store
.enqueue_projection_task(&sample_insert("bad-mark"))
.await
.unwrap();
let err = store
.mark_projection_task_failed(id, 1, ProjectionTaskStatus::Completed, "oops")
.await
.expect_err("must reject COMPLETED on the failure path");
match err {
SystemStoreError::InvalidInput(msg) => {
assert!(msg.contains("FAILED or DEAD_LETTER"));
}
other => panic!("expected InvalidInput, got: {other}"),
}
}
}