Skip to main content

moirai_core/executor/
config.rs

1//! Configuration settings for executor behavior.
2
3use crate::platform::String;
4
5use super::placement::WorkerPlacement;
6
7// Memory pool size constants
8const KILOBYTE: usize = 1024;
9const MEGABYTE: usize = 1024 * KILOBYTE;
10/// Default capacity for the small object allocation pool.
11pub const SMALL_POOL_SIZE: usize = 64 * KILOBYTE;
12/// Default capacity for the medium object allocation pool.
13pub const MEDIUM_POOL_SIZE: usize = MEGABYTE;
14/// Default capacity for the large object allocation pool.
15pub const LARGE_POOL_SIZE: usize = 16 * MEGABYTE;
16
17/// Default aggregate bound for worker admission queues (tasks, not bytes).
18/// Sized for burst absorption across all workers before producers observe
19/// backpressure.
20pub const DEFAULT_GLOBAL_QUEUE_CAPACITY: usize = 8192;
21/// Default initial slot count for each worker's resizable local priority queue.
22///
23/// This is a retained-storage policy, not an admission bound. The Chase-Lev
24/// queues grow when full and normalize this value to a supported power of two.
25pub const DEFAULT_LOCAL_QUEUE_INITIAL_CAPACITY: usize = 128;
26
27/// Configuration settings for executor behavior and performance characteristics.
28///
29/// This struct encapsulates all tunable parameters that affect executor operation,
30/// including thread pool sizes, queue capacities, and various performance optimizations.
31#[allow(clippy::module_name_repetitions)]
32pub struct ExecutorConfig {
33    /// Number of worker threads for parallel tasks
34    pub worker_threads: usize,
35    /// Number of threads dedicated to async tasks
36    pub async_threads: usize,
37    /// Maximum aggregate size of the workers' external admission queues.
38    ///
39    /// Executor construction partitions this bound across workers without
40    /// exceeding it. The value must supply at least two slots per worker.
41    pub max_global_queue_size: usize,
42    /// Initial slot count for each resizable per-worker local priority queue.
43    ///
44    /// Values below the deque minimum normalize upward. Local queues grow when
45    /// full; [`Self::max_global_queue_size`] is the external admission bound.
46    pub local_queue_initial_capacity: usize,
47    /// Thread name prefix for worker threads
48    pub thread_name_prefix: String,
49    /// Whether each worker is confined to one logical processor.
50    ///
51    /// [`WorkerPlacement::Pinned`] can fail executor construction; see
52    /// [`ExecutorError::WorkerPlacementFailed`](crate::error::ExecutorError::WorkerPlacementFailed).
53    pub worker_placement: WorkerPlacement,
54    /// Whether to enable metrics collection
55    #[cfg(feature = "metrics")]
56    pub enable_metrics: bool,
57    /// Task preemption configuration
58    pub preemption: PreemptionConfig,
59    /// Memory management configuration
60    pub memory: MemoryConfig,
61    /// Task cleanup configuration
62    pub cleanup: CleanupConfig,
63}
64
65impl Default for ExecutorConfig {
66    fn default() -> Self {
67        Self {
68            worker_threads: super::logical_parallelism(),
69            async_threads: (super::logical_parallelism() / 4).max(1),
70            max_global_queue_size: DEFAULT_GLOBAL_QUEUE_CAPACITY,
71            local_queue_initial_capacity: DEFAULT_LOCAL_QUEUE_INITIAL_CAPACITY,
72            thread_name_prefix: "moirai-worker".into(),
73            worker_placement: WorkerPlacement::default(),
74            #[cfg(feature = "metrics")]
75            enable_metrics: true,
76            preemption: PreemptionConfig::default(),
77            memory: MemoryConfig::default(),
78            cleanup: CleanupConfig::default(),
79        }
80    }
81}
82
83/// Configuration for task preemption.
84#[derive(Debug, Clone)]
85pub struct PreemptionConfig {
86    /// Whether to enable cooperative preemption
87    pub enabled: bool,
88    /// Time slice for each task before preemption (microseconds)
89    pub time_slice_us: u64,
90    /// Whether to preempt based on priority
91    pub priority_based: bool,
92    /// Minimum execution time before preemption (microseconds)
93    pub min_execution_time_us: u64,
94}
95
96impl Default for PreemptionConfig {
97    fn default() -> Self {
98        Self {
99            enabled: true,
100            time_slice_us: 10_000, // 10ms
101            priority_based: true,
102            min_execution_time_us: 1_000, // 1ms
103        }
104    }
105}
106
107/// Configuration for memory management.
108#[derive(Debug, Clone)]
109pub struct MemoryConfig {
110    /// Whether to use memory pools
111    pub use_memory_pools: bool,
112    /// Size of small object pool (bytes)
113    pub small_pool_size: usize,
114    /// Size of medium object pool (bytes)
115    pub medium_pool_size: usize,
116    /// Size of large object pool (bytes)
117    pub large_pool_size: usize,
118    /// Whether to track memory usage per task
119    pub track_per_task_memory: bool,
120}
121
122impl Default for MemoryConfig {
123    fn default() -> Self {
124        Self {
125            use_memory_pools: true,
126            small_pool_size: SMALL_POOL_SIZE,
127            medium_pool_size: MEDIUM_POOL_SIZE,
128            large_pool_size: LARGE_POOL_SIZE,
129            track_per_task_memory: cfg!(feature = "metrics"),
130        }
131    }
132}
133
134/// Configuration for completed-task retention.
135///
136/// The executor releases the state of finished tasks in whole blocks of 1,024
137/// tasks as it registers new ones, so retained memory follows the spawn rate
138/// without a background thread. A block is released once every task in it has
139/// finished and either its newest completion is older than
140/// `task_retention_duration` or more than `max_retained_tasks` finished tasks
141/// are retained. One long-running task therefore keeps its own block resident.
142///
143/// A released task still reports as completed to `wait_for_task` and accepts
144/// `cancel_task` as a no-op; only its status and statistics are gone, so
145/// `task_status` and `task_stats` return `None`.
146#[derive(Debug, Clone)]
147pub struct CleanupConfig {
148    /// How long finished-task metadata stays observable.
149    ///
150    /// # Default: 5 minutes
151    pub task_retention_duration: core::time::Duration,
152
153    /// Whether the executor releases finished-task metadata at all.
154    ///
155    /// When disabled every task's metadata is retained for the executor's
156    /// lifetime and memory grows with the number of tasks spawned.
157    /// # Default: true
158    pub enable_automatic_cleanup: bool,
159
160    /// Most finished tasks retained regardless of age, rounded up to whole
161    /// blocks of 1,024 tasks.
162    ///
163    /// This is the hard bound on retained memory for a workload that finishes
164    /// tasks faster than `task_retention_duration` elapses.
165    /// # Default: 10,000 tasks
166    pub max_retained_tasks: usize,
167}
168
169impl Default for CleanupConfig {
170    fn default() -> Self {
171        Self {
172            task_retention_duration: core::time::Duration::from_mins(5),
173            enable_automatic_cleanup: true,
174            max_retained_tasks: 10_000,
175        }
176    }
177}