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 was cancelled
132    ///
133    /// # Guarantees
134    /// - Task did not complete normally
135    /// - Cancellation was requested and honored
136    /// - Resources have been cleaned up
137    Cancelled,
138
139    /// Task failed with an error
140    ///
141    /// # Guarantees
142    /// - Task encountered an unrecoverable error
143    /// - Error information is available via task handle
144    /// - Resources have been cleaned up
145    Failed,
146}
147
148impl core::fmt::Display for TaskStatus {
149    fn fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
150        match self {
151            Self::Queued => write!(f, "Queued"),
152            Self::Running => write!(f, "Running"),
153            Self::Completed => write!(f, "Completed"),
154            Self::Cancelled => write!(f, "Cancelled"),
155            Self::Failed => write!(f, "Failed"),
156        }
157    }
158}
159
160/// Detailed statistics about a specific task.
161///
162/// Task statistics report the lifecycle timing the executor actually
163/// tracks: spawn, start, and completion timestamps from a monotonic
164/// high-resolution clock, plus the derived execution duration.
165#[derive(Debug, Clone)]
166pub struct TaskStats {
167    /// Task identifier
168    pub id: TaskId,
169    /// Current status
170    pub status: TaskStatus,
171    /// Priority level
172    pub priority: Priority,
173    /// When the task was spawned
174    pub spawn_time: Instant,
175    /// When the task started executing (if started)
176    pub start_time: Option<Instant>,
177    /// When the task completed (if completed)
178    pub completion_time: Option<Instant>,
179    /// Total CPU time used (nanoseconds), derived from the lifecycle
180    /// start/completion timestamps.
181    pub cpu_time_ns: u64,
182}
183
184impl TaskStats {
185    /// Returns the total execution time of the task, if available.
186    ///
187    /// # Returns
188    /// `Some(duration)` if the task has completed execution, `None` if still running or queued.
189    #[must_use]
190    pub fn execution_time(&self) -> Option<core::time::Duration> {
191        match (&self.start_time, &self.completion_time) {
192            (Some(start), Some(end)) => Some(end.duration_since(*start)),
193            _ => None,
194        }
195    }
196
197    /// Returns the time the task spent in the queue before execution.
198    ///
199    /// # Returns
200    /// - `Some(duration_since_spawn)` if task is still queued
201    /// - `Some(queue_duration)` if task has started execution
202    /// - `None` if timing information is unavailable
203    #[must_use]
204    pub fn queue_time(&self) -> Option<core::time::Duration> {
205        match &self.start_time {
206            Some(start) => Some(start.duration_since(self.spawn_time)),
207            None => Some(Instant::now().duration_since(self.spawn_time)),
208        }
209    }
210
211    /// Returns whether the task is currently active (queued or running).
212    #[must_use]
213    pub fn is_active(&self) -> bool {
214        matches!(self.status, TaskStatus::Queued | TaskStatus::Running)
215    }
216
217    /// Returns whether the task has reached a terminal state.
218    #[must_use]
219    pub fn is_finished(&self) -> bool {
220        matches!(
221            self.status,
222            TaskStatus::Completed | TaskStatus::Cancelled | TaskStatus::Failed
223        )
224    }
225}