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
11pub 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 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 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 pub fn next(&mut self) -> Option<RuntimeSignal> {
168 self.queue.pop()
169 }
170
171 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 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 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}