1use reifydb_codec::{key::encoded::EncodedKeyRange, row::pod::EncodedPodRow};
5use reifydb_core::{
6 interface::catalog::{
7 id::QueueId,
8 queue::{
9 QueueItemState, QueueItemStatus, QueuePartitionCounters, decode_queue_item_state,
10 decode_queue_partition_counters, encode_queue_item_state, encode_queue_partition_counters,
11 },
12 },
13 key::queue::{QueueDueKey, QueueItemStateKey, QueueKeyActiveKey, QueuePartitionKey},
14};
15use reifydb_value::{
16 Result,
17 value::{datetime::DateTime, row_number::RowNumber},
18};
19use tracing::debug;
20
21use crate::{
22 change::{QueueAckTransition, QueueRowAck},
23 queue::chain::{ChainHead, ChainOverlay, chain_add, chain_peek, chain_remove},
24 single::{SingleTransaction, write::SingleWriteTransaction},
25};
26
27pub struct QueueAdmission {
28 pub row: RowNumber,
29 pub key_hash: Option<u64>,
30 pub not_before: Option<DateTime>,
31}
32
33struct TransitionEffect {
34 requeued: bool,
35 blocked_delta: i64,
36}
37
38fn partition_ranges(queue: QueueId, partition: u16) -> Vec<EncodedKeyRange> {
39 vec![
40 QueueItemStateKey::partition_scan(queue, partition).encode(),
41 QueueDueKey::partition_scan(queue, partition).encode(),
42 QueueKeyActiveKey::partition_scan(queue, partition).encode(),
43 ]
44}
45
46pub fn admit_ready_items(
47 single: &SingleTransaction,
48 queue: QueueId,
49 partition: u16,
50 items: &[QueueAdmission],
51) -> Result<u64> {
52 if items.is_empty() {
53 return Ok(0);
54 }
55
56 let lock_key = QueuePartitionKey::new(queue, partition);
57 let mut tx = single.begin_command_ranged([&lock_key.encode()], partition_ranges(queue, partition))?;
58
59 let mut overlay = ChainOverlay::default();
60 let mut admitted = 0u64;
61 let mut blocked_delta = 0i64;
62 for item in items {
63 let state_key = QueueItemStateKey::new(queue, partition, item.row);
64 if tx.contains_key(&state_key)? {
65 continue;
66 }
67
68 let mut state = QueueItemState::ready(item.not_before);
69 state.key_hash = item.key_hash.unwrap_or_default();
70
71 if let Some(key_hash) = item.key_hash {
72 match chain_peek(single, &overlay, queue, partition, key_hash)? {
73 ChainHead::Empty => {}
74 ChainHead::Single(_) => {
75 state.status = QueueItemStatus::Parked;
76 blocked_delta += 1;
77 }
78 ChainHead::Multiple(_) => state.status = QueueItemStatus::Parked,
79 }
80 chain_add(&mut tx, &mut overlay, queue, partition, key_hash, item.row)?;
81 }
82
83 tx.set(&state_key, encode_queue_item_state(&state))?;
84 if state.status == QueueItemStatus::Ready {
85 expose_due(&mut tx, queue, partition, item.row, &state)?;
86 }
87 admitted += 1;
88 }
89
90 if admitted > 0 {
91 let mut counters = read_counters(&mut tx, &lock_key)?;
92 counters.depth += admitted;
93 counters.blocked_keys = counters.blocked_keys.saturating_add_signed(blocked_delta);
94 tx.set(&lock_key, encode_queue_partition_counters(&counters))?;
95 }
96
97 tx.commit()?;
98
99 Ok(admitted)
100}
101
102pub fn apply_ack_transitions(
103 single: &SingleTransaction,
104 queue: QueueId,
105 partition: u16,
106 acks: &[QueueRowAck],
107) -> Result<u64> {
108 if acks.is_empty() {
109 return Ok(0);
110 }
111
112 let lock_key = QueuePartitionKey::new(queue, partition);
113 let mut tx = single.begin_command_ranged([&lock_key.encode()], partition_ranges(queue, partition))?;
114
115 let mut overlay = ChainOverlay::default();
116 let mut applied = 0u64;
117 let mut requeued = 0u64;
118 let mut blocked_delta = 0i64;
119 for ack in acks {
120 let state_key = QueueItemStateKey::new(queue, partition, ack.row_number);
121 let Some(stored) = tx.get(&state_key)? else {
122 debug!(queue = queue.0, partition, item = ack.row_number.0, "ack has no item state");
123 continue;
124 };
125 let Some(mut state) = decode_queue_item_state(EncodedPodRow::view(&stored.bytes)) else {
126 continue;
127 };
128
129 if state.status != QueueItemStatus::Leased || state.attempt != ack.attempt {
130 debug!(
131 queue = queue.0,
132 partition,
133 item = ack.row_number.0,
134 attempt = ack.attempt,
135 "ack no longer matches the lease it was issued for"
136 );
137 continue;
138 }
139
140 let effect = apply_state_transition(
141 single,
142 &mut tx,
143 &mut overlay,
144 TransitionTarget::for_ack(queue, partition, ack),
145 &mut state,
146 &ack.transition,
147 )?;
148 if effect.requeued {
149 requeued += 1;
150 }
151 blocked_delta += effect.blocked_delta;
152 applied += 1;
153 }
154
155 if applied > 0 {
156 adjust_counters(&mut tx, &lock_key, applied, requeued, blocked_delta)?;
157 }
158
159 tx.commit()?;
160
161 Ok(applied)
162}
163
164pub struct ExpiredLease {
165 pub row: RowNumber,
166 pub attempt: u32,
167 pub key_hash: Option<u64>,
168 pub lease_deadline: DateTime,
169}
170
171pub fn apply_reap_transition(
172 single: &SingleTransaction,
173 queue: QueueId,
174 partition: u16,
175 lease: &ExpiredLease,
176 transition: &QueueAckTransition,
177 now: DateTime,
178) -> Result<bool> {
179 let lock_key = QueuePartitionKey::new(queue, partition);
180 let mut tx = single.begin_command_ranged([&lock_key.encode()], partition_ranges(queue, partition))?;
181
182 let state_key = QueueItemStateKey::new(queue, partition, lease.row);
183 let Some(stored) = tx.get(&state_key)? else {
184 return Ok(false);
185 };
186 let Some(mut state) = decode_queue_item_state(EncodedPodRow::view(&stored.bytes)) else {
187 return Ok(false);
188 };
189
190 if state.status != QueueItemStatus::Leased
191 || state.attempt != lease.attempt
192 || state.lease_deadline != Some(lease.lease_deadline)
193 || lease.lease_deadline > now
194 {
195 debug!(
196 queue = queue.0,
197 partition,
198 item = lease.row.0,
199 attempt = lease.attempt,
200 "the lease moved between the reaper's scan and its compare-and-set"
201 );
202 return Ok(false);
203 }
204
205 let mut overlay = ChainOverlay::default();
206 let effect = apply_state_transition(
207 single,
208 &mut tx,
209 &mut overlay,
210 TransitionTarget::for_lease(queue, partition, lease),
211 &mut state,
212 transition,
213 )?;
214 adjust_counters(&mut tx, &lock_key, 1, u64::from(effect.requeued), effect.blocked_delta)?;
215 tx.commit()?;
216
217 Ok(true)
218}
219
220pub enum ReplayOutcome {
221 Ready,
222 Parked,
223 Unknown,
224 Unreadable,
225 NotDead(QueueItemStatus),
226}
227
228pub fn apply_replay_transition(
229 single: &SingleTransaction,
230 queue: QueueId,
231 partition: u16,
232 row: RowNumber,
233 key_hash: Option<u64>,
234) -> Result<ReplayOutcome> {
235 let lock_key = QueuePartitionKey::new(queue, partition);
236 let mut tx = single.begin_command_ranged([&lock_key.encode()], partition_ranges(queue, partition))?;
237
238 let state_key = QueueItemStateKey::new(queue, partition, row);
239 let Some(stored) = tx.get(&state_key)? else {
240 return Ok(ReplayOutcome::Unknown);
241 };
242 let Some(mut state) = decode_queue_item_state(EncodedPodRow::view(&stored.bytes)) else {
243 return Ok(ReplayOutcome::Unreadable);
244 };
245
246 if state.status != QueueItemStatus::Dead {
247 return Ok(ReplayOutcome::NotDead(state.status));
248 }
249
250 state.status = QueueItemStatus::Ready;
251 state.budget_base = state.attempt;
252 state.backoff_until = None;
253 state.lease_deadline = None;
254
255 let mut overlay = ChainOverlay::default();
256 let mut blocked_delta = 0i64;
257 if let Some(key_hash) = key_hash {
258 match chain_peek(single, &overlay, queue, partition, key_hash)? {
259 ChainHead::Empty => {}
260 ChainHead::Single(_) => {
261 state.status = QueueItemStatus::Parked;
262 blocked_delta += 1;
263 }
264 ChainHead::Multiple(_) => state.status = QueueItemStatus::Parked,
265 }
266 chain_add(&mut tx, &mut overlay, queue, partition, key_hash, row)?;
267 }
268
269 tx.set(&state_key, encode_queue_item_state(&state))?;
270 if state.status == QueueItemStatus::Ready {
271 expose_due(&mut tx, queue, partition, row, &state)?;
272 }
273
274 let mut counters = read_counters(&mut tx, &lock_key)?;
275 counters.depth += 1;
276 counters.blocked_keys = counters.blocked_keys.saturating_add_signed(blocked_delta);
277 tx.set(&lock_key, encode_queue_partition_counters(&counters))?;
278
279 tx.commit()?;
280
281 Ok(if state.status == QueueItemStatus::Parked {
282 ReplayOutcome::Parked
283 } else {
284 ReplayOutcome::Ready
285 })
286}
287
288pub fn remove_item_states(
289 single: &SingleTransaction,
290 queue: QueueId,
291 partition: u16,
292 rows: &[RowNumber],
293) -> Result<u64> {
294 if rows.is_empty() {
295 return Ok(0);
296 }
297
298 let lock_key = QueuePartitionKey::new(queue, partition);
299 let mut tx = single.begin_command_ranged([&lock_key.encode()], partition_ranges(queue, partition))?;
300
301 let mut removed = 0u64;
302 for row in rows {
303 let state_key = QueueItemStateKey::new(queue, partition, *row);
304 let Some(stored) = tx.get(&state_key)? else {
305 continue;
306 };
307 let Some(state) = decode_queue_item_state(EncodedPodRow::view(&stored.bytes)) else {
308 tx.remove(&state_key)?;
309 removed += 1;
310 continue;
311 };
312
313 if state.status != QueueItemStatus::Done && state.status != QueueItemStatus::Dead {
314 debug!(
315 queue = queue.0,
316 partition,
317 item = row.0,
318 "retention skipped an item that stopped being terminal under it"
319 );
320 continue;
321 }
322
323 tx.remove(&state_key)?;
324 removed += 1;
325 }
326
327 tx.commit()?;
328
329 Ok(removed)
330}
331
332struct TransitionTarget {
333 queue: QueueId,
334 partition: u16,
335 row: RowNumber,
336 key_hash: Option<u64>,
337}
338
339impl TransitionTarget {
340 fn for_ack(queue: QueueId, partition: u16, ack: &QueueRowAck) -> Self {
341 Self {
342 queue,
343 partition,
344 row: ack.row_number,
345 key_hash: ack.key_hash,
346 }
347 }
348
349 fn for_lease(queue: QueueId, partition: u16, lease: &ExpiredLease) -> Self {
350 Self {
351 queue,
352 partition,
353 row: lease.row,
354 key_hash: lease.key_hash,
355 }
356 }
357}
358
359fn apply_state_transition(
360 single: &SingleTransaction,
361 tx: &mut SingleWriteTransaction<'_>,
362 overlay: &mut ChainOverlay,
363 target: TransitionTarget,
364 state: &mut QueueItemState,
365 transition: &QueueAckTransition,
366) -> Result<TransitionEffect> {
367 let TransitionTarget {
368 queue,
369 partition,
370 row,
371 key_hash,
372 } = target;
373
374 state.lease_deadline = None;
375
376 let terminal = match transition {
377 QueueAckTransition::Done => {
378 state.status = QueueItemStatus::Done;
379 true
380 }
381 QueueAckTransition::Dead => {
382 state.status = QueueItemStatus::Dead;
383 true
384 }
385 QueueAckTransition::Retry {
386 backoff_until,
387 } => {
388 state.status = QueueItemStatus::Ready;
389 state.backoff_until = Some(*backoff_until);
390 expose_due(tx, queue, partition, row, state)?;
391 false
392 }
393 };
394
395 tx.set(&QueueItemStateKey::new(queue, partition, row), encode_queue_item_state(state))?;
396
397 let blocked_delta = match (terminal, key_hash) {
398 (true, Some(key_hash)) => {
399 chain_remove(tx, overlay, queue, partition, key_hash, row)?;
400 promote_next(single, tx, overlay, queue, partition, key_hash)?
401 }
402 _ => 0,
403 };
404
405 Ok(TransitionEffect {
406 requeued: !terminal,
407 blocked_delta,
408 })
409}
410
411fn promote_next(
412 single: &SingleTransaction,
413 tx: &mut SingleWriteTransaction<'_>,
414 overlay: &ChainOverlay,
415 queue: QueueId,
416 partition: u16,
417 key_hash: u64,
418) -> Result<i64> {
419 let (successor, blocked_delta) = match chain_peek(single, overlay, queue, partition, key_hash)? {
420 ChainHead::Empty => return Ok(0),
421 ChainHead::Single(row) => (row, -1),
422 ChainHead::Multiple(row) => (row, 0),
423 };
424
425 let state_key = QueueItemStateKey::new(queue, partition, successor);
426 let Some(stored) = tx.get(&state_key)? else {
427 debug!(queue = queue.0, partition, item = successor.0, "the successor of a key has no item state");
428 return Ok(blocked_delta);
429 };
430 let Some(mut state) = decode_queue_item_state(EncodedPodRow::view(&stored.bytes)) else {
431 return Ok(blocked_delta);
432 };
433
434 if state.status != QueueItemStatus::Parked {
435 debug!(
436 queue = queue.0,
437 partition,
438 item = successor.0,
439 "the successor of a key was not parked when its predecessor finished"
440 );
441 return Ok(blocked_delta);
442 }
443
444 state.status = QueueItemStatus::Ready;
445 tx.set(&state_key, encode_queue_item_state(&state))?;
446 expose_due(tx, queue, partition, successor, &state)?;
447
448 Ok(blocked_delta)
449}
450
451fn expose_due(
452 tx: &mut SingleWriteTransaction<'_>,
453 queue: QueueId,
454 partition: u16,
455 row: RowNumber,
456 state: &QueueItemState,
457) -> Result<()> {
458 tx.set(&QueueDueKey::new(queue, partition, state.due(), row), EncodedPodRow::new(&[]).into_bytes())
459}
460
461fn read_counters(tx: &mut SingleWriteTransaction<'_>, lock_key: &QueuePartitionKey) -> Result<QueuePartitionCounters> {
462 Ok(tx.get(lock_key)?
463 .map(|stored| decode_queue_partition_counters(EncodedPodRow::view(&stored.bytes)))
464 .unwrap_or_default())
465}
466
467fn adjust_counters(
468 tx: &mut SingleWriteTransaction<'_>,
469 lock_key: &QueuePartitionKey,
470 applied: u64,
471 requeued: u64,
472 blocked_delta: i64,
473) -> Result<()> {
474 let mut counters = read_counters(tx, lock_key)?;
475 counters.in_flight = counters.in_flight.saturating_sub(applied);
476 counters.depth += requeued;
477 counters.blocked_keys = counters.blocked_keys.saturating_add_signed(blocked_delta);
478 tx.set(lock_key, encode_queue_partition_counters(&counters))?;
479
480 Ok(())
481}