1use std::cmp::Ordering;
2use std::collections::BinaryHeap;
3
4use compact_str::CompactString;
5
6use crate::types::signal::{RuntimeSignal, Urgency};
7
8#[derive(Clone)]
10struct PrioritizedSignal {
11 urgency: Urgency,
12 timestamp_ms: u64,
13 deadline_escalated: bool,
14 dedupe_keys: Vec<CompactString>,
15 signal: RuntimeSignal,
16}
17
18#[derive(Debug, Clone)]
19pub(crate) struct QueuedSignalRuntimeState {
20 pub signal: RuntimeSignal,
21 pub deadline_escalated: bool,
22 pub dedupe_keys: Vec<CompactString>,
23}
24
25impl PartialEq for PrioritizedSignal {
26 fn eq(&self, other: &Self) -> bool {
27 self.signal.id == other.signal.id
28 }
29}
30impl Eq for PrioritizedSignal {}
31
32impl PartialOrd for PrioritizedSignal {
33 fn partial_cmp(&self, other: &Self) -> Option<Ordering> {
34 Some(self.cmp(other))
35 }
36}
37
38impl Ord for PrioritizedSignal {
39 fn cmp(&self, other: &Self) -> Ordering {
40 self.urgency
41 .cmp(&other.urgency)
42 .then_with(|| other.timestamp_ms.cmp(&self.timestamp_ms))
43 .then_with(|| self.signal.id.cmp(&other.signal.id))
44 }
45}
46
47pub(super) struct SignalQueue {
49 heap: BinaryHeap<PrioritizedSignal>,
50 max_size: usize,
51}
52
53pub(super) struct QueueAdmission {
54 pub(super) admitted: bool,
55 pub(super) displaced: Option<RuntimeSignal>,
56 pub(super) displaced_dedupe_keys: Vec<CompactString>,
57}
58
59impl SignalQueue {
60 pub(super) fn new(max_size: usize) -> Self {
61 Self {
62 heap: BinaryHeap::new(),
63 max_size,
64 }
65 }
66
67 #[cfg(test)]
69 pub(super) fn push(&mut self, signal: RuntimeSignal) -> bool {
70 self.admit(signal).admitted
71 }
72
73 #[cfg(test)]
77 pub(super) fn admit(&mut self, signal: RuntimeSignal) -> QueueAdmission {
78 self.admit_with_deadline_state(signal, false)
79 }
80
81 pub(super) fn admit_with_deadline_state(
82 &mut self,
83 signal: RuntimeSignal,
84 deadline_escalated: bool,
85 ) -> QueueAdmission {
86 if let Some(key) = signal.coalesce_key.as_ref() {
87 let existing_id = self
88 .heap
89 .iter()
90 .find(|queued| queued.signal.coalesce_key.as_ref() == Some(key))
91 .map(|queued| queued.signal.id.clone());
92 if let Some(existing_id) = existing_id {
93 let mut retained = BinaryHeap::with_capacity(self.heap.len());
94 for mut queued in self.heap.drain() {
95 if queued.signal.id == existing_id {
96 queued.signal.coalesced_count = queued
97 .signal
98 .coalesced_count
99 .max(1)
100 .saturating_add(signal.coalesced_count.max(1));
101 queued.signal.urgency = queued.signal.urgency.max(signal.urgency);
102 queued.urgency = queued.signal.urgency;
103 queued.signal.deadline_ms =
104 earliest_deadline(queued.signal.deadline_ms, signal.deadline_ms);
105 queued.deadline_escalated |= deadline_escalated;
106 if let Some(dedupe_key) = signal.dedupe_key.as_ref() {
107 if !queued.dedupe_keys.contains(dedupe_key) {
108 queued.dedupe_keys.push(dedupe_key.clone());
109 }
110 }
111 }
112 retained.push(queued);
113 }
114 self.heap = retained;
115 return QueueAdmission {
116 admitted: true,
117 displaced: None,
118 displaced_dedupe_keys: Vec::new(),
119 };
120 }
121 }
122
123 if self.heap.len() >= self.max_size {
124 let lowest = self.heap.iter().map(|queued| queued.urgency).min();
125 if lowest.is_none_or(|urgency| signal.urgency <= urgency) {
126 return QueueAdmission {
127 admitted: false,
128 displaced: None,
129 displaced_dedupe_keys: Vec::new(),
130 };
131 }
132
133 let lowest = lowest.expect("a full queue has a lowest urgency");
134 let displaced_id = self
135 .heap
136 .iter()
137 .filter(|queued| queued.urgency == lowest)
138 .max_by(|left, right| {
139 left.timestamp_ms
140 .cmp(&right.timestamp_ms)
141 .then_with(|| left.signal.id.cmp(&right.signal.id))
142 })
143 .map(|queued| queued.signal.id.clone())
144 .expect("a full queue has a displacement candidate");
145 let mut displaced = None;
146 let mut displaced_dedupe_keys = Vec::new();
147 let mut retained = BinaryHeap::with_capacity(self.heap.len());
148 for queued in self.heap.drain() {
149 if queued.signal.id == displaced_id {
150 displaced = Some(queued.signal);
151 displaced_dedupe_keys = queued.dedupe_keys;
152 } else {
153 retained.push(queued);
154 }
155 }
156 self.heap = retained;
157
158 let urgency = signal.urgency;
159 let timestamp_ms = signal.timestamp_ms;
160 let dedupe_keys = signal.dedupe_key.iter().cloned().collect();
161 self.heap.push(PrioritizedSignal {
162 urgency,
163 timestamp_ms,
164 deadline_escalated,
165 dedupe_keys,
166 signal,
167 });
168 return QueueAdmission {
169 admitted: true,
170 displaced,
171 displaced_dedupe_keys,
172 };
173 }
174
175 let urgency = signal.urgency;
176 let timestamp_ms = signal.timestamp_ms;
177 let dedupe_keys = signal.dedupe_key.iter().cloned().collect();
178 self.heap.push(PrioritizedSignal {
179 urgency,
180 timestamp_ms,
181 deadline_escalated,
182 dedupe_keys,
183 signal,
184 });
185 QueueAdmission {
186 admitted: true,
187 displaced: None,
188 displaced_dedupe_keys: Vec::new(),
189 }
190 }
191
192 pub(super) fn expire(
194 &mut self,
195 now_ms: u64,
196 ttl_ms: Option<u64>,
197 ) -> Vec<(RuntimeSignal, Vec<CompactString>)> {
198 let Some(ttl_ms) = ttl_ms else {
199 return Vec::new();
200 };
201 let mut expired = Vec::new();
202 let mut retained = BinaryHeap::with_capacity(self.heap.len());
203 for queued in self.heap.drain() {
204 let has_timestamp = queued.timestamp_ms > 0;
205 let reached_expiry = now_ms >= queued.timestamp_ms.saturating_add(ttl_ms);
206 if has_timestamp && reached_expiry {
207 expired.push((queued.signal, queued.dedupe_keys));
208 } else {
209 retained.push(queued);
210 }
211 }
212 self.heap = retained;
213 expired
214 }
215
216 pub(super) fn escalate_deadlines(&mut self, now_ms: u64) {
218 let mut rebuilt = BinaryHeap::with_capacity(self.heap.len());
219 for mut queued in self.heap.drain() {
220 let due = queued
221 .signal
222 .deadline_ms
223 .is_some_and(|deadline_ms| now_ms >= deadline_ms);
224 if due && !queued.deadline_escalated {
225 queued.signal.urgency = escalate_one_tier(queued.signal.urgency);
226 queued.urgency = queued.signal.urgency;
227 queued.deadline_escalated = true;
228 }
229 rebuilt.push(queued);
230 }
231 self.heap = rebuilt;
232 }
233
234 pub(super) fn pop(&mut self) -> Option<RuntimeSignal> {
235 self.heap.pop().map(|ps| ps.signal)
236 }
237
238 pub(super) fn len(&self) -> usize {
239 self.heap.len()
240 }
241
242 pub(super) fn checkpoint_entries(&self) -> Vec<QueuedSignalRuntimeState> {
243 let mut heap = self.heap.clone();
244 let mut entries = Vec::with_capacity(heap.len());
245 while let Some(queued) = heap.pop() {
246 entries.push(QueuedSignalRuntimeState {
247 signal: queued.signal,
248 deadline_escalated: queued.deadline_escalated,
249 dedupe_keys: queued.dedupe_keys,
250 });
251 }
252 entries
253 }
254
255 pub(super) fn restore_entries(
256 &mut self,
257 entries: Vec<QueuedSignalRuntimeState>,
258 ) -> Result<(), String> {
259 if entries.len() > self.max_size {
260 return Err(format!(
261 "checkpoint carries {} queued signals for capacity {}",
262 entries.len(),
263 self.max_size
264 ));
265 }
266 let mut ids = std::collections::HashSet::with_capacity(entries.len());
267 let mut heap = BinaryHeap::with_capacity(entries.len());
268 for entry in entries {
269 if !ids.insert(entry.signal.id.clone()) {
270 return Err(format!(
271 "checkpoint carries duplicate queued signal id {:?}",
272 entry.signal.id
273 ));
274 }
275 heap.push(PrioritizedSignal {
276 urgency: entry.signal.urgency,
277 timestamp_ms: entry.signal.timestamp_ms,
278 deadline_escalated: entry.deadline_escalated,
279 dedupe_keys: entry.dedupe_keys,
280 signal: entry.signal,
281 });
282 }
283 self.heap = heap;
284 Ok(())
285 }
286}
287
288fn earliest_deadline(left: Option<u64>, right: Option<u64>) -> Option<u64> {
289 match (left, right) {
290 (Some(left), Some(right)) => Some(left.min(right)),
291 (Some(deadline), None) | (None, Some(deadline)) => Some(deadline),
292 (None, None) => None,
293 }
294}
295
296fn escalate_one_tier(urgency: Urgency) -> Urgency {
297 match urgency {
298 Urgency::Low => Urgency::Normal,
299 Urgency::Normal => Urgency::High,
300 Urgency::High | Urgency::Critical => Urgency::Critical,
301 }
302}
303
304#[cfg(test)]
305mod tests {
306 use super::*;
307 use crate::types::signal::{SignalSource, SignalType};
308
309 #[test]
310 fn higher_urgency_dequeued_first() {
311 let mut q = SignalQueue::new(10);
312 q.push(
313 RuntimeSignal::new(SignalSource::Cron, SignalType::Event, Urgency::Low, "low")
314 .with_timestamp(1),
315 );
316 q.push(
317 RuntimeSignal::new(
318 SignalSource::Gateway,
319 SignalType::Alert,
320 Urgency::Critical,
321 "crit",
322 )
323 .with_timestamp(2),
324 );
325 q.push(
326 RuntimeSignal::new(
327 SignalSource::Cron,
328 SignalType::Event,
329 Urgency::Normal,
330 "norm",
331 )
332 .with_timestamp(3),
333 );
334
335 assert_eq!(q.pop().unwrap().urgency, Urgency::Critical);
336 assert_eq!(q.pop().unwrap().urgency, Urgency::Normal);
337 assert_eq!(q.pop().unwrap().urgency, Urgency::Low);
338 }
339
340 #[test]
341 fn respects_max_size() {
342 let mut q = SignalQueue::new(1);
343 assert!(
344 q.push(
345 RuntimeSignal::new(SignalSource::Cron, SignalType::Event, Urgency::Low, "a")
346 .with_timestamp(1)
347 )
348 );
349 assert!(
350 !q.push(
351 RuntimeSignal::new(SignalSource::Cron, SignalType::Event, Urgency::Low, "b")
352 .with_timestamp(2)
353 )
354 );
355 }
356
357 #[test]
358 fn same_urgency_older_first() {
359 let mut q = SignalQueue::new(10);
360 q.push(
361 RuntimeSignal::new(
362 SignalSource::Cron,
363 SignalType::Event,
364 Urgency::Normal,
365 "newer",
366 )
367 .with_timestamp(100),
368 );
369 q.push(
370 RuntimeSignal::new(
371 SignalSource::Cron,
372 SignalType::Event,
373 Urgency::Normal,
374 "older",
375 )
376 .with_timestamp(1),
377 );
378
379 assert_eq!(q.pop().unwrap().summary.as_str(), "older");
380 }
381
382 #[test]
383 fn full_queue_accepts_strictly_higher_urgency_and_preserves_oldest_lowest() {
384 let mut q = SignalQueue::new(2);
385 assert!(
386 q.push(
387 RuntimeSignal::new(SignalSource::Cron, SignalType::Event, Urgency::Low, "old")
388 .with_timestamp(1)
389 )
390 );
391 assert!(
392 q.push(
393 RuntimeSignal::new(SignalSource::Cron, SignalType::Event, Urgency::Low, "new")
394 .with_timestamp(2)
395 )
396 );
397
398 let admission = q.admit(
399 RuntimeSignal::new(
400 SignalSource::Gateway,
401 SignalType::Alert,
402 Urgency::Critical,
403 "critical",
404 )
405 .with_timestamp(3),
406 );
407 assert!(admission.admitted);
408 assert_eq!(
409 admission
410 .displaced
411 .as_ref()
412 .map(|signal| signal.summary.as_str()),
413 Some("new")
414 );
415
416 assert_eq!(q.pop().unwrap().summary.as_str(), "critical");
417 assert_eq!(q.pop().unwrap().summary.as_str(), "old");
418 }
419
420 #[test]
421 fn expired_entries_are_removed_before_capacity_is_evaluated() {
422 let mut q = SignalQueue::new(1);
423 assert!(
424 q.push(
425 RuntimeSignal::new(
426 SignalSource::Cron,
427 SignalType::Event,
428 Urgency::Critical,
429 "stale"
430 )
431 .with_timestamp(10)
432 )
433 );
434
435 let expired = q.expire(30, Some(10));
436 let admission = q.admit(
437 RuntimeSignal::new(SignalSource::Cron, SignalType::Event, Urgency::Low, "fresh")
438 .with_timestamp(30),
439 );
440
441 assert!(admission.admitted);
442 assert!(admission.displaced.is_none());
443 assert_eq!(expired.len(), 1);
444 assert_eq!(expired[0].0.summary.as_str(), "stale");
445 assert_eq!(q.pop().unwrap().summary.as_str(), "fresh");
446 }
447
448 #[test]
449 fn coalescing_keeps_first_identity_and_combines_policy_inputs() {
450 let mut q = SignalQueue::new(1);
451 let first = RuntimeSignal::new(
452 SignalSource::Cron,
453 SignalType::Event,
454 Urgency::Normal,
455 "first",
456 )
457 .with_timestamp(10)
458 .with_deadline(200)
459 .with_coalesce("batch");
460 let first_id = first.id.clone();
461 assert!(q.admit(first).admitted);
462
463 let second = RuntimeSignal::new(
464 SignalSource::Cron,
465 SignalType::Event,
466 Urgency::High,
467 "second",
468 )
469 .with_timestamp(20)
470 .with_deadline(100)
471 .with_coalesce("batch");
472 let admission = q.admit(second);
473
474 assert!(admission.admitted);
475 assert!(admission.displaced.is_none());
476 assert_eq!(q.len(), 1);
477 let merged = q.pop().unwrap();
478 assert_eq!(merged.id, first_id);
479 assert_eq!(merged.urgency, Urgency::High);
480 assert_eq!(merged.deadline_ms, Some(100));
481 assert_eq!(merged.coalesced_count, 2);
482 }
483}