Skip to main content

deepstrike_core/signals/
queue.rs

1use std::cmp::Ordering;
2use std::collections::BinaryHeap;
3
4use compact_str::CompactString;
5
6use crate::types::signal::{RuntimeSignal, Urgency};
7
8/// Wrapper for priority ordering: higher urgency first, then older timestamp first.
9#[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
47/// Priority queue for runtime signals. Internal to the signals module.
48pub(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    /// Admit using the queue's deterministic capacity policy without TTL cleanup.
68    #[cfg(test)]
69    pub(super) fn push(&mut self, signal: RuntimeSignal) -> bool {
70        self.admit(signal).admitted
71    }
72
73    /// At capacity, a new signal may replace only a strictly lower-urgency entry. Among the lowest
74    /// urgency, the newest entry is displaced so an older waiter never loses its place to an
75    /// equal-priority arrival. The router calls [`Self::expire`] before admission.
76    #[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    /// Remove entries whose timestamp plus the configured TTL has elapsed.
193    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    /// Promote each due queued signal at most once, then rebuild priority ordering.
217    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}