use crate::budget::Budget;
use std::collections::VecDeque;
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, Ordering};
use yo_common::Result;
use yo_shard::Epochs;
use yo_shard::spsc::Receiver;
pub const BATCH_MAX: usize = 64;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum Flow {
Next,
Break,
}
pub trait Engine {
type Work;
fn key_hash(&self, work: &Self::Work) -> Option<u64>;
fn prefetch(&self, work: &Self::Work, hash: u64);
fn run(&mut self, work: Self::Work, hash: Option<u64>) -> Flow;
fn flush(&mut self);
fn submit_io(&mut self) -> Result<()> {
Ok(())
}
fn drain_io(&mut self) -> Result<()> {
Ok(())
}
fn maintain(&mut self, budget: &mut Budget) {
let _ = budget;
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
pub struct Turn {
pub commands: usize,
pub broke: bool,
pub carried: usize,
pub maintained: u32,
}
impl Turn {
#[must_use]
pub const fn is_idle(&self) -> bool {
self.commands == 0 && self.carried == 0
}
}
pub struct Reactor<E: Engine> {
engine: E,
id: usize,
epochs: Arc<Epochs>,
lanes: Vec<Receiver<E::Work>>,
pending: VecDeque<E::Work>,
hashes: Vec<Option<u64>>,
lane: usize,
budget: u32,
turns: u64,
commands: u64,
batches: u64,
full: u64,
breaks: u64,
idle: u64,
}
impl<E: Engine> Reactor<E> {
pub fn new(engine: E, id: usize, epochs: Arc<Epochs>, lanes: Vec<Receiver<E::Work>>) -> Self {
assert!(id < epochs.len(), "shard {id} has no epoch slot");
Reactor {
engine,
id,
epochs,
lanes,
pending: VecDeque::with_capacity(BATCH_MAX),
hashes: Vec::with_capacity(BATCH_MAX),
lane: 0,
budget: crate::MAINTENANCE_UNITS,
turns: 0,
commands: 0,
batches: 0,
full: 0,
breaks: 0,
idle: 0,
}
}
pub fn inline(engine: E) -> Self {
Reactor::new(engine, 0, Epochs::new(1), Vec::new())
}
#[must_use]
pub fn with_maintenance(mut self, units: u32) -> Self {
self.budget = units;
self
}
pub const fn engine(&self) -> &E {
&self.engine
}
pub const fn engine_mut(&mut self) -> &mut E {
&mut self.engine
}
#[must_use]
pub const fn id(&self) -> usize {
self.id
}
#[must_use]
pub const fn turns(&self) -> u64 {
self.turns
}
#[must_use]
pub const fn commands(&self) -> u64 {
self.commands
}
#[must_use]
pub const fn batches(&self) -> u64 {
self.batches
}
#[must_use]
pub const fn full_batches(&self) -> u64 {
self.full
}
#[must_use]
pub const fn breaks(&self) -> u64 {
self.breaks
}
#[must_use]
pub const fn idle_turns(&self) -> u64 {
self.idle
}
#[must_use]
pub fn carried(&self) -> usize {
self.pending.len()
}
pub fn tick(&mut self) -> Result<Turn> {
self.turns += 1;
self.engine.submit_io()?;
if self.pending.is_empty() {
self.fill();
if !self.pending.is_empty() {
self.batches += 1;
if self.pending.len() == BATCH_MAX {
self.full += 1;
}
}
}
let mut turn = Turn::default();
if self.pending.is_empty() {
self.idle += 1;
} else {
self.epochs.enter(self.id);
self.hashes.clear();
for w in &self.pending {
let h = self.engine.key_hash(w);
if let Some(h) = h {
self.engine.prefetch(w, h);
}
self.hashes.push(h);
}
let mut n = 0;
while let Some(w) = self.pending.pop_front() {
let h = self.hashes[n];
n += 1;
if self.engine.run(w, h) == Flow::Break {
turn.broke = true;
self.breaks += 1;
break;
}
}
self.epochs.leave(self.id);
self.engine.flush();
self.commands += n as u64;
turn.commands = n;
turn.carried = self.pending.len();
}
self.engine.drain_io()?;
if self.budget > 0 {
let mut budget = Budget::new(self.budget);
self.engine.maintain(&mut budget);
turn.maintained = budget.spent();
}
Ok(turn)
}
pub fn run_until(&mut self, stop: &AtomicBool) -> Result<()> {
const SPINS: u32 = 128;
let mut idle = 0u32;
loop {
if !self.tick()?.is_idle() {
idle = 0;
continue;
}
if stop.load(Ordering::Acquire) {
if self.tick()?.is_idle() {
return Ok(());
}
continue;
}
idle += 1;
if idle < SPINS {
std::hint::spin_loop();
} else {
std::thread::yield_now();
idle = 0;
}
}
}
pub fn execute(&mut self, work: E::Work) -> Flow {
self.turns += 1;
self.commands += 1;
self.epochs.enter(self.id);
let hash = self.engine.key_hash(&work);
if let Some(h) = hash {
self.engine.prefetch(&work, h);
}
let flow = self.engine.run(work, hash);
self.epochs.leave(self.id);
flow
}
pub fn execute_all<I>(&mut self, work: I) -> usize
where
I: IntoIterator<Item = E::Work>,
{
self.pending.extend(work);
if self.pending.is_empty() {
return 0;
}
self.turns += 1;
self.epochs.enter(self.id);
self.hashes.clear();
for w in &self.pending {
let h = self.engine.key_hash(w);
if let Some(h) = h {
self.engine.prefetch(w, h);
}
self.hashes.push(h);
}
let mut n = 0;
while let Some(w) = self.pending.pop_front() {
let h = self.hashes[n];
n += 1;
if self.engine.run(w, h) == Flow::Break {
self.breaks += 1;
break;
}
}
self.pending.clear();
self.epochs.leave(self.id);
self.commands += n as u64;
n
}
fn fill(&mut self) {
if self.lanes.is_empty() {
return;
}
let n = self.lanes.len();
let mut room = BATCH_MAX;
let mut at = self.lane;
loop {
let mut took = 0;
for _ in 0..n {
if room == 0 {
self.lane = at;
return;
}
if let Some(w) = self.lanes[at].pop() {
self.pending.push_back(w);
room -= 1;
took += 1;
}
at = (at + 1) % n;
}
if took == 0 {
self.lane = at;
return;
}
}
}
}
impl<E: Engine> std::fmt::Debug for Reactor<E> {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("Reactor")
.field("id", &self.id)
.field("lanes", &self.lanes.len())
.field("turns", &self.turns)
.field("commands", &self.commands)
.field("batches", &self.batches)
.field("full_batches", &self.full)
.field("breaks", &self.breaks)
.field("idle_turns", &self.idle)
.field("carried", &self.pending.len())
.finish()
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::cell::RefCell;
use yo_common::{Code, Error};
use yo_shard::spsc::{Sender, lane};
#[derive(Debug, Clone, PartialEq, Eq)]
enum Step {
Submit,
Prefetch(u64),
Run(u64),
Flush,
Drain,
Maintain,
}
const KEYLESS: u64 = u64::MAX;
struct Recorder {
steps: RefCell<Vec<Step>>,
break_on: Option<u64>,
fail_submit: bool,
fail_drain: bool,
maintenance_item: u32,
}
impl Recorder {
fn new() -> Recorder {
Recorder {
steps: RefCell::new(Vec::new()),
break_on: None,
fail_submit: false,
fail_drain: false,
maintenance_item: 0,
}
}
fn push(&self, step: Step) {
self.steps.borrow_mut().push(step);
}
fn steps(&self) -> Vec<Step> {
self.steps.borrow().clone()
}
fn runs(&self) -> Vec<u64> {
self.steps
.borrow()
.iter()
.filter_map(|s| match s {
Step::Run(v) => Some(*v),
_ => None,
})
.collect()
}
fn count(&self, want: &Step) -> usize {
self.steps.borrow().iter().filter(|s| *s == want).count()
}
}
impl Engine for Recorder {
type Work = u64;
fn key_hash(&self, work: &u64) -> Option<u64> {
if *work == KEYLESS { None } else { Some(*work) }
}
fn prefetch(&self, _work: &u64, hash: u64) {
self.push(Step::Prefetch(hash));
}
fn run(&mut self, work: u64, hash: Option<u64>) -> Flow {
assert_eq!(
hash,
if work == KEYLESS { None } else { Some(work) },
"the second walk gets the hash the first walk took"
);
self.push(Step::Run(work));
if self.break_on == Some(work) {
return Flow::Break;
}
Flow::Next
}
fn flush(&mut self) {
self.push(Step::Flush);
}
fn submit_io(&mut self) -> Result<()> {
self.push(Step::Submit);
if self.fail_submit {
return Err(Error::new(Code::Io, "submit said no"));
}
Ok(())
}
fn drain_io(&mut self) -> Result<()> {
self.push(Step::Drain);
if self.fail_drain {
return Err(Error::new(Code::Io, "drain said no"));
}
Ok(())
}
fn maintain(&mut self, budget: &mut Budget) {
self.push(Step::Maintain);
if self.maintenance_item == 0 {
return;
}
while budget.spend(self.maintenance_item) {}
}
}
fn wired(lanes: usize) -> (Reactor<Recorder>, Vec<Sender<u64>>, Arc<Epochs>) {
let mut rxs = Vec::new();
let mut txs = Vec::new();
for _ in 0..lanes {
let (tx, rx) = lane(1024);
txs.push(tx);
rxs.push(rx);
}
let epochs = Epochs::new(1);
let r = Reactor::new(Recorder::new(), 0, Arc::clone(&epochs), rxs);
(r, txs, epochs)
}
fn walk_boundary(steps: &[Step]) -> (usize, usize) {
let last_prefetch = steps
.iter()
.rposition(|s| matches!(s, Step::Prefetch(_)))
.expect("nothing was prefetched");
let first_run = steps
.iter()
.position(|s| matches!(s, Step::Run(_)))
.expect("nothing ran");
(last_prefetch, first_run)
}
#[test]
fn the_stages_run_in_the_order_the_spec_lists_them() {
let (mut r, tx, _e) = wired(1);
tx[0].push(7).unwrap();
let turn = r.tick().unwrap();
assert_eq!(turn.commands, 1);
assert_eq!(
r.engine().steps(),
vec![
Step::Submit,
Step::Prefetch(7),
Step::Run(7),
Step::Flush,
Step::Drain,
Step::Maintain,
]
);
}
#[test]
fn every_command_is_prefetched_before_any_of_them_runs() {
let (mut r, tx, _e) = wired(1);
for i in 0..8 {
tx[0].push(i).unwrap();
}
r.tick().unwrap();
let (last_prefetch, first_run) = walk_boundary(&r.engine().steps());
assert!(
last_prefetch < first_run,
"the two walks overlapped, which makes the prefetch distance one"
);
assert_eq!(r.engine().runs(), (0..8).collect::<Vec<_>>());
}
#[test]
fn a_batch_stops_at_sixty_four() {
let (mut r, tx, _e) = wired(1);
for i in 0..200 {
tx[0].push(i).unwrap();
}
assert_eq!(r.tick().unwrap().commands, BATCH_MAX);
assert_eq!(r.full_batches(), 1);
assert_eq!(r.tick().unwrap().commands, BATCH_MAX);
assert_eq!(r.tick().unwrap().commands, BATCH_MAX);
assert_eq!(r.tick().unwrap().commands, 200 - 3 * BATCH_MAX);
assert_eq!(r.batches(), 4);
assert_eq!(r.full_batches(), 3);
assert_eq!(r.commands(), 200);
assert_eq!(r.engine().runs(), (0..200).collect::<Vec<_>>());
}
#[test]
fn a_break_leaves_the_rest_of_the_batch_for_the_next_turn() {
let (mut r, tx, _e) = wired(1);
for i in 0..10 {
tx[0].push(i).unwrap();
}
r.engine_mut().break_on = Some(3);
let turn = r.tick().unwrap();
assert!(turn.broke);
assert_eq!(turn.commands, 4, "the command that broke it still ran");
assert_eq!(turn.carried, 6);
assert_eq!(r.engine().count(&Step::Flush), 1, "the replies still went");
r.engine_mut().break_on = None;
let turn = r.tick().unwrap();
assert_eq!(turn.commands, 6);
assert_eq!(turn.carried, 0);
assert_eq!(r.engine().runs(), (0..10).collect::<Vec<_>>());
assert_eq!(r.breaks(), 1);
assert_eq!(r.batches(), 1, "a broken batch is one batch, not two");
}
#[test]
fn the_epoch_moves_once_per_batch_and_not_once_per_command() {
let (mut r, tx, epochs) = wired(1);
for i in 0..10 {
tx[0].push(i).unwrap();
}
let before = epochs.get(0);
r.tick().unwrap();
assert_eq!(
epochs.get(0),
before + 2,
"one enter and one leave for ten commands"
);
assert_eq!(r.commands(), 10);
}
#[test]
fn an_idle_turn_does_not_touch_the_epoch_or_the_replies() {
let (mut r, _tx, epochs) = wired(1);
let before = epochs.get(0);
let turn = r.tick().unwrap();
assert!(turn.is_idle());
assert_eq!(epochs.get(0), before, "an idle shard holds nothing");
assert_eq!(r.engine().count(&Step::Flush), 0);
assert_eq!(
r.engine().count(&Step::Drain),
1,
"completions still get picked up"
);
assert_eq!(
r.engine().count(&Step::Maintain),
1,
"and background work still runs"
);
assert_eq!(r.idle_turns(), 1);
assert_eq!(r.batches(), 0);
}
#[test]
fn work_comes_off_every_lane_rather_than_the_first_one() {
let (mut r, tx, _e) = wired(4);
for (i, t) in tx.iter().enumerate() {
for j in 0..4u64 {
t.push(i as u64 * 10 + j).unwrap();
}
}
let turn = r.tick().unwrap();
assert_eq!(turn.commands, 16);
let runs = r.engine().runs();
for lane in 0..4u64 {
let from_lane = runs.iter().filter(|v| **v / 10 == lane).count();
assert_eq!(from_lane, 4, "lane {lane} was skipped or drained twice");
}
assert_eq!(
runs[..4].iter().map(|v| v / 10).collect::<Vec<_>>(),
vec![0, 1, 2, 3],
"round robin, so the first four are one from each lane"
);
}
#[test]
fn a_lane_that_never_stops_cannot_starve_the_others() {
let (mut r, tx, _e) = wired(2);
for i in 0..200 {
tx[0].push(i).unwrap();
}
tx[1].push(9_999).unwrap();
r.tick().unwrap();
assert!(
r.engine().runs().contains(&9_999),
"the quiet lane waited behind a whole batch of the busy one"
);
}
#[test]
fn maintenance_gets_a_budget_and_stops_when_it_is_spent() {
let (mut r, _tx, _e) = wired(1);
r.engine_mut().maintenance_item = 100;
let turn = r.tick().unwrap();
assert!(
(crate::MAINTENANCE_UNITS..crate::MAINTENANCE_UNITS + 100).contains(&turn.maintained),
"spent {}",
turn.maintained
);
let mut r = r.with_maintenance(0);
let before = r.engine().count(&Step::Maintain);
let turn = r.tick().unwrap();
assert_eq!(turn.maintained, 0);
assert_eq!(
r.engine().count(&Step::Maintain),
before,
"a zero budget is no slice at all rather than an empty one"
);
}
#[test]
fn a_failed_submit_stops_the_turn_before_the_batch() {
let (mut r, tx, _e) = wired(1);
tx[0].push(1).unwrap();
r.engine_mut().fail_submit = true;
assert_eq!(r.tick().unwrap_err().code(), Code::Io);
assert_eq!(r.engine().runs(), Vec::<u64>::new());
assert_eq!(r.carried(), 0, "nothing was drained, so nothing is held");
r.engine_mut().fail_submit = false;
assert_eq!(r.tick().unwrap().commands, 1);
}
#[test]
fn a_failed_drain_still_ran_the_batch_and_flushed_it() {
let (mut r, tx, _e) = wired(1);
tx[0].push(1).unwrap();
r.engine_mut().fail_drain = true;
assert_eq!(r.tick().unwrap_err().code(), Code::Io);
assert_eq!(r.engine().runs(), vec![1]);
assert_eq!(r.engine().count(&Step::Flush), 1);
}
#[test]
fn a_command_with_no_key_is_run_without_a_hash() {
let (mut r, tx, _e) = wired(1);
tx[0].push(KEYLESS).unwrap();
let turn = r.tick().unwrap();
assert_eq!(turn.commands, 1);
assert_eq!(
r.engine().count(&Step::Prefetch(KEYLESS)),
0,
"there is nothing to warm for a command with no key"
);
}
#[test]
fn inline_execution_takes_the_same_walk_as_the_loop() {
let mut r = Reactor::inline(Recorder::new());
assert_eq!(r.execute(9), Flow::Next);
assert_eq!(
r.engine().steps(),
vec![Step::Prefetch(9), Step::Run(9)],
"no submit, no flush and no maintenance, and the same two calls"
);
assert_eq!(r.commands(), 1);
}
#[test]
fn inline_batches_pay_for_the_epoch_once() {
let mut r = Reactor::inline(Recorder::new());
assert_eq!(r.execute_all(0..20), 20);
assert_eq!(r.engine().runs(), (0..20).collect::<Vec<_>>());
let (last_prefetch, first_run) = walk_boundary(&r.engine().steps());
assert!(last_prefetch < first_run, "inline gets the two walks too");
}
#[test]
fn an_inline_batch_stops_at_a_break() {
let mut r = Reactor::inline(Recorder::new());
r.engine_mut().break_on = Some(2);
assert_eq!(r.execute_all(0..10), 3);
assert_eq!(r.engine().runs(), vec![0, 1, 2]);
assert_eq!(r.carried(), 0, "inline holds nothing over to next time");
r.engine_mut().break_on = None;
assert_eq!(r.execute_all(100..103), 3);
assert_eq!(r.engine().runs(), vec![0, 1, 2, 100, 101, 102]);
}
#[test]
fn an_empty_inline_batch_is_not_a_turn() {
let mut r = Reactor::inline(Recorder::new());
assert_eq!(r.execute_all(Vec::new()), 0);
assert_eq!(r.turns(), 0);
assert!(r.engine().steps().is_empty());
}
#[test]
fn run_until_returns_when_the_flag_is_set_and_the_lanes_are_dry() {
let (mut r, tx, _e) = wired(1);
for i in 0..300 {
tx[0].push(i).unwrap();
}
let stop = AtomicBool::new(true);
r.run_until(&stop).unwrap();
assert_eq!(r.commands(), 300, "the flag does not drop queued work");
assert!(r.turns() >= 5);
}
#[test]
fn a_reactor_says_what_it_has_been_doing() {
let (mut r, tx, _e) = wired(1);
tx[0].push(1).unwrap();
r.tick().unwrap();
let said = format!("{r:?}");
assert!(said.contains("commands: 1"), "{said}");
assert!(said.contains("lanes: 1"), "{said}");
}
}