Skip to main content

reifydb_engine/queue/
interceptor.rs

1// SPDX-License-Identifier: Apache-2.0
2// Copyright (c) 2026 ReifyDB
3
4use std::collections::BTreeMap;
5
6use reifydb_core::{common::CommitVersion, interface::catalog::id::QueueId};
7use reifydb_runtime::context::clock::Clock;
8use reifydb_transaction::{
9	change::{QueueAckTransition, QueueRowAck, RowChange},
10	interceptor::transaction::{PostCommitContext, PostCommitInterceptor},
11	queue::scheduling::{QueueAdmission, admit_ready_items, apply_ack_transitions},
12	single::SingleTransaction,
13};
14use reifydb_value::value::datetime::DateTime;
15use tracing::{error, instrument};
16
17use crate::{Result, queue::wake::QueueWakeRegistry};
18
19pub struct QueueSchedulingInterceptor {
20	single: SingleTransaction,
21	wake: QueueWakeRegistry,
22	clock: Clock,
23}
24
25impl QueueSchedulingInterceptor {
26	pub fn new(single: SingleTransaction, wake: QueueWakeRegistry, clock: Clock) -> Self {
27		Self {
28			single,
29			wake,
30			clock,
31		}
32	}
33
34	#[instrument(
35		name = "queue::interceptor::enqueue",
36		level = "debug",
37		skip_all,
38		fields(queue = queue.0, partition = partition, items = items.len())
39	)]
40	fn admit(&self, queue: QueueId, partition: u16, items: &[QueueAdmission], version: CommitVersion) {
41		if let Err(err) = admit_ready_items(&self.single, queue, partition, items) {
42			error!(
43				queue = queue.0,
44				partition,
45				version = version.0,
46				items = items.len(),
47				error = %err,
48				"queue scheduling handoff failed; hydration will recover these items at next boot"
49			);
50			return;
51		}
52
53		let now = self.clock.now();
54		self.wake.nudge(queue, items.iter().filter(|item| is_due(item.not_before, now)).count());
55	}
56
57	#[instrument(
58		name = "queue::interceptor::ack",
59		level = "debug",
60		skip_all,
61		fields(queue = queue.0, partition = partition, items = items.len())
62	)]
63	fn ack(&self, queue: QueueId, partition: u16, items: &[QueueRowAck], version: CommitVersion) {
64		if let Err(err) = apply_ack_transitions(&self.single, queue, partition, items) {
65			error!(
66				queue = queue.0,
67				partition,
68				version = version.0,
69				items = items.len(),
70				error = %err,
71				"queue ack transition failed; the lease will expire and the item is redelivered"
72			);
73			return;
74		}
75
76		let now = self.clock.now();
77		self.wake.nudge(queue, items.iter().filter(|ack| releases_work(ack, now)).count());
78	}
79}
80
81fn is_due(not_before: Option<DateTime>, now: DateTime) -> bool {
82	not_before.is_none_or(|instant| instant <= now)
83}
84
85fn releases_work(ack: &QueueRowAck, now: DateTime) -> bool {
86	match &ack.transition {
87		QueueAckTransition::Retry {
88			backoff_until,
89		} => *backoff_until <= now,
90		QueueAckTransition::Done | QueueAckTransition::Dead => ack.key_hash.is_some(),
91	}
92}
93
94impl PostCommitInterceptor for QueueSchedulingInterceptor {
95	fn intercept(&self, ctx: &mut PostCommitContext) -> Result<()> {
96		if ctx.version == CommitVersion(0) || ctx.row_changes.is_empty() {
97			return Ok(());
98		}
99
100		let mut admissions: BTreeMap<(QueueId, u16), Vec<QueueAdmission>> = BTreeMap::new();
101		let mut acks: BTreeMap<(QueueId, u16), Vec<QueueRowAck>> = BTreeMap::new();
102		for change in &ctx.row_changes {
103			match change {
104				RowChange::QueueInsert(insertion) => {
105					admissions.entry((insertion.queue_id, insertion.partition)).or_default().push(
106						QueueAdmission {
107							row: insertion.row_number,
108							key_hash: insertion.key_hash,
109							not_before: insertion.not_before,
110						},
111					);
112				}
113				RowChange::QueueAck(ack) => {
114					acks.entry((ack.queue_id, ack.partition)).or_default().push(ack.clone());
115				}
116				RowChange::TableInsert(_) => {}
117			}
118		}
119
120		for ((queue, partition), items) in admissions {
121			self.admit(queue, partition, &items, ctx.version);
122		}
123
124		for ((queue, partition), items) in acks {
125			self.ack(queue, partition, &items, ctx.version);
126		}
127
128		Ok(())
129	}
130}