use alloc::collections::VecDeque;
use core::time::Duration;
use std::collections::{HashMap, HashSet};
use std::time::{SystemTime, UNIX_EPOCH};
use routers_network::Entry;
use thiserror::Error;
use tokio::time::Instant;
use crate::bus::adapter::AckHandle;
use crate::event::{Payload, VehicleId};
use crate::orchestrator::admission;
use crate::protocol::ids::{JobId, ObservationId};
use crate::protocol::job::JobIdentity;
use crate::store::checkpoint::VehicleCheckpoint;
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub struct PublishedAtMicros(i64);
impl PublishedAtMicros {
#[must_use]
pub fn from_system_time(value: SystemTime) -> Option<Self> {
let micros = value.duration_since(UNIX_EPOCH).ok()?.as_micros();
i64::try_from(micros).ok().map(Self)
}
#[must_use]
pub const fn from_unix_micros(value: i64) -> Option<Self> {
if value < 0 { None } else { Some(Self(value)) }
}
#[must_use]
pub fn freshness_target_after(self, budget: Duration) -> i64 {
let budget_us = i64::try_from(budget.as_micros()).unwrap_or(i64::MAX);
self.0.saturating_add(budget_us)
}
}
#[allow(clippy::large_enum_variant)]
#[derive(Clone, Debug)]
pub enum CheckpointState<E: Entry> {
Unloaded,
Absent,
Present(VehicleCheckpoint<E>),
}
impl<E: Entry> CheckpointState<E> {
#[must_use]
pub fn present(&self) -> Option<&VehicleCheckpoint<E>> {
match self {
CheckpointState::Present(cp) => Some(cp),
CheckpointState::Unloaded | CheckpointState::Absent => None,
}
}
#[must_use]
pub fn is_loaded(&self) -> bool {
!matches!(self, CheckpointState::Unloaded)
}
}
pub struct PendingObservation<H: AckHandle> {
pub id: ObservationId,
pub payload: Payload,
pub published_at: PublishedAtMicros,
pub handle: H,
pub received: Instant,
}
#[derive(Debug)]
pub enum JobReservation {
Admitted(admission::Permit),
SyntheticTerminal,
}
#[derive(Debug)]
pub struct ActiveJob {
pub id: JobId,
pub identity: JobIdentity,
pub observation: ObservationId,
pub bytes: u64,
pub reservation: JobReservation,
pub dispatched: Instant,
}
#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
pub struct Depth {
pub vehicles: usize,
pub pending: usize,
pub active: usize,
}
pub struct VehicleState<E: Entry, H: AckHandle> {
pub checkpoint: CheckpointState<E>,
pub pending: VecDeque<PendingObservation<H>>,
pub active: Option<ActiveJob>,
pub committing: bool,
pub last_touch: Instant,
}
impl<E: Entry, H: AckHandle> VehicleState<E, H> {
fn new(now: Instant) -> Self {
Self {
checkpoint: CheckpointState::Unloaded,
pending: VecDeque::new(),
active: None,
committing: false,
last_touch: now,
}
}
fn is_eligible(&self) -> bool {
self.active.is_none() && !self.committing && !self.pending.is_empty()
}
fn is_idle(&self) -> bool {
self.pending.is_empty() && self.active.is_none() && !self.committing
}
}
#[must_use = "an enqueue result may carry an ack handle that must be acknowledged"]
pub enum Enqueue<H: AckHandle> {
Queued {
depth: usize,
},
Coalesced(H),
Committed(H),
Overflow {
handle: H,
limit: usize,
},
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub enum EnqueueReadiness {
Ready,
Committed,
Coalesced,
Backpressured {
limit: usize,
},
}
pub struct Finished<H: AckHandle> {
pub observation: PendingObservation<H>,
pub job: ActiveJob,
}
#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
pub struct SchedulerStats {
pub vehicles: usize,
pub pending: usize,
pub active: usize,
pub ready: usize,
pub committing: usize,
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub struct SchedulerConfig {
pub idle_ttl: Duration,
pub pending_limit: usize,
}
impl Default for SchedulerConfig {
fn default() -> Self {
Self {
idle_ttl: Duration::from_secs(10 * 60),
pending_limit: 256,
}
}
}
#[derive(Clone, Copy, Debug, PartialEq, Eq, Error)]
pub enum SchedulerError {
#[error("vehicle is not tracked by the scheduler")]
UnknownVehicle,
#[error("vehicle already has an active job")]
AlreadyActive,
#[error("vehicle has no active job")]
NoActiveJob,
#[error("vehicle is already committing")]
AlreadyCommitting,
#[error("vehicle has no head observation to finish")]
NoHead,
}
pub struct Scheduler<E: Entry, H: AckHandle> {
config: SchedulerConfig,
vehicles: HashMap<VehicleId, VehicleState<E, H>>,
ready: VecDeque<VehicleId>,
ready_set: HashSet<VehicleId>,
absent_checkpoint: CheckpointState<E>,
}
impl<E: Entry, H: AckHandle> Scheduler<E, H> {
#[must_use]
pub fn new(config: SchedulerConfig) -> Self {
Self {
config,
vehicles: HashMap::new(),
ready: VecDeque::new(),
ready_set: HashSet::new(),
absent_checkpoint: CheckpointState::Unloaded,
}
}
#[must_use]
pub fn config(&self) -> &SchedulerConfig {
&self.config
}
fn mark_ready(&mut self, vehicle: VehicleId) {
if self.ready_set.insert(vehicle) {
self.ready.push_back(vehicle);
}
}
pub fn enqueue(&mut self, obs: PendingObservation<H>, now: Instant) -> Enqueue<H> {
let vehicle = obs.payload.vehicle_id;
match self.enqueue_readiness(&obs) {
EnqueueReadiness::Committed => return Enqueue::Committed(obs.handle),
EnqueueReadiness::Coalesced => return Enqueue::Coalesced(obs.handle),
EnqueueReadiness::Backpressured { limit } => {
return Enqueue::Overflow {
handle: obs.handle,
limit,
};
}
EnqueueReadiness::Ready => {}
}
let state = self
.vehicles
.entry(vehicle)
.or_insert_with(|| VehicleState::new(now));
state.last_touch = now;
state.pending.push_back(obs);
let depth = state.pending.len();
let eligible = state.is_eligible();
if eligible {
self.mark_ready(vehicle);
}
Enqueue::Queued { depth }
}
#[must_use]
pub fn enqueue_readiness(&self, obs: &PendingObservation<H>) -> EnqueueReadiness {
let Some(state) = self.vehicles.get(&obs.payload.vehicle_id) else {
return EnqueueReadiness::Ready;
};
if let CheckpointState::Present(cp) = &state.checkpoint
&& obs.id.sequence <= cp.last_input.sequence
{
return EnqueueReadiness::Committed;
}
if state
.active
.as_ref()
.is_some_and(|active| active.observation == obs.id)
|| state.pending.iter().any(|pending| pending.id == obs.id)
{
return EnqueueReadiness::Coalesced;
}
if state.pending.len() >= self.config.pending_limit {
return EnqueueReadiness::Backpressured {
limit: self.config.pending_limit,
};
}
EnqueueReadiness::Ready
}
pub fn next_ready(&mut self) -> Option<VehicleId> {
while let Some(vehicle) = self.ready.pop_front() {
self.ready_set.remove(&vehicle);
if self
.vehicles
.get(&vehicle)
.is_some_and(|state| state.is_eligible())
{
return Some(vehicle);
}
}
None
}
#[must_use]
pub fn head(&self, vehicle: VehicleId) -> Option<&PendingObservation<H>> {
self.vehicles.get(&vehicle).and_then(|s| s.pending.front())
}
#[must_use]
pub fn checkpoint(&self, vehicle: VehicleId) -> &CheckpointState<E> {
self.vehicles
.get(&vehicle)
.map_or(&self.absent_checkpoint, |s| &s.checkpoint)
}
pub fn set_checkpoint(&mut self, vehicle: VehicleId, checkpoint: CheckpointState<E>) {
if let Some(state) = self.vehicles.get_mut(&vehicle) {
state.checkpoint = checkpoint;
}
}
pub fn activate(&mut self, vehicle: VehicleId, job: ActiveJob) -> Result<(), SchedulerError> {
let state = self
.vehicles
.get_mut(&vehicle)
.ok_or(SchedulerError::UnknownVehicle)?;
if state.active.is_some() {
return Err(SchedulerError::AlreadyActive);
}
state.active = Some(job);
Ok(())
}
#[must_use]
pub fn state(&self, vehicle: VehicleId) -> Option<&VehicleState<E, H>> {
self.vehicles.get(&vehicle)
}
#[must_use]
pub fn active(&self, vehicle: VehicleId) -> Option<&ActiveJob> {
self.vehicles.get(&vehicle).and_then(|s| s.active.as_ref())
}
pub fn active_mut(&mut self, vehicle: VehicleId) -> Option<&mut ActiveJob> {
self.vehicles
.get_mut(&vehicle)
.and_then(|s| s.active.as_mut())
}
pub fn begin_commit(&mut self, vehicle: VehicleId) -> Result<(), SchedulerError> {
let state = self
.vehicles
.get_mut(&vehicle)
.ok_or(SchedulerError::UnknownVehicle)?;
if state.active.is_none() {
return Err(SchedulerError::NoActiveJob);
}
if state.committing {
return Err(SchedulerError::AlreadyCommitting);
}
state.committing = true;
Ok(())
}
pub fn finish(
&mut self,
vehicle: VehicleId,
now: Instant,
) -> Result<Finished<H>, SchedulerError> {
let state = self
.vehicles
.get_mut(&vehicle)
.ok_or(SchedulerError::UnknownVehicle)?;
let job = state.active.take().ok_or(SchedulerError::NoActiveJob)?;
let Some(observation) = state.pending.pop_front() else {
state.active = Some(job);
return Err(SchedulerError::NoHead);
};
state.committing = false;
state.last_touch = now;
let eligible = state.is_eligible();
if eligible {
self.mark_ready(vehicle);
}
Ok(Finished { observation, job })
}
pub fn abandon(&mut self, vehicle: VehicleId) -> Option<ActiveJob> {
let state = self.vehicles.get_mut(&vehicle)?;
let job = state.active.take()?;
let eligible = state.is_eligible();
if eligible {
self.mark_ready(vehicle);
}
Some(job)
}
pub fn rollback_commit(&mut self, vehicle: VehicleId) -> Option<ActiveJob> {
let state = self.vehicles.get_mut(&vehicle)?;
state.committing = false;
let job = state.active.take()?;
let eligible = state.is_eligible();
if eligible {
self.mark_ready(vehicle);
}
Some(job)
}
pub fn evict_idle(&mut self, now: Instant) -> Vec<VehicleId> {
let ttl = self.config.idle_ttl;
let mut evicted = Vec::new();
self.vehicles.retain(|&vehicle, state| {
let expired = now.saturating_duration_since(state.last_touch) >= ttl;
if state.is_idle() && expired {
evicted.push(vehicle);
false
} else {
true
}
});
evicted
}
#[must_use]
pub fn stats(&self) -> SchedulerStats {
let mut stats = SchedulerStats {
vehicles: self.vehicles.len(),
ready: self.ready_set.len(),
..SchedulerStats::default()
};
for state in self.vehicles.values() {
stats.pending += state.pending.len();
if state.active.is_some() {
stats.active += 1;
}
if state.committing {
stats.committing += 1;
}
}
stats
}
#[must_use]
pub fn oldest_pending(&self, now: Instant) -> Option<Duration> {
self.vehicles
.values()
.filter_map(|s| s.pending.front())
.map(|p| p.received)
.min()
.map(|oldest| now.saturating_duration_since(oldest))
}
#[must_use]
pub fn depth(&self) -> Depth {
Depth {
vehicles: self.vehicles.len(),
pending: self.vehicles.values().map(|s| s.pending.len()).sum(),
active: self
.vehicles
.values()
.filter(|s| s.active.is_some())
.count(),
}
}
}
#[cfg(test)]
mod tests {
use alloc::sync::Arc;
use std::sync::Mutex;
use chrono::{DateTime, Utc};
use geo::Point;
use routers_network::mock::MockEntryId;
use routers_transition::matcher::Trip;
use super::*;
use crate::orchestrator::admission::{Admission, AdmissionConfig};
use crate::protocol::ids::{GraphVersion, RegionId, Revision, SCHEMA_VERSION, SegmentId};
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
enum AckOp {
Ack(u64),
Nak(u64),
}
#[derive(Clone)]
struct TestAck {
seq: u64,
log: Arc<Mutex<Vec<AckOp>>>,
}
impl AckHandle for TestAck {
async fn ack(self) -> anyhow::Result<()> {
self.log.lock().unwrap().push(AckOp::Ack(self.seq));
Ok(())
}
async fn nak(self, _delay: Option<Duration>) -> anyhow::Result<()> {
self.log.lock().unwrap().push(AckOp::Nak(self.seq));
Ok(())
}
fn sequence(&self) -> u64 {
self.seq
}
fn deliveries(&self) -> u32 {
1
}
}
struct Fixture {
log: Arc<Mutex<Vec<AckOp>>>,
admission: Admission,
region: RegionId,
base: Instant,
}
impl Fixture {
fn new() -> Self {
let region = RegionId::new("r1").unwrap();
let admission = Admission::new(AdmissionConfig::default(), core::iter::once(®ion));
Self {
log: Arc::new(Mutex::new(Vec::new())),
admission,
region,
base: Instant::now(),
}
}
fn at(&self, secs: u64) -> Instant {
self.base + Duration::from_secs(secs)
}
fn obs(&self, vehicle: u64, seq: u64, received: Instant) -> PendingObservation<TestAck> {
PendingObservation {
id: ObservationId {
partition: 7,
sequence: seq,
},
payload: Payload {
vehicle_id: VehicleId(vehicle),
timestamp: DateTime::<Utc>::from_timestamp(0, 0).unwrap(),
point: Point::new(0.0, 0.0),
},
published_at: PublishedAtMicros::from_unix_micros(0)
.expect("zero is a valid publication time"),
handle: TestAck {
seq,
log: self.log.clone(),
},
received,
}
}
fn identity(&self, vehicle: u64, seq: u64) -> JobIdentity {
JobIdentity {
schema: SCHEMA_VERSION,
vehicle_id: VehicleId(vehicle),
observation: ObservationId {
partition: 7,
sequence: seq,
},
base: None,
graph: GraphVersion::new("g1").unwrap(),
region: self.region.clone(),
}
}
fn job(&self, vehicle: u64, seq: u64, now: Instant) -> ActiveJob {
let identity = self.identity(vehicle, seq);
let id = identity.local_decision_id();
ActiveJob {
id,
identity,
observation: ObservationId {
partition: 7,
sequence: seq,
},
bytes: 100,
reservation: JobReservation::Admitted(
self.admission.try_admit(&self.region, 100).unwrap(),
),
dispatched: now,
}
}
fn present_checkpoint(&self, last_seq: u64) -> CheckpointState<MockEntryId> {
CheckpointState::Present(VehicleCheckpoint {
trip: Trip::new(),
last_input: ObservationId {
partition: 7,
sequence: last_seq,
},
revision: Revision(last_seq),
segment: SegmentId(1),
finalized_through: None,
graph: GraphVersion::new("g1").unwrap(),
schema: SCHEMA_VERSION,
region: self.region.clone(),
routing_version: 1,
})
}
}
fn scheduler(_fx: &Fixture) -> Scheduler<MockEntryId, TestAck> {
Scheduler::new(SchedulerConfig::default())
}
#[track_caller]
fn feed(
sched: &mut Scheduler<MockEntryId, TestAck>,
obs: PendingObservation<TestAck>,
now: Instant,
) {
assert!(matches!(sched.enqueue(obs, now), Enqueue::Queued { .. }));
}
#[track_caller]
fn expect_queued(result: Enqueue<TestAck>) -> usize {
match result {
Enqueue::Queued { depth } => depth,
Enqueue::Coalesced(_) => panic!("expected Queued, got Coalesced"),
Enqueue::Committed(_) => panic!("expected Queued, got Committed"),
Enqueue::Overflow { .. } => panic!("expected Queued, got Overflow"),
}
}
#[test]
fn fifo_is_preserved_per_vehicle() {
let fx = Fixture::new();
let mut sched = scheduler(&fx);
let now = fx.at(0);
assert_eq!(expect_queued(sched.enqueue(fx.obs(1, 10, now), now)), 1);
assert_eq!(expect_queued(sched.enqueue(fx.obs(1, 11, now), now)), 2);
assert_eq!(expect_queued(sched.enqueue(fx.obs(1, 12, now), now)), 3);
assert_eq!(sched.head(VehicleId(1)).unwrap().id.sequence, 10);
assert_eq!(sched.next_ready(), Some(VehicleId(1)));
assert_eq!(sched.next_ready(), None);
sched.activate(VehicleId(1), fx.job(1, 10, now)).unwrap();
let finished = sched.finish(VehicleId(1), now).unwrap();
assert_eq!(finished.observation.id.sequence, 10);
assert_eq!(sched.head(VehicleId(1)).unwrap().id.sequence, 11);
assert_eq!(sched.next_ready(), Some(VehicleId(1)));
}
#[test]
fn ready_queue_is_fair_across_vehicles() {
let fx = Fixture::new();
let mut sched = scheduler(&fx);
let now = fx.at(0);
feed(&mut sched, fx.obs(1, 1, now), now);
feed(&mut sched, fx.obs(2, 1, now), now);
feed(&mut sched, fx.obs(3, 1, now), now);
assert_eq!(sched.next_ready(), Some(VehicleId(1)));
assert_eq!(sched.next_ready(), Some(VehicleId(2)));
assert_eq!(sched.next_ready(), Some(VehicleId(3)));
assert_eq!(sched.next_ready(), None);
sched.activate(VehicleId(1), fx.job(1, 1, now)).unwrap();
feed(&mut sched, fx.obs(1, 2, now), now);
feed(&mut sched, fx.obs(2, 2, now), now);
sched.finish(VehicleId(1), now).unwrap();
assert_eq!(sched.next_ready(), Some(VehicleId(2)));
assert_eq!(sched.next_ready(), Some(VehicleId(1)));
}
#[test]
fn duplicate_pending_observation_is_coalesced() {
let fx = Fixture::new();
let mut sched = scheduler(&fx);
let now = fx.at(0);
feed(&mut sched, fx.obs(1, 5, now), now);
let again = sched.enqueue(fx.obs(1, 5, now), now);
match again {
Enqueue::Coalesced(handle) => assert_eq!(handle.sequence(), 5),
_ => panic!("expected Coalesced for a duplicate pending id"),
}
assert_eq!(sched.stats().pending, 1);
}
#[test]
fn duplicate_active_observation_is_coalesced() {
let fx = Fixture::new();
let mut sched = scheduler(&fx);
let now = fx.at(0);
feed(&mut sched, fx.obs(1, 5, now), now);
sched.next_ready();
sched.activate(VehicleId(1), fx.job(1, 5, now)).unwrap();
let again = sched.enqueue(fx.obs(1, 5, now), now);
match again {
Enqueue::Coalesced(handle) => assert_eq!(handle.sequence(), 5),
_ => panic!("expected Coalesced for a duplicate active id"),
}
assert_eq!(sched.stats().pending, 1);
}
#[test]
fn already_committed_observation_is_suppressed() {
let fx = Fixture::new();
let mut sched = scheduler(&fx);
let now = fx.at(0);
feed(&mut sched, fx.obs(1, 21, now), now);
sched.set_checkpoint(VehicleId(1), fx.present_checkpoint(20));
match sched.enqueue(fx.obs(1, 20, now), now) {
Enqueue::Committed(handle) => assert_eq!(handle.sequence(), 20),
_ => panic!("expected Committed at the boundary"),
}
match sched.enqueue(fx.obs(1, 19, now), now) {
Enqueue::Committed(handle) => assert_eq!(handle.sequence(), 19),
_ => panic!("expected Committed below the boundary"),
}
assert!(matches!(
sched.enqueue(fx.obs(1, 22, now), now),
Enqueue::Queued { .. }
));
}
#[test]
fn suppression_needs_a_loaded_checkpoint() {
let fx = Fixture::new();
let mut sched = scheduler(&fx);
let now = fx.at(0);
assert!(matches!(
sched.enqueue(fx.obs(1, 1, now), now),
Enqueue::Queued { .. }
));
}
#[test]
fn pending_overflow_refuses_and_returns_the_handle() {
let fx = Fixture::new();
let mut sched = Scheduler::<MockEntryId, TestAck>::new(SchedulerConfig {
pending_limit: 2,
..SchedulerConfig::default()
});
let now = fx.at(0);
assert!(matches!(
sched.enqueue(fx.obs(1, 1, now), now),
Enqueue::Queued { depth: 1 }
));
assert!(matches!(
sched.enqueue(fx.obs(1, 2, now), now),
Enqueue::Queued { depth: 2 }
));
match sched.enqueue(fx.obs(1, 3, now), now) {
Enqueue::Overflow { handle, limit } => {
assert_eq!(handle.sequence(), 3);
assert_eq!(limit, 2);
}
_ => panic!("expected Overflow at the limit"),
}
assert_eq!(sched.stats().pending, 2);
}
#[test]
fn activate_twice_is_rejected() {
let fx = Fixture::new();
let mut sched = scheduler(&fx);
let now = fx.at(0);
feed(&mut sched, fx.obs(1, 1, now), now);
sched.next_ready();
sched.activate(VehicleId(1), fx.job(1, 1, now)).unwrap();
assert_eq!(
sched.activate(VehicleId(1), fx.job(1, 1, now)),
Err(SchedulerError::AlreadyActive)
);
}
#[test]
fn activate_unknown_vehicle_is_rejected() {
let fx = Fixture::new();
let mut sched = scheduler(&fx);
let now = fx.at(0);
assert_eq!(
sched.activate(VehicleId(99), fx.job(99, 1, now)),
Err(SchedulerError::UnknownVehicle)
);
}
#[test]
fn begin_commit_transitions_and_errors() {
let fx = Fixture::new();
let mut sched = scheduler(&fx);
let now = fx.at(0);
feed(&mut sched, fx.obs(1, 1, now), now);
assert_eq!(
sched.begin_commit(VehicleId(1)),
Err(SchedulerError::NoActiveJob)
);
sched.next_ready();
sched.activate(VehicleId(1), fx.job(1, 1, now)).unwrap();
assert_eq!(sched.begin_commit(VehicleId(1)), Ok(()));
assert_eq!(
sched.begin_commit(VehicleId(1)),
Err(SchedulerError::AlreadyCommitting)
);
}
#[test]
fn finish_without_active_job_is_rejected() {
let fx = Fixture::new();
let mut sched = scheduler(&fx);
let now = fx.at(0);
feed(&mut sched, fx.obs(1, 1, now), now);
assert!(matches!(
sched.finish(VehicleId(1), now),
Err(SchedulerError::NoActiveJob)
));
}
#[test]
fn finish_pops_head_and_clears_committing() {
let fx = Fixture::new();
let mut sched = scheduler(&fx);
let now = fx.at(0);
feed(&mut sched, fx.obs(1, 1, now), now);
sched.next_ready();
sched.activate(VehicleId(1), fx.job(1, 1, now)).unwrap();
sched.begin_commit(VehicleId(1)).unwrap();
let finished = sched.finish(VehicleId(1), now).unwrap();
assert_eq!(finished.observation.id.sequence, 1);
assert_eq!(finished.job.observation.sequence, 1);
let stats = sched.stats();
assert_eq!(stats.active, 0);
assert_eq!(stats.committing, 0);
assert_eq!(stats.pending, 0);
assert_eq!(sched.next_ready(), None);
}
#[test]
fn abandon_keeps_the_head_and_re_readies() {
let fx = Fixture::new();
let mut sched = scheduler(&fx);
let now = fx.at(0);
feed(&mut sched, fx.obs(1, 1, now), now);
sched.next_ready();
sched.activate(VehicleId(1), fx.job(1, 1, now)).unwrap();
let job = sched.abandon(VehicleId(1)).expect("a job to abandon");
assert_eq!(job.observation.sequence, 1);
assert_eq!(sched.head(VehicleId(1)).unwrap().id.sequence, 1);
assert_eq!(sched.next_ready(), Some(VehicleId(1)));
assert!(sched.abandon(VehicleId(1)).is_none());
}
#[test]
fn rollback_commit_keeps_the_raw_head_and_re_readies_it() {
let fx = Fixture::new();
let mut sched = scheduler(&fx);
let now = fx.at(0);
feed(&mut sched, fx.obs(1, 1, now), now);
sched.next_ready();
sched.activate(VehicleId(1), fx.job(1, 1, now)).unwrap();
sched.begin_commit(VehicleId(1)).unwrap();
let job = sched
.rollback_commit(VehicleId(1))
.expect("active commit rolls back");
assert_eq!(job.observation.sequence, 1);
assert_eq!(sched.head(VehicleId(1)).unwrap().id.sequence, 1);
assert_eq!(sched.next_ready(), Some(VehicleId(1)));
assert!(!sched.state(VehicleId(1)).unwrap().committing);
}
#[test]
fn eviction_only_reclaims_idle_vehicles() {
let fx = Fixture::new();
let mut sched = scheduler(&fx);
let ttl = sched.config().idle_ttl;
let t0 = fx.at(0);
feed(&mut sched, fx.obs(1, 1, t0), t0);
feed(&mut sched, fx.obs(2, 1, t0), t0);
sched.next_ready();
sched.next_ready();
sched.activate(VehicleId(2), fx.job(2, 1, t0)).unwrap();
feed(&mut sched, fx.obs(3, 1, t0), t0);
sched.next_ready();
sched.activate(VehicleId(3), fx.job(3, 1, t0)).unwrap();
sched.begin_commit(VehicleId(3)).unwrap();
feed(&mut sched, fx.obs(4, 1, t0), t0);
sched.next_ready();
sched.activate(VehicleId(4), fx.job(4, 1, t0)).unwrap();
sched.finish(VehicleId(4), t0).unwrap();
assert!(sched.evict_idle(t0).is_empty());
let after = t0 + ttl;
let evicted = sched.evict_idle(after);
assert_eq!(evicted, vec![VehicleId(4)]);
let stats = sched.stats();
assert_eq!(stats.vehicles, 3);
assert_eq!(stats.pending, 3);
}
#[test]
fn oldest_pending_picks_the_oldest_head() {
let fx = Fixture::new();
let mut sched = scheduler(&fx);
feed(&mut sched, fx.obs(1, 1, fx.at(0)), fx.at(0));
feed(&mut sched, fx.obs(1, 2, fx.at(20)), fx.at(20));
feed(&mut sched, fx.obs(2, 1, fx.at(5)), fx.at(5));
let now = fx.at(30);
assert_eq!(sched.oldest_pending(now), Some(Duration::from_secs(30)));
let empty = Scheduler::<MockEntryId, TestAck>::new(SchedulerConfig::default());
assert_eq!(empty.oldest_pending(now), None);
}
#[test]
fn checkpoint_accessor_is_total() {
let fx = Fixture::new();
let mut sched = scheduler(&fx);
let now = fx.at(0);
assert!(matches!(
sched.checkpoint(VehicleId(1)),
CheckpointState::Unloaded
));
feed(&mut sched, fx.obs(1, 1, now), now);
sched.set_checkpoint(VehicleId(1), CheckpointState::Absent);
assert!(matches!(
sched.checkpoint(VehicleId(1)),
CheckpointState::Absent
));
sched.set_checkpoint(VehicleId(2), CheckpointState::Absent);
assert!(matches!(
sched.checkpoint(VehicleId(2)),
CheckpointState::Unloaded
));
}
#[test]
fn returned_handle_can_be_acknowledged() {
let fx = Fixture::new();
let mut sched = scheduler(&fx);
let now = fx.at(0);
feed(&mut sched, fx.obs(1, 5, now), now);
let Enqueue::Coalesced(handle) = sched.enqueue(fx.obs(1, 5, now), now) else {
panic!("expected Coalesced");
};
futures::executor::block_on(handle.ack()).unwrap();
assert_eq!(&*fx.log.lock().unwrap(), &[AckOp::Ack(5)]);
}
}