Skip to main content

deepstrike_core/signals/
router.rs

1use std::collections::{HashSet, VecDeque};
2
3use compact_str::CompactString;
4
5use super::attention::UrgencyBasedPolicy;
6use super::queue::{QueuedSignalRuntimeState, SignalQueue};
7use crate::scheduler::tcb::TaskLifecycle;
8use crate::types::policy::SignalDisposition;
9use crate::types::signal::RuntimeSignal;
10
11/// Signal router: dedup set + urgency-based attention + bounded priority queue.
12pub struct SignalRouter {
13    seen: HashSet<CompactString>,
14    seen_order: VecDeque<CompactString>,
15    dedupe_capacity: usize,
16    queue: SignalQueue,
17    attention: UrgencyBasedPolicy,
18    ttl_ms: Option<u64>,
19    deadline_escalation: bool,
20}
21
22#[derive(Debug, Clone)]
23pub(crate) struct SignalRouterRuntimeState {
24    pub queued: Vec<QueuedSignalRuntimeState>,
25    pub seen_order: Vec<CompactString>,
26}
27
28#[derive(Debug, Clone, PartialEq)]
29pub struct SignalRouteOutcome {
30    pub disposition: SignalDisposition,
31    pub displaced_signal_id: Option<String>,
32    pub expired_signal_ids: Vec<String>,
33}
34
35impl SignalRouter {
36    pub const DEFAULT_DEDUPE_CAPACITY: usize = 256;
37
38    pub fn new(max_queue_size: usize) -> Self {
39        Self::with_policy(max_queue_size, None, false)
40    }
41
42    pub fn with_policy(
43        max_queue_size: usize,
44        ttl_ms: Option<u64>,
45        deadline_escalation: bool,
46    ) -> Self {
47        Self {
48            seen: HashSet::with_capacity(Self::DEFAULT_DEDUPE_CAPACITY),
49            seen_order: VecDeque::with_capacity(Self::DEFAULT_DEDUPE_CAPACITY),
50            dedupe_capacity: Self::DEFAULT_DEDUPE_CAPACITY,
51            queue: SignalQueue::new(max_queue_size),
52            attention: UrgencyBasedPolicy,
53            ttl_ms,
54            deadline_escalation,
55        }
56    }
57
58    /// Ingest a signal. Returns the disposition after dedup + attention evaluation.
59    /// `Queue` dispositions are buffered; if the queue is full, returns `Dropped`
60    /// so the SDK can apply backpressure or surface the loss to telemetry.
61    /// All other dispositions are returned directly to the caller.
62    pub fn ingest(&mut self, signal: RuntimeSignal, lifecycle: TaskLifecycle) -> SignalDisposition {
63        let now_ms = signal.timestamp_ms;
64        self.ingest_at(signal, lifecycle, now_ms).disposition
65    }
66
67    pub fn ingest_at(
68        &mut self,
69        mut signal: RuntimeSignal,
70        lifecycle: TaskLifecycle,
71        now_ms: u64,
72    ) -> SignalRouteOutcome {
73        let expired_signal_ids = self.expire(now_ms);
74        let dedupe_key = signal.dedupe_key.clone();
75        if let Some(ref key) = dedupe_key {
76            if self.seen.contains(key) {
77                return SignalRouteOutcome {
78                    disposition: SignalDisposition::Ignore,
79                    displaced_signal_id: None,
80                    expired_signal_ids,
81                };
82            }
83        }
84
85        let deadline_escalated = self.deadline_escalation
86            && signal
87                .deadline_ms
88                .is_some_and(|deadline_ms| now_ms >= deadline_ms);
89        if deadline_escalated {
90            signal.urgency = escalate_one_tier(signal.urgency);
91        }
92
93        let disposition = self.attention.evaluate(&signal, lifecycle);
94
95        if disposition == SignalDisposition::Queue {
96            let admission = self
97                .queue
98                .admit_with_deadline_state(signal, deadline_escalated);
99            for key in &admission.displaced_dedupe_keys {
100                self.release_dedupe_key(key);
101            }
102            let displaced_signal_id = admission
103                .displaced
104                .as_ref()
105                .map(|displaced| displaced.id.to_string());
106            if !admission.admitted {
107                return SignalRouteOutcome {
108                    disposition: SignalDisposition::Dropped,
109                    displaced_signal_id: None,
110                    expired_signal_ids,
111                };
112            }
113            if let Some(key) = dedupe_key {
114                self.commit_dedupe(key);
115            }
116            return SignalRouteOutcome {
117                disposition,
118                displaced_signal_id,
119                expired_signal_ids,
120            };
121        }
122
123        if let Some(key) = dedupe_key {
124            self.commit_dedupe(key);
125        }
126
127        SignalRouteOutcome {
128            disposition,
129            displaced_signal_id: None,
130            expired_signal_ids,
131        }
132    }
133
134    /// Expire queued signals against the journaled clock before either admission or delivery.
135    pub fn expire(&mut self, now_ms: u64) -> Vec<String> {
136        let expired = self.queue.expire(now_ms, self.ttl_ms);
137        if self.deadline_escalation {
138            self.queue.escalate_deadlines(now_ms);
139        }
140        for (_, dedupe_keys) in &expired {
141            for key in dedupe_keys {
142                self.release_dedupe_key(key);
143            }
144        }
145        expired
146            .into_iter()
147            .map(|(signal, _)| signal.id.to_string())
148            .collect()
149    }
150
151    fn commit_dedupe(&mut self, key: CompactString) {
152        if self.seen_order.len() == self.dedupe_capacity {
153            if let Some(expired) = self.seen_order.pop_front() {
154                self.seen.remove(&expired);
155            }
156        }
157        self.seen.insert(key.clone());
158        self.seen_order.push_back(key);
159    }
160
161    fn release_dedupe_key(&mut self, key: &CompactString) {
162        self.seen.remove(key);
163        self.seen_order.retain(|seen_key| seen_key != key);
164    }
165
166    /// Pull next queued signal.
167    pub fn next(&mut self) -> Option<RuntimeSignal> {
168        self.queue.pop()
169    }
170
171    /// Number of queued signals.
172    pub fn depth(&self) -> usize {
173        self.queue.len()
174    }
175
176    pub(crate) fn checkpoint_state(&self) -> SignalRouterRuntimeState {
177        SignalRouterRuntimeState {
178            queued: self.queue.checkpoint_entries(),
179            seen_order: self.seen_order.iter().cloned().collect(),
180        }
181    }
182
183    pub(crate) fn restore_state(&mut self, state: SignalRouterRuntimeState) -> Result<(), String> {
184        if state.seen_order.len() > self.dedupe_capacity {
185            return Err(format!(
186                "checkpoint carries {} signal dedupe keys for capacity {}",
187                state.seen_order.len(),
188                self.dedupe_capacity
189            ));
190        }
191        let mut seen = HashSet::with_capacity(state.seen_order.len());
192        for key in &state.seen_order {
193            if !seen.insert(key.clone()) {
194                return Err(format!(
195                    "checkpoint carries duplicate signal dedupe key {key:?}"
196                ));
197            }
198        }
199        let mut queued_keys = HashSet::new();
200        for queued in &state.queued {
201            if let Some(primary) = queued.signal.dedupe_key.as_ref()
202                && !queued.dedupe_keys.contains(primary)
203            {
204                return Err(format!(
205                    "queued signal {:?} omits its primary dedupe key {primary:?}",
206                    queued.signal.id
207                ));
208            }
209            for key in &queued.dedupe_keys {
210                if !seen.contains(key) {
211                    return Err(format!(
212                        "queued signal {:?} carries uncommitted dedupe key {key:?}",
213                        queued.signal.id
214                    ));
215                }
216                if !queued_keys.insert(key.clone()) {
217                    return Err(format!(
218                        "signal dedupe key {key:?} belongs to more than one queue entry"
219                    ));
220                }
221            }
222        }
223        self.queue.restore_entries(state.queued)?;
224        self.seen = seen;
225        self.seen_order = state.seen_order.into();
226        Ok(())
227    }
228
229    /// Clear the dedup set (call at session boundaries to prevent unbounded growth).
230    pub fn clear_dedup(&mut self) {
231        self.seen.clear();
232        self.seen_order.clear();
233    }
234
235    #[cfg(test)]
236    fn dedupe_len(&self) -> usize {
237        self.seen.len()
238    }
239}
240
241fn escalate_one_tier(urgency: crate::types::signal::Urgency) -> crate::types::signal::Urgency {
242    use crate::types::signal::Urgency;
243    match urgency {
244        Urgency::Low => Urgency::Normal,
245        Urgency::Normal => Urgency::High,
246        Urgency::High | Urgency::Critical => Urgency::Critical,
247    }
248}
249
250#[cfg(test)]
251mod tests {
252    use super::*;
253    use crate::scheduler::tcb::TaskLifecycle;
254    use crate::types::signal::{SignalSource, SignalType, Urgency};
255
256    #[test]
257    fn deduplicates_signals() {
258        let mut router = SignalRouter::new(100);
259        let sig = RuntimeSignal::new(
260            SignalSource::Cron,
261            SignalType::Event,
262            Urgency::Normal,
263            "tick",
264        )
265        .with_dedupe("cron-tick-1");
266
267        let d1 = router.ingest(sig.clone(), TaskLifecycle::Running);
268        assert_ne!(d1, SignalDisposition::Ignore);
269
270        let d2 = router.ingest(sig, TaskLifecycle::Running);
271        assert_eq!(d2, SignalDisposition::Ignore);
272    }
273
274    #[test]
275    fn normal_signal_queued() {
276        let mut router = SignalRouter::new(100);
277        let sig = RuntimeSignal::new(
278            SignalSource::Cron,
279            SignalType::Event,
280            Urgency::Normal,
281            "job",
282        );
283
284        let d = router.ingest(sig, TaskLifecycle::Running);
285        assert_eq!(d, SignalDisposition::Queue);
286        assert_eq!(router.depth(), 1);
287        assert!(router.next().is_some());
288    }
289
290    #[test]
291    fn interrupt_signals_not_queued() {
292        let mut router = SignalRouter::new(100);
293        let sig = RuntimeSignal::new(
294            SignalSource::Gateway,
295            SignalType::Alert,
296            Urgency::Critical,
297            "fire",
298        );
299
300        let d = router.ingest(sig, TaskLifecycle::Running);
301        assert_eq!(d, SignalDisposition::InterruptNow);
302        assert_eq!(router.depth(), 0);
303    }
304
305    #[test]
306    fn full_queue_drops_signal() {
307        let mut router = SignalRouter::new(1);
308        let s1 = RuntimeSignal::new(
309            SignalSource::Cron,
310            SignalType::Event,
311            Urgency::Normal,
312            "first",
313        );
314        let s2 = RuntimeSignal::new(
315            SignalSource::Cron,
316            SignalType::Event,
317            Urgency::Normal,
318            "second",
319        );
320
321        assert_eq!(
322            router.ingest(s1, TaskLifecycle::Running),
323            SignalDisposition::Queue
324        );
325        assert_eq!(
326            router.ingest(s2, TaskLifecycle::Running),
327            SignalDisposition::Dropped
328        );
329    }
330
331    #[test]
332    fn clear_dedup_allows_reingest() {
333        let mut router = SignalRouter::new(100);
334        let sig = RuntimeSignal::new(
335            SignalSource::Cron,
336            SignalType::Event,
337            Urgency::Normal,
338            "tick",
339        )
340        .with_dedupe("key-1");
341
342        router.ingest(sig.clone(), TaskLifecycle::Running);
343        assert_eq!(
344            router.ingest(sig.clone(), TaskLifecycle::Running),
345            SignalDisposition::Ignore
346        );
347
348        router.clear_dedup();
349        assert_ne!(
350            router.ingest(sig, TaskLifecycle::Running),
351            SignalDisposition::Ignore
352        );
353    }
354
355    #[test]
356    fn dedupe_window_is_bounded_and_expires_oldest_key() {
357        let mut router = SignalRouter::new(1);
358        for index in 0..=SignalRouter::DEFAULT_DEDUPE_CAPACITY {
359            let signal =
360                RuntimeSignal::new(SignalSource::Cron, SignalType::Event, Urgency::Low, "tick")
361                    .with_dedupe(format!("key-{index}"));
362            assert_ne!(
363                router.ingest(signal, TaskLifecycle::Running),
364                SignalDisposition::Ignore
365            );
366        }
367
368        assert_eq!(router.dedupe_len(), SignalRouter::DEFAULT_DEDUPE_CAPACITY);
369        let expired =
370            RuntimeSignal::new(SignalSource::Cron, SignalType::Event, Urgency::Low, "tick")
371                .with_dedupe("key-0");
372        assert_ne!(
373            router.ingest(expired, TaskLifecycle::Running),
374            SignalDisposition::Ignore
375        );
376    }
377
378    #[test]
379    fn dropped_signal_does_not_commit_its_dedupe_key() {
380        let mut router = SignalRouter::new(1);
381        let admitted = RuntimeSignal::new(
382            SignalSource::Cron,
383            SignalType::Event,
384            Urgency::Normal,
385            "admitted",
386        );
387        let retryable = RuntimeSignal::new(
388            SignalSource::Cron,
389            SignalType::Event,
390            Urgency::Normal,
391            "retryable",
392        )
393        .with_dedupe("retryable-key");
394
395        assert_eq!(
396            router.ingest(admitted, TaskLifecycle::Running),
397            SignalDisposition::Queue
398        );
399        assert_eq!(
400            router.ingest(retryable.clone(), TaskLifecycle::Running),
401            SignalDisposition::Dropped
402        );
403        assert!(router.next().is_some());
404        assert_eq!(
405            router.ingest(retryable, TaskLifecycle::Running),
406            SignalDisposition::Queue
407        );
408    }
409
410    #[test]
411    fn ttl_cleanup_precedes_urgency_displacement() {
412        let mut router = SignalRouter::with_policy(1, Some(10), false);
413        let fresh = RuntimeSignal::new(
414            SignalSource::Cron,
415            SignalType::Event,
416            Urgency::Normal,
417            "fresh",
418        )
419        .with_timestamp(30);
420
421        // Terminal forces the critical signal into the pending queue so TTL can be tested.
422        let stale_queued = RuntimeSignal::new(
423            SignalSource::Gateway,
424            SignalType::Alert,
425            Urgency::Critical,
426            "stale queued",
427        )
428        .with_timestamp(10);
429        assert_eq!(
430            router
431                .ingest_at(
432                    stale_queued,
433                    TaskLifecycle::Done(crate::types::result::TerminationReason::Completed),
434                    10,
435                )
436                .disposition,
437            SignalDisposition::Queue
438        );
439
440        let outcome = router.ingest_at(fresh, TaskLifecycle::Running, 30);
441        assert_eq!(outcome.disposition, SignalDisposition::Queue);
442        assert_eq!(outcome.expired_signal_ids.len(), 1);
443        assert!(outcome.displaced_signal_id.is_none());
444    }
445
446    #[test]
447    fn expiration_releases_dedupe_key_before_redelivery_is_checked() {
448        let mut router = SignalRouter::with_policy(1, Some(10), false);
449        let first = RuntimeSignal::new(
450            SignalSource::Cron,
451            SignalType::Event,
452            Urgency::Normal,
453            "first lease",
454        )
455        .with_timestamp(10)
456        .with_dedupe("leased-work");
457        assert_eq!(
458            router
459                .ingest_at(first, TaskLifecycle::Running, 10)
460                .disposition,
461            SignalDisposition::Queue
462        );
463
464        let redelivery = RuntimeSignal::new(
465            SignalSource::Cron,
466            SignalType::Event,
467            Urgency::Normal,
468            "redelivery",
469        )
470        .with_timestamp(30)
471        .with_dedupe("leased-work");
472        let outcome = router.ingest_at(redelivery, TaskLifecycle::Running, 30);
473
474        assert_eq!(outcome.disposition, SignalDisposition::Queue);
475        assert_eq!(outcome.expired_signal_ids.len(), 1);
476    }
477
478    #[test]
479    fn displacement_releases_the_evicted_signals_dedupe_key() {
480        let mut router = SignalRouter::new(2);
481        let old = RuntimeSignal::new(
482            SignalSource::Cron,
483            SignalType::Event,
484            Urgency::Low,
485            "old low",
486        )
487        .with_timestamp(1)
488        .with_dedupe("old-low");
489        let newest = RuntimeSignal::new(
490            SignalSource::Cron,
491            SignalType::Event,
492            Urgency::Low,
493            "new low",
494        )
495        .with_timestamp(2)
496        .with_dedupe("new-low");
497        let newest_id = newest.id.to_string();
498        let critical = RuntimeSignal::new(
499            SignalSource::Gateway,
500            SignalType::Alert,
501            Urgency::Critical,
502            "critical",
503        )
504        .with_timestamp(3);
505        let terminal = TaskLifecycle::Done(crate::types::result::TerminationReason::Completed);
506        assert_eq!(router.ingest(old, terminal), SignalDisposition::Queue);
507        assert_eq!(router.ingest(newest, terminal), SignalDisposition::Queue);
508
509        let outcome = router.ingest_at(critical, terminal, 3);
510        assert_eq!(outcome.disposition, SignalDisposition::Queue);
511        assert_eq!(
512            outcome.displaced_signal_id.as_deref(),
513            Some(newest_id.as_str())
514        );
515
516        let redelivery = RuntimeSignal::new(
517            SignalSource::Cron,
518            SignalType::Event,
519            Urgency::Low,
520            "new low redelivery",
521        )
522        .with_timestamp(4)
523        .with_dedupe("new-low");
524        assert_ne!(
525            router.ingest(redelivery, TaskLifecycle::Running),
526            SignalDisposition::Ignore
527        );
528    }
529
530    #[test]
531    fn due_deadline_escalates_exactly_one_urgency_tier_when_enabled() {
532        let mut router = SignalRouter::with_policy(4, None, true);
533        let due = RuntimeSignal::new(
534            SignalSource::Gateway,
535            SignalType::Event,
536            Urgency::Normal,
537            "due work",
538        )
539        .with_timestamp(10)
540        .with_deadline(20);
541
542        let outcome = router.ingest_at(due, TaskLifecycle::Running, 20);
543
544        assert_eq!(outcome.disposition, SignalDisposition::Interrupt);
545        assert_eq!(router.depth(), 0);
546    }
547
548    #[test]
549    fn deadline_is_inert_when_escalation_policy_is_disabled() {
550        let mut router = SignalRouter::with_policy(4, None, false);
551        let due = RuntimeSignal::new(
552            SignalSource::Gateway,
553            SignalType::Event,
554            Urgency::Normal,
555            "due work",
556        )
557        .with_timestamp(10)
558        .with_deadline(20);
559
560        let outcome = router.ingest_at(due, TaskLifecycle::Running, 20);
561
562        assert_eq!(outcome.disposition, SignalDisposition::Queue);
563        assert_eq!(router.next().unwrap().urgency, Urgency::Normal);
564    }
565
566    #[test]
567    fn queued_signals_coalesce_without_consuming_capacity_or_dedupe_semantics() {
568        let mut router = SignalRouter::with_policy(1, None, false);
569        let first = RuntimeSignal::new(
570            SignalSource::Cron,
571            SignalType::Event,
572            Urgency::Normal,
573            "first sample",
574        )
575        .with_timestamp(10)
576        .with_deadline(200)
577        .with_coalesce("telemetry")
578        .with_dedupe("event-1");
579        let second = RuntimeSignal::new(
580            SignalSource::Cron,
581            SignalType::Event,
582            Urgency::Normal,
583            "second sample",
584        )
585        .with_timestamp(20)
586        .with_deadline(100)
587        .with_coalesce("telemetry");
588
589        assert_eq!(
590            router.ingest(first, TaskLifecycle::Running),
591            SignalDisposition::Queue
592        );
593        assert_eq!(
594            router.ingest(second, TaskLifecycle::Running),
595            SignalDisposition::Queue
596        );
597        assert_eq!(router.depth(), 1);
598
599        let duplicate = RuntimeSignal::new(
600            SignalSource::Cron,
601            SignalType::Event,
602            Urgency::Normal,
603            "dedupe still wins",
604        )
605        .with_timestamp(30)
606        .with_coalesce("telemetry")
607        .with_dedupe("event-1");
608        assert_eq!(
609            router.ingest(duplicate, TaskLifecycle::Running),
610            SignalDisposition::Ignore
611        );
612
613        let merged = router.next().unwrap();
614        assert_eq!(merged.summary.as_str(), "first sample");
615        assert_eq!(merged.timestamp_ms, 10);
616        assert_eq!(merged.deadline_ms, Some(100));
617        assert_eq!(merged.coalesced_count, 2);
618    }
619
620    #[test]
621    fn expiration_releases_every_dedupe_key_merged_into_a_coalesced_entry() {
622        let mut router = SignalRouter::with_policy(1, Some(10), false);
623        let first = RuntimeSignal::new(
624            SignalSource::Cron,
625            SignalType::Event,
626            Urgency::Normal,
627            "first",
628        )
629        .with_timestamp(10)
630        .with_coalesce("batch")
631        .with_dedupe("event-1");
632        let second = RuntimeSignal::new(
633            SignalSource::Cron,
634            SignalType::Event,
635            Urgency::Normal,
636            "second",
637        )
638        .with_timestamp(11)
639        .with_coalesce("batch")
640        .with_dedupe("event-2");
641
642        assert_eq!(
643            router.ingest(first, TaskLifecycle::Running),
644            SignalDisposition::Queue
645        );
646        assert_eq!(
647            router.ingest(second, TaskLifecycle::Running),
648            SignalDisposition::Queue
649        );
650        assert_eq!(router.expire(30).len(), 1);
651
652        for key in ["event-1", "event-2"] {
653            let redelivery = RuntimeSignal::new(
654                SignalSource::Cron,
655                SignalType::Event,
656                Urgency::Normal,
657                "redelivery",
658            )
659            .with_timestamp(30)
660            .with_dedupe(key);
661            assert_eq!(
662                router.ingest(redelivery, TaskLifecycle::Running),
663                SignalDisposition::Queue
664            );
665            router.next();
666        }
667    }
668
669    #[test]
670    fn runtime_signal_wire_rejects_removed_topic_field() {
671        let encoded = serde_json::json!({
672            "id": uuid::Uuid::nil(),
673            "source": "custom",
674            "signal_type": "event",
675            "urgency": "normal",
676            "summary": "signal",
677            "payload": null,
678            "topic": "removed-field",
679            "timestamp_ms": 1
680        });
681
682        assert!(serde_json::from_value::<RuntimeSignal>(encoded).is_err());
683    }
684}