use std::collections::BTreeMap;
use reifydb_codec::{key::encoded::EncodedKeyRange, row::pod::EncodedPodRow};
use reifydb_core::{
interface::{
catalog::queue::{
Queue, QueueItemState, QueueItemStatus, QueuePartitionCounters, decode_queue_item_state,
decode_queue_partition_counters, encode_queue_item_state, encode_queue_partition_counters,
},
store::{SingleVersionGet, SingleVersionRange, SingleVersionRow},
},
key::{
any::TaggedKey,
queue::{QueueDueKey, QueueItemStateKey, QueueKeyActiveKey, QueuePartitionKey},
},
};
use reifydb_engine::queue::hydrate::hydrate_queues;
use reifydb_test_harness::engine::TestEngine;
use reifydb_transaction::{single::write::SingleWriteTransaction, transaction::Transaction};
use reifydb_value::value::{datetime::DateTime, row_number::RowNumber};
fn engine_with_queue(declaration: &str) -> TestEngine {
let t = TestEngine::new();
t.admin("CREATE NAMESPACE test");
t.admin(declaration);
t
}
fn find_queue(t: &TestEngine, name: &str) -> Queue {
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()
}
const SCAN_LIMIT: u64 = 8192;
fn scan(t: &TestEngine, range: EncodedKeyRange) -> Vec<SingleVersionRow> {
let store = t.inner().single().read_store();
let batch = SingleVersionRange::range_batch(&store, range, SCAN_LIMIT).unwrap();
assert!(!batch.has_more, "the scan limit is too small to observe the whole keyspace under test");
batch.items
}
fn states(t: &TestEngine, queue: &Queue) -> BTreeMap<RowNumber, (u16, QueueItemState)> {
scan(t, QueueItemStateKey::queue_scan(queue.id).encode())
.iter()
.map(|item| {
let key = QueueItemStateKey::decode(&item.key).unwrap();
(key.row, (key.partition, decode_queue_item_state(EncodedPodRow::view(&item.bytes)).unwrap()))
})
.collect()
}
fn dues(t: &TestEngine, queue: &Queue) -> BTreeMap<RowNumber, QueueDueKey> {
scan(t, QueueDueKey::queue_scan(queue.id).encode())
.iter()
.map(|item| {
let key = QueueDueKey::decode(&item.key).unwrap();
(key.row, key)
})
.collect()
}
fn counters(t: &TestEngine, queue: &Queue, partition: u16) -> QueuePartitionCounters {
let store = t.inner().single().read_store();
SingleVersionGet::get(&store, &QueuePartitionKey::encoded(queue.id, partition))
.unwrap()
.map(|stored| decode_queue_partition_counters(EncodedPodRow::view(&stored.bytes)))
.unwrap_or_default()
}
fn total_depth(t: &TestEngine, queue: &Queue) -> u64 {
(0..queue.partitions()).map(|partition| counters(t, queue, partition).depth).sum()
}
fn keys_in(t: &TestEngine, range: EncodedKeyRange) -> Vec<TaggedKey> {
scan(t, range).iter().map(|item| TaggedKey::decode(&item.key).unwrap()).collect()
}
fn with_partition<F>(t: &TestEngine, queue: &Queue, partition: u16, f: F)
where
F: FnOnce(&mut SingleWriteTransaction<'_>),
{
let single = t.inner().single();
let lock_key = QueuePartitionKey::encoded(queue.id, partition);
let mut tx = single
.begin_command_ranged(
[&lock_key],
vec![
QueueItemStateKey::partition_scan(queue.id, partition).encode(),
QueueDueKey::partition_scan(queue.id, partition).encode(),
QueueKeyActiveKey::partition_scan(queue.id, partition).encode(),
],
)
.unwrap();
f(&mut tx);
tx.commit().unwrap();
}
fn crash_before_handoff(t: &TestEngine, queue: &Queue) {
for partition in 0..queue.partitions() {
let state_keys = keys_in(t, QueueItemStateKey::partition_scan(queue.id, partition).encode());
let due_keys = keys_in(t, QueueDueKey::partition_scan(queue.id, partition).encode());
let chain_keys = keys_in(t, QueueKeyActiveKey::partition_scan(queue.id, partition).encode());
if state_keys.is_empty() && due_keys.is_empty() {
continue;
}
with_partition(t, queue, partition, |tx| {
for key in state_keys.iter().chain(due_keys.iter()).chain(chain_keys.iter()) {
tx.remove(key).unwrap();
}
tx.remove(&QueuePartitionKey::new(queue.id, partition)).unwrap();
});
}
}
fn forget_item(t: &TestEngine, queue: &Queue, row: RowNumber) {
let (partition, _) = states(t, queue)[&row];
let due = dues(t, queue)[&row].clone();
let mut counters = counters(t, queue, partition);
counters.depth -= 1;
with_partition(t, queue, partition, |tx| {
tx.remove(&QueueItemStateKey::new(queue.id, partition, row)).unwrap();
tx.remove(&due).unwrap();
tx.set(&QueuePartitionKey::new(queue.id, partition), encode_queue_partition_counters(&counters))
.unwrap();
});
}
#[test]
fn test_hydration_recreates_a_lost_handoff() {
let t = engine_with_queue("CREATE QUEUE test::jobs { id: int4 } WITH { fifo: { partitions: 1 } }");
t.command("INSERT test::jobs [{ id: 1 }, { id: 2 }, { id: 3 }]");
let queue = find_queue(&t, "jobs");
crash_before_handoff(&t, &queue);
assert!(states(&t, &queue).is_empty(), "the crash image must have no scheduling state left");
assert_eq!(hydrate_queues(t.inner()).unwrap(), 3);
let states = states(&t, &queue);
assert_eq!(states.len(), 3);
for (partition, state) in states.values() {
assert_eq!(*partition, 0);
assert_eq!(state.status, QueueItemStatus::Ready, "recovered items must be claimable, not parked");
assert_eq!(state.attempt, 0);
}
assert_eq!(dues(&t, &queue).len(), 3);
assert_eq!(total_depth(&t, &queue), 3);
}
#[test]
fn test_hydration_recovers_not_before_from_the_stored_row() {
let t = engine_with_queue("CREATE QUEUE test::jobs { id: int4 } WITH { fifo: { partitions: 1 } }");
t.command(r#"INSERT test::jobs [{ id: 1 }] WITH { not_before: datetime::from_epoch_millis(1700000000000) }"#);
let queue = find_queue(&t, "jobs");
let expected = DateTime::from_nanos(1_700_000_000_000 * 1_000_000);
crash_before_handoff(&t, &queue);
assert_eq!(hydrate_queues(t.inner()).unwrap(), 1);
let states = states(&t, &queue);
let (_, state) = states.values().next().unwrap();
assert_eq!(state.not_before, Some(expected), "the recovered state must carry the stored not_before");
assert_eq!(dues(&t, &queue).values().next().unwrap().due, expected, "and so must the due index entry");
}
#[test]
fn test_hydration_admits_only_the_items_that_are_missing() {
let t = engine_with_queue("CREATE QUEUE test::jobs { id: int4 } WITH { fifo: { partitions: 1 } }");
t.command("INSERT test::jobs [{ id: 1 }, { id: 2 }, { id: 3 }]");
let queue = find_queue(&t, "jobs");
let lost = *states(&t, &queue).keys().next().unwrap();
forget_item(&t, &queue, lost);
assert_eq!(total_depth(&t, &queue), 2);
assert_eq!(hydrate_queues(t.inner()).unwrap(), 1, "only the forgotten item may be admitted");
assert_eq!(states(&t, &queue).len(), 3);
assert_eq!(dues(&t, &queue).len(), 3, "the surviving items must not gain a second due entry");
assert_eq!(total_depth(&t, &queue), 3);
}
#[test]
fn test_hydration_is_idempotent() {
let t = engine_with_queue("CREATE QUEUE test::jobs { id: int4 } WITH { fifo: { partitions: 4 } }");
t.command("INSERT test::jobs [{ id: 1 }, { id: 2 }, { id: 3 }]");
let queue = find_queue(&t, "jobs");
let before = states(&t, &queue);
assert_eq!(hydrate_queues(t.inner()).unwrap(), 0, "a healthy store must admit nothing");
assert_eq!(hydrate_queues(t.inner()).unwrap(), 0);
assert_eq!(states(&t, &queue), before);
assert_eq!(dues(&t, &queue).len(), 3);
assert_eq!(total_depth(&t, &queue), 3);
}
#[test]
fn test_hydration_does_not_re_admit_an_item_that_reached_a_terminal_status() {
let t = engine_with_queue("CREATE QUEUE test::jobs { id: int4 } WITH { fifo: { partitions: 1 } }");
t.command("INSERT test::jobs [{ id: 1 }, { id: 2 }]");
let queue = find_queue(&t, "jobs");
let done = *states(&t, &queue).keys().next().unwrap();
let due = dues(&t, &queue)[&done].clone();
let mut counters = counters(&t, &queue, 0);
counters.depth -= 1;
with_partition(&t, &queue, 0, |tx| {
let mut state = QueueItemState::ready(None);
state.status = QueueItemStatus::Done;
tx.set(&QueueItemStateKey::new(queue.id, 0, done), encode_queue_item_state(&state)).unwrap();
tx.remove(&due).unwrap();
tx.set(&QueuePartitionKey::new(queue.id, 0), encode_queue_partition_counters(&counters)).unwrap();
});
assert_eq!(hydrate_queues(t.inner()).unwrap(), 0, "a terminal item must not be re-admitted");
assert_eq!(states(&t, &queue)[&done].1.status, QueueItemStatus::Done, "its status must not be reset");
assert!(!dues(&t, &queue).contains_key(&done), "and it must not reappear in the due index");
assert_eq!(total_depth(&t, &queue), 1);
}
#[test]
fn test_hydration_recovers_a_queue_that_outgrows_one_scan_batch() {
let t = engine_with_queue("CREATE QUEUE test::jobs { id: int4 } WITH { fifo: { partitions: 1 } }");
const ITEMS: usize = 1100;
for chunk in 0..11 {
let rows: Vec<String> = (0..100).map(|i| format!("{{ id: {} }}", chunk * 100 + i)).collect();
t.command(&format!("INSERT test::jobs [{}]", rows.join(", ")));
}
let queue = find_queue(&t, "jobs");
assert_eq!(states(&t, &queue).len(), ITEMS);
crash_before_handoff(&t, &queue);
assert!(states(&t, &queue).is_empty());
assert_eq!(
hydrate_queues(t.inner()).unwrap(),
ITEMS as u64,
"every item must be admitted, not just the first batch"
);
assert_eq!(states(&t, &queue).len(), ITEMS);
assert_eq!(dues(&t, &queue).len(), ITEMS);
assert_eq!(total_depth(&t, &queue), ITEMS as u64);
}
#[test]
fn test_hydration_recomputes_the_original_partition_of_an_ordered_item() {
let t = engine_with_queue(
"CREATE QUEUE test::jobs { id: int4, tenant: utf8 } WITH { fifo: { partitions: 8, ordered_by: tenant } }",
);
t.command(
r#"INSERT test::jobs [{ id: 1, tenant: "a" }, { id: 2, tenant: "b" }, { id: 3, tenant: "c" }, { id: 4, tenant: "a" }]"#,
);
let queue = find_queue(&t, "jobs");
let before: BTreeMap<RowNumber, u16> = states(&t, &queue).iter().map(|(row, (p, _))| (*row, *p)).collect();
crash_before_handoff(&t, &queue);
assert_eq!(hydrate_queues(t.inner()).unwrap(), 4);
let after: BTreeMap<RowNumber, u16> = states(&t, &queue).iter().map(|(row, (p, _))| (*row, *p)).collect();
assert_eq!(after, before, "every recovered item must land in the partition it was enqueued to");
let dues = dues(&t, &queue);
assert_eq!(dues.len(), 3, "the second item of tenant a is parked behind its sibling, so it has no due entry");
for (row, partition) in &before {
let Some(due) = dues.get(row) else {
continue;
};
assert_eq!(due.partition, *partition, "the due entry must follow the item's partition");
}
}