Skip to main content

reifydb_transaction/queue/
scheduling.rs

1// SPDX-License-Identifier: Apache-2.0
2// Copyright (c) 2026 ReifyDB
3
4use 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}