Skip to main content

moirai_core/executor/
manager.rs

1//! Task management and status tracking.
2
3use crate::error::ExecutorResult;
4use crate::platform::Instant;
5use crate::{Priority, TaskId};
6
7/// Task management and monitoring capabilities.
8///
9/// This trait provides operations for managing and monitoring running tasks.
10/// It follows the Interface Segregation Principle by separating management
11/// concerns from spawning concerns.
12///
13/// # Behavior Guarantees
14/// - All operations are thread-safe and non-blocking where possible
15/// - Task state is eventually consistent across all observers
16/// - Cancellation is cooperative and may not be immediate
17/// - Statistics are updated atomically and consistently
18///
19/// # Performance Characteristics
20/// - Status queries: O(1) lookup time, < 50ns typical latency
21/// - Cancellation: O(1) operation, cooperative completion
22/// - Statistics: Atomic operations, minimal overhead
23pub trait TaskManager: Send + Sync + 'static {
24    /// Cancels a running task by its ID.
25    ///
26    /// # Arguments
27    /// * `id` - The unique identifier of the task to cancel
28    ///
29    /// # Returns
30    /// `Ok(())` if the task was successfully cancelled or was already completed.
31    ///
32    /// # Errors
33    /// Returns `TaskError` in the following cases:
34    /// - `NotFound` if no task with the given ID exists
35    /// - `InvalidState` if the task cannot be cancelled (e.g., already completed)
36    /// - `SystemError` if the cancellation operation fails due to internal errors
37    fn cancel_task(&self, id: TaskId) -> ExecutorResult<()>;
38
39    /// Get the current status of a task.
40    ///
41    /// # Behavior Guarantees
42    /// - Returns None if task ID is not found
43    /// - Status is eventually consistent across threads
44    /// - Completed tasks may be garbage collected after timeout
45    /// - Status transitions are monotonic (no backwards moves)
46    ///
47    /// # Performance Characteristics
48    /// - O(1) lookup time using hash table
49    /// - Latency: < 50ns for status query
50    /// - Memory: Minimal overhead for status tracking
51    /// - Non-blocking: Never blocks calling thread
52    fn task_status(&self, id: TaskId) -> Option<TaskStatus>;
53
54    /// Wait for a task to complete.
55    ///
56    /// Returns a future that resolves when the task completes or the timeout expires.
57    /// This enables async/await patterns for task coordination.
58    ///
59    /// # Arguments
60    /// - `id`: The task ID to wait for
61    /// - `timeout`: Optional timeout duration
62    ///
63    /// # Returns
64    /// A future that resolves to:
65    /// - `Ok(())` when the task completes successfully
66    /// - `Err(TaskError::Timeout)` if the timeout expires
67    /// - `Err(TaskError::NotFound)` if the task doesn't exist
68    ///
69    /// # Performance
70    /// - Immediate return: < 10ns if already complete
71    /// - Waiting overhead: Event-driven, no busy polling
72    /// - Memory: Minimal waker chain overhead
73    fn wait_for_task(
74        &self,
75        id: TaskId,
76        timeout: Option<core::time::Duration>,
77    ) -> impl core::future::Future<Output = ExecutorResult<()>> + Send;
78
79    /// Get statistics about task execution.
80    ///
81    /// # Behavior Guarantees
82    /// - Returns None if task ID is not found or stats not enabled
83    /// - Statistics are eventually consistent
84    /// - Timing measurements use high-resolution monotonic clock
85    /// - Memory usage tracking depends on executor configuration
86    ///
87    /// # Performance Characteristics
88    /// - Lookup: O(1) hash table access
89    /// - Overhead: ~100 bytes per task when metrics enabled
90    /// - Collection cost: < 5% runtime overhead when enabled
91    fn task_stats(&self, id: TaskId) -> Option<TaskStats>;
92}
93
94/// Status of a task within the executor.
95///
96/// Task status transitions follow a strict state machine:
97/// Queued → Running → (Completed | Cancelled | Failed)
98///
99/// # State Transitions
100/// - Queued: Initial state when task is spawned
101/// - Running: Task is currently executing on a worker thread
102/// - Completed: Task finished successfully
103/// - Cancelled: Task was cancelled before or during execution
104/// - Failed: Task encountered an error or panic
105#[derive(Debug, Clone, Copy, PartialEq, Eq)]
106pub enum TaskStatus {
107    /// Task is queued but not yet started
108    ///
109    /// # Guarantees
110    /// - Task will eventually transition to Running
111    /// - Cancellation is possible in this state
112    /// - Memory has been allocated for task execution
113    Queued,
114
115    /// Task is currently running
116    ///
117    /// # Guarantees
118    /// - Task is actively executing on a worker thread
119    /// - Cancellation is cooperative in this state
120    /// - Progress is being made toward completion
121    Running,
122
123    /// Task completed successfully
124    ///
125    /// # Guarantees
126    /// - Task result is available via task handle
127    /// - No further state transitions possible
128    /// - Resources have been cleaned up
129    Completed,
130
131    /// Task ended without a result
132    ///
133    /// Reached when a cancel request is honored, or when the executor discards
134    /// the task before it finishes (for example at shutdown). The task's handle
135    /// resolves to `TaskError::Cancelled` in both cases.
136    ///
137    /// # Guarantees
138    /// - Task did not complete normally
139    /// - Cancellation was requested and honored, or the task was discarded
140    /// - Resources have been cleaned up
141    Cancelled,
142
143    /// Task failed with an error
144    ///
145    /// # Guarantees
146    /// - Task encountered an unrecoverable error
147    /// - Error information is available via task handle
148    /// - Resources have been cleaned up
149    Failed,
150}
151
152impl core::fmt::Display for TaskStatus {
153    fn fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
154        match self {
155            Self::Queued => write!(f, "Queued"),
156            Self::Running => write!(f, "Running"),
157            Self::Completed => write!(f, "Completed"),
158            Self::Cancelled => write!(f, "Cancelled"),
159            Self::Failed => write!(f, "Failed"),
160        }
161    }
162}
163
164/// Detailed statistics about a specific task.
165///
166/// Task statistics report the lifecycle timing the executor actually
167/// tracks: spawn, start, and completion timestamps from a monotonic
168/// high-resolution clock, plus the derived execution duration.
169#[derive(Debug, Clone)]
170pub struct TaskStats {
171    /// Task identifier
172    pub id: TaskId,
173    /// Current status
174    pub status: TaskStatus,
175    /// Priority level
176    pub priority: Priority,
177    /// When the task was spawned
178    pub spawn_time: Instant,
179    /// When the task started executing (if started)
180    pub start_time: Option<Instant>,
181    /// When the task completed (if completed)
182    pub completion_time: Option<Instant>,
183    /// Total CPU time used (nanoseconds), derived from the lifecycle
184    /// start/completion timestamps.
185    pub cpu_time_ns: u64,
186}
187
188impl TaskStats {
189    /// Returns the total execution time of the task, if available.
190    ///
191    /// # Returns
192    /// `Some(duration)` if the task has completed execution, `None` if still running or queued.
193    #[must_use]
194    pub fn execution_time(&self) -> Option<core::time::Duration> {
195        match (&self.start_time, &self.completion_time) {
196            (Some(start), Some(end)) => Some(end.duration_since(*start)),
197            _ => None,
198        }
199    }
200
201    /// Returns the time the task spent in the queue before execution.
202    ///
203    /// # Returns
204    /// - `Some(duration_since_spawn)` if task is still queued
205    /// - `Some(queue_duration)` if task has started execution
206    /// - `None` if timing information is unavailable
207    #[must_use]
208    pub fn queue_time(&self) -> Option<core::time::Duration> {
209        match &self.start_time {
210            Some(start) => Some(start.duration_since(self.spawn_time)),
211            None => Some(Instant::now().duration_since(self.spawn_time)),
212        }
213    }
214
215    /// Returns whether the task is currently active (queued or running).
216    #[must_use]
217    pub fn is_active(&self) -> bool {
218        matches!(self.status, TaskStatus::Queued | TaskStatus::Running)
219    }
220
221    /// Returns whether the task has reached a terminal state.
222    #[must_use]
223    pub fn is_finished(&self) -> bool {
224        matches!(
225            self.status,
226            TaskStatus::Completed | TaskStatus::Cancelled | TaskStatus::Failed
227        )
228    }
229}