1use std::{
2 collections::{HashMap, VecDeque},
3 sync::LazyLock,
4};
5
6use dora_message::{
7 config::{DEFAULT_QUEUE_SIZE, QueuePolicy},
8 daemon_to_node::NodeEvent,
9 id::DataId,
10 metadata::{
11 GOAL_ID, GOAL_STATUS, MetadataParameters, REQUEST_ID, carries_pattern_correlation,
12 get_string_param,
13 },
14};
15
16use super::thread::EventItem;
17
18fn event_parameters(event: &EventItem) -> Option<&MetadataParameters> {
25 match event {
26 EventItem::NodeEvent {
27 event: NodeEvent::Input { metadata, .. },
28 ..
29 } => Some(&metadata.parameters),
30 EventItem::ZenohInput { metadata, .. } => Some(&metadata.parameters),
31 _ => None,
32 }
33}
34
35fn is_correlated(event: &EventItem) -> bool {
43 event_parameters(event).is_some_and(carries_pattern_correlation)
50}
51
52fn is_stop(event: &EventItem) -> bool {
60 matches!(
61 event,
62 EventItem::NodeEvent {
63 event: NodeEvent::Stop,
64 ..
65 }
66 )
67}
68
69enum Eviction {
71 RemoveAt(usize),
73 DropIncoming,
77 DropCorrelatedLoud(usize),
81}
82
83fn select_eviction(queue: &VecDeque<EventItem>, incoming: &EventItem) -> Eviction {
90 if let Some(idx) = queue.iter().position(|e| !is_correlated(e) && !is_stop(e)) {
92 return Eviction::RemoveAt(idx);
93 }
94 if !is_correlated(incoming) && !is_stop(incoming) {
97 return Eviction::DropIncoming;
98 }
99 match queue.iter().position(|e| !is_stop(e)) {
104 Some(idx) => Eviction::DropCorrelatedLoud(idx),
105 None => Eviction::DropIncoming,
106 }
107}
108
109fn log_correlation_drop(event_id: &DataId, dropped: &EventItem) {
113 let Some(params) = event_parameters(dropped) else {
114 return;
115 };
116 let request_id = get_string_param(params, REQUEST_ID);
117 let goal_id = get_string_param(params, GOAL_ID);
118 let goal_status = get_string_param(params, GOAL_STATUS);
119 tracing::error!(
120 input = %event_id,
121 ?request_id,
122 ?goal_id,
123 ?goal_status,
124 "queue full of correlated messages; dropping oldest correlation. \
125 This breaks the service/action request-response contract. \
126 Consider increasing queue_size or switching this input to \
127 `queue_policy: backpressure`."
128 );
129}
130pub(crate) const NON_INPUT_EVENT: &str = "dora.non_input_event";
131
132static NON_INPUT_EVENT_ID: LazyLock<DataId> =
135 LazyLock::new(|| DataId::from(NON_INPUT_EVENT.to_string()));
136
137#[derive(Debug)]
209pub struct Scheduler {
210 last_used: VecDeque<DataId>,
212 event_queues: HashMap<DataId, (usize, VecDeque<EventItem>)>,
214 queue_policies: HashMap<DataId, QueuePolicy>,
216 dropped: HashMap<DataId, u64>,
218}
219
220impl Scheduler {
221 pub(crate) fn with_policies(
222 event_queues: HashMap<DataId, (usize, VecDeque<EventItem>)>,
223 queue_policies: HashMap<DataId, QueuePolicy>,
224 ) -> Self {
225 let topic = VecDeque::from_iter(
226 event_queues
227 .keys()
228 .filter(|t| **t != *NON_INPUT_EVENT_ID)
229 .cloned(),
230 );
231 Self {
232 last_used: topic,
233 event_queues,
234 queue_policies,
235 dropped: HashMap::new(),
236 }
237 }
238
239 pub fn drain_drop_counts(&mut self) -> HashMap<DataId, u64> {
241 std::mem::take(&mut self.dropped)
242 }
243
244 pub(crate) fn add_event(&mut self, event: EventItem) {
245 let (event_id, should_flush) = match &event {
246 EventItem::NodeEvent {
247 event: NodeEvent::Input { id, metadata, .. },
248 ..
249 } => {
250 let flush = dora_message::metadata::get_bool_param(
251 &metadata.parameters,
252 dora_message::metadata::FLUSH,
253 ) == Some(true);
254 (id, flush)
255 }
256 EventItem::ZenohInput { id, metadata, .. } => {
257 let flush = dora_message::metadata::get_bool_param(
258 &metadata.parameters,
259 dora_message::metadata::FLUSH,
260 ) == Some(true);
261 (id, flush)
262 }
263 _ => (&*NON_INPUT_EVENT_ID, false),
264 };
265
266 if should_flush && let Some((_size, queue)) = self.event_queues.get_mut(event_id) {
280 let before = queue.len();
281 queue.retain(|e| is_correlated(e) || is_stop(e));
282 let drained = before - queue.len();
283 if drained > 0 {
284 tracing::debug!(
285 "Flushed {drained} queued event(s) for input `{event_id}` (flush signal)"
286 );
287 }
288 if !queue.is_empty() {
289 tracing::debug!(
290 input = %event_id,
291 preserved = queue.len(),
292 "flush signal retained correlated (request_id/goal_id) events"
293 );
294 }
295 }
296
297 if !self.event_queues.contains_key(event_id) {
304 tracing::warn!(
305 "no queue config for input `{event_id}`, using default size {DEFAULT_QUEUE_SIZE}"
306 );
307 self.last_used.push_back(event_id.clone());
308 self.event_queues
309 .insert(event_id.clone(), (DEFAULT_QUEUE_SIZE, Default::default()));
310 }
311 let Some((size, queue)) = self.event_queues.get_mut(event_id) else {
312 return;
314 };
315
316 let policy = self
317 .queue_policies
318 .get(event_id)
319 .copied()
320 .unwrap_or_default();
321
322 let cap = policy.effective_cap(*size);
323 if queue.len() >= cap {
324 if policy == QueuePolicy::Backpressure {
325 tracing::error!(
326 "Backpressure input `{event_id}` hit hard cap ({cap}), \
327 dropping oldest to prevent OOM"
328 );
329 } else {
330 tracing::warn!("Discarding event for input `{event_id}` due to queue size limit");
331 }
332 *self.dropped.entry(event_id.clone()).or_insert(0) += 1;
333 match select_eviction(queue, &event) {
334 Eviction::RemoveAt(idx) => {
335 queue.remove(idx);
336 }
337 Eviction::DropIncoming => {
338 return;
341 }
342 Eviction::DropCorrelatedLoud(idx) => {
343 if let Some(dropped) = queue.remove(idx) {
344 log_correlation_drop(event_id, &dropped);
345 }
346 }
347 }
348 }
349 queue.push_back(event);
350 }
351
352 pub(crate) fn next(&mut self) -> Option<EventItem> {
353 if let Some((_size, queue)) = self.event_queues.get_mut(&*NON_INPUT_EVENT_ID)
355 && let Some(event) = queue.pop_front()
356 {
357 return Some(event);
358 }
359
360 for index in 0..self.last_used.len() {
364 let id = &self.last_used[index];
365 if let Some((_size, queue)) = self.event_queues.get_mut(id)
366 && let Some(event) = queue.pop_front()
367 {
368 if let Some(id) = self.last_used.remove(index) {
370 self.last_used.push_back(id);
371 }
372 return Some(event);
373 }
374 }
375
376 None
377 }
378
379 pub(crate) fn is_empty(&self) -> bool {
380 self.event_queues
381 .iter()
382 .all(|(_id, (_size, queue))| queue.is_empty())
383 }
384}
385
386#[cfg(test)]
387mod tests {
388 use super::*;
389 use crate::uhlc;
390 use dora_message::{
391 daemon_to_node::NodeEvent,
392 metadata::{FLUSH, Metadata, MetadataParameters, Parameter},
393 };
394
395 fn make_input(id: &str, params: MetadataParameters) -> EventItem {
396 let ts = uhlc::HLC::default().new_timestamp();
397 let metadata = Metadata::from_parameters(ts, params);
398 EventItem::NodeEvent {
399 event: NodeEvent::Input {
400 id: DataId::from(id.to_string()),
401 metadata: std::sync::Arc::new(metadata),
402 data: None,
403 },
404 }
405 }
406
407 fn make_stop() -> EventItem {
408 EventItem::NodeEvent {
409 event: NodeEvent::Stop,
410 }
411 }
412
413 fn make_input_closed(id: &str) -> EventItem {
414 EventItem::NodeEvent {
415 event: NodeEvent::InputClosed {
416 id: DataId::from(id.to_string()),
417 },
418 }
419 }
420
421 fn make_scheduler(audio_capacity: usize) -> (Scheduler, DataId) {
422 let id = DataId::from("audio".to_string());
423 let mut queues = HashMap::new();
424 queues.insert(id.clone(), (audio_capacity, VecDeque::new()));
425 queues.insert(
426 DataId::from(NON_INPUT_EVENT.to_string()),
427 (10, VecDeque::new()),
428 );
429 (Scheduler::with_policies(queues, HashMap::new()), id)
430 }
431
432 #[test]
433 fn flush_clears_older_queued_events() {
434 let (mut sched, id) = make_scheduler(10);
435
436 sched.add_event(make_input("audio", MetadataParameters::new()));
437 sched.add_event(make_input("audio", MetadataParameters::new()));
438 sched.add_event(make_input("audio", MetadataParameters::new()));
439 assert_eq!(sched.event_queues[&id].1.len(), 3);
440
441 let mut flush_params = MetadataParameters::new();
443 flush_params.insert(FLUSH.into(), Parameter::Bool(true));
444 sched.add_event(make_input("audio", flush_params));
445
446 assert_eq!(sched.event_queues[&id].1.len(), 1);
447 }
448
449 #[test]
450 fn non_flush_does_not_clear_queue() {
451 let (mut sched, id) = make_scheduler(10);
452
453 sched.add_event(make_input("audio", MetadataParameters::new()));
454 sched.add_event(make_input("audio", MetadataParameters::new()));
455 sched.add_event(make_input("audio", MetadataParameters::new()));
456 assert_eq!(sched.event_queues[&id].1.len(), 3);
457 }
458
459 #[test]
460 fn flush_false_does_not_clear_queue() {
461 let (mut sched, id) = make_scheduler(10);
462
463 sched.add_event(make_input("audio", MetadataParameters::new()));
464 sched.add_event(make_input("audio", MetadataParameters::new()));
465
466 let mut params = MetadataParameters::new();
467 params.insert(FLUSH.into(), Parameter::Bool(false));
468 sched.add_event(make_input("audio", params));
469
470 assert_eq!(sched.event_queues[&id].1.len(), 3);
471 }
472
473 #[test]
474 fn flush_with_queue_size_one_retains_flush_message() {
475 let (mut sched, id) = make_scheduler(1);
476
477 sched.add_event(make_input("audio", MetadataParameters::new()));
478 assert_eq!(sched.event_queues[&id].1.len(), 1);
479
480 let mut flush_params = MetadataParameters::new();
482 flush_params.insert(FLUSH.into(), Parameter::Bool(true));
483 sched.add_event(make_input("audio", flush_params));
484
485 assert_eq!(sched.event_queues[&id].1.len(), 1);
487 }
488
489 #[test]
490 fn drop_oldest_tracks_drop_count() {
491 let (mut sched, id) = make_scheduler(2);
492
493 sched.add_event(make_input("audio", MetadataParameters::new()));
495 sched.add_event(make_input("audio", MetadataParameters::new()));
496 assert_eq!(sched.event_queues[&id].1.len(), 2);
497
498 sched.add_event(make_input("audio", MetadataParameters::new()));
500 sched.add_event(make_input("audio", MetadataParameters::new()));
501 sched.add_event(make_input("audio", MetadataParameters::new()));
502
503 assert_eq!(sched.event_queues[&id].1.len(), 2);
505
506 let counts = sched.drain_drop_counts();
508 assert_eq!(counts.get(&id), Some(&3));
509
510 let counts = sched.drain_drop_counts();
512 assert!(counts.is_empty());
513 }
514
515 #[test]
518 fn flush_retains_correlated_events() {
519 let (mut sched, id) = make_scheduler(10);
522
523 sched.add_event(make_input("audio", with_request_id("req-1")));
524 sched.add_event(make_input("audio", MetadataParameters::new()));
525 sched.add_event(make_input("audio", MetadataParameters::new()));
526 assert_eq!(sched.event_queues[&id].1.len(), 3);
527
528 let mut flush_params = MetadataParameters::new();
529 flush_params.insert(FLUSH.into(), Parameter::Bool(true));
530 sched.add_event(make_input("audio", flush_params));
531
532 let queue = &sched.event_queues[&id].1;
533 assert_eq!(queue.len(), 2);
535 assert!(
536 queue
537 .iter()
538 .any(|e| request_id_of(e).as_deref() == Some("req-1")),
539 "service response with request_id was wiped by flush"
540 );
541 }
542
543 #[test]
544 fn flush_retains_goal_id_events() {
545 let (mut sched, id) = make_scheduler(10);
547
548 let mut goal_params = MetadataParameters::new();
549 goal_params.insert(GOAL_ID.into(), Parameter::String("goal-7".to_string()));
550 sched.add_event(make_input("audio", goal_params));
551 sched.add_event(make_input("audio", MetadataParameters::new()));
552
553 let mut flush_params = MetadataParameters::new();
554 flush_params.insert(FLUSH.into(), Parameter::Bool(true));
555 sched.add_event(make_input("audio", flush_params));
556
557 let queue = &sched.event_queues[&id].1;
558 assert_eq!(queue.len(), 2);
560 let has_goal = queue.iter().any(|e| {
561 let EventItem::NodeEvent {
562 event: NodeEvent::Input { metadata, .. },
563 ..
564 } = e
565 else {
566 return false;
567 };
568 get_string_param(&metadata.parameters, GOAL_ID) == Some("goal-7")
569 });
570 assert!(has_goal, "action result with goal_id was wiped by flush");
571 }
572
573 #[test]
574 fn flush_with_all_correlated_queue_keeps_everything() {
575 let (mut sched, id) = make_scheduler(10);
578
579 sched.add_event(make_input("audio", with_request_id("req-1")));
580 sched.add_event(make_input("audio", with_request_id("req-2")));
581 sched.add_event(make_input("audio", with_request_id("req-3")));
582
583 let mut flush_params = MetadataParameters::new();
584 flush_params.insert(FLUSH.into(), Parameter::Bool(true));
585 sched.add_event(make_input("audio", flush_params));
586
587 assert_eq!(sched.event_queues[&id].1.len(), 4);
589 }
590
591 fn request_id_of(event: &EventItem) -> Option<String> {
595 let EventItem::NodeEvent {
596 event: NodeEvent::Input { metadata, .. },
597 ..
598 } = event
599 else {
600 return None;
601 };
602 get_string_param(&metadata.parameters, REQUEST_ID).map(|s| s.to_string())
603 }
604
605 fn with_request_id(id: &str) -> MetadataParameters {
606 let mut params = MetadataParameters::new();
607 params.insert(REQUEST_ID.into(), Parameter::String(id.to_string()));
608 params
609 }
610
611 #[test]
612 fn drop_oldest_preserves_correlated_when_non_correlated_present() {
613 let (mut sched, id) = make_scheduler(3);
616
617 sched.add_event(make_input("audio", with_request_id("req-1")));
618 sched.add_event(make_input("audio", MetadataParameters::new()));
619 sched.add_event(make_input("audio", MetadataParameters::new()));
620 sched.add_event(make_input("audio", MetadataParameters::new()));
621
622 let queue = &sched.event_queues[&id].1;
623 assert_eq!(queue.len(), 3);
624 assert!(
626 queue
627 .iter()
628 .any(|e| request_id_of(e).as_deref() == Some("req-1")),
629 "correlated message was dropped even though non-correlated events were available"
630 );
631 }
632
633 #[test]
634 fn drop_oldest_drops_middle_non_correlated_to_save_front_correlated() {
635 let (mut sched, id) = make_scheduler(3);
638
639 sched.add_event(make_input("audio", with_request_id("req-1")));
640 sched.add_event(make_input("audio", MetadataParameters::new()));
641 sched.add_event(make_input("audio", with_request_id("req-2")));
642 sched.add_event(make_input("audio", MetadataParameters::new()));
643
644 let queue = &sched.event_queues[&id].1;
645 assert_eq!(queue.len(), 3);
646 assert!(
647 queue
648 .iter()
649 .any(|e| request_id_of(e).as_deref() == Some("req-1"))
650 );
651 assert!(
652 queue
653 .iter()
654 .any(|e| request_id_of(e).as_deref() == Some("req-2"))
655 );
656 }
657
658 #[test]
659 fn drop_oldest_drops_incoming_if_queue_is_fully_correlated_and_incoming_is_not() {
660 let (mut sched, id) = make_scheduler(2);
663
664 sched.add_event(make_input("audio", with_request_id("req-1")));
665 sched.add_event(make_input("audio", with_request_id("req-2")));
666 sched.add_event(make_input("audio", MetadataParameters::new()));
667
668 let queue = &sched.event_queues[&id].1;
669 assert_eq!(queue.len(), 2);
670 let ids: Vec<_> = queue.iter().filter_map(request_id_of).collect();
671 assert_eq!(ids, vec!["req-1".to_string(), "req-2".to_string()]);
672
673 let counts = sched.drain_drop_counts();
675 assert_eq!(counts.get(&id), Some(&1));
676 }
677
678 #[test]
679 fn drop_oldest_drops_front_loudly_when_both_queue_and_incoming_are_correlated() {
680 let (mut sched, id) = make_scheduler(2);
683
684 sched.add_event(make_input("audio", with_request_id("req-1")));
685 sched.add_event(make_input("audio", with_request_id("req-2")));
686 sched.add_event(make_input("audio", with_request_id("req-3")));
687
688 let queue = &sched.event_queues[&id].1;
689 assert_eq!(queue.len(), 2);
690 let ids: Vec<_> = queue.iter().filter_map(request_id_of).collect();
691 assert_eq!(ids, vec!["req-2".to_string(), "req-3".to_string()]);
692 }
693
694 #[test]
695 fn drop_oldest_goal_id_is_also_preserved() {
696 let (mut sched, id) = make_scheduler(2);
698
699 let mut goal_params = MetadataParameters::new();
700 goal_params.insert(GOAL_ID.into(), Parameter::String("goal-42".to_string()));
701
702 sched.add_event(make_input("audio", goal_params));
703 sched.add_event(make_input("audio", MetadataParameters::new()));
704 sched.add_event(make_input("audio", MetadataParameters::new()));
705
706 let queue = &sched.event_queues[&id].1;
707 assert_eq!(queue.len(), 2);
708 let has_goal = queue.iter().any(|e| {
709 let EventItem::NodeEvent {
710 event: NodeEvent::Input { metadata, .. },
711 ..
712 } = e
713 else {
714 return false;
715 };
716 get_string_param(&metadata.parameters, GOAL_ID) == Some("goal-42")
717 });
718 assert!(
719 has_goal,
720 "goal-42 was dropped despite having non-correlated events to drop"
721 );
722 }
723
724 #[test]
731 fn stop_survives_non_input_queue_overflow() {
732 let (mut sched, _id) = make_scheduler(10);
734
735 sched.add_event(make_stop());
736 for i in 0..100 {
737 sched.add_event(make_input_closed(&format!("in-{i}")));
738 }
739
740 let non_input = &sched.event_queues[&*NON_INPUT_EVENT_ID].1;
741 assert_eq!(non_input.len(), 10, "non-input queue must stay bounded");
742 assert!(
743 non_input.iter().any(is_stop),
744 "the Stop event must survive non-input queue overflow"
745 );
746 }
747
748 #[test]
752 fn incoming_stop_is_admitted_into_a_full_non_input_queue() {
753 let (mut sched, _id) = make_scheduler(10);
754
755 for i in 0..10 {
756 sched.add_event(make_input_closed(&format!("in-{i}")));
757 }
758 sched.add_event(make_stop());
759
760 let non_input = &sched.event_queues[&*NON_INPUT_EVENT_ID].1;
761 assert_eq!(non_input.len(), 10, "non-input queue must stay bounded");
762 assert!(
763 non_input.iter().any(is_stop),
764 "an incoming Stop must be admitted into a full non-input queue"
765 );
766 }
767
768 #[test]
773 fn flush_retains_stop_when_targeting_the_non_input_queue() {
774 let (mut sched, _id) = make_scheduler(10);
775
776 sched.add_event(make_stop());
777 sched.add_event(make_input_closed("x"));
778
779 let mut flush_params = MetadataParameters::new();
780 flush_params.insert(FLUSH.into(), Parameter::Bool(true));
781 sched.add_event(make_input(NON_INPUT_EVENT, flush_params));
782
783 let non_input = &sched.event_queues[&*NON_INPUT_EVENT_ID].1;
784 assert!(
785 non_input.iter().any(is_stop),
786 "flush must not evict the Stop shutdown signal"
787 );
788 }
789
790 #[test]
791 fn backpressure_policy_prevents_drops() {
792 let id = DataId::from("commands".to_string());
793 let mut queues = HashMap::new();
794 queues.insert(id.clone(), (2, VecDeque::new()));
795 queues.insert(
796 DataId::from(NON_INPUT_EVENT.to_string()),
797 (10, VecDeque::new()),
798 );
799 let policies = HashMap::from([(id.clone(), QueuePolicy::Backpressure)]);
800 let mut sched = Scheduler::with_policies(queues, policies);
801
802 sched.add_event(make_input("commands", MetadataParameters::new()));
804 sched.add_event(make_input("commands", MetadataParameters::new()));
805 sched.add_event(make_input("commands", MetadataParameters::new()));
806 sched.add_event(make_input("commands", MetadataParameters::new()));
807
808 assert_eq!(sched.event_queues[&id].1.len(), 4);
810
811 let counts = sched.drain_drop_counts();
813 assert!(counts.is_empty());
814 }
815
816 fn make_zenoh_input(id: &str, params: MetadataParameters) -> EventItem {
819 let ts = uhlc::HLC::default().new_timestamp();
820 let metadata = Metadata::from_parameters(ts, params);
821 use dora_arrow_convert::IntoArrow;
822 EventItem::ZenohInput {
823 id: DataId::from(id.to_string()),
824 metadata: std::sync::Arc::new(metadata),
825 data: dora_arrow_convert::internal::into_array_ref(().into_arrow()).to_data(),
826 }
827 }
828
829 fn request_id_of_zenoh(event: &EventItem) -> Option<String> {
831 let EventItem::ZenohInput { metadata, .. } = event else {
832 return None;
833 };
834 get_string_param(&metadata.parameters, REQUEST_ID).map(|s| s.to_string())
835 }
836
837 #[test]
838 fn zenoh_drop_oldest_drops_front_loudly_when_both_queue_and_incoming_are_correlated() {
839 let (mut sched, id) = make_scheduler(2);
846
847 sched.add_event(make_zenoh_input("audio", with_request_id("req-1")));
848 sched.add_event(make_zenoh_input("audio", with_request_id("req-2")));
849 sched.add_event(make_zenoh_input("audio", with_request_id("req-3")));
850
851 let queue = &sched.event_queues[&id].1;
852 assert_eq!(queue.len(), 2);
853 let ids: Vec<_> = queue.iter().filter_map(request_id_of_zenoh).collect();
854 assert_eq!(ids, vec!["req-2".to_string(), "req-3".to_string()]);
855 }
856
857 #[test]
858 fn zenoh_drop_oldest_preserves_correlated_when_non_correlated_present() {
859 let (mut sched, id) = make_scheduler(3);
862
863 sched.add_event(make_zenoh_input("audio", with_request_id("req-1")));
864 sched.add_event(make_zenoh_input("audio", MetadataParameters::new()));
865 sched.add_event(make_zenoh_input("audio", MetadataParameters::new()));
866 sched.add_event(make_zenoh_input("audio", MetadataParameters::new()));
867
868 let queue = &sched.event_queues[&id].1;
869 assert_eq!(queue.len(), 3);
870 assert!(
871 queue
872 .iter()
873 .any(|e| request_id_of_zenoh(e).as_deref() == Some("req-1")),
874 "correlated ZenohInput was dropped even though non-correlated events were available"
875 );
876 }
877}