use std::{collections::VecDeque, sync::Arc};
use tokio::time::Instant;
use crate::{TaskSpec, core::deferred_drop::OwnedTask, identity::TaskId};
pub(in crate::controller::engine) struct PendingSubmission {
pub(in crate::controller::engine) id: TaskId,
pub(in crate::controller::engine) task_name: Arc<str>,
pub(in crate::controller::engine) owned: OwnedTask<TaskSpec>,
}
impl PendingSubmission {
pub(in crate::controller::engine) fn new(
id: TaskId,
task_name: Arc<str>,
owned: OwnedTask<TaskSpec>,
) -> Self {
Self {
id,
task_name,
owned,
}
}
#[cfg(test)]
pub(in crate::controller::engine) fn task_spec(&self) -> &TaskSpec {
&self.owned.value
}
}
pub(in crate::controller::engine) struct SlotState {
phase: SlotPhase,
pub(in crate::controller::engine) queue: VecDeque<PendingSubmission>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(in crate::controller::engine) enum SlotPhase {
Idle,
Admitting {
owner: TaskId,
since: Instant,
},
CancelPendingAdmission {
owner: TaskId,
requested_at: Instant,
},
Running {
owner: TaskId,
started_at: Instant,
},
Terminating {
owner: TaskId,
requested_at: Instant,
},
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(in crate::controller::engine) enum ReplaceAction {
RemoveNow(TaskId),
WaitForAdmission,
AlreadyRequested,
Idle,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(in crate::controller::engine) enum AdmissionTransition {
Running,
RemoveNow(TaskId),
Stale,
}
impl SlotPhase {
pub(in crate::controller::engine) fn owner_id(self) -> Option<TaskId> {
match self {
Self::Idle => None,
Self::Admitting { owner, .. }
| Self::CancelPendingAdmission { owner, .. }
| Self::Running { owner, .. }
| Self::Terminating { owner, .. } => Some(owner),
}
}
pub(in crate::controller::engine) fn label(self) -> &'static str {
match self {
Self::Idle => "idle",
Self::Admitting { .. } => "admitting",
Self::Running { .. } => "running",
Self::CancelPendingAdmission { .. } | Self::Terminating { .. } => "terminating",
}
}
}
impl SlotState {
pub(in crate::controller::engine) fn new() -> Self {
Self {
phase: SlotPhase::Idle,
queue: VecDeque::new(),
}
}
pub(in crate::controller::engine) fn phase(&self) -> SlotPhase {
self.phase
}
pub(in crate::controller::engine) fn owner_id(&self) -> Option<TaskId> {
self.phase.owner_id()
}
pub(in crate::controller::engine) fn is_idle(&self) -> bool {
matches!(self.phase, SlotPhase::Idle)
}
pub(in crate::controller::engine) fn status_label(&self) -> &'static str {
self.phase.label()
}
pub(in crate::controller::engine) fn begin_admission(
&mut self,
owner: TaskId,
since: Instant,
) -> bool {
if !self.is_idle() {
return false;
}
self.phase = SlotPhase::Admitting { owner, since };
true
}
pub(in crate::controller::engine) fn request_replacement(
&mut self,
requested_at: Instant,
) -> ReplaceAction {
match self.phase {
SlotPhase::Idle => ReplaceAction::Idle,
SlotPhase::Admitting { owner, .. } => {
self.phase = SlotPhase::CancelPendingAdmission {
owner,
requested_at,
};
ReplaceAction::WaitForAdmission
}
SlotPhase::Running { owner, .. } => {
self.phase = SlotPhase::Terminating {
owner,
requested_at,
};
ReplaceAction::RemoveNow(owner)
}
SlotPhase::CancelPendingAdmission { .. } | SlotPhase::Terminating { .. } => {
ReplaceAction::AlreadyRequested
}
}
}
pub(in crate::controller::engine) fn confirm_admission(
&mut self,
owner: TaskId,
started_at: Instant,
) -> AdmissionTransition {
match self.phase {
SlotPhase::Admitting { owner: current, .. } if current == owner => {
self.phase = SlotPhase::Running { owner, started_at };
AdmissionTransition::Running
}
SlotPhase::CancelPendingAdmission {
owner: current,
requested_at,
} if current == owner => {
self.phase = SlotPhase::Terminating {
owner,
requested_at,
};
AdmissionTransition::RemoveNow(owner)
}
_ => AdmissionTransition::Stale,
}
}
pub(in crate::controller::engine) fn reject_admission(&mut self, owner: TaskId) -> bool {
let matches_current = matches!(
self.phase,
SlotPhase::Admitting { owner: current, .. }
| SlotPhase::CancelPendingAdmission { owner: current, .. }
if current == owner
);
if matches_current {
self.phase = SlotPhase::Idle;
}
matches_current
}
pub(in crate::controller::engine) fn complete_owner(&mut self, owner: TaskId) -> bool {
let matches_current = matches!(
self.phase,
SlotPhase::Running { owner: current, .. }
| SlotPhase::Terminating { owner: current, .. }
if current == owner
);
if matches_current {
self.phase = SlotPhase::Idle;
}
matches_current
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn new_slot_is_idle_with_empty_queue() {
let slot = SlotState::new();
assert_eq!(slot.phase(), SlotPhase::Idle);
assert_eq!(slot.owner_id(), None);
assert!(slot.queue.is_empty());
}
#[test]
fn every_occupied_phase_carries_its_owner() {
let owner = TaskId::next();
let now = Instant::now();
for phase in [
SlotPhase::Admitting { owner, since: now },
SlotPhase::CancelPendingAdmission {
owner,
requested_at: now,
},
SlotPhase::Running {
owner,
started_at: now,
},
SlotPhase::Terminating {
owner,
requested_at: now,
},
] {
assert_eq!(phase.owner_id(), Some(owner));
}
assert_eq!(SlotPhase::Idle.owner_id(), None);
}
#[test]
fn early_replace_waits_for_admission_then_enters_real_termination() {
let owner = TaskId::next();
let now = Instant::now();
let mut slot = SlotState::new();
assert!(slot.begin_admission(owner, now));
assert_eq!(
slot.request_replacement(now),
ReplaceAction::WaitForAdmission
);
assert!(matches!(
slot.phase(),
SlotPhase::CancelPendingAdmission { owner: id, .. } if id == owner
));
assert!(
!slot.complete_owner(owner),
"completion cannot release an admission that is still pending"
);
assert_eq!(
slot.confirm_admission(owner, now),
AdmissionTransition::RemoveNow(owner)
);
assert!(matches!(
slot.phase(),
SlotPhase::Terminating { owner: id, .. } if id == owner
));
assert!(slot.complete_owner(owner));
assert!(slot.is_idle());
}
#[test]
fn stale_results_do_not_mutate_current_owner() {
let owner = TaskId::next();
let stale = TaskId::next();
let now = Instant::now();
let mut slot = SlotState::new();
assert!(slot.begin_admission(owner, now));
assert_eq!(
slot.confirm_admission(stale, now),
AdmissionTransition::Stale
);
assert!(!slot.reject_admission(stale));
assert_eq!(slot.owner_id(), Some(owner));
assert!(matches!(slot.phase(), SlotPhase::Admitting { .. }));
}
#[test]
fn queue_push_pop_fifo() {
let mut slot = SlotState::new();
let pending = |name: &str| {
let task_spec = make_spec(name);
let retained = task_spec.task().clone();
let reservation = crate::core::deferred_drop::test_reservation();
PendingSubmission::new(
TaskId::next(),
Arc::from(name),
OwnedTask::new(task_spec, retained, reservation),
)
};
slot.queue.push_back(pending("a"));
slot.queue.push_back(pending("b"));
slot.queue.push_back(pending("c"));
assert_eq!(slot.queue.len(), 3);
assert_eq!(slot.queue.pop_front().unwrap().task_spec().name(), "a");
assert_eq!(slot.queue.pop_front().unwrap().task_spec().name(), "b");
assert_eq!(slot.queue.pop_front().unwrap().task_spec().name(), "c");
assert!(slot.queue.is_empty());
}
fn make_spec(name: &str) -> TaskSpec {
use crate::TaskContext;
use crate::{TaskFn, TaskRef};
let task: TaskRef = TaskFn::arc(|_ctx: TaskContext| async { Ok(()) });
TaskSpec::once(name, task)
}
}