use reifydb_codec::row::{pod::EncodedPodRow, queue_attempt::EncodedQueueAttemptRow};
use reifydb_core::{
interface::{
catalog::{
id::QueueId,
queue::{
AttemptOutcome, QueueAttemptRecord, QueueItemState, QueueItemStatus,
QueuePartitionCounters, decode_queue_attempt, decode_queue_item_state,
decode_queue_partition_counters,
},
},
store::{SingleVersionGet, SingleVersionRange},
},
key::{
any::TaggedKey,
queue::{QueueAttemptKey, QueueDueKey, QueueItemStateKey, QueuePartitionKey},
},
};
use reifydb_test_harness::engine::TestEngine;
use reifydb_transaction::{
change::{QueueAckTransition, QueueRowAck},
queue::scheduling::apply_ack_transitions,
transaction::Transaction,
};
use reifydb_value::value::{Value, frame::frame::Frame, row_number::RowNumber};
fn engine_with_queue(declaration: &str) -> TestEngine {
let t = TestEngine::new();
t.admin("CREATE NAMESPACE test");
t.admin(declaration);
t
}
fn queue_id(t: &TestEngine, name: &str) -> QueueId {
let catalog = t.inner().catalog();
let mut query_txn = t.inner().begin_query(TestEngine::identity()).unwrap();
let mut txn = Transaction::Query(&mut query_txn);
let namespace = catalog.find_namespace_by_name(&mut txn, "test").unwrap().unwrap();
catalog.find_queue_by_name(&mut txn, namespace.id(), name).unwrap().unwrap().id
}
fn states(t: &TestEngine, queue: QueueId) -> Vec<(QueueItemStateKey, QueueItemState)> {
let store = t.inner().single().read_store();
SingleVersionRange::range_batch(&store, QueueItemStateKey::queue_scan(queue).encode(), 1024)
.unwrap()
.items
.iter()
.map(|item| {
(
QueueItemStateKey::decode(&item.key).unwrap(),
decode_queue_item_state(EncodedPodRow::view(&item.bytes)).unwrap(),
)
})
.collect()
}
fn state_of(t: &TestEngine, queue: QueueId) -> QueueItemState {
states(t, queue).into_iter().next().expect("the queue must hold exactly one item").1
}
fn dues(t: &TestEngine, queue: QueueId) -> Vec<QueueDueKey> {
let store = t.inner().single().read_store();
SingleVersionRange::range_batch(&store, QueueDueKey::queue_scan(queue).encode(), 1024)
.unwrap()
.items
.iter()
.map(|item| QueueDueKey::decode(&item.key).unwrap())
.collect()
}
fn counters(t: &TestEngine, queue: QueueId, partition: u16) -> QueuePartitionCounters {
let store = t.inner().single().read_store();
SingleVersionGet::get(&store, &QueuePartitionKey::encoded(queue, partition))
.unwrap()
.map(|stored| decode_queue_partition_counters(EncodedPodRow::view(&stored.bytes)))
.unwrap_or_default()
}
fn attempts(t: &TestEngine, queue: QueueId) -> Vec<(QueueAttemptKey, QueueAttemptRecord)> {
let mut query_txn = t.inner().begin_query(TestEngine::identity()).unwrap();
let mut txn = Transaction::Query(&mut query_txn);
let mut stream = txn
.range(QueueAttemptKey::queue_scan(queue), reifydb_transaction::multi::RangeScope::All, 1024)
.unwrap();
let mut out = Vec::new();
while let Some(item) = stream.next() {
let item = item.unwrap();
out.push((
match item.key {
TaggedKey::QueueAttempt(key) => key,
other => panic!("queue attempt scan yielded {other:?}"),
},
decode_queue_attempt(EncodedQueueAttemptRow::view(&item.bytes)).unwrap(),
));
}
out
}
fn claim_one(t: &TestEngine, worker: &str) -> String {
let frames = t.command(&format!(r#"CALL queue::claim("{worker}", "test::jobs", 1, duration::seconds(30))"#));
token_of(&frames)
}
fn token_of(frames: &[Frame]) -> String {
let frame = frames.first().expect("claim must return a frame");
assert_eq!(frame.row_count(), 1, "expected exactly one claimed item");
match frame.columns.iter().find(|c| c.name == "token").unwrap().data.get_value(0) {
Value::Utf8(t) => t,
other => panic!("token must be Utf8, got {other:?}"),
}
}
fn status_of(frames: &[Frame]) -> String {
match frames[0].columns.iter().find(|c| c.name == "status").unwrap().data.get_value(0) {
Value::Utf8(s) => s,
other => panic!("status must be Utf8, got {other:?}"),
}
}
fn ack(t: &TestEngine, token: &str) -> String {
status_of(&t.command(&format!(r#"CALL queue::ack("{token}")"#)))
}
fn fail(t: &TestEngine, token: &str) -> String {
status_of(&t.command(&format!(r#"CALL queue::fail("{token}", none)"#)))
}
fn kill(t: &TestEngine, token: &str) -> String {
status_of(&t.command(&format!(r#"CALL queue::kill("{token}", none)"#)))
}
fn claimable(t: &TestEngine) -> usize {
TestEngine::row_count(&t.command(r#"CALL queue::claim("probe", "test::jobs", 10, duration::seconds(30))"#))
}
const ONE_PARTITION: &str = "CREATE QUEUE test::jobs { id: int4 } WITH { fifo: { partitions: 1 } }";
#[test]
fn test_an_ok_ack_finishes_the_item_for_good() {
let t = engine_with_queue(ONE_PARTITION);
t.command("INSERT test::jobs [{ id: 1 }]");
let queue = queue_id(&t, "jobs");
let token = claim_one(&t, "w1");
assert_eq!(ack(&t, &token), "ok");
let state = state_of(&t, queue);
assert_eq!(state.status, QueueItemStatus::Done);
assert_eq!(state.lease_deadline, None, "a finished item must not keep a live lease");
assert_eq!(counters(&t, queue, 0).in_flight, 0);
assert_eq!(counters(&t, queue, 0).depth, 0, "depth already fell at claim time");
assert_eq!(claimable(&t), 0, "a done item must never be delivered again");
}
#[test]
fn test_a_fail_ack_records_what_the_worker_reported() {
let t = engine_with_queue(ONE_PARTITION);
t.command("INSERT test::jobs [{ id: 1 }]");
let queue = queue_id(&t, "jobs");
t.mock_clock().set_nanos(4_200);
let token = claim_one(&t, "worker-a");
t.command(&format!(r#"CALL queue::fail("{token}", "gateway timeout")"#));
let recorded = attempts(&t, queue);
assert_eq!(recorded.len(), 1);
let (key, record) = &recorded[0];
assert_eq!(key.attempt, 1);
assert_eq!(record.worker, "worker-a");
assert_eq!(record.outcome, AttemptOutcome::Err);
assert_eq!(record.response.as_deref(), Some("gateway timeout"));
assert_eq!(record.lost, false, "only the step-5 reaper writes lost attempts");
assert_eq!(record.anomaly, None);
}
#[test]
fn test_an_err_ack_requeues_the_item_until_the_retry_budget_is_spent() {
let t = engine_with_queue(
"CREATE QUEUE test::jobs { id: int4 } WITH { fifo: { partitions: 1 }, retry: { attempts: 2 } }",
);
t.command("INSERT test::jobs [{ id: 1 }]");
let queue = queue_id(&t, "jobs");
assert_eq!(fail(&t, &claim_one(&t, "w1")), "ok");
assert_eq!(state_of(&t, queue).status, QueueItemStatus::Ready, "attempt 1 of 2 must be retried");
assert_eq!(dues(&t, queue).len(), 1, "a retried item needs a due entry or no scan will find it");
assert_eq!(counters(&t, queue, 0).depth, 1);
assert_eq!(counters(&t, queue, 0).in_flight, 0);
assert_eq!(claimable(&t), 0, "a backed-off retry must not be redelivered before its delay elapses");
t.mock_clock().advance_millis(10_000);
let second = claim_one(&t, "w1");
assert_eq!(fail(&t, &second), "ok");
assert_eq!(state_of(&t, queue).status, QueueItemStatus::Dead, "the budget is spent at attempt 2");
assert_eq!(claimable(&t), 0);
assert_eq!(counters(&t, queue, 0).in_flight, 0);
}
#[test]
fn test_an_err_ack_places_the_due_entry_at_the_backoff_instant() {
let t = engine_with_queue(
"CREATE QUEUE test::jobs { id: int4 } WITH { fifo: { partitions: 1 }, retry: { attempts: 5 } }",
);
t.command("INSERT test::jobs [{ id: 1 }]");
let queue = queue_id(&t, "jobs");
t.mock_clock().set_millis(50_000);
assert_eq!(fail(&t, &claim_one(&t, "w1")), "ok");
let state = state_of(&t, queue);
assert_eq!(
state.backoff_until.map(|b| b.to_nanos()),
Some(60_000_000_000),
"the default 10s backoff must land 10s after the ack"
);
let due = dues(&t, queue);
assert_eq!(due.len(), 1);
assert_eq!(
due[0].due.to_nanos(),
60_000_000_000,
"the due entry and backoff_until must name the same instant or claim removes the wrong key"
);
}
#[test]
fn test_consecutive_failures_double_the_wait() {
let t = engine_with_queue(
"CREATE QUEUE test::jobs { id: int4 } WITH { fifo: { partitions: 1 }, retry: { attempts: 5 } }",
);
t.command("INSERT test::jobs [{ id: 1 }]");
let queue = queue_id(&t, "jobs");
t.mock_clock().set_millis(0);
assert_eq!(fail(&t, &claim_one(&t, "w1")), "ok");
assert_eq!(dues(&t, queue)[0].due.to_nanos(), 10_000_000_000, "first failure waits one base interval");
t.mock_clock().set_millis(10_000);
assert_eq!(fail(&t, &claim_one(&t, "w1")), "ok");
assert_eq!(dues(&t, queue)[0].due.to_nanos(), 30_000_000_000, "the second failure waits 20s, not another 10s");
}
#[test]
fn test_a_retry_does_not_overwrite_the_user_declared_not_before() {
let t = engine_with_queue(
"CREATE QUEUE test::jobs { id: int4 } WITH { fifo: { partitions: 1 }, retry: { attempts: 5 } }",
);
t.command(r#"INSERT test::jobs [{ id: 1 }] WITH { not_before: datetime::from_epoch_millis(10000) }"#);
let queue = queue_id(&t, "jobs");
t.mock_clock().set_millis(10_000);
assert_eq!(fail(&t, &claim_one(&t, "w1")), "ok");
let state = state_of(&t, queue);
assert_eq!(
state.not_before.map(|n| n.to_nanos()),
Some(10_000_000_000),
"the declared not_before must survive a retry untouched"
);
assert_eq!(state.backoff_until.map(|b| b.to_nanos()), Some(20_000_000_000));
}
#[test]
fn test_a_dead_ack_buries_the_item_immediately() {
let t = engine_with_queue(
"CREATE QUEUE test::jobs { id: int4 } WITH { fifo: { partitions: 1 }, retry: { attempts: 5 } }",
);
t.command("INSERT test::jobs [{ id: 1 }]");
let queue = queue_id(&t, "jobs");
assert_eq!(kill(&t, &claim_one(&t, "w1")), "ok");
assert_eq!(state_of(&t, queue).status, QueueItemStatus::Dead);
assert_eq!(claimable(&t), 0, "a dead item must not be retried despite an unspent budget");
}
#[test]
fn test_a_repeated_ack_is_a_no_op() {
let t = engine_with_queue(ONE_PARTITION);
t.command("INSERT test::jobs [{ id: 1 }]");
let queue = queue_id(&t, "jobs");
let token = claim_one(&t, "w1");
assert_eq!(ack(&t, &token), "ok");
let after_first = counters(&t, queue, 0);
assert_eq!(ack(&t, &token), "repeat");
assert_eq!(attempts(&t, queue).len(), 1, "a repeat must not write a second attempt record");
assert_eq!(counters(&t, queue, 0), after_first, "a repeat must not move the counters");
assert_eq!(state_of(&t, queue).status, QueueItemStatus::Done);
}
#[test]
fn test_the_first_outcome_wins_and_the_conflicting_one_is_recorded() {
let t = engine_with_queue(ONE_PARTITION);
t.command("INSERT test::jobs [{ id: 1 }]");
let queue = queue_id(&t, "jobs");
let token = claim_one(&t, "w1");
ack(&t, &token);
assert_eq!(fail(&t, &token), "stale");
let (_, record) = attempts(&t, queue).into_iter().next().unwrap();
assert_eq!(record.outcome, AttemptOutcome::Ok, "the first outcome must stand");
assert!(record.anomaly.unwrap().contains("conflicting late ack"), "the contradiction must be recorded");
assert_eq!(state_of(&t, queue).status, QueueItemStatus::Done, "a conflicting ack must not transition");
}
#[test]
fn test_an_ack_for_an_attempt_that_is_not_live_is_recorded_but_never_transitions() {
let t = engine_with_queue(ONE_PARTITION);
t.command("INSERT test::jobs [{ id: 1 }]");
let queue = queue_id(&t, "jobs");
let live = claim_one(&t, "w1");
let forged = live.replace(":1:w1", ":7:w1");
assert_eq!(ack(&t, &forged), "stale");
let (key, record) = attempts(&t, queue).into_iter().next().unwrap();
assert_eq!(key.attempt, 7);
assert!(record.anomaly.unwrap().contains("stale"));
assert_eq!(state_of(&t, queue).status, QueueItemStatus::Leased, "the real lease must be untouched");
assert_eq!(state_of(&t, queue).attempt, 1);
}
#[test]
fn test_an_ack_after_the_item_was_already_finished_does_not_resurrect_it() {
let t = engine_with_queue(ONE_PARTITION);
t.command("INSERT test::jobs [{ id: 1 }]");
let queue = queue_id(&t, "jobs");
let token = claim_one(&t, "w1");
ack(&t, &token);
let before = counters(&t, queue, 0);
let other_attempt = token.replace(":1:w1", ":2:w1");
assert_eq!(fail(&t, &other_attempt), "stale");
assert_eq!(state_of(&t, queue).status, QueueItemStatus::Done);
assert_eq!(counters(&t, queue, 0), before);
}
fn apply_ack(t: &TestEngine, queue: QueueId, row: RowNumber, attempt: u32, transition: QueueAckTransition) -> u64 {
apply_ack_transitions(
t.inner().single(),
queue,
0,
&[QueueRowAck {
queue_id: queue,
partition: 0,
key_hash: None,
row_number: row,
attempt,
transition,
}],
)
.unwrap()
}
#[test]
fn test_the_interceptor_refuses_an_ack_whose_attempt_no_longer_matches_the_lease() {
let t = engine_with_queue(ONE_PARTITION);
t.command("INSERT test::jobs [{ id: 1 }]");
let queue = queue_id(&t, "jobs");
claim_one(&t, "w1");
let row = states(&t, queue)[0].0.row;
let before = counters(&t, queue, 0);
let applied = apply_ack(&t, queue, row, 99, QueueAckTransition::Done);
assert_eq!(applied, 0, "a mismatched attempt must apply no transition");
assert_eq!(state_of(&t, queue).status, QueueItemStatus::Leased);
assert_eq!(counters(&t, queue, 0), before, "a refused ack must not move the counters");
}
#[test]
fn test_the_interceptor_refuses_an_ack_against_an_item_that_is_no_longer_leased() {
let t = engine_with_queue(ONE_PARTITION);
t.command("INSERT test::jobs [{ id: 1 }]");
let queue = queue_id(&t, "jobs");
let token = claim_one(&t, "w1");
let row = states(&t, queue)[0].0.row;
ack(&t, &token);
let before = counters(&t, queue, 0);
let applied = apply_ack(&t, queue, row, 1, QueueAckTransition::Done);
assert_eq!(applied, 0, "an item that is already Done must not transition again");
assert_eq!(counters(&t, queue, 0), before);
}
#[test]
fn test_a_malformed_token_is_rejected_with_queue_003() {
let t = engine_with_queue(ONE_PARTITION);
let err = t.command_err(r#"CALL queue::ack("not-a-token")"#);
assert!(err.contains("QUEUE_003"), "{err}");
}
#[test]
fn test_an_ack_is_rejected_outside_a_command_transaction() {
let t = engine_with_queue(ONE_PARTITION);
let err = t.query_err(r#"CALL queue::ack("qt1:1:0:1:1:w1")"#);
assert!(err.contains("must run in a command transaction"), "{err}");
}