use std::{
collections::{HashMap, VecDeque},
sync::LazyLock,
};
use dora_message::{
config::{DEFAULT_QUEUE_SIZE, QueuePolicy},
daemon_to_node::NodeEvent,
id::DataId,
metadata::{GOAL_ID, GOAL_STATUS, REQUEST_ID, get_string_param},
};
use super::thread::EventItem;
fn is_correlated(event: &EventItem) -> bool {
let params = match event {
EventItem::NodeEvent {
event: NodeEvent::Input { metadata, .. },
..
} => &metadata.parameters,
EventItem::ZenohInput { metadata, .. } => &metadata.parameters,
_ => return false,
};
params.contains_key(REQUEST_ID)
|| params.contains_key(GOAL_ID)
|| params.contains_key(GOAL_STATUS)
}
enum Eviction {
RemoveAt(usize),
DropIncoming,
DropFrontLoud,
}
fn select_eviction(queue: &VecDeque<EventItem>, incoming: &EventItem) -> Eviction {
if let Some(idx) = queue.iter().position(|e| !is_correlated(e)) {
return Eviction::RemoveAt(idx);
}
if !is_correlated(incoming) {
return Eviction::DropIncoming;
}
Eviction::DropFrontLoud
}
fn log_correlation_drop(event_id: &DataId, dropped: &EventItem) {
let params = match dropped {
EventItem::NodeEvent {
event: NodeEvent::Input { metadata, .. },
..
} => &metadata.parameters,
EventItem::ZenohInput { metadata, .. } => &metadata.parameters,
_ => return,
};
let request_id = get_string_param(params, REQUEST_ID);
let goal_id = get_string_param(params, GOAL_ID);
let goal_status = get_string_param(params, GOAL_STATUS);
tracing::error!(
input = %event_id,
?request_id,
?goal_id,
?goal_status,
"queue full of correlated messages; dropping oldest correlation. \
This breaks the service/action request-response contract. \
Consider increasing queue_size or switching this input to \
`queue_policy: backpressure`."
);
}
pub(crate) const NON_INPUT_EVENT: &str = "dora.non_input_event";
static NON_INPUT_EVENT_ID: LazyLock<DataId> =
LazyLock::new(|| DataId::from(NON_INPUT_EVENT.to_string()));
#[derive(Debug)]
pub struct Scheduler {
last_used: VecDeque<DataId>,
event_queues: HashMap<DataId, (usize, VecDeque<EventItem>)>,
queue_policies: HashMap<DataId, QueuePolicy>,
dropped: HashMap<DataId, u64>,
}
impl Scheduler {
pub(crate) fn with_policies(
event_queues: HashMap<DataId, (usize, VecDeque<EventItem>)>,
queue_policies: HashMap<DataId, QueuePolicy>,
) -> Self {
let topic = VecDeque::from_iter(
event_queues
.keys()
.filter(|t| **t != *NON_INPUT_EVENT_ID)
.cloned(),
);
Self {
last_used: topic,
event_queues,
queue_policies,
dropped: HashMap::new(),
}
}
pub fn drain_drop_counts(&mut self) -> HashMap<DataId, u64> {
std::mem::take(&mut self.dropped)
}
pub(crate) fn add_event(&mut self, event: EventItem) {
let (event_id, should_flush) = match &event {
EventItem::NodeEvent {
event: NodeEvent::Input { id, metadata, .. },
..
} => {
let flush = dora_message::metadata::get_bool_param(
&metadata.parameters,
dora_message::metadata::FLUSH,
) == Some(true);
(id, flush)
}
EventItem::ZenohInput { id, metadata, .. } => {
let flush = dora_message::metadata::get_bool_param(
&metadata.parameters,
dora_message::metadata::FLUSH,
) == Some(true);
(id, flush)
}
_ => (&*NON_INPUT_EVENT_ID, false),
};
if should_flush && let Some((_size, queue)) = self.event_queues.get_mut(event_id) {
let before = queue.len();
queue.retain(is_correlated);
let drained = before - queue.len();
if drained > 0 {
tracing::debug!(
"Flushed {drained} queued event(s) for input `{event_id}` (flush signal)"
);
}
if !queue.is_empty() {
tracing::debug!(
input = %event_id,
preserved = queue.len(),
"flush signal retained correlated (request_id/goal_id) events"
);
}
}
if !self.event_queues.contains_key(event_id) {
tracing::warn!(
"no queue config for input `{event_id}`, using default size {DEFAULT_QUEUE_SIZE}"
);
self.last_used.push_back(event_id.clone());
self.event_queues
.insert(event_id.clone(), (DEFAULT_QUEUE_SIZE, Default::default()));
}
let Some((size, queue)) = self.event_queues.get_mut(event_id) else {
return;
};
let policy = self
.queue_policies
.get(event_id)
.copied()
.unwrap_or_default();
let cap = policy.effective_cap(*size);
if queue.len() >= cap {
if policy == QueuePolicy::Backpressure {
tracing::error!(
"Backpressure input `{event_id}` hit hard cap ({cap}), \
dropping oldest to prevent OOM"
);
} else {
tracing::warn!("Discarding event for input `{event_id}` due to queue size limit");
}
*self.dropped.entry(event_id.clone()).or_insert(0) += 1;
match select_eviction(queue, &event) {
Eviction::RemoveAt(idx) => {
queue.remove(idx);
}
Eviction::DropIncoming => {
return;
}
Eviction::DropFrontLoud => {
if let Some(front) = queue.pop_front() {
log_correlation_drop(event_id, &front);
}
}
}
}
queue.push_back(event);
}
pub(crate) fn next(&mut self) -> Option<EventItem> {
if let Some((_size, queue)) = self.event_queues.get_mut(&*NON_INPUT_EVENT_ID)
&& let Some(event) = queue.pop_front()
{
return Some(event);
}
for index in 0..self.last_used.len() {
let id = &self.last_used[index];
if let Some((_size, queue)) = self.event_queues.get_mut(id)
&& let Some(event) = queue.pop_front()
{
if let Some(id) = self.last_used.remove(index) {
self.last_used.push_back(id);
}
return Some(event);
}
}
None
}
pub(crate) fn is_empty(&self) -> bool {
self.event_queues
.iter()
.all(|(_id, (_size, queue))| queue.is_empty())
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::uhlc;
use dora_message::{
daemon_to_node::NodeEvent,
metadata::{FLUSH, Metadata, MetadataParameters, Parameter},
};
fn make_input(id: &str, params: MetadataParameters) -> EventItem {
let ts = uhlc::HLC::default().new_timestamp();
let metadata = Metadata::from_parameters(ts, params);
EventItem::NodeEvent {
event: NodeEvent::Input {
id: DataId::from(id.to_string()),
metadata: std::sync::Arc::new(metadata),
data: None,
},
}
}
fn make_scheduler(audio_capacity: usize) -> (Scheduler, DataId) {
let id = DataId::from("audio".to_string());
let mut queues = HashMap::new();
queues.insert(id.clone(), (audio_capacity, VecDeque::new()));
queues.insert(
DataId::from(NON_INPUT_EVENT.to_string()),
(10, VecDeque::new()),
);
(Scheduler::with_policies(queues, HashMap::new()), id)
}
#[test]
fn flush_clears_older_queued_events() {
let (mut sched, id) = make_scheduler(10);
sched.add_event(make_input("audio", MetadataParameters::new()));
sched.add_event(make_input("audio", MetadataParameters::new()));
sched.add_event(make_input("audio", MetadataParameters::new()));
assert_eq!(sched.event_queues[&id].1.len(), 3);
let mut flush_params = MetadataParameters::new();
flush_params.insert(FLUSH.into(), Parameter::Bool(true));
sched.add_event(make_input("audio", flush_params));
assert_eq!(sched.event_queues[&id].1.len(), 1);
}
#[test]
fn non_flush_does_not_clear_queue() {
let (mut sched, id) = make_scheduler(10);
sched.add_event(make_input("audio", MetadataParameters::new()));
sched.add_event(make_input("audio", MetadataParameters::new()));
sched.add_event(make_input("audio", MetadataParameters::new()));
assert_eq!(sched.event_queues[&id].1.len(), 3);
}
#[test]
fn flush_false_does_not_clear_queue() {
let (mut sched, id) = make_scheduler(10);
sched.add_event(make_input("audio", MetadataParameters::new()));
sched.add_event(make_input("audio", MetadataParameters::new()));
let mut params = MetadataParameters::new();
params.insert(FLUSH.into(), Parameter::Bool(false));
sched.add_event(make_input("audio", params));
assert_eq!(sched.event_queues[&id].1.len(), 3);
}
#[test]
fn flush_with_queue_size_one_retains_flush_message() {
let (mut sched, id) = make_scheduler(1);
sched.add_event(make_input("audio", MetadataParameters::new()));
assert_eq!(sched.event_queues[&id].1.len(), 1);
let mut flush_params = MetadataParameters::new();
flush_params.insert(FLUSH.into(), Parameter::Bool(true));
sched.add_event(make_input("audio", flush_params));
assert_eq!(sched.event_queues[&id].1.len(), 1);
}
#[test]
fn drop_oldest_tracks_drop_count() {
let (mut sched, id) = make_scheduler(2);
sched.add_event(make_input("audio", MetadataParameters::new()));
sched.add_event(make_input("audio", MetadataParameters::new()));
assert_eq!(sched.event_queues[&id].1.len(), 2);
sched.add_event(make_input("audio", MetadataParameters::new()));
sched.add_event(make_input("audio", MetadataParameters::new()));
sched.add_event(make_input("audio", MetadataParameters::new()));
assert_eq!(sched.event_queues[&id].1.len(), 2);
let counts = sched.drain_drop_counts();
assert_eq!(counts.get(&id), Some(&3));
let counts = sched.drain_drop_counts();
assert!(counts.is_empty());
}
#[test]
fn flush_retains_correlated_events() {
let (mut sched, id) = make_scheduler(10);
sched.add_event(make_input("audio", with_request_id("req-1")));
sched.add_event(make_input("audio", MetadataParameters::new()));
sched.add_event(make_input("audio", MetadataParameters::new()));
assert_eq!(sched.event_queues[&id].1.len(), 3);
let mut flush_params = MetadataParameters::new();
flush_params.insert(FLUSH.into(), Parameter::Bool(true));
sched.add_event(make_input("audio", flush_params));
let queue = &sched.event_queues[&id].1;
assert_eq!(queue.len(), 2);
assert!(
queue
.iter()
.any(|e| request_id_of(e).as_deref() == Some("req-1")),
"service response with request_id was wiped by flush"
);
}
#[test]
fn flush_retains_goal_id_events() {
let (mut sched, id) = make_scheduler(10);
let mut goal_params = MetadataParameters::new();
goal_params.insert(GOAL_ID.into(), Parameter::String("goal-7".to_string()));
sched.add_event(make_input("audio", goal_params));
sched.add_event(make_input("audio", MetadataParameters::new()));
let mut flush_params = MetadataParameters::new();
flush_params.insert(FLUSH.into(), Parameter::Bool(true));
sched.add_event(make_input("audio", flush_params));
let queue = &sched.event_queues[&id].1;
assert_eq!(queue.len(), 2);
let has_goal = queue.iter().any(|e| {
let EventItem::NodeEvent {
event: NodeEvent::Input { metadata, .. },
..
} = e
else {
return false;
};
get_string_param(&metadata.parameters, GOAL_ID) == Some("goal-7")
});
assert!(has_goal, "action result with goal_id was wiped by flush");
}
#[test]
fn flush_with_all_correlated_queue_keeps_everything() {
let (mut sched, id) = make_scheduler(10);
sched.add_event(make_input("audio", with_request_id("req-1")));
sched.add_event(make_input("audio", with_request_id("req-2")));
sched.add_event(make_input("audio", with_request_id("req-3")));
let mut flush_params = MetadataParameters::new();
flush_params.insert(FLUSH.into(), Parameter::Bool(true));
sched.add_event(make_input("audio", flush_params));
assert_eq!(sched.event_queues[&id].1.len(), 4);
}
fn request_id_of(event: &EventItem) -> Option<String> {
let EventItem::NodeEvent {
event: NodeEvent::Input { metadata, .. },
..
} = event
else {
return None;
};
get_string_param(&metadata.parameters, REQUEST_ID).map(|s| s.to_string())
}
fn with_request_id(id: &str) -> MetadataParameters {
let mut params = MetadataParameters::new();
params.insert(REQUEST_ID.into(), Parameter::String(id.to_string()));
params
}
#[test]
fn drop_oldest_preserves_correlated_when_non_correlated_present() {
let (mut sched, id) = make_scheduler(3);
sched.add_event(make_input("audio", with_request_id("req-1")));
sched.add_event(make_input("audio", MetadataParameters::new()));
sched.add_event(make_input("audio", MetadataParameters::new()));
sched.add_event(make_input("audio", MetadataParameters::new()));
let queue = &sched.event_queues[&id].1;
assert_eq!(queue.len(), 3);
assert!(
queue
.iter()
.any(|e| request_id_of(e).as_deref() == Some("req-1")),
"correlated message was dropped even though non-correlated events were available"
);
}
#[test]
fn drop_oldest_drops_middle_non_correlated_to_save_front_correlated() {
let (mut sched, id) = make_scheduler(3);
sched.add_event(make_input("audio", with_request_id("req-1")));
sched.add_event(make_input("audio", MetadataParameters::new()));
sched.add_event(make_input("audio", with_request_id("req-2")));
sched.add_event(make_input("audio", MetadataParameters::new()));
let queue = &sched.event_queues[&id].1;
assert_eq!(queue.len(), 3);
assert!(
queue
.iter()
.any(|e| request_id_of(e).as_deref() == Some("req-1"))
);
assert!(
queue
.iter()
.any(|e| request_id_of(e).as_deref() == Some("req-2"))
);
}
#[test]
fn drop_oldest_drops_incoming_if_queue_is_fully_correlated_and_incoming_is_not() {
let (mut sched, id) = make_scheduler(2);
sched.add_event(make_input("audio", with_request_id("req-1")));
sched.add_event(make_input("audio", with_request_id("req-2")));
sched.add_event(make_input("audio", MetadataParameters::new()));
let queue = &sched.event_queues[&id].1;
assert_eq!(queue.len(), 2);
let ids: Vec<_> = queue.iter().filter_map(request_id_of).collect();
assert_eq!(ids, vec!["req-1".to_string(), "req-2".to_string()]);
let counts = sched.drain_drop_counts();
assert_eq!(counts.get(&id), Some(&1));
}
#[test]
fn drop_oldest_drops_front_loudly_when_both_queue_and_incoming_are_correlated() {
let (mut sched, id) = make_scheduler(2);
sched.add_event(make_input("audio", with_request_id("req-1")));
sched.add_event(make_input("audio", with_request_id("req-2")));
sched.add_event(make_input("audio", with_request_id("req-3")));
let queue = &sched.event_queues[&id].1;
assert_eq!(queue.len(), 2);
let ids: Vec<_> = queue.iter().filter_map(request_id_of).collect();
assert_eq!(ids, vec!["req-2".to_string(), "req-3".to_string()]);
}
#[test]
fn drop_oldest_goal_id_is_also_preserved() {
let (mut sched, id) = make_scheduler(2);
let mut goal_params = MetadataParameters::new();
goal_params.insert(GOAL_ID.into(), Parameter::String("goal-42".to_string()));
sched.add_event(make_input("audio", goal_params));
sched.add_event(make_input("audio", MetadataParameters::new()));
sched.add_event(make_input("audio", MetadataParameters::new()));
let queue = &sched.event_queues[&id].1;
assert_eq!(queue.len(), 2);
let has_goal = queue.iter().any(|e| {
let EventItem::NodeEvent {
event: NodeEvent::Input { metadata, .. },
..
} = e
else {
return false;
};
get_string_param(&metadata.parameters, GOAL_ID) == Some("goal-42")
});
assert!(
has_goal,
"goal-42 was dropped despite having non-correlated events to drop"
);
}
#[test]
fn backpressure_policy_prevents_drops() {
let id = DataId::from("commands".to_string());
let mut queues = HashMap::new();
queues.insert(id.clone(), (2, VecDeque::new()));
queues.insert(
DataId::from(NON_INPUT_EVENT.to_string()),
(10, VecDeque::new()),
);
let policies = HashMap::from([(id.clone(), QueuePolicy::Backpressure)]);
let mut sched = Scheduler::with_policies(queues, policies);
sched.add_event(make_input("commands", MetadataParameters::new()));
sched.add_event(make_input("commands", MetadataParameters::new()));
sched.add_event(make_input("commands", MetadataParameters::new()));
sched.add_event(make_input("commands", MetadataParameters::new()));
assert_eq!(sched.event_queues[&id].1.len(), 4);
let counts = sched.drain_drop_counts();
assert!(counts.is_empty());
}
fn make_zenoh_input(id: &str, params: MetadataParameters) -> EventItem {
let ts = uhlc::HLC::default().new_timestamp();
let metadata = Metadata::from_parameters(ts, params);
use dora_arrow_convert::IntoArrow;
EventItem::ZenohInput {
id: DataId::from(id.to_string()),
metadata: std::sync::Arc::new(metadata),
data: ().into_arrow().into(),
}
}
fn request_id_of_zenoh(event: &EventItem) -> Option<String> {
let EventItem::ZenohInput { metadata, .. } = event else {
return None;
};
get_string_param(&metadata.parameters, REQUEST_ID).map(|s| s.to_string())
}
#[test]
fn zenoh_drop_oldest_drops_front_loudly_when_both_queue_and_incoming_are_correlated() {
let (mut sched, id) = make_scheduler(2);
sched.add_event(make_zenoh_input("audio", with_request_id("req-1")));
sched.add_event(make_zenoh_input("audio", with_request_id("req-2")));
sched.add_event(make_zenoh_input("audio", with_request_id("req-3")));
let queue = &sched.event_queues[&id].1;
assert_eq!(queue.len(), 2);
let ids: Vec<_> = queue.iter().filter_map(request_id_of_zenoh).collect();
assert_eq!(ids, vec!["req-2".to_string(), "req-3".to_string()]);
}
#[test]
fn zenoh_drop_oldest_preserves_correlated_when_non_correlated_present() {
let (mut sched, id) = make_scheduler(3);
sched.add_event(make_zenoh_input("audio", with_request_id("req-1")));
sched.add_event(make_zenoh_input("audio", MetadataParameters::new()));
sched.add_event(make_zenoh_input("audio", MetadataParameters::new()));
sched.add_event(make_zenoh_input("audio", MetadataParameters::new()));
let queue = &sched.event_queues[&id].1;
assert_eq!(queue.len(), 3);
assert!(
queue
.iter()
.any(|e| request_id_of_zenoh(e).as_deref() == Some("req-1")),
"correlated ZenohInput was dropped even though non-correlated events were available"
);
}
}