reifydb_engine/queue/
interceptor.rs1use 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}