use core::convert::Infallible;
use core::marker::PhantomData;
use std::collections::{BTreeMap, VecDeque};
use std::time::Duration;
use crate::{
Actions, Address, Behavior, Births, Crash, CreationRejection, Delivery, Exit, Never, Own,
Proxy, ProxyCommand, Recipient, RestartPolicy, SendAlgebra, SendInput, Strategy,
SupervisionEvent, Supervisor, SupervisorSends, User, WorkerCreationResolved, WorkerStopped,
delegate_transition,
};
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
pub struct JobId(pub u64);
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
pub struct AssignmentId(pub u64);
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct PoolAssignment<J> {
pub assignment: AssignmentId,
pub job: JobId,
pub payload: J,
}
#[derive(Clone, PartialEq, Eq)]
pub enum PoolMessage<A: Address, J, R> {
Submit {
job: JobId,
payload: J,
reply_to: Recipient<A, PoolResponse<J, R, A>>,
},
Completed {
worker: A::Nonce,
assignment: AssignmentId,
result: R,
},
}
#[derive(Clone, PartialEq, Eq)]
pub enum KeyedPoolMessage<A: Address, K, J, R> {
Submit {
key: K,
job: JobId,
payload: J,
reply_to: Recipient<A, PoolResponse<J, R, A>>,
},
Completed {
worker: A::Nonce,
assignment: AssignmentId,
result: R,
},
Rebalance {
key: K,
worker: A::Nonce,
},
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum PoolRejection {
BacklogFull,
AffinityUnavailable,
}
#[derive(Clone, PartialEq, Eq)]
pub enum PoolInterruption<A: Address> {
WorkerStopped {
worker: A::Nonce,
outcome: Result<Exit<A>, Crash>,
},
NoRecoverableWorkers,
AffinityRetired {
worker: A::Nonce,
reason: WorkerRetirement,
},
}
#[derive(Clone, PartialEq, Eq)]
pub enum PoolResponse<J, R, A: Address> {
Accepted {
job: JobId,
},
Rejected {
job: JobId,
payload: J,
reason: PoolRejection,
},
Completed {
job: JobId,
result: R,
},
Interrupted {
job: JobId,
payload: J,
reason: PoolInterruption<A>,
},
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum InterruptionPolicy {
Fail,
Retry,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum WorkerPhase {
Installing,
Idle,
Assigned {
assignment: AssignmentId,
job: JobId,
},
Retired {
reason: WorkerRetirement,
},
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum WorkerRetirement {
CreationRejected(CreationRejection),
ReplacementUnavailable,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum PoolConfigError<N> {
NoWorkers,
DuplicateWorker(N),
}
pub trait AffinitySelector<K, N> {
fn select(&self, key: &K) -> N;
}
impl<K, N, F> AffinitySelector<K, N> for F
where
F: Fn(&K) -> N,
{
fn select(&self, key: &K) -> N {
self(key)
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum PoolError<N> {
UnknownWorker(N),
CompletionForUnavailableWorker {
worker: N,
phase: WorkerPhase,
},
StaleCompletion {
worker: N,
expected: AssignmentId,
received: AssignmentId,
},
WorkerStoppedWhileUnavailable {
worker: N,
phase: WorkerPhase,
},
CreationResolvedWhileUnavailable {
worker: N,
phase: WorkerPhase,
},
RebalanceToRetiredWorker {
worker: N,
reason: WorkerRetirement,
},
}
struct AcceptedJob<A: Address, J, R> {
id: JobId,
payload: J,
reply_to: Recipient<A, PoolResponse<J, R, A>>,
interruption: Option<PoolInterruption<A>>,
target: Option<A::Nonce>,
}
struct QueuedJob<A: Address, J, R> {
accepted: AcceptedJob<A, J, R>,
dispatch_payload: J,
}
enum SlotState<A: Address, J, R> {
Installing,
Idle,
Assigned {
assignment: AssignmentId,
job: AcceptedJob<A, J, R>,
},
Retired {
reason: WorkerRetirement,
},
}
struct Slot<A: Address, J, R> {
nonce: A::Nonce,
state: SlotState<A, J, R>,
}
struct PlannedDispatch {
slot_position: usize,
job_position: usize,
}
enum Admission {
Accepted,
Rejected,
}
pub type PoolEvent<A, J, R> = SupervisionEvent<User<A, PoolMessage<A, J, R>>>;
pub type KeyedPoolEvent<A, K, J, R> = SupervisionEvent<User<A, KeyedPoolMessage<A, K, J, R>>>;
pub struct PoolBehaviorSends<A: Address, J, R, C: Behavior<Addr = A>> {
pub responses: Vec<Delivery<A, PoolResponse<J, R, A>>>,
pub assignments: Vec<Delivery<A, ProxyCommand<C>>>,
}
impl<A: Address, J, R, C: Behavior<Addr = A>> SendAlgebra for PoolBehaviorSends<A, J, R, C> {
fn empty() -> Self {
Self {
responses: Vec::new(),
assignments: Vec::new(),
}
}
fn append(&mut self, mut other: Self) {
self.responses.append(&mut other.responses);
self.assignments.append(&mut other.assignments);
}
}
impl<A: Address, J, R, C: Behavior<Addr = A>> SendInput<Delivery<A, PoolResponse<J, R, A>>, Own>
for PoolBehaviorSends<A, J, R, C>
{
fn emit(&mut self, input: Delivery<A, PoolResponse<J, R, A>>) {
self.responses.push(input);
}
}
impl<A: Address, J, R, C: Behavior<Addr = A>> SendInput<Delivery<A, ProxyCommand<C>>, Own>
for PoolBehaviorSends<A, J, R, C>
{
fn emit(&mut self, input: Delivery<A, ProxyCommand<C>>) {
self.assignments.push(input);
}
}
type KernelSends<A, J, R, C> = PoolBehaviorSends<A, J, R, C>;
pub type PoolSends<A, J, R, C> = SupervisorSends<A, KernelSends<A, J, R, C>, C>;
pub type PoolActions<A, J, R, C> = Actions<A, Never, PoolSends<A, J, R, C>, Births<Proxy<C>>>;
struct PoolKernel<A: Address, J, R, C>(PhantomData<fn(A, J, R, C)>);
impl<A: Address, J, R, C> PoolKernel<A, J, R, C> {
const fn new() -> Self {
Self(PhantomData)
}
}
impl<A, J, R, C> Behavior for PoolKernel<A, J, R, C>
where
A: Address,
C: Behavior<Addr = A, Msg = PoolAssignment<J>, Ph = Never>,
{
type Addr = A;
type Msg = PoolMessage<A, J, R>;
type Event = User<A, PoolMessage<A, J, R>>;
type Sends = KernelSends<A, J, R, C>;
type Ph = Never;
type Error = Infallible;
type Birth = Births<C>;
fn init(&mut self) -> crate::BehaviorActed<Self> {
Ok(Actions::cont())
}
fn transition(&mut self, _event: Self::Event) -> crate::BehaviorActed<Self> {
Ok(Actions::cont())
}
}
type PoolSupervisor<A, J, R, C> = Supervisor<PoolKernel<A, J, R, C>, C>;
pub struct WorkerPool<A: Address, J, R, C>
where
C: Behavior<Addr = A, Msg = PoolAssignment<J>, Ph = Never>,
{
supervisor: PoolSupervisor<A, J, R, C>,
slots: Vec<Slot<A, J, R>>,
backlog: VecDeque<QueuedJob<A, J, R>>,
backlog_capacity: usize,
next_assignment: u64,
interruption: InterruptionPolicy,
}
impl<A, J, R, C> WorkerPool<A, J, R, C>
where
A: Address,
C: Behavior<Addr = A, Msg = PoolAssignment<J>, Ph = Never>,
{
#[allow(
clippy::too_many_arguments,
reason = "the arguments expose the complete pool policy"
)]
pub fn new(
nonces: fn(usize) -> A::Nonce,
count: usize,
build: fn(usize) -> C,
backlog_capacity: usize,
interruption: InterruptionPolicy,
restart_policy: RestartPolicy,
max_restarts: u32,
restart_window: Duration,
) -> Result<Self, PoolConfigError<A::Nonce>> {
if count == 0 {
return Err(PoolConfigError::NoWorkers);
}
let mut slots = Vec::with_capacity(count);
for index in 0..count {
let nonce = nonces(index);
if slots.iter().any(|slot: &Slot<A, J, R>| slot.nonce == nonce) {
return Err(PoolConfigError::DuplicateWorker(nonce));
}
slots.push(Slot {
nonce,
state: SlotState::Installing,
});
}
Ok(Self {
supervisor: Supervisor::new(
PoolKernel::new(),
nonces,
count,
build,
Strategy::OneForOne,
restart_policy,
max_restarts,
restart_window,
),
slots,
backlog: VecDeque::new(),
backlog_capacity,
next_assignment: 0,
interruption,
})
}
#[must_use]
pub fn backlog_len(&self) -> usize {
self.backlog.len()
}
#[must_use]
pub fn worker_phase(&self, worker: A::Nonce) -> Option<WorkerPhase> {
self.slots
.iter()
.find(|slot| slot.nonce == worker)
.map(|slot| match &slot.state {
SlotState::Installing => WorkerPhase::Installing,
SlotState::Idle => WorkerPhase::Idle,
SlotState::Assigned { assignment, job } => WorkerPhase::Assigned {
assignment: *assignment,
job: job.id,
},
SlotState::Retired { reason } => WorkerPhase::Retired { reason: *reason },
})
}
fn slot_position(&self, worker: A::Nonce) -> Result<usize, PoolError<A::Nonce>> {
self.slots
.iter()
.position(|slot| slot.nonce == worker)
.ok_or(PoolError::UnknownWorker(worker))
}
}
impl<A, J, R, C> WorkerPool<A, J, R, C>
where
A: Address,
A::Nonce: From<u64>,
J: Clone,
C: Behavior<Addr = A, Msg = PoolAssignment<J>, Ph = Never>,
{
fn supervisor_transition(&mut self, event: PoolEvent<A, J, R>) -> PoolActions<A, J, R, C> {
match delegate_transition(&mut self.supervisor, event) {
Ok(actions) => actions,
Err(never) => match never {},
}
}
fn submit(
&mut self,
job: JobId,
payload: J,
reply_to: Recipient<A, PoolResponse<J, R, A>>,
actions: &mut PoolActions<A, J, R, C>,
) {
let can_dispatch = self
.slots
.iter()
.any(|slot| matches!(slot.state, SlotState::Idle));
if !can_dispatch && self.backlog.len() == self.backlog_capacity {
actions.sends.behavior.send::<_, Own>(Delivery::new(
reply_to,
PoolResponse::Rejected {
job,
payload,
reason: PoolRejection::BacklogFull,
},
));
return;
}
let dispatch_payload = payload.clone();
self.backlog.push_back(QueuedJob {
accepted: AcceptedJob {
id: job,
payload,
reply_to,
interruption: None,
target: None,
},
dispatch_payload,
});
actions
.sends
.behavior
.send::<_, Own>(Delivery::new(reply_to, PoolResponse::Accepted { job }));
}
fn submit_to(
&mut self,
target: A::Nonce,
job: JobId,
payload: J,
reply_to: Recipient<A, PoolResponse<J, R, A>>,
actions: &mut PoolActions<A, J, R, C>,
) -> Admission {
let Some(slot) = self.slots.iter().find(|slot| slot.nonce == target) else {
actions.sends.behavior.send::<_, Own>(Delivery::new(
reply_to,
PoolResponse::Rejected {
job,
payload,
reason: PoolRejection::AffinityUnavailable,
},
));
return Admission::Rejected;
};
if matches!(slot.state, SlotState::Retired { .. }) {
actions.sends.behavior.send::<_, Own>(Delivery::new(
reply_to,
PoolResponse::Rejected {
job,
payload,
reason: PoolRejection::AffinityUnavailable,
},
));
return Admission::Rejected;
}
let can_dispatch = matches!(slot.state, SlotState::Idle);
if !can_dispatch && self.backlog.len() == self.backlog_capacity {
actions.sends.behavior.send::<_, Own>(Delivery::new(
reply_to,
PoolResponse::Rejected {
job,
payload,
reason: PoolRejection::BacklogFull,
},
));
return Admission::Rejected;
}
let dispatch_payload = payload.clone();
self.backlog.push_back(QueuedJob {
accepted: AcceptedJob {
id: job,
payload,
reply_to,
interruption: None,
target: Some(target),
},
dispatch_payload,
});
actions
.sends
.behavior
.send::<_, Own>(Delivery::new(reply_to, PoolResponse::Accepted { job }));
Admission::Accepted
}
fn complete(
&mut self,
worker: A::Nonce,
assignment: AssignmentId,
result: R,
actions: &mut PoolActions<A, J, R, C>,
) -> Result<(), PoolError<A::Nonce>> {
let position = self.slot_position(worker)?;
let phase = self
.worker_phase(worker)
.expect("position proves the slot exists");
let SlotState::Assigned {
assignment: expected,
..
} = &self.slots[position].state
else {
return Err(PoolError::CompletionForUnavailableWorker { worker, phase });
};
if *expected != assignment {
return Err(PoolError::StaleCompletion {
worker,
expected: *expected,
received: assignment,
});
}
let SlotState::Assigned { job, .. } =
core::mem::replace(&mut self.slots[position].state, SlotState::Idle)
else {
unreachable!("the state was proven assigned")
};
actions.sends.behavior.send::<_, Own>(Delivery::new(
job.reply_to,
PoolResponse::Completed {
job: job.id,
result,
},
));
Ok(())
}
fn worker_stopped(
&mut self,
stopped: &WorkerStopped<A>,
responses: &mut Vec<Delivery<A, PoolResponse<J, R, A>>>,
) -> Result<(), PoolError<A::Nonce>> {
let position = self.slot_position(stopped.proxy)?;
let phase = self
.worker_phase(stopped.proxy)
.expect("position proves the slot exists");
if matches!(phase, WorkerPhase::Installing | WorkerPhase::Retired { .. }) {
return Err(PoolError::WorkerStoppedWhileUnavailable {
worker: stopped.proxy,
phase,
});
}
if self.interruption == InterruptionPolicy::Retry {
if let SlotState::Assigned { job, .. } = &self.slots[position].state {
let dispatch_payload = job.payload.clone();
let SlotState::Assigned { mut job, .. } =
core::mem::replace(&mut self.slots[position].state, SlotState::Installing)
else {
unreachable!("the assigned state was matched before committing retry")
};
job.interruption = Some(PoolInterruption::WorkerStopped {
worker: stopped.proxy,
outcome: stopped.outcome,
});
self.backlog.push_front(QueuedJob {
accepted: job,
dispatch_payload,
});
} else {
self.slots[position].state = SlotState::Installing;
}
return Ok(());
}
let previous = core::mem::replace(&mut self.slots[position].state, SlotState::Installing);
if let SlotState::Assigned { job, .. } = previous {
responses.push(Delivery::new(
job.reply_to,
PoolResponse::Interrupted {
job: job.id,
payload: job.payload,
reason: PoolInterruption::WorkerStopped {
worker: stopped.proxy,
outcome: stopped.outcome,
},
},
));
}
Ok(())
}
fn fail_backlog_if_irrecoverable(&mut self, actions: &mut PoolActions<A, J, R, C>) {
if self
.slots
.iter()
.any(|slot| !matches!(slot.state, SlotState::Retired { .. }))
{
return;
}
for queued in self.backlog.drain(..) {
let job = queued.accepted;
actions.sends.behavior.send::<_, Own>(Delivery::new(
job.reply_to,
PoolResponse::Interrupted {
job: job.id,
payload: job.payload,
reason: job
.interruption
.unwrap_or(PoolInterruption::NoRecoverableWorkers),
},
));
}
}
fn fail_jobs_for_retired_slot(
&mut self,
worker: A::Nonce,
reason: WorkerRetirement,
actions: &mut PoolActions<A, J, R, C>,
) {
let mut retained = VecDeque::with_capacity(self.backlog.len());
while let Some(queued) = self.backlog.pop_front() {
if queued.accepted.target == Some(worker) {
let job = queued.accepted;
actions.sends.behavior.send::<_, Own>(Delivery::new(
job.reply_to,
PoolResponse::Interrupted {
job: job.id,
payload: job.payload,
reason: job
.interruption
.unwrap_or(PoolInterruption::AffinityRetired { worker, reason }),
},
));
} else {
retained.push_back(queued);
}
}
self.backlog = retained;
}
fn creation_resolved(
&mut self,
resolved: &WorkerCreationResolved<A::Nonce>,
) -> Result<(), PoolError<A::Nonce>> {
let position = self.slot_position(resolved.proxy)?;
let phase = self
.worker_phase(resolved.proxy)
.expect("position proves the slot exists");
if !matches!(phase, WorkerPhase::Installing) {
return Err(PoolError::CreationResolvedWhileUnavailable {
worker: resolved.proxy,
phase,
});
}
self.slots[position].state = match resolved.result {
Ok(()) => SlotState::Idle,
Err(rejection) => SlotState::Retired {
reason: WorkerRetirement::CreationRejected(rejection),
},
};
Ok(())
}
fn dispatch(&mut self, actions: &mut PoolActions<A, J, R, C>) {
let mut selected_jobs = Vec::new();
let mut plan = Vec::new();
for (slot_position, slot) in self.slots.iter().enumerate() {
if !matches!(slot.state, SlotState::Idle) {
continue;
}
let Some(job_position) = self.backlog.iter().enumerate().find_map(|(position, job)| {
(!selected_jobs.contains(&position)
&& job
.accepted
.target
.is_none_or(|target| target == slot.nonce))
.then_some(position)
}) else {
continue;
};
selected_jobs.push(job_position);
plan.push(PlannedDispatch {
slot_position,
job_position,
});
}
let count =
u64::try_from(plan.len()).expect("a pool cannot contain more than u64::MAX slots");
let next_assignment = self
.next_assignment
.checked_add(count)
.expect("pool assignment identifiers exhausted");
let mut selected_by_position = BTreeMap::new();
for planned in plan {
selected_by_position.insert(planned.job_position, planned.slot_position);
}
let mut selected_by_slot: Vec<Option<QueuedJob<A, J, R>>> = std::iter::repeat_with(|| None)
.take(self.slots.len())
.collect();
let mut remaining = VecDeque::new();
for (position, queued) in self.backlog.drain(..).enumerate() {
if let Some(slot_position) = selected_by_position.remove(&position) {
selected_by_slot[slot_position] = Some(queued);
} else {
remaining.push_back(queued);
}
}
self.backlog = remaining;
for (offset, (slot_position, queued)) in selected_by_slot
.into_iter()
.enumerate()
.filter_map(|(slot_position, queued)| queued.map(|queued| (slot_position, queued)))
.enumerate()
{
let payload = queued.dispatch_payload;
let job = queued.accepted;
let assignment = AssignmentId(
self.next_assignment
+ u64::try_from(offset).expect("offset is bounded by the checked plan length"),
);
let nonce = self.slots[slot_position].nonce;
let job_id = job.id;
self.slots[slot_position].state = SlotState::Assigned { assignment, job };
actions.sends.behavior.send::<_, Own>(Delivery::new(
Recipient::child(nonce),
ProxyCommand::Forward(PoolAssignment {
assignment,
job: job_id,
payload,
}),
));
}
self.next_assignment = next_assignment;
}
}
impl<A, J, R, C> Behavior for WorkerPool<A, J, R, C>
where
A: Address,
A::Nonce: From<u64>,
J: Clone,
C: Behavior<Addr = A, Msg = PoolAssignment<J>, Ph = Never>,
{
type Addr = A;
type Msg = PoolMessage<A, J, R>;
type Event = PoolEvent<A, J, R>;
type Sends = PoolSends<A, J, R, C>;
type Ph = Never;
type Error = PoolError<A::Nonce>;
type Birth = Births<Proxy<C>>;
fn init(&mut self) -> crate::BehaviorActed<Self> {
match self.supervisor.init() {
Ok(actions) => Ok(actions),
Err(never) => match never {},
}
}
fn transition(&mut self, event: Self::Event) -> crate::BehaviorActed<Self> {
match event {
SupervisionEvent::Inner(User {
message:
PoolMessage::Submit {
job,
payload,
reply_to,
},
..
}) => {
let mut actions = Actions::cont();
self.submit(job, payload, reply_to, &mut actions);
self.dispatch(&mut actions);
Ok(actions)
}
SupervisionEvent::Inner(User {
message:
PoolMessage::Completed {
worker,
assignment,
result,
},
..
}) => {
let mut actions = Actions::cont();
self.complete(worker, assignment, result, &mut actions)?;
self.dispatch(&mut actions);
Ok(actions)
}
SupervisionEvent::WorkerStopped(stopped) => {
let proxy = stopped.proxy;
let mut responses = Vec::new();
self.worker_stopped(&stopped, &mut responses)?;
let mut actions =
self.supervisor_transition(SupervisionEvent::WorkerStopped(stopped));
actions.sends.behavior.responses.extend(responses);
let replacement_requested = actions
.sends
.replacement_commands
.iter()
.any(|delivery| delivery.to.route() == crate::Route::Child(proxy));
if !replacement_requested {
let position = self.slot_position(proxy)?;
let reason = WorkerRetirement::ReplacementUnavailable;
self.slots[position].state = SlotState::Retired { reason };
self.fail_jobs_for_retired_slot(proxy, reason, &mut actions);
}
self.dispatch(&mut actions);
self.fail_backlog_if_irrecoverable(&mut actions);
Ok(actions)
}
SupervisionEvent::WorkerCreationResolved(resolved) => {
let proxy = resolved.proxy;
self.creation_resolved(&resolved)?;
let mut actions =
self.supervisor_transition(SupervisionEvent::WorkerCreationResolved(resolved));
if let Some(WorkerPhase::Retired { reason }) = self.worker_phase(proxy) {
self.fail_jobs_for_retired_slot(proxy, reason, &mut actions);
}
self.dispatch(&mut actions);
self.fail_backlog_if_irrecoverable(&mut actions);
Ok(actions)
}
SupervisionEvent::ChildStopped(stopped) => {
Ok(self.supervisor_transition(SupervisionEvent::ChildStopped(stopped)))
}
SupervisionEvent::CreationResolved(resolved) => {
Ok(self.supervisor_transition(SupervisionEvent::CreationResolved(resolved)))
}
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::{MailAddr, NoBirths, Route};
#[derive(Clone, Copy)]
struct TestWorker;
impl Behavior for TestWorker {
type Addr = MailAddr;
type Msg = PoolAssignment<u8>;
type Event = User<MailAddr, PoolAssignment<u8>>;
type Sends = Vec<Delivery<MailAddr, Never>>;
type Ph = Never;
type Error = Never;
type Birth = NoBirths;
fn init(&mut self) -> crate::BehaviorActed<Self> {
Ok(Actions::cont())
}
fn transition(&mut self, _: Self::Event) -> crate::BehaviorActed<Self> {
Ok(Actions::cont())
}
}
fn test_worker(_: usize) -> TestWorker {
TestWorker
}
#[test]
fn one_dispatch_batch_preserves_fifo_jobs_across_index_removal() {
let mut pool = WorkerPool::new(
|index| u64::try_from(index).unwrap(),
2,
test_worker,
3,
InterruptionPolicy::Fail,
RestartPolicy::Permanent,
1,
Duration::from_secs(1),
)
.unwrap();
pool.init().unwrap();
for job in 1..=3 {
pool.transition(SupervisionEvent::Inner(User::new(
MailAddr(90),
PoolMessage::Submit {
job: JobId(job),
payload: u8::try_from(job).unwrap(),
reply_to: Recipient::global(MailAddr(91)),
},
)))
.unwrap();
}
pool.slots[0].state = SlotState::Idle;
pool.slots[1].state = SlotState::Idle;
let mut actions: PoolActions<MailAddr, u8, (), TestWorker> = Actions::cont();
pool.dispatch(&mut actions);
let assignments = &actions.sends.behavior.assignments;
assert_eq!(assignments.len(), 2);
for (index, expected_job) in [JobId(1), JobId(2)].into_iter().enumerate() {
assert_eq!(
assignments[index].to.route(),
Route::Child(u64::try_from(index).unwrap())
);
let ProxyCommand::Forward(assignment) = &assignments[index].message else {
panic!("pool dispatches with Forward");
};
assert_eq!(
assignment.assignment,
AssignmentId(u64::try_from(index).unwrap())
);
assert_eq!(assignment.job, expected_job);
}
assert_eq!(pool.backlog.len(), 1);
assert_eq!(pool.backlog[0].accepted.id, JobId(3));
}
}
pub struct KeyedWorkerPool<A: Address, K, J, R, C, S>
where
K: Eq,
C: Behavior<Addr = A, Msg = PoolAssignment<J>, Ph = Never>,
S: AffinitySelector<K, A::Nonce>,
{
pool: WorkerPool<A, J, R, C>,
bindings: Vec<(K, A::Nonce)>,
selector: S,
}
impl<A, K, J, R, C, S> KeyedWorkerPool<A, K, J, R, C, S>
where
A: Address,
K: Eq,
C: Behavior<Addr = A, Msg = PoolAssignment<J>, Ph = Never>,
S: AffinitySelector<K, A::Nonce>,
{
#[allow(
clippy::too_many_arguments,
reason = "the arguments expose the complete pool and affinity policy"
)]
pub fn new(
nonces: fn(usize) -> A::Nonce,
count: usize,
build: fn(usize) -> C,
backlog_capacity: usize,
interruption: InterruptionPolicy,
restart_policy: RestartPolicy,
max_restarts: u32,
restart_window: Duration,
selector: S,
) -> Result<Self, PoolConfigError<A::Nonce>> {
Ok(Self {
pool: WorkerPool::new(
nonces,
count,
build,
backlog_capacity,
interruption,
restart_policy,
max_restarts,
restart_window,
)?,
bindings: Vec::new(),
selector,
})
}
#[must_use]
pub fn affinity(&self, key: &K) -> Option<A::Nonce> {
self.bindings
.iter()
.find_map(|(bound, worker)| (bound == key).then_some(*worker))
}
#[must_use]
pub fn backlog_len(&self) -> usize {
self.pool.backlog_len()
}
#[must_use]
pub fn worker_phase(&self, worker: A::Nonce) -> Option<WorkerPhase> {
self.pool.worker_phase(worker)
}
fn rebalance(&mut self, key: K, worker: A::Nonce) -> Result<(), PoolError<A::Nonce>> {
let position = self.pool.slot_position(worker)?;
if let SlotState::Retired { reason } = self.pool.slots[position].state {
return Err(PoolError::RebalanceToRetiredWorker { worker, reason });
}
if let Some((_, bound)) = self.bindings.iter_mut().find(|(bound, _)| *bound == key) {
*bound = worker;
} else {
self.bindings.push((key, worker));
}
Ok(())
}
}
impl<A, K, J, R, C, S> Behavior for KeyedWorkerPool<A, K, J, R, C, S>
where
A: Address,
A::Nonce: From<u64>,
K: Eq,
J: Clone,
C: Behavior<Addr = A, Msg = PoolAssignment<J>, Ph = Never>,
S: AffinitySelector<K, A::Nonce>,
{
type Addr = A;
type Msg = KeyedPoolMessage<A, K, J, R>;
type Event = KeyedPoolEvent<A, K, J, R>;
type Sends = PoolSends<A, J, R, C>;
type Ph = Never;
type Error = PoolError<A::Nonce>;
type Birth = Births<Proxy<C>>;
fn init(&mut self) -> crate::BehaviorActed<Self> {
self.pool.init()
}
fn transition(&mut self, event: Self::Event) -> crate::BehaviorActed<Self> {
match event {
SupervisionEvent::Inner(User {
message:
KeyedPoolMessage::Submit {
key,
job,
payload,
reply_to,
},
..
}) => {
let existing = self.affinity(&key);
let target = existing.unwrap_or_else(|| self.selector.select(&key));
let mut actions = Actions::cont();
let admission = self
.pool
.submit_to(target, job, payload, reply_to, &mut actions);
match (admission, existing) {
(Admission::Accepted, None) => self.bindings.push((key, target)),
(Admission::Accepted | Admission::Rejected, Some(_))
| (Admission::Rejected, None) => {}
}
self.pool.dispatch(&mut actions);
Ok(actions)
}
SupervisionEvent::Inner(User {
message:
KeyedPoolMessage::Completed {
worker,
assignment,
result,
},
..
}) => {
let mut actions = Actions::cont();
self.pool
.complete(worker, assignment, result, &mut actions)?;
self.pool.dispatch(&mut actions);
Ok(actions)
}
SupervisionEvent::Inner(User {
message: KeyedPoolMessage::Rebalance { key, worker },
..
}) => {
self.rebalance(key, worker)?;
Ok(Actions::cont())
}
SupervisionEvent::WorkerStopped(stopped) => self
.pool
.transition(SupervisionEvent::WorkerStopped(stopped)),
SupervisionEvent::WorkerCreationResolved(resolved) => self
.pool
.transition(SupervisionEvent::WorkerCreationResolved(resolved)),
SupervisionEvent::ChildStopped(stopped) => self
.pool
.transition(SupervisionEvent::ChildStopped(stopped)),
SupervisionEvent::CreationResolved(resolved) => self
.pool
.transition(SupervisionEvent::CreationResolved(resolved)),
}
}
}