use super::test_fixtures::{insert_pending, insert_pending_op, open_temp_queue};
use super::*;
#[test]
fn dequeue_does_not_claim_pair_keys_under_memory_bindings() {
let (conn, path) = open_temp_queue();
let _pair = insert_pending_op(&conn, "pair:21560:159670", "entity_pair", "EntityConnect");
assert!(matches!(
dequeue_next_pending(&conn, "MemoryBindings", "", "").unwrap(),
DequeueOutcome::Empty
));
match dequeue_next_pending(&conn, "EntityConnect", "", "").unwrap() {
DequeueOutcome::Claimed(row) => {
assert_eq!(row.item_key, "pair:21560:159670");
}
DequeueOutcome::Empty => panic!("expected EC claim"),
}
let _ = std::fs::remove_file(&path);
}
#[test]
fn dequeue_filters_by_operation_isolation() {
let (conn, path) = open_temp_queue();
let _ed_id = insert_pending_op(&conn, "aeroporto-guarulhos", "entity", "EntityDescriptions");
let mb_id = insert_pending_op(&conn, "yt-memory", "memory", "MemoryBindings");
match dequeue_next_pending(&conn, "MemoryBindings", "", "").unwrap() {
DequeueOutcome::Claimed(row) => {
assert_eq!(row.id, mb_id);
assert_eq!(row.item_key, "yt-memory");
assert_eq!(row.operation, "MemoryBindings");
}
DequeueOutcome::Empty => panic!("expected MB claim"),
}
assert!(matches!(
dequeue_next_pending(&conn, "MemoryBindings", "", "").unwrap(),
DequeueOutcome::Empty
));
match dequeue_next_pending(&conn, "EntityDescriptions", "", "").unwrap() {
DequeueOutcome::Claimed(row) => {
assert_eq!(row.item_key, "aeroporto-guarulhos");
assert_eq!(row.operation, "EntityDescriptions");
}
DequeueOutcome::Empty => panic!("expected ED claim"),
}
let _ = std::fs::remove_file(&path);
}
#[test]
fn dequeue_next_pending_distinguishes_empty_from_claimed() {
let (conn, path) = open_temp_queue();
let id = insert_pending(&conn, "mem-dequeue");
let claimed =
dequeue_next_pending(&conn, "MemoryBindings", "", "").expect("dequeue must succeed");
match claimed {
DequeueOutcome::Claimed(row) => {
assert_eq!(row.id, id);
assert_eq!(row.item_key, "mem-dequeue");
assert_eq!(row.operation, "MemoryBindings");
}
DequeueOutcome::Empty => panic!("expected a claimed row"),
}
let empty =
dequeue_next_pending(&conn, "MemoryBindings", "", "").expect("dequeue must succeed");
assert!(matches!(empty, DequeueOutcome::Empty));
let _ = std::fs::remove_file(&path);
}
#[test]
fn dequeue_next_pending_isolates_by_namespace() {
let (conn, path) = open_temp_queue();
conn.execute(
"INSERT INTO queue (namespace, item_key, item_type, status, operation)
VALUES (\"global\", \"chunk:1\", \"chunk\", \"pending\", \"ReEmbed\")",
[],
)
.unwrap();
conn.execute(
"INSERT INTO queue (namespace, item_key, item_type, status, operation)
VALUES (\"ai-sdd\", \"entity:x\", \"entity\", \"pending\", \"ReEmbed\")",
[],
)
.unwrap();
match dequeue_next_pending(&conn, "ReEmbed", "ai-sdd", "").unwrap() {
DequeueOutcome::Claimed(row) => {
assert_eq!(row.item_key, "entity:x");
}
DequeueOutcome::Empty => panic!("expected ai-sdd claim"),
}
let still: i64 = conn
.query_row(
"SELECT COUNT(*) FROM queue WHERE namespace=\"global\" AND status=\"pending\"",
[],
|r| r.get(0),
)
.unwrap();
assert_eq!(still, 1);
assert!(matches!(
dequeue_next_pending(&conn, "ReEmbed", "ai-research", "").unwrap(),
DequeueOutcome::Empty
));
let _ = std::fs::remove_file(&path);
}
#[test]
fn dequeue_skips_future_retry_and_dead() {
let (conn, path) = open_temp_queue();
let eligible = insert_pending(&conn, "mem-eligible");
let waiting = insert_pending(&conn, "mem-waiting");
conn.execute(
"UPDATE queue SET next_retry_at=datetime('now', '+3600 seconds') WHERE id=?1",
rusqlite::params![waiting],
)
.unwrap();
let dead = insert_pending(&conn, "mem-dead");
conn.execute(
"UPDATE queue SET status='dead' WHERE id=?1",
rusqlite::params![dead],
)
.unwrap();
let claimed: Option<i64> = conn
.query_row(
"UPDATE queue SET status='processing', attempt=attempt+1 \
WHERE id = (SELECT id FROM queue WHERE status='pending' \
AND (next_retry_at IS NULL OR next_retry_at <= datetime('now')) \
ORDER BY id LIMIT 1) \
RETURNING id",
[],
|r| r.get(0),
)
.ok();
assert_eq!(claimed, Some(eligible));
let second: Option<i64> = conn
.query_row(
"UPDATE queue SET status='processing', attempt=attempt+1 \
WHERE id = (SELECT id FROM queue WHERE status='pending' \
AND (next_retry_at IS NULL OR next_retry_at <= datetime('now')) \
ORDER BY id LIMIT 1) \
RETURNING id",
[],
|r| r.get(0),
)
.ok();
assert_eq!(second, None);
let _ = std::fs::remove_file(&path);
}
#[test]
fn fresh_processing_claim_is_preserved() {
let (conn, path) = open_temp_queue();
let id = insert_pending(&conn, "mem-fresh");
conn.execute(
"UPDATE queue SET status='processing', claimed_at = CAST(strftime('%s','now') AS INTEGER) WHERE id=?1",
rusqlite::params![id],
)
.unwrap();
let reset = reset_stale_processing_claims(&conn, 1800).unwrap();
assert_eq!(reset, 0, "a fresh claim must not be reset");
let status: String = conn
.query_row(
"SELECT status FROM queue WHERE id=?1",
rusqlite::params![id],
|r| r.get(0),
)
.unwrap();
assert_eq!(status, "processing");
let _ = std::fs::remove_file(&path);
}
#[test]
fn heartbeat_updates_claimed_at() {
let (conn, path) = open_temp_queue();
let id = insert_pending(&conn, "mem-heartbeat");
conn.execute(
"UPDATE queue SET status='processing', claimed_at = CAST(strftime('%s','now') AS INTEGER) - 7200 WHERE id=?1",
rusqlite::params![id],
)
.unwrap();
heartbeat(&conn, id).unwrap();
let claimed_at: Option<i64> = conn
.query_row(
"SELECT claimed_at FROM queue WHERE id=?1",
rusqlite::params![id],
|r| r.get(0),
)
.unwrap();
let claimed_at = claimed_at.expect("claimed_at must be set after heartbeat");
let now: i64 = conn
.query_row("SELECT CAST(strftime('%s','now') AS INTEGER)", [], |r| {
r.get::<_, i64>(0)
})
.unwrap();
assert!(
now - claimed_at <= 5,
"claimed_at must be within 5s of now after heartbeat"
);
let reset = reset_stale_processing_claims(&conn, 1800).unwrap();
assert_eq!(reset, 0, "fresh claim survives the sweep after heartbeat");
let _ = std::fs::remove_file(&path);
}
#[test]
fn legacy_unscoped_rows_are_not_claimed() {
let (conn, path) = open_temp_queue();
conn.execute(
"INSERT INTO queue (item_key, item_type, status, operation) \
VALUES ('orphan-legacy', 'memory', 'pending', 'LegacyUnscoped')",
[],
)
.unwrap();
assert!(matches!(
dequeue_next_pending(&conn, "MemoryBindings", "", "").unwrap(),
DequeueOutcome::Empty
));
let _ = std::fs::remove_file(&path);
}
#[test]
fn stale_processing_claim_is_reset_after_threshold() {
let (conn, path) = open_temp_queue();
let id = insert_pending(&conn, "mem-stale");
conn.execute(
"UPDATE queue SET status='processing', claimed_at = CAST(strftime('%s','now') AS INTEGER) - 7200 WHERE id=?1",
rusqlite::params![id],
)
.unwrap();
let reset = reset_stale_processing_claims(&conn, 1800).unwrap();
assert_eq!(reset, 1, "a stale claim older than the threshold is reset");
let (status, claimed_at): (String, Option<i64>) = conn
.query_row(
"SELECT status, claimed_at FROM queue WHERE id=?1",
rusqlite::params![id],
|r| Ok((r.get(0)?, r.get(1)?)),
)
.unwrap();
assert_eq!(status, "pending");
assert!(claimed_at.is_none(), "claimed_at must be cleared on reset");
let _ = std::fs::remove_file(&path);
}
#[test]
fn validate_claim_rejects_wrong_type_and_key_shape() {
let pair = ClaimedRow {
id: 1,
item_key: "pair:1:2".into(),
item_type: "entity_pair".into(),
operation: "MemoryBindings".into(),
attempt: 1,
};
match validate_claim(&pair, "MemoryBindings", "memory") {
ClaimCheck::SkipWrongType { reason } => {
assert!(reason.contains("wrong_"), "reason={reason}");
}
other => panic!("expected SkipWrongType, got {other:?}"),
}
let entity = ClaimedRow {
id: 2,
item_key: "aeroporto-guarulhos".into(),
item_type: "entity".into(),
operation: "EntityDescriptions".into(),
attempt: 1,
};
assert_eq!(
validate_claim(&entity, "MemoryBindings", "memory"),
ClaimCheck::RequeueWrongOp
);
assert_eq!(
validate_claim(&entity, "EntityDescriptions", "entity"),
ClaimCheck::Ok
);
assert!(is_non_memory_key_shape("pair:1:2"));
assert!(is_non_memory_key_shape("entity:99"));
assert!(!is_non_memory_key_shape("plain-memory-name"));
}
#[test]
fn with_busy_retry_bounds_dequeue_under_sustained_contention() {
let (conn, path) = open_temp_queue();
insert_pending(&conn, "mem-busy");
conn.pragma_update(None, "busy_timeout", 0i64)
.expect("busy_timeout override must succeed");
let blocker = Connection::open(&path).expect("blocker connection must open");
blocker
.execute_batch("BEGIN EXCLUSIVE;")
.expect("exclusive lock must be acquired");
let calls = std::sync::Arc::new(std::sync::atomic::AtomicU32::new(0));
let calls_clone = std::sync::Arc::clone(&calls);
const ATTEMPTS: u32 = 5;
const BASE_DELAY_MS: u64 = 1;
let result: Result<DequeueOutcome, AppError> =
crate::storage::utils::with_busy_retry_policy(ATTEMPTS, BASE_DELAY_MS, || {
calls_clone.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
dequeue_next_pending(&conn, "MemoryBindings", "", "")
});
assert!(
matches!(result, Err(AppError::DbBusy(_))),
"sustained SQLITE_BUSY must convert to DbBusy, not hang or silently report Empty"
);
assert_eq!(
calls.load(std::sync::atomic::Ordering::SeqCst),
ATTEMPTS,
"must attempt exactly the declared budget, never retry unbounded"
);
blocker
.execute_batch("ROLLBACK;")
.expect("releasing the exclusive lock must succeed");
let _ = std::fs::remove_file(&path);
}