use std::num::NonZeroUsize;
use std::sync::atomic::AtomicBool;
use std::sync::Arc;
use thiserror::Error;
use crate::{
CpuDomainExecutor, CpuDomainExecutorCapabilities, CpuDomainId, CpuPlacementGuarantee, CpuSet,
ResolvedCpuPlacement,
};
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub enum CpuDomainOwnership {
Managed,
ExternalManaged,
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub enum CpuAdmissionMode {
CooperativeCpuSet,
CallerManaged,
}
#[derive(Debug)]
enum CpuDomainAdmission {
CooperativeCpuSet {
placement: ResolvedCpuPlacement,
guarantee: CpuPlacementGuarantee,
},
CallerManaged {
active: Arc<AtomicBool>,
},
}
#[derive(Clone, Copy, Debug, Eq, Error, PartialEq)]
pub enum ExternalCpuDomainError {
#[error("external CPU domain placement must contain at least one CPU")]
EmptyPlacementCpuSet,
#[error("external CPU domain executor must report at least one worker")]
ZeroExecutorWorkers,
#[error(
"external CPU domain thread budget {thread_budget} exceeds executor worker count {worker_count}"
)]
ThreadBudgetExceedsWorkerCount {
thread_budget: usize,
worker_count: usize,
},
}
#[derive(Debug)]
pub(crate) struct CpuResourceDomain {
id: CpuDomainId,
admission: CpuDomainAdmission,
executor: Arc<dyn CpuDomainExecutor>,
thread_budget: NonZeroUsize,
ownership: CpuDomainOwnership,
}
impl CpuResourceDomain {
pub(crate) fn new(
id: CpuDomainId,
placement: ResolvedCpuPlacement,
executor: Arc<dyn CpuDomainExecutor>,
thread_budget: NonZeroUsize,
placement_guarantee: CpuPlacementGuarantee,
ownership: CpuDomainOwnership,
) -> Self {
Self {
id,
admission: CpuDomainAdmission::CooperativeCpuSet {
placement,
guarantee: placement_guarantee,
},
executor,
thread_budget,
ownership,
}
}
fn new_caller_managed(
id: CpuDomainId,
executor: Arc<dyn CpuDomainExecutor>,
thread_budget: NonZeroUsize,
) -> Self {
Self {
id,
admission: CpuDomainAdmission::CallerManaged {
active: Arc::new(AtomicBool::new(false)),
},
executor,
thread_budget,
ownership: CpuDomainOwnership::ExternalManaged,
}
}
pub(crate) fn id(&self) -> CpuDomainId {
self.id
}
pub(crate) fn admission_mode(&self) -> CpuAdmissionMode {
match self.admission {
CpuDomainAdmission::CooperativeCpuSet { .. } => CpuAdmissionMode::CooperativeCpuSet,
CpuDomainAdmission::CallerManaged { .. } => CpuAdmissionMode::CallerManaged,
}
}
pub(crate) fn placement(&self) -> Option<&ResolvedCpuPlacement> {
match &self.admission {
CpuDomainAdmission::CooperativeCpuSet { placement, .. } => Some(placement),
CpuDomainAdmission::CallerManaged { .. } => None,
}
}
pub(crate) fn cpus(&self) -> Option<&CpuSet> {
self.placement().map(ResolvedCpuPlacement::cpus)
}
pub(crate) fn caller_managed_active(&self) -> Option<Arc<AtomicBool>> {
match &self.admission {
CpuDomainAdmission::CallerManaged { active } => Some(Arc::clone(active)),
CpuDomainAdmission::CooperativeCpuSet { .. } => None,
}
}
pub(crate) fn executor(&self) -> &Arc<dyn CpuDomainExecutor> {
&self.executor
}
pub(crate) fn thread_budget(&self) -> NonZeroUsize {
self.thread_budget
}
pub(crate) fn placement_guarantee(&self) -> Option<CpuPlacementGuarantee> {
match self.admission {
CpuDomainAdmission::CooperativeCpuSet { guarantee, .. } => Some(guarantee),
CpuDomainAdmission::CallerManaged { .. } => None,
}
}
pub(crate) fn ownership(&self) -> CpuDomainOwnership {
self.ownership
}
pub(crate) fn executor_capabilities(&self) -> CpuDomainExecutorCapabilities {
self.executor().capabilities()
}
}
#[derive(Debug)]
pub struct ExternalCpuDomain {
domain: CpuResourceDomain,
}
impl ExternalCpuDomain {
pub fn new(
id: CpuDomainId,
placement: ResolvedCpuPlacement,
executor: Arc<dyn CpuDomainExecutor>,
thread_budget: NonZeroUsize,
placement_guarantee: CpuPlacementGuarantee,
) -> Result<Self, ExternalCpuDomainError> {
let worker_count = executor.capabilities().worker_count.get();
validate_external_domain_config(Some(placement.cpus().len()), worker_count, thread_budget)?;
Ok(Self {
domain: CpuResourceDomain::new(
id,
placement,
executor,
thread_budget,
placement_guarantee,
CpuDomainOwnership::ExternalManaged,
),
})
}
pub fn new_caller_managed(
id: CpuDomainId,
executor: Arc<dyn CpuDomainExecutor>,
thread_budget: NonZeroUsize,
) -> Result<Self, ExternalCpuDomainError> {
let worker_count = executor.capabilities().worker_count.get();
validate_external_domain_config(None, worker_count, thread_budget)?;
Ok(Self {
domain: CpuResourceDomain::new_caller_managed(id, executor, thread_budget),
})
}
pub fn id(&self) -> CpuDomainId {
self.domain.id()
}
pub fn placement(&self) -> Option<&ResolvedCpuPlacement> {
self.domain.placement()
}
pub fn cpus(&self) -> Option<&CpuSet> {
self.domain.cpus()
}
pub fn admission_mode(&self) -> CpuAdmissionMode {
self.domain.admission_mode()
}
pub fn thread_budget(&self) -> NonZeroUsize {
self.domain.thread_budget()
}
pub fn placement_guarantee(&self) -> Option<CpuPlacementGuarantee> {
self.domain.placement_guarantee()
}
pub fn ownership(&self) -> CpuDomainOwnership {
self.domain.ownership()
}
pub fn executor_capabilities(&self) -> CpuDomainExecutorCapabilities {
self.domain.executor_capabilities()
}
}
impl From<ExternalCpuDomain> for CpuResourceDomain {
fn from(domain: ExternalCpuDomain) -> Self {
domain.domain
}
}
fn validate_external_domain_config(
cpu_count: Option<usize>,
worker_count: usize,
thread_budget: NonZeroUsize,
) -> Result<(), ExternalCpuDomainError> {
if cpu_count == Some(0) {
return Err(ExternalCpuDomainError::EmptyPlacementCpuSet);
}
if worker_count == 0 {
return Err(ExternalCpuDomainError::ZeroExecutorWorkers);
}
if thread_budget.get() > worker_count {
return Err(ExternalCpuDomainError::ThreadBudgetExceedsWorkerCount {
thread_budget: thread_budget.get(),
worker_count,
});
}
Ok(())
}
#[cfg(test)]
mod tests;