Skip to main content

moirai_core/task/
mod.rs

1//! # Task Abstraction Layer
2//!
3//! This module provides the core task abstractions for the Moirai concurrency library.
4//! All task types are designed to be zero-cost abstractions that compile away to optimal code.
5//!
6//! ## Safety Guarantees
7//!
8//! - **Memory Safety**: All task operations are memory-safe by construction
9//! - **Data Race Freedom**: Rust's ownership system prevents data races
10//! - **Resource Cleanup**: Automatic resource cleanup on task completion or panic
11//! - **Type Safety**: Generic type system ensures compile-time correctness
12//!
13//! ## Performance Characteristics
14//!
15//! - **Task Creation**: O(1) constant time with zero allocations for simple closures
16//! - **Task Execution**: Zero-cost abstractions compile to direct function calls
17//! - **Memory Overhead**: < 64 bytes per task for metadata and context
18//! - **Cache Efficiency**: Task data structures are cache-line aligned
19//!
20//! ## Examples
21//!
22//! ### Basic Task Creation
23//!
24//! ```rust
25//! use moirai_core::{Task, TaskBuilder, Priority};
26//!
27//! // Simple closure task
28//! let task = TaskBuilder::new()
29//!     .priority(Priority::Normal)
30//!     .name("computation")
31//!     .build(|| {
32//!         (1..=100).sum::<i32>()
33//!     });
34//!
35//! assert_eq!(task.execute(), 5050);
36//! ```
37//!
38//! ### Task Chaining and Composition
39//!
40//! ```rust,ignore
41//! use moirai_core::{TaskBuilder, TaskExt, Task};
42//!
43//! let base_task = TaskBuilder::new().build(|| 21);
44//!
45//! // Chain operations
46//! let doubled = base_task.then(|x| x * 2);
47//! let result = doubled.execute();
48//! assert_eq!(result, 42);
49//!
50//! // Map transformations
51//! let mapped = TaskBuilder::new().build(|| "hello")
52//!     .map(|s| s.to_uppercase());
53//! assert_eq!(mapped.execute(), "HELLO");
54//! ```
55//!
56//! ### Error Handling
57//!
58//! `Task::execute` has no error channel: fallible work returns a `Result` as
59//! its `Output` and callers branch on the value. Panic recovery is the
60//! executor's responsibility (the hybrid executor catches unwinds at the job
61//! boundary), not a task-combinator concern.
62
63#![cfg_attr(test, allow(clippy::unwrap_used, reason = "test scope"))]
64
65/// Maximum generic spin attempts before falling back to blocking
66pub(crate) const MAX_SPIN_ATTEMPTS: usize = 64;
67
68/// Task builder types: `TaskBuilder`, `BaseTask`, `Closure`, `Chained`, `Mapped`, `Parameterized`, `Group`, `Spawner`.
69pub(crate) mod builder;
70/// Task extension trait and combinator types: `TaskExt`, `ContextualTask`.
71pub(crate) mod ext;
72/// `TaskFuture<T>`: fused one-shot future that runs a task synchronously on first poll.
73pub(crate) mod future;
74/// Task handle and result-slot types: `TaskHandle`, `TaskResultSender`, `BlockingResultWait`, `ResultWaitPolicy`.
75pub(crate) mod handle;
76/// Core identity and context types: `TaskId`, `Priority`, `TaskContext`.
77pub(crate) mod id_and_context;
78/// Core `Task` trait definition and `Box<T: Task>` delegation impl.
79pub(crate) mod traits;
80
81// ── Public API re-exports ─────────────────────────────────────────────────────
82
83pub use builder::{BaseTask, Chained, Closure, Group, Mapped, Parameterized, Spawner, TaskBuilder};
84pub use ext::{ContextualTask, TaskExt};
85pub use future::TaskFuture;
86pub use handle::TaskHandle;
87pub use id_and_context::{Priority, TaskContext, TaskId};
88pub use traits::Task;
89
90#[cfg(feature = "std")]
91pub use handle::{BlockingResultWait, ResultWaitPolicy, TaskResultSender};
92
93#[cfg(all(feature = "std", feature = "result-diagnostics"))]
94pub use handle::{
95    diagnostic_result_slot_complete_waiting, diagnostic_result_slot_ready_take,
96    diagnostic_result_slot_register_waiter, diagnostic_result_slot_spin_miss,
97};
98
99// ── Tests ─────────────────────────────────────────────────────────────────────
100
101#[cfg(test)]
102mod tests {
103    use super::*;
104    use crate::TaskBuilder;
105
106    #[test]
107    fn test_task_future() {
108        let id = TaskId::new(1);
109        let task = TaskBuilder::new().with_id(id).build(|| 42);
110        let future = TaskFuture::new(task, TaskContext::new(id));
111
112        assert_eq!(future.context().id, id);
113    }
114
115    fn poll_once<F: core::future::Future + Unpin>(future: &mut F) -> core::task::Poll<F::Output> {
116        // No-op waker: TaskFuture never returns Pending, so the waker is unused.
117        let waker = core::task::Waker::noop();
118        let mut cx = core::task::Context::from_waker(waker);
119        core::pin::Pin::new(future).poll(&mut cx)
120    }
121
122    #[test]
123    fn task_future_resolves_task_output_on_first_poll() {
124        let id = TaskId::new(2);
125        let task = TaskBuilder::new().with_id(id).build(|| 21 * 2);
126        let mut future = TaskFuture::new(task, TaskContext::new(id));
127
128        assert_eq!(poll_once(&mut future), core::task::Poll::Ready(42));
129    }
130
131    #[test]
132    #[should_panic(expected = "TaskFuture polled after completion")]
133    fn task_future_panics_when_polled_after_completion() {
134        let id = TaskId::new(3);
135        let task = TaskBuilder::new().with_id(id).build(|| 1);
136        let mut future = TaskFuture::new(task, TaskContext::new(id));
137
138        assert_eq!(poll_once(&mut future), core::task::Poll::Ready(1));
139        let _ = poll_once(&mut future); // fused: must panic, not hang as Pending
140    }
141
142    #[test]
143    fn test_task_composition() {
144        let id = TaskId::new(1);
145        let task = TaskBuilder::new().with_id(id).build(|| 10);
146
147        // Test map combinator
148        let mapped = task.map(|x| x * 2);
149        assert_eq!(mapped.execute(), 20);
150    }
151
152    #[test]
153    fn test_task_group() {
154        let mut group = Group::new(TaskId::new(1));
155
156        let task1 = TaskBuilder::new().with_id(TaskId::new(2)).build(|| 42);
157
158        let task2 = TaskBuilder::new().with_id(TaskId::new(3)).build(|| 24);
159
160        // Wrap tasks in closures for the group
161        group.add_task(|| {
162            let _ = task1.execute();
163        });
164        group.add_task(|| {
165            let _ = task2.execute();
166        });
167
168        assert_eq!(group.len(), 2);
169        assert!(!group.is_empty());
170
171        // Execute the group
172        group.execute();
173    }
174
175    #[test]
176    fn test_parameterized_task() {
177        let id = TaskId::new(1);
178        let task = Parameterized::new(|x: i32| x * 3, 7, TaskContext::new(id));
179
180        assert_eq!(task.execute(), 21);
181    }
182
183    #[test]
184    fn test_spawner_task() {
185        let id = TaskId::new(1);
186        let spawner = Spawner::new(
187            || {
188                // This would spawn other tasks in a real implementation
189            },
190            TaskContext::new(id),
191        );
192
193        spawner.execute(); // Should not panic
194    }
195
196    #[test]
197    fn task_handle_returns_sent_result() {
198        let (handle, sender) = TaskHandle::new_pending(TaskId::new(10));
199
200        assert!(!handle.is_finished());
201        sender.send(Ok(42usize));
202        assert!(handle.is_finished());
203        assert_eq!(handle.join(), Some(Ok(42)));
204    }
205
206    #[test]
207    fn task_handle_ready_returns_stored_result() {
208        let handle = TaskHandle::ready(TaskId::new(12), Ok(84usize));
209
210        assert!(handle.is_finished());
211        assert_eq!(handle.id(), TaskId::new(12));
212        assert_eq!(handle.join(), Some(Ok(84)));
213    }
214
215    #[test]
216    fn task_handle_reports_cancelled_when_sender_drops() {
217        let (handle, sender) = TaskHandle::<usize>::new_pending(TaskId::new(11));
218
219        drop(sender);
220
221        assert!(handle.is_finished());
222        assert_eq!(handle.join(), Some(Err(crate::error::TaskError::Cancelled)));
223    }
224
225    #[test]
226    fn task_handle_waits_for_cross_thread_completion() {
227        let (handle, sender) = TaskHandle::new_pending(TaskId::new(13));
228
229        let worker = std::thread::spawn(move || {
230            sender.send(Ok(168usize));
231        });
232
233        assert_eq!(handle.join(), Some(Ok(168)));
234        worker.join().unwrap();
235    }
236
237    #[test]
238    fn task_handle_parks_until_delayed_completion() {
239        let (handle, sender) = TaskHandle::new_pending(TaskId::new(14));
240        let (ready_tx, ready_rx) = std::sync::mpsc::sync_channel(0);
241        let (release_tx, release_rx) = std::sync::mpsc::sync_channel(0);
242
243        let worker = std::thread::spawn(move || {
244            ready_tx.send(()).expect("worker publishes readiness");
245            release_rx.recv().expect("worker receives release");
246            sender.send(Ok(336usize));
247        });
248
249        ready_rx
250            .recv_timeout(std::time::Duration::from_secs(1))
251            .expect("worker reaches the delayed-completion gate");
252        release_tx.send(()).expect("release delayed completion");
253        assert_eq!(handle.join(), Some(Ok(336)));
254        worker.join().unwrap();
255    }
256
257    #[test]
258    fn result_wait_policy_is_zero_sized_and_const_bounded() {
259        assert_eq!(core::mem::size_of::<BlockingResultWait>(), 0);
260        assert_eq!(BlockingResultWait::SPIN_ATTEMPTS, MAX_SPIN_ATTEMPTS);
261    }
262}