use alloc::collections::VecDeque;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct Priority {
pub urgency: u8,
pub incremental: bool,
}
impl Default for Priority {
fn default() -> Self {
Self {
urgency: 3,
incremental: false,
}
}
}
impl Priority {
pub fn parse(s: &[u8]) -> Option<Self> {
let mut out = Self::default();
let mut i = 0usize;
while i < s.len() {
while i < s.len() && (s[i] == b' ' || s[i] == b',' || s[i] == b'\t') {
i += 1;
}
let start = i;
while i < s.len() && s[i] != b',' && s[i] != b' ' && s[i] != b'\t' {
i += 1;
}
let tok = &s[start..i];
if tok.is_empty() {
continue;
}
match tok {
b"i" => out.incremental = true,
_ => {
if let Some(eq) = tok.iter().position(|&c| c == b'=') {
let k = &tok[..eq];
let v = &tok[eq + 1..];
if k == b"u" {
let v = core::str::from_utf8(v).ok()?;
let u: u8 = v.parse().ok()?;
if u <= 7 {
out.urgency = u;
}
}
}
}
}
}
Some(out)
}
}
impl core::fmt::Display for Priority {
fn fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
core::write!(f, "u={}", self.urgency)?;
if self.incremental {
f.write_str(", i")?;
}
Ok(())
}
}
#[derive(Default)]
struct Bucket {
non_incremental: VecDeque<u32>,
incremental: VecDeque<u32>,
deficit: u32,
quantum: u32,
}
impl Bucket {
fn new(quantum: u32) -> Self {
Self {
non_incremental: VecDeque::new(),
incremental: VecDeque::new(),
deficit: 0,
quantum,
}
}
}
pub struct Scheduler {
buckets: [Bucket; 8],
count: usize,
pub rounds: u64,
}
impl Scheduler {
pub fn new(quantum: u32) -> Self {
let quantum = quantum.max(1);
Self {
buckets: [
Bucket::new(quantum),
Bucket::new(quantum),
Bucket::new(quantum),
Bucket::new(quantum),
Bucket::new(quantum),
Bucket::new(quantum),
Bucket::new(quantum),
Bucket::new(quantum),
],
count: 0,
rounds: 0,
}
}
#[inline]
pub fn len(&self) -> usize {
self.count
}
#[inline]
pub fn is_empty(&self) -> bool {
self.count == 0
}
pub fn add(&mut self, stream_id: u32, p: Priority) {
let b = &mut self.buckets[p.urgency as usize];
if p.incremental {
b.incremental.push_back(stream_id);
} else {
b.non_incremental.push_back(stream_id);
}
self.count += 1;
}
pub fn update(&mut self, stream_id: u32, old: Priority, new: Priority) {
if old == new {
return;
}
self.remove(stream_id);
self.add(stream_id, new);
}
pub fn remove(&mut self, stream_id: u32) {
for b in self.buckets.iter_mut() {
if let Some(pos) = b.non_incremental.iter().position(|&s| s == stream_id) {
b.non_incremental.remove(pos);
self.count -= 1;
return;
}
if let Some(pos) = b.incremental.iter().position(|&s| s == stream_id) {
b.incremental.remove(pos);
self.count -= 1;
return;
}
}
}
pub fn next(&mut self, want: usize) -> Option<u32> {
for u in 0..8 {
let b = &mut self.buckets[u];
if let Some(sid) = b.incremental.pop_front() {
b.incremental.push_back(sid);
return Some(sid);
}
if !b.non_incremental.is_empty() && b.deficit.saturating_add(want as u32) <= b.quantum {
b.deficit = b.deficit.saturating_add(want as u32);
return b.non_incremental.pop_front();
}
}
let any = self.buckets.iter().any(|b| !b.non_incremental.is_empty());
if any {
self.rounds += 1;
for b in self.buckets.iter_mut() {
b.deficit = 0;
}
for u in 0..8 {
let b = &mut self.buckets[u];
if let Some(sid) = b.incremental.pop_front() {
b.incremental.push_back(sid);
return Some(sid);
}
if !b.non_incremental.is_empty()
&& b.deficit.saturating_add(want as u32) <= b.quantum
{
b.deficit = b.deficit.saturating_add(want as u32);
return b.non_incremental.pop_front();
}
}
}
None
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn parse_priority_values() {
assert_eq!(Priority::parse(b""), Some(Priority::default()));
assert_eq!(
Priority::parse(b"u=0"),
Some(Priority {
urgency: 0,
incremental: false
})
);
assert_eq!(
Priority::parse(b"u=5, i"),
Some(Priority {
urgency: 5,
incremental: true
})
);
assert_eq!(
Priority::parse(b"i, u=1"),
Some(Priority {
urgency: 1,
incremental: true
})
);
assert_eq!(Priority::parse(b"u=9"), Some(Priority::default())); assert_eq!(
Priority::parse(b"x=y, u=7"),
Some(Priority {
urgency: 7,
incremental: false
})
);
}
#[test]
fn urgency_ordering() {
let mut s = Scheduler::new(16_384);
s.add(
10,
Priority {
urgency: 3,
incremental: false,
},
);
s.add(
11,
Priority {
urgency: 0,
incremental: false,
},
);
s.add(
12,
Priority {
urgency: 7,
incremental: false,
},
);
assert_eq!(s.next(100), Some(11));
assert_eq!(s.next(100), Some(10));
assert_eq!(s.next(100), Some(12));
assert_eq!(s.next(100), None);
}
#[test]
fn incremental_round_robin() {
let mut s = Scheduler::new(16_384);
s.add(
1,
Priority {
urgency: 3,
incremental: true,
},
);
s.add(
2,
Priority {
urgency: 3,
incremental: true,
},
);
s.add(
3,
Priority {
urgency: 0,
incremental: false,
},
);
assert_eq!(s.next(100), Some(3));
assert_eq!(s.next(100), Some(1));
assert_eq!(s.next(100), Some(2));
assert_eq!(s.next(100), Some(1));
}
#[test]
fn no_starvation_across_urgencies() {
let mut s = Scheduler::new(1000);
s.add(
1,
Priority {
urgency: 0,
incremental: false,
},
);
s.add(
2,
Priority {
urgency: 7,
incremental: false,
},
);
let mut served: alloc::vec::Vec<u32> = alloc::vec::Vec::new();
for _ in 0..16 {
if let Some(sid) = s.next(100) {
served.push(sid);
if sid == 1 {
s.add(
1,
Priority {
urgency: 0,
incremental: false,
},
);
}
}
}
assert_eq!(&served[..10], &[1, 1, 1, 1, 1, 1, 1, 1, 1, 1]);
assert_eq!(served[10], 2);
}
#[test]
fn priority_update_moves_bucket() {
let mut s = Scheduler::new(16_384);
s.add(
1,
Priority {
urgency: 3,
incremental: false,
},
);
s.add(
2,
Priority {
urgency: 7,
incremental: false,
},
);
s.update(
2,
Priority {
urgency: 7,
incremental: false,
},
Priority {
urgency: 0,
incremental: false,
},
);
assert_eq!(s.next(100), Some(2));
assert_eq!(s.next(100), Some(1));
}
}