Skip to main content

pi_async_rt/rt/
multi_thread.rs

1//! # 多线程运行时
2//!
3//! - [ComputationalTaskPool]\: 计算型的多线程任务池,适合用于Cpu密集型的应用,
4//!   不支持运行时伸缩
5//! - [StealableTaskPool]\:
6//!   可窃取的多线程任务池,适合用于block较多的应用,支持运行时伸缩
7//! - [MultiTaskRuntime]\: 异步多线程任务运行时,支持运行时线程伸缩
8//! - [MultiTaskRuntimeBuilder]\: 异步多线程任务运行时构建器
9//!
10//! [ComputationalTaskPool]: struct.ComputationalTaskPool.html
11//! [StealableTaskPool]: struct.StealableTaskPool.html
12//! [MultiTaskRuntime]: struct.MultiTaskRuntime.html
13//! [MultiTaskRuntimeBuilder]: struct.MultiTaskRuntimeBuilder.html
14//!
15//! # Examples
16//!
17//! ```
18//! use pi_async_rt::rt::{AsyncRuntime, AsyncRuntimeExt};
19//! use pi_async_rt::rt::multi_thread::{MultiTaskRuntime, MultiTaskRuntimeBuilder, StealableTaskPool};
20//!
21//! let pool = StealableTaskPool::with(4,100000,[1, 254],3000);
22//! let builer = MultiTaskRuntimeBuilder::new(pool)
23//!     .set_timer_interval(1)
24//!     .init_worker_size(4)
25//!     .set_worker_limit(4, 4);
26//! let rt = builer.build();
27//! let _ = rt.spawn(async move {});
28//! ```
29//!
30//! # Concurrency and worker ownership
31//!
32//! The runtime shares one task pool between all worker threads. Pool-wide state must therefore use
33//! thread-safe containers, while each worker-local stack, queue selector, and deque worker may only
34//! be accessed by the OS thread bound to that exact pool and worker slot. The implementation binds
35//! a private `{thread_id, pool pointer}` context when a worker starts and validates it before every
36//! owner-only access. A wrong-pool owner operation panics before touching worker-local state;
37//! cross-runtime `spawn_local` remains a supported public-queue fallback.
38//!
39//! This validation is O(1), performs one thread-local read and pointer comparison, does not allocate,
40//! clone, lock, spin, block, poll user code, or call into V8/FFI. See the standard regression target
41//! `tests/stealable_task_pool_concurrency.rs`.
42
43use std::sync::Arc;
44use std::vec::IntoIter;
45use std::time::Duration;
46use std::future::Future;
47use std::cell::{Cell, UnsafeCell};
48use std::marker::PhantomData;
49use std::io::{Error, ErrorKind, Result};
50use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
51use std::task::{Context, Poll, Waker};
52use std::thread::{self, Builder};
53
54use async_stream::stream;
55use crossbeam_channel::{bounded, Sender};
56use crossbeam_deque::{Injector, Steal, Stealer, Worker};
57use crossbeam_queue::{ArrayQueue, SegQueue};
58use crossbeam_utils::atomic::AtomicCell;
59use st3::{StealError,
60          fifo::{Worker as FIFOWorker, Stealer as FIFOStealer}};
61use flume::bounded as async_bounded;
62use futures::{
63    future::{BoxFuture, FutureExt},
64    stream::{BoxStream, Stream, StreamExt},
65    task::waker_ref,
66    TryFuture,
67};
68use parking_lot::{Condvar, Mutex};
69use rand::{Rng, thread_rng};
70use num_cpus;
71use wrr::IWRRSelector;
72use quanta::{Clock, Instant as QInstant};
73use log::warn;
74
75use super::{
76    PI_ASYNC_LOCAL_THREAD_ASYNC_RUNTIME, PI_ASYNC_THREAD_LOCAL_ID, DEFAULT_MAX_HIGH_PRIORITY_BOUNDED, DEFAULT_HIGH_PRIORITY_BOUNDED, DEFAULT_MAX_LOW_PRIORITY_BOUNDED, alloc_rt_uid, local_async_runtime, AsyncMapReduce, AsyncPipelineResult, AsyncRuntime,
77    AsyncRuntimeExt, AsyncTask, AsyncTaskPool, AsyncTaskPoolExt, AsyncTaskTimerByNotCancel, AsyncTimingTask,
78    AsyncWait, AsyncWaitAny, AsyncWaitAnyCallback, AsyncWaitTimeout, LocalAsyncWaitTimeout, LocalAsyncRuntime, TaskId, TaskHandle, YieldNow, prune_stale_waiting_workers, register_waiting_worker, wake_waiting_worker
79};
80
81/*
82* 默认的初始工作者数量
83*/
84#[cfg(not(target_arch = "wasm32"))]
85const DEFAULT_INIT_WORKER_SIZE: usize = 2;
86#[cfg(target_arch = "wasm32")]
87const DEFAULT_INIT_WORKER_SIZE: usize = 1;
88
89/*
90* 默认的工作者线程名称前缀
91*/
92const DEFAULT_WORKER_THREAD_PREFIX: &str = "Default-Multi-RT";
93
94/*
95* 默认的线程栈大小
96*/
97const DEFAULT_THREAD_STACK_SIZE: usize = 1024 * 1024;
98
99/*
100* 默认的工作者线程空闲休眠时长,单位ms
101*/
102const DEFAULT_WORKER_THREAD_SLEEP_TIME: u64 = 10;
103
104/*
105* 默认的运行时空闲休眠时长,单位ms,运行时空闲是指绑定当前运行时的队列为空,且定时器内未到期的任务为空
106*/
107const DEFAULT_RUNTIME_SLEEP_TIME: u64 = 1000;
108
109/*
110* 默认的最大权重
111*/
112const DEFAULT_MAX_WEIGHT: u8 = 254;
113
114/*
115* 默认的最小权重
116*/
117const DEFAULT_MIN_WEIGHT: u8 = 1;
118
119/*
120* multi-thread worker id中worker index的掩码
121*/
122const MULTI_THREAD_WORKER_ID_MASK: usize = 0xffffffff;
123
124/// 当前 OS 线程绑定的 multi-thread runtime worker 上下文。
125///
126/// 说明:
127/// - `thread_id` 保持既有 packed runtime-id/worker-index 表示;未绑定线程为 `usize::MAX`。
128/// - `pool` 是当前 worker 持有的 `Arc<P>` 数据地址经类型擦除后的 raw pointer。
129/// - raw pointer 只做地址相等比较,永不解引用、释放或转换回引用,也不拥有 pool。
130///
131/// 业务边界:
132/// - 只用于本模块两个内置 multi-thread pool 的 worker-local owner 校验。
133/// - 不替代全局 `PI_ASYNC_THREAD_LOCAL_ID`,后者仍服务既有 runtime 公共逻辑。
134/// - 不允许据此从错误 runtime/pool 访问 worker-local stack、deque、selector 或 waker。
135///
136/// 性能与副作用:
137/// - Copy-only 固定大小状态;读取、worker id 提取和地址比较均为 O(1),每线程空间 O(1)。
138/// - 上下文读取是纯操作;绑定不是纯操作,会修改当前线程 TLS,且对同一绑定值幂等。
139/// - 不分配、不 clone、不加锁、不自旋、不阻塞、不执行 I/O、future、回调或唤醒。
140///
141/// 安全性:
142/// - TLS `Cell` 只被所属 OS 线程访问,因此无需跨线程同步。
143/// - worker 闭包在整个工作循环持有 runtime `Arc`,比较期间 pool 数据地址有效。
144/// - 内存安全和线程安全不依赖 raw pointer 解引用;错误地址只会导致 fail-fast panic。
145/// - 不触达 V8/FFI,不跨 await 保存 guard,异步安全。
146#[derive(Clone, Copy)]
147struct MultiThreadWorkerContext {
148    thread_id: usize,
149    pool: *const (),
150}
151
152impl MultiThreadWorkerContext {
153    const UNBOUND: Self = MultiThreadWorkerContext {
154        thread_id: usize::MAX,
155        pool: std::ptr::null(),
156    };
157
158    /// 返回当前上下文所属 runtime id。
159    ///
160    /// 只对 packed id 做位移,O(1)、纯函数、幂等、无副作用且不阻塞。调用方必须先使用该值
161    /// 判断 task 是否属于当前 runtime,才能决定是 local owner 路径还是 public fallback。
162    #[inline]
163    const fn runtime_id(self) -> usize {
164        self.thread_id >> 32
165    }
166
167    /// 校验当前线程属于 `pool`,并返回 owner worker index。
168    ///
169    /// 参数:
170    /// - `pool`:即将访问 owner-only worker 状态的具体 pool;只借用,不保存、不 clone。
171    ///
172    /// 返回:
173    /// - pool 地址匹配时返回 packed id 的低 32 位 worker index。
174    ///
175    /// Panic:
176    /// - 当前线程未绑定 multi-thread worker,或绑定到另一个 pool 时,在访问任何 owner-only
177    ///   状态之前 panic。该 panic 表示调用方违反 `try_pop`/local worker owner 前置条件。
178    ///
179    /// 性能和安全:
180    /// - O(1) 时间、O(1) 空间、纯只读、幂等;正常路径仅一次 raw pointer 比较和位与。
181    /// - 无锁、无分配、无 clone、无阻塞;raw pointer 永不解引用,因此比较本身内存安全。
182    #[inline]
183    fn owner_worker_id<P>(self, pool: &P) -> usize {
184        let expected = pool as *const P as *const ();
185        if self.pool != expected {
186            panic!(
187                "Multi-thread task pool owner mismatch: owner-only worker state requires the worker bound to this pool"
188            );
189        }
190
191        self.thread_id & MULTI_THREAD_WORKER_ID_MASK
192    }
193}
194
195thread_local! {
196    /// 当前 OS 线程的 multi-thread worker/pool 绑定。
197    ///
198    /// 默认未绑定;只由 `bind_multi_thread_worker_context` 在 worker 启动时写入。`Cell` 不会
199    /// 跨线程共享,无锁、无分配,也不拥有 raw pool pointer 指向的对象。
200    static PI_ASYNC_MULTI_THREAD_WORKER_CONTEXT: Cell<MultiThreadWorkerContext>
201        = Cell::new(MultiThreadWorkerContext::UNBOUND);
202}
203
204/// 读取当前 multi-thread worker 上下文。
205///
206/// 返回 Copy snapshot,O(1)、无分配、无锁、无阻塞。正常线程生命周期中该操作是纯只读且
207/// 幂等;仅在线程 TLS 已销毁的非法调用阶段 panic。它是 owner helper 的内部 fail-fast
208/// 边界,不替代或改变公开 `get_thread_id` 使用的既有 TLS。
209#[inline]
210fn current_multi_thread_worker_context() -> MultiThreadWorkerContext {
211    match PI_ASYNC_MULTI_THREAD_WORKER_CONTEXT.try_with(|context| context.get()) {
212        Ok(context) => context,
213        Err(e) => {
214            panic!(
215                "Get multi-thread worker context failed, thread: {:?}, reason: {:?}",
216                thread::current(),
217                e
218            );
219        },
220    }
221}
222
223///
224/// 计算型的工作者任务队列
225///
226struct ComputationalTaskQueue<O: Default + 'static> {
227    stack: Worker<Arc<AsyncTask<ComputationalTaskPool<O>, O>>>,     //工作者任务栈
228    queue: SegQueue<Arc<AsyncTask<ComputationalTaskPool<O>, O>>>,   //工作者任务队列
229    thread_waker: Arc<(AtomicBool, Mutex<()>, Condvar)>,            //工作者线程的唤醒器
230}
231
232impl<O: Default + 'static> ComputationalTaskQueue<O> {
233    //构建计算型的工作者任务队列
234    pub fn new(thread_waker: Arc<(AtomicBool, Mutex<()>, Condvar)>) -> Self {
235        let stack = Worker::new_lifo();
236        let queue = SegQueue::new();
237
238        ComputationalTaskQueue {
239            stack,
240            queue,
241            thread_waker,
242        }
243    }
244
245    //获取计算型的工作者任务队列的任务数量
246    pub fn len(&self) -> usize {
247        self.stack.len() + self.queue.len()
248    }
249}
250
251/// 计算型多线程任务池,适合 CPU 密集型应用,不支持运行时伸缩。
252///
253/// # 使用方式
254///
255/// 通过 [`ComputationalTaskPool::new`] 创建与 runtime worker slot 数匹配的 pool,再交给
256/// [`MultiTaskRuntimeBuilder`]。外部线程可以使用 runtime 的 `spawn`/`spawn_local`;直接
257/// `try_pop`、`try_pop_all` 或取得当前 worker waker 只允许在拥有该 pool 的 worker 上执行。
258///
259/// ```
260/// use pi_async_rt::rt::{AsyncRuntime, multi_thread::{ComputationalTaskPool, MultiTaskRuntimeBuilder}};
261///
262/// let pool = ComputationalTaskPool::new(2);
263/// let runtime = MultiTaskRuntimeBuilder::new(pool)
264///     .init_worker_size(2)
265///     .set_worker_limit(2, 2)
266///     .build();
267/// runtime.spawn(async {}).unwrap();
268/// ```
269///
270/// # 业务与错误边界
271///
272/// - `push` 允许任意线程调用;`push_local` 在非所属 runtime 线程上保持公共队列 fallback。
273/// - owner-only pop/waker 操作在未绑定线程或其它 pool 的 worker 上调用会 panic,并保证在访问
274///   worker-local `crossbeam_deque::Worker` 前失败。
275/// - pool 不定义任务结果、取消、timeout 或 runtime 关闭语义。
276///
277/// # 性能与安全
278///
279/// - push/pop 快路径为 O(1);总长度为原子计数近似值,空间为 O(W + N)。
280/// - owner 校验一次 TLS 读取和一次指针比较,无分配、clone、锁、自旋或阻塞。
281/// - 类型不是纯值容器:push/pop 会修改队列和统计,非幂等;不会在内部 poll 用户 future。
282/// - pool 可在线程间共享;worker-local stack 只由通过 pool identity 校验的 owner worker 访问。
283/// - 不持有跨 await guard,不触达 V8/FFI。专项入口:
284///   `tests/stealable_task_pool_concurrency.rs` 的 owner/cross-runtime 用例。
285pub struct ComputationalTaskPool<O: Default + 'static> {
286    workers: Vec<ComputationalTaskQueue<O>>, //工作者的任务队列列表
287    waits: Option<Arc<ArrayQueue<Arc<(AtomicBool, Mutex<()>, Condvar)>>>>, //待唤醒的工作者唤醒器队列
288    consume_count: Arc<AtomicUsize>,                                       //任务消费计数
289    produce_count: Arc<AtomicUsize>,                                       //任务生产计数
290}
291
292// SAFETY: shared queue/counters/waits provide their own synchronization. The non-Sync local
293// crossbeam `Worker` is accessed only after `MultiThreadWorkerContext::owner_worker_id` proves the
294// caller is the worker bound to this exact pool; other threads use the `SegQueue` path. The worker
295// closure retains the runtime Arc for the full access lifetime. The invariant is covered by
296// `tests/stealable_task_pool_concurrency.rs`.
297unsafe impl<O: Default + 'static> Send for ComputationalTaskPool<O> {}
298// SAFETY: see the field-by-field owner and lifetime proof above. No method exposes a reference to
299// the local `Worker`, and wrong-pool owner operations panic before indexing or touching it.
300unsafe impl<O: Default + 'static> Sync for ComputationalTaskPool<O> {}
301
302impl<O: Default + 'static> Default for ComputationalTaskPool<O> {
303    fn default() -> Self {
304        #[cfg(not(target_arch = "wasm32"))]
305        let core_len = num_cpus::get(); //工作者任务池数据等于本机逻辑核数
306        #[cfg(target_arch = "wasm32")]
307        let core_len = 1; //工作者任务池数据等于1
308        ComputationalTaskPool::new(core_len)
309    }
310}
311
312impl<O: Default + 'static> AsyncTaskPool<O> for ComputationalTaskPool<O> {
313    type Pool = ComputationalTaskPool<O>;
314
315    /// 返回当前线程既有的 packed runtime/worker id。
316    ///
317    /// worker 线程返回 runtime id 与 worker index 的组合值,未绑定线程保持既有
318    /// `usize::MAX`。本方法不校验 pool owner,调用目的、返回语义和 panic 边界均未改变。
319    /// O(1)、纯只读、幂等、无分配/锁/阻塞/唤醒;只读取当前线程 TLS,线程和内存安全。
320    #[inline]
321    fn get_thread_id(&self) -> usize {
322        match PI_ASYNC_THREAD_LOCAL_ID.try_with(move |thread_id| unsafe {
323            // SAFETY: this reads only the current OS thread's TLS UnsafeCell and returns a Copy id.
324            *thread_id.get()
325        }) {
326            Err(e) => {
327                panic!(
328                    "Get thread id failed, thread: {:?}, reason: {:?}",
329                    thread::current(),
330                    e
331                );
332            }
333            Ok(id) => id,
334        }
335    }
336
337    #[inline]
338    fn len(&self) -> usize {
339        if let Some(len) = self
340            .produce_count
341            .load(Ordering::Relaxed)
342            .checked_sub(self.consume_count.load(Ordering::Relaxed))
343        {
344            len
345        } else {
346            0
347        }
348    }
349
350    #[inline]
351    fn push(&self, task: Arc<AsyncTask<Self::Pool, O>>) -> Result<()> {
352        let index = self.produce_count.fetch_add(1, Ordering::Relaxed) % self.workers.len();
353        self.workers[index].queue.push(task);
354        Ok(())
355    }
356
357    /// 将任务优先提交给所属 worker 的本地 queue,否则保持公共 queue fallback。
358    ///
359    /// `task` 的 `Arc` 所有权被 queue 接收;本实现当前成功返回 `Ok(())`。当 task owner 与当前
360    /// runtime 不同(含外部线程)时不触碰 owner-only 状态;owner 相同但 pool 身份错误时在
361    /// queue 访问前 panic。O(1)、零额外分配/clone/锁/阻塞,不 poll 用户 future。该操作非纯、
362    /// 非幂等;合法调用域内线程安全、内存安全和异步安全,不触达 V8/FFI。
363    #[inline]
364    fn push_local(&self, task: Arc<AsyncTask<Self::Pool, O>>) -> Result<()> {
365        let context = current_multi_thread_worker_context();
366        let rt_uid = task.owner();
367        if context.runtime_id() == rt_uid {
368            //当前是运行时所在线程
369            let worker = &self.workers[context.owner_worker_id(self)];
370            worker.queue.push(task);
371
372            self.produce_count.fetch_add(1, Ordering::Relaxed);
373            Ok(())
374        } else {
375            //当前不是运行时所在线程
376            self.push(task)
377        }
378    }
379
380    /// 按既有优先级语义提交任务,并保护最高优先级 owner-only stack。
381    ///
382    /// `priority` 的阈值和 fallback 顺序不变;`task` 被成功转移后返回 `Ok(())`。只有最高优先级
383    /// 且 task 属于当前 runtime 时执行 exact-pool owner 校验,错误 pool 在 stack 访问前 panic。
384    /// O(1)、无新增分配/clone/锁/阻塞;非纯、非幂等,不执行任务或用户回调,线程/异步/内存
385    /// 安全边界与 [`ComputationalTaskPool::push_local`] 相同。
386    #[inline]
387    fn push_priority(&self,
388                     priority: usize,
389                     task: Arc<AsyncTask<Self::Pool, O>>) -> Result<()> {
390        if priority >= DEFAULT_MAX_HIGH_PRIORITY_BOUNDED {
391            //最高优先级
392            let context = current_multi_thread_worker_context();
393            let rt_uid = task.owner();
394            if context.runtime_id() == rt_uid {
395                let worker = &self.workers[context.owner_worker_id(self)];
396                worker.stack.push(task);
397
398                self.produce_count.fetch_add(1, Ordering::Relaxed);
399                Ok(())
400            } else {
401                self.push(task)
402            }
403        } else if priority >= DEFAULT_HIGH_PRIORITY_BOUNDED {
404            //高优先级
405            self.push_local(task)
406        } else {
407            //低优先级
408            self.push(task)
409        }
410    }
411
412    #[inline]
413    fn push_keep(&self, task: Arc<AsyncTask<Self::Pool, O>>) -> Result<()> {
414        self.push_priority(DEFAULT_HIGH_PRIORITY_BOUNDED, task)
415    }
416
417    /// 从当前 exact pool 的 owner worker slot 尝试取一个任务。
418    ///
419    /// 返回任务的 `Arc` 或空队列时的 `None`。只能由该 pool 绑定的 worker 调用;外部线程或
420    /// wrong-pool worker 会在索引及 local stack 访问前 panic。快路径 O(1)、零分配、零 clone、
421    /// 无锁/自旋/阻塞;消费队列使其非纯、非幂等。本方法不 poll/drop 任务,不触达 V8/FFI,
422    /// 在 owner 契约内线程、内存与异步安全。
423    #[inline]
424    fn try_pop(&self) -> Option<Arc<AsyncTask<Self::Pool, O>>> {
425        let id = current_multi_thread_worker_context().owner_worker_id(self);
426        let worker = &self.workers[id];
427        let task = worker.stack.pop();
428        if task.is_some() {
429            //指定工作者的任务栈有任务,则立即返回任务
430            self.consume_count.fetch_add(1, Ordering::Relaxed);
431            return task;
432        }
433
434        let task = worker.queue.pop();
435        if task.is_some() {
436            self.consume_count.fetch_add(1, Ordering::Relaxed);
437        }
438
439        task
440    }
441
442    /// 反复执行 owner-checked `try_pop`,返回当前可取得任务的 owned iterator。
443    ///
444    /// 空池返回空 iterator;wrong-pool 调用在首次 pop 前 panic。时间 O(N)、临时空间 O(N),会
445    /// 按既有实现分配 `Vec`,其预分配容量来自近似 `len()`。非纯、非幂等;不阻塞、不 poll
446    /// 用户 future,owner 契约内线程/内存/异步安全。该既有批量分配行为不属于本轮修改。
447    #[inline]
448    fn try_pop_all(&self) -> IntoIter<Arc<AsyncTask<Self::Pool, O>>> {
449        let mut tasks = Vec::with_capacity(self.len());
450        while let Some(task) = self.try_pop() {
451            tasks.push(task);
452        }
453
454        tasks.into_iter()
455    }
456
457    #[inline]
458    fn get_thread_waker(&self) -> Option<&Arc<(AtomicBool, Mutex<()>, Condvar)>> {
459        //多线程任务运行时不支持此方法
460        None
461    }
462}
463
464impl<O: Default + 'static> AsyncTaskPoolExt<O> for ComputationalTaskPool<O> {
465    #[inline]
466    fn set_waits(&mut self, waits: Arc<ArrayQueue<Arc<(AtomicBool, Mutex<()>, Condvar)>>>) {
467        self.waits = Some(waits);
468    }
469
470    #[inline]
471    fn get_waits(&self) -> Option<&Arc<ArrayQueue<Arc<(AtomicBool, Mutex<()>, Condvar)>>>> {
472        self.waits.as_ref()
473    }
474
475    #[inline]
476    fn worker_len(&self) -> usize {
477        self.workers.len()
478    }
479
480    /// clone 当前 exact pool/worker 的休眠唤醒器。
481    ///
482    /// 合法 owner 返回 `Some(Arc<...>)`;wrong-pool 或未绑定线程在索引前 panic。O(1),会按
483    /// 既有语义执行一次 `Arc::clone`,但 owner 校验自身不 clone/分配/加锁/阻塞。该只读查询
484    /// 不唤醒线程,本身幂等但增加引用计数;返回 Arc 的释放由调用方负责,线程和内存安全。
485    #[inline]
486    fn clone_thread_waker(&self) -> Option<Arc<(AtomicBool, Mutex<()>, Condvar)>> {
487        let worker = &self.workers[current_multi_thread_worker_context().owner_worker_id(self)];
488        Some(worker.thread_waker.clone())
489    }
490}
491
492impl<O: Default + 'static> ComputationalTaskPool<O> {
493    /// 构建指定 worker slot 数量的计算型多线程任务池。
494    ///
495    /// # 参数与返回
496    ///
497    /// - `size`:请求的 worker slot 数;小于默认初始 worker 数时保持旧行为,提升到默认值。
498    /// - 返回独立 pool;尚未绑定 runtime/waits,也不会启动线程或执行 future。
499    ///
500    /// # 性能与副作用
501    ///
502    /// 时间和空间复杂度均为 O(W),W 为实际 slot 数。该函数会分配 worker/queue/waker,但不
503    /// 阻塞、不执行 I/O、不创建线程。它不是纯函数;每次调用创建不同 pool,因此非幂等。
504    ///
505    /// # 安全与边界
506    ///
507    /// 返回值可安全移动并由 builder 在线程间共享。调用方必须让 builder 的最大 worker 数
508    /// 不超过 slot 数;builder 会再次收敛该边界。owner-only API 的线程/pool 限制见类型文档。
509    pub fn new(mut size: usize) -> Self {
510        if size < DEFAULT_INIT_WORKER_SIZE {
511            //工作者数量过少,则设置为默认的工作者数量
512            size = DEFAULT_INIT_WORKER_SIZE;
513        }
514
515        let mut workers = Vec::with_capacity(size);
516        for _ in 0..size {
517            let thread_waker = Arc::new((AtomicBool::new(false), Mutex::new(()), Condvar::new()));
518            let worker = ComputationalTaskQueue::new(thread_waker);
519            workers.push(worker);
520        }
521        let consume_count = Arc::new(AtomicUsize::new(0));
522        let produce_count = Arc::new(AtomicUsize::new(0));
523
524        ComputationalTaskPool {
525            workers,
526            waits: None,
527            consume_count,
528            produce_count,
529        }
530    }
531}
532
533///
534/// 可窃取的混合任务队列
535///
536struct StealableTaskQueue<O: Default + 'static> {
537    stack:          UnsafeCell<Option<Arc<AsyncTask<StealableTaskPool<O>, O>>>>,    //工作者任务栈
538    internal:       FIFOWorker<Arc<AsyncTask<StealableTaskPool<O>, O>>>,            //工作者本地内部任务队列,可窃取
539    external:       Worker<Arc<AsyncTask<StealableTaskPool<O>, O>>>,                //工作者本地外部任务队列,可窃取
540    selector:       UnsafeCell<IWRRSelector<2>>,                                    //工作者任务队列选择器
541    thread_waker:   Arc<(AtomicBool, Mutex<()>, Condvar)>,                          //工作者线程的唤醒器
542}
543
544impl<O: Default + 'static> StealableTaskQueue<O> {
545    // 构建可窃取的混合任务队列,允许设置初始的栈和队列的初始容量,并自动设置栈和队列的容量
546    // 栈和队列的容量是初始容量的最小二次方,例如初始容量为0,则容量为1
547    pub fn new(
548        init_queue_capacity: usize,
549        thread_waker: Arc<(AtomicBool, Mutex<()>, Condvar)>,
550    ) -> (Self,
551          FIFOStealer<Arc<AsyncTask<StealableTaskPool<O>, O>>>,
552          Stealer<Arc<AsyncTask<StealableTaskPool<O>, O>>>) {
553        let stack = UnsafeCell::new(None);
554        let internal = FIFOWorker::new(init_queue_capacity);
555        let external = Worker::new_fifo();
556        let internal_stealer = internal.stealer();
557        let external_stealer = external.stealer();
558        let selector = UnsafeCell::new(IWRRSelector::new([2, 1]));
559
560        (
561            StealableTaskQueue {
562                stack,
563                internal,
564                external,
565                selector,
566                thread_waker,
567            },
568            internal_stealer,
569            external_stealer
570        )
571    }
572
573    // 获取栈容量
574    pub const fn stack_capacity(&self) -> usize {
575        1
576    }
577
578    // 获取本地内部任务队列容量
579    pub fn internal_capacity(&self) -> usize {
580        self.internal.capacity()
581    }
582
583    // 获取剩余的本地内部任务队列容量,不准确
584    pub fn remaining_internal_capacity(&self) -> usize {
585        self.internal.spare_capacity()
586    }
587
588    /// 返回当前 worker single-item stack 的近似精确长度(0 或 1)。
589    ///
590    /// 该私有 helper 的强前置条件是调用方已经通过 exact pool owner 校验;它随后只读取当前
591    /// owner 的 `UnsafeCell<Option<_>>`。O(1)、零分配/clone/锁/阻塞,纯只读且幂等,不执行
592    /// future 或回调。若绕过 owner 前置条件并发调用会破坏 `unsafe impl Sync` 的安全不变量,
593    /// 因而所有生产入口必须从 `push_priority` 的已校验 local 分支到达。
594    #[inline]
595    pub fn stack_len(&self) -> usize {
596        unsafe {
597            // SAFETY: callers reach this private queue only through the exact pool's owner-checked
598            // local branch. No other thread reads or writes this worker slot's single-item stack.
599            if (&*self.stack.get()).is_some() {
600                1
601            } else {
602                0
603            }
604        }
605    }
606
607    // 获取本地内部任务队列长度
608    pub fn internal_len(&self) -> usize {
609        self
610            .internal_capacity()
611            .checked_sub(self.remaining_internal_capacity())
612            .unwrap_or(0)
613    }
614
615    // 获取本地外部任务队列长度
616    pub fn external_len(&self) -> usize {
617        self.external.len()
618    }
619}
620
621/// 可窃取的混合多线程任务池。
622///
623/// # 使用方式
624///
625/// pool 将 external/public 任务和 runtime worker 内产生的 internal/local 任务分开保存,并由
626/// 每个 worker 的 `IWRRSelector` 按既有 best-effort 权重选择顺序。典型用法:
627///
628/// ```
629/// use pi_async_rt::rt::{AsyncRuntime, multi_thread::{MultiTaskRuntimeBuilder, StealableTaskPool}};
630///
631/// let pool = StealableTaskPool::with(4, 4096, [1, 1], 3000);
632/// let runtime = MultiTaskRuntimeBuilder::new(pool)
633///     .init_worker_size(4)
634///     .set_worker_limit(4, 4)
635///     .build();
636/// runtime.spawn(async {}).unwrap();
637/// ```
638///
639/// # 调度与业务边界
640///
641/// - `push` 可从任意线程进入 public injector;所属 worker 的 `push_local` 优先进入 internal
642///   queue,跨 runtime 调用保持 public fallback。
643/// - `try_pop`/`try_pop_all` 和 worker waker 获取是 owner-only 操作,只允许拥有该 exact pool
644///   的 worker 调用;错误 pool/外部线程会在访问 worker-local 状态前 panic。
645/// - pool-level 流量计数和刷新时间只用于近似调度启发式,不发布任务对象,也不承诺所有 worker
646///   selector 同时或最终应用同一权重。
647/// - `weights` 保留既有 API/存储语义;本版本不改变其历史行为,不在此处定义新的公平性保证。
648/// - timeout、取消、任务结果、worker sleep/wake 和 runtime 关闭由其它层负责。
649///
650/// # 性能、纯度和副作用
651///
652/// - local/public pop 快路径为 O(1);尝试其它 W-1 个 stealer 的最坏时间为 O(W)。空间
653///   O(W + N),N 为排队任务数。
654/// - 每次 weighted pop 读取一次 `AtomicCell<QInstant>`;到刷新窗口时写入一次。x86_64 验收
655///   目标要求该类型 lock-free;其它目标保证正确性但不承诺无内部锁。
656/// - owner 校验 O(1),一次 TLS Copy 读取和 raw pointer 比较;不分配、不 clone、不锁、不阻塞。
657/// - push/pop/刷新均非纯且非幂等,会修改队列、统计或当前 worker selector;不会在 pool 内
658///   poll 用户 future、执行回调、I/O 或 V8/FFI。
659///
660/// # 线程、内存与异步安全
661///
662/// pool-wide 状态使用并发容器/原子;`stack`、deque worker 和 selector 只由 owner worker
663/// 访问。raw pool pointer 仅比较不解引用,worker 闭包持有 runtime Arc 保证生命周期。没有锁
664/// guard 跨 await,也不引入引用环。专项和 TSan 入口:
665/// `tests/stealable_task_pool_concurrency.rs`。
666pub struct StealableTaskPool<O: Default + 'static> {
667    public:                         Injector<Arc<AsyncTask<StealableTaskPool<O>, O>>>,          //公共的任务池
668    workers:                        Vec<StealableTaskQueue<O>>,                                 //工作者的任务队列列表
669    internal_stealers:              Vec<FIFOStealer<Arc<AsyncTask<StealableTaskPool<O>, O>>>>,  //工作者任务队列的本地内部任务窃取者
670    external_stealers:              Vec<Stealer<Arc<AsyncTask<StealableTaskPool<O>, O>>>>,      //工作者任务队列的本地外部任务窃取者
671    internal_consume:               AtomicUsize,                                                //内部任务消费计数
672    internal_produce:               AtomicUsize,                                                //内部任务生产计数
673    internal_traffic_statistics:    AtomicUsize,                                                //内部任务流量统计
674    external_consume:               AtomicUsize,                                                //外部任务消费计数
675    external_produce:               AtomicUsize,                                                //外部任务生产计数
676    external_traffic_statistics:    AtomicUsize,                                                //外部任务流量统计
677    weights:                        [u8; 2],                                                    //工作者任务队列的权重
678    clock:                          Clock,                                                      //任务池的时钟
679    interval:                       usize,                                                      //整理的间隔时长,单位ms
680    last_time:                      AtomicCell<QInstant>,                                       //线程安全的上一次整理时间
681    waits:                          Option<Arc<ArrayQueue<Arc<(AtomicBool, Mutex<()>, Condvar)>>>>, //待唤醒的工作者唤醒器队列
682}
683
684// SAFETY: public injector, stealers, atomic counters, AtomicCell timestamp, clock and waits are safe
685// to share. Each queue's stack, local workers and selector remain owner-only: every production path
686// that reads or mutates them first validates the current TLS pool pointer and worker index. Stealers
687// are the only cross-worker view of local queues. The worker closure retains the runtime Arc for the
688// entire loop. This invariant is exercised by the owner, exact-once and TSan standard tests.
689unsafe impl<O: Default + 'static> Send for StealableTaskPool<O> {}
690// SAFETY: see the field-by-field synchronization and owner proof above. Wrong-pool safe API calls
691// panic before worker lookup or UnsafeCell access; no owner-only reference escapes from the pool.
692unsafe impl<O: Default + 'static> Sync for StealableTaskPool<O> {}
693
694impl<O: Default + 'static> Default for StealableTaskPool<O> {
695    fn default() -> Self {
696        StealableTaskPool::new()
697    }
698}
699
700impl<O: Default + 'static> AsyncTaskPool<O> for StealableTaskPool<O> {
701    type Pool = StealableTaskPool<O>;
702
703    /// 返回当前线程既有的 packed runtime/worker id。
704    ///
705    /// worker 线程返回 runtime id 与 worker index 的组合值,未绑定线程保持既有
706    /// `usize::MAX`。本方法不执行新增 owner 校验;O(1)、纯只读、幂等、无分配/锁/阻塞,
707    /// 只读取当前线程 TLS,调用目的、可见语义及线程/内存安全边界均保持不变。
708    #[inline]
709    fn get_thread_id(&self) -> usize {
710        match PI_ASYNC_THREAD_LOCAL_ID.try_with(move |thread_id| unsafe {
711            // SAFETY: this reads only the current OS thread's TLS UnsafeCell and returns a Copy id.
712            *thread_id.get()
713        }) {
714            Err(e) => {
715                panic!(
716                    "Get thread id failed, thread: {:?}, reason: {:?}",
717                    thread::current(),
718                    e
719                );
720            }
721            Ok(id) => id,
722        }
723    }
724
725    #[inline]
726    fn len(&self) -> usize {
727        self.internal_produce
728            .load(Ordering::Relaxed)
729            .checked_sub(self.internal_consume.load(Ordering::Relaxed))
730            .unwrap_or(0)
731            +
732            self.external_produce
733                .load(Ordering::Relaxed)
734                .checked_sub(self.external_consume.load(Ordering::Relaxed))
735                .unwrap_or(0)
736    }
737
738    #[inline]
739    fn push(&self, task: Arc<AsyncTask<Self::Pool, O>>) -> Result<()> {
740        self.public.push(task);
741
742        self
743            .external_produce
744            .fetch_add(1, Ordering::Relaxed);
745        Ok(())
746    }
747
748    /// 将任务提交到所属 worker 的 internal queue,无法使用本地路径时进入 public injector。
749    ///
750    /// `task` 所有权被成功转移并返回 `Ok(())`。外部线程或跨 runtime 调用保持 public fallback;
751    /// task owner 与当前 runtime 相同但 pool 身份错误时,在读取 internal queue 前 panic。正常
752    /// local/public 路径分摊 O(1)、零新增分配/clone/锁/阻塞;非纯、非幂等,不 poll 用户
753    /// future。合法调用域内线程安全、内存安全、异步安全,不触达 V8/FFI。
754    #[inline]
755    fn push_local(&self, task: Arc<AsyncTask<Self::Pool, O>>) -> Result<()> {
756        let context = current_multi_thread_worker_context();
757        let rt_uid = task.owner();
758        if context.runtime_id() == rt_uid {
759            //当前是运行时所在线程
760            let worker = &self.workers[context.owner_worker_id(self)];
761            if worker.remaining_internal_capacity() > 0 {
762                //本地内部任务队列有空闲容量,则立即将任务加入本地内部任务队列
763                let _ = worker.internal.push(task);
764
765                self
766                    .internal_produce
767                    .fetch_add(1, Ordering::Relaxed);
768                Ok(())
769            } else {
770                //本地内部任务队列没有空闲容量,则立即将任务加入公共任务池
771                self.push(task)
772            }
773        } else {
774            //当前不是运行时所在线程
775            self.push(task)
776        }
777    }
778
779    /// 按既有优先级规则提交任务,并保护最高优先级 owner-only single-item stack。
780    ///
781    /// `priority` 的阈值、internal/public fallback 和返回值保持不变。最高优先级本地分支先
782    /// 校验 exact pool,错误上下文在 stack/queue 访问前 panic;其它 runtime 继续 public
783    /// fallback。分摊 O(1),无新增分配/clone/锁/阻塞;非纯、非幂等,不执行任务、回调、
784    /// I/O 或 V8/FFI,owner 契约内线程/内存/异步安全。
785    #[inline]
786    fn push_priority(&self,
787                     priority: usize,
788                     task: Arc<AsyncTask<Self::Pool, O>>) -> Result<()> {
789        if priority >= DEFAULT_MAX_HIGH_PRIORITY_BOUNDED {
790            //最高优先级
791            let context = current_multi_thread_worker_context();
792            let rt_uid = task.owner();
793            if context.runtime_id() == rt_uid {
794                //当前是运行时所在线程
795                let worker = &self.workers[context.owner_worker_id(self)];
796                if worker.stack_len() < 1 {
797                    //本地任务栈有空闲容量,则立即将任务加入本地任务栈
798                    unsafe {
799                        // SAFETY: context.owner_worker_id(self) above proved this exact pool and
800                        // worker slot are owned by the current OS thread. The stack never escapes.
801                        *worker.stack.get() = Some(task);
802                    }
803                } else if worker.remaining_internal_capacity() > 0 {
804                    //本地内部任务队列有空闲容量,则立即将任务加入本地内部任务队列
805                    let _ = worker.internal.push(task);
806                } else {
807                    //本地任务栈和本地内部任务队列都没有空闲容量,则立即将任务加入公共任务池
808                    return self.push(task);
809                }
810
811                self
812                    .internal_produce
813                    .fetch_add(1, Ordering::Relaxed);
814                Ok(())
815            } else {
816                //当前不是运行时所在线程
817                self.push(task)
818            }
819        } else if priority >= DEFAULT_HIGH_PRIORITY_BOUNDED {
820            //高优先级
821            self.push_local(task)
822        } else {
823            //低优先级
824            self.push(task)
825        }
826    }
827
828    #[inline]
829    fn push_keep(&self, task: Arc<AsyncTask<Self::Pool, O>>) -> Result<()> {
830        self.push_priority(DEFAULT_HIGH_PRIORITY_BOUNDED, task)
831    }
832
833    /// 从当前 exact pool 的 owner worker slot 按既有权重和 steal 顺序尝试取一个任务。
834    ///
835    /// 返回任务 `Arc` 或所有候选队列均空时的 `None`。只能由所属 worker 调用;wrong-pool 或
836    /// 外部线程在索引、stack、selector 及 local deque 访问前 panic。local 快路径 O(1),尝试
837    /// W-1 个 stealer 最坏 O(W);原实现的空队列 steal 路径可能构造 O(W) 临时 `Vec`,本轮不
838    /// 改该行为。新增 owner guard 零分配/clone/锁/阻塞。消费和统计更新使本方法非纯、非幂等;
839    /// 不 poll 用户 future,owner 契约内线程/内存/异步安全。
840    #[inline]
841    fn try_pop(&self) -> Option<Arc<AsyncTask<Self::Pool, O>>> {
842        let id = current_multi_thread_worker_context().owner_worker_id(self);
843        let worker = &self.workers[id];
844        let task = unsafe {
845            // SAFETY: owner_worker_id validated the exact pool before indexing. Only this worker
846            // accesses its stack, so taking the Option through UnsafeCell is exclusive.
847            (&mut *worker
848                .stack
849                .get())
850                .take()
851        };
852        if task.is_some() {
853            //指定工作者的任务栈有任务,则立即返回任务
854            return task;
855        }
856
857        //从指定工作者的任务队列中弹出任务
858        try_pop_by_weight(self, worker, id)
859    }
860
861    /// 反复执行 owner-checked `try_pop`,返回当前可取得任务的 owned iterator。
862    ///
863    /// 空池返回空 iterator;wrong-pool 调用在首次 pop 前 panic。除每次 pop 的既有复杂度外,
864    /// 汇总 N 个任务需 O(N) 时间和 O(N) `Vec` 空间;会分配但不阻塞、不执行任务。该操作非纯、
865    /// 非幂等,owner 契约内线程/内存/异步安全;批量收集语义和分配行为未被本轮改变。
866    #[inline]
867    fn try_pop_all(&self) -> IntoIter<Arc<AsyncTask<Self::Pool, O>>> {
868        let mut tasks = Vec::with_capacity(self.len());
869        while let Some(task) = self.try_pop() {
870            tasks.push(task);
871        }
872
873        tasks.into_iter()
874    }
875
876    #[inline]
877    fn get_thread_waker(&self) -> Option<&Arc<(AtomicBool, Mutex<()>, Condvar)>> {
878        //多线程任务运行时不支持此方法
879        None
880    }
881}
882
883// 获取指定数字的MSB
884const fn get_msb(n: usize) -> usize {
885    usize::BITS as usize - n.leading_zeros() as usize
886}
887
888/// 使用既有近似流量统计刷新当前 owner worker 的 selector,并按权重尝试取任务。
889///
890/// 参数:
891/// - `pool`:所有 worker 共享的 exact `StealableTaskPool`;只借用,不 clone/保存。
892/// - `local_worker`:已经由调用方 owner 校验选出的当前 worker slot。
893/// - `local_worker_id`:上述 slot 的索引,用于偷取时排除自己。
894///
895/// 返回:按原 selector、fallback 和 steal 顺序取得一个任务,或队列均为空时返回 `None`。
896/// 任务所有权随 `Arc` 返回给工作循环;本函数不 poll、drop 或执行任务。
897///
898/// 并发与边界:
899/// - `last_time` 是 pool-wide `AtomicCell<QInstant>`;所有 worker 可并发 load/store,不存在普通
900///   共享读写。允许多个 worker 在同一近似窗口刷新,保持旧 best-effort 行为,故不使用 CAS。
901/// - traffic counters 是 Relaxed 近似统计,不承担任务对象发布;交错采样只影响启发式权重。
902/// - selector 仍只修改 `local_worker`,其 UnsafeCell 安全性依赖调用方已完成 exact pool owner
903///   校验。该前置条件由私有调用链 `StealableTaskPool::try_pop` 保证。
904/// - interval 非零由构造器保证;quanta 的 `duration_since` 对倒退采样饱和为零。
905///
906/// 性能、纯度与安全:
907/// - 常规选择 O(1);steal 最坏 O(W),W 为 worker 数;本轮新增原子访问 O(1)、零分配。
908/// - 非纯、非幂等:可能更新统计、selector、刷新时间并消费队列任务。
909/// - 本实现不显式加锁、自旋等待或阻塞,不执行 I/O/回调/V8/FFI,不跨 await 保存状态。
910/// - x86_64 上 `AtomicCell<QInstant>` 的 lock-free 前提由专项测试固定;其它 target 的
911///   crossbeam fallback 可能使用内部同步,只保证线程安全正确性,不承诺无锁性能。
912fn try_pop_by_weight<O: Default + 'static>(pool: &StealableTaskPool<O>,
913                                           local_worker: &StealableTaskQueue<O>,
914                                           local_worker_id: usize)
915                                           -> Option<Arc<AsyncTask<StealableTaskPool<O>, O>>> {
916    unsafe {
917        // SAFETY: this private helper is called only after try_pop validates the current thread as
918        // owner of `local_worker` in this exact pool. Consequently selector has one mutable owner;
919        // pool-wide timestamp/counters are atomic and stealing uses dedicated thread-safe stealers.
920        let duration = pool
921            .clock
922            .recent()
923            .duration_since(pool.last_time.load())
924            .as_millis() as usize;
925        if duration >= pool.interval {
926            //开始整理外部任务队列和内部任务队列的任务数量,并更新权重
927            let new_external_traffic_statistics = pool
928                .external_produce
929                .load(Ordering::Relaxed);
930            let new_internal_traffic_statistics = pool
931                .internal_produce
932                .load(Ordering::Relaxed);
933
934            //获取外部任务增量和内部任务增量
935            let external_delta = if new_external_traffic_statistics == 0 {
936                //上次整理到本次整理之间,外部任务数量为空,则增量为1
937                1
938            } else {
939                //上次整理到本次整理之间,外部任务数量不为空,则计算两次整理之间的外部任务数量的增量
940                new_external_traffic_statistics
941                    .checked_sub(pool
942                        .external_traffic_statistics
943                        .load(Ordering::Relaxed))
944                    .unwrap_or(1)
945            };
946            pool
947                .external_traffic_statistics
948                .store(new_external_traffic_statistics, Ordering::Relaxed); //更新外部任务流量统计
949            let internal_delta = if new_internal_traffic_statistics == 0 {
950                //上次整理到本次整理之间,内部任务数量为空,则增量为1
951                1
952            } else {
953                //上次整理到本次整理之间,内部任务数量不为空,则计算两次整理之间的内部任务数量的增量
954                new_internal_traffic_statistics
955                    .checked_sub(pool
956                        .internal_traffic_statistics
957                        .load(Ordering::Relaxed))
958                    .unwrap_or(1)
959            };
960            pool
961                .internal_traffic_statistics
962                .store(new_internal_traffic_statistics, Ordering::Relaxed); //更新内部任务流量统计
963
964            //更新外部任务队列和内部任务队列的权重
965            let selector = &mut *local_worker.selector.get();
966            if external_delta > internal_delta {
967                //内部任务增量较小
968                let msb = get_msb(internal_delta);
969                let internal_weight
970                    = (internal_delta >> msb.checked_sub(2).unwrap_or(0)).max(1);
971                let external_weight
972                    = ((external_delta >> msb).min(DEFAULT_MAX_WEIGHT as usize)).max(1);
973
974                selector.change_weight(0, external_weight as u8);
975                selector.change_weight(1, internal_weight as u8);
976            } else if external_delta < internal_delta {
977                //外部任务增量较小
978                let msb = get_msb(external_delta);
979                let external_weight
980                    = (external_delta >> msb.checked_sub(2).unwrap_or(0)).max(1);
981                let internal_weight
982                    = ((internal_delta >> msb).min(DEFAULT_MAX_WEIGHT as usize)).max(1);
983
984                selector.change_weight(0, external_weight as u8);
985                selector.change_weight(1, internal_weight as u8);
986            } else {
987                //外部任务和内部任务增量相同
988                selector.change_weight(0, 1);
989                selector.change_weight(1, 1);
990            }
991
992            pool.last_time.store(pool.clock.recent()); //线程安全地更新上一次整理时间
993        }
994
995        //根据权重选择从指定的任务队列弹出任务
996        match (&mut *local_worker.selector.get()).select() {
997            0 => {
998                //弹出外部任务
999                let task = try_pop_external(pool, local_worker, local_worker_id);
1000                if task.is_some() {
1001                    task
1002                } else {
1003                    //当前没有外部任务,则尝试弹出内部任务
1004                    try_pop_internal(pool, local_worker, local_worker_id)
1005                }
1006            },
1007            _ => {
1008                //弹出内部任务
1009                let task = try_pop_internal(pool, local_worker, local_worker_id);
1010                if task.is_some() {
1011                    task
1012                } else {
1013                    //当前没有内部任务,则尝试弹出外部任务
1014                    try_pop_external(pool, local_worker, local_worker_id)
1015                }
1016            },
1017        }
1018    }
1019}
1020
1021// 尝试弹出内部任务队列的任务
1022#[inline]
1023fn try_pop_internal<O: Default + 'static>(pool: &StealableTaskPool<O>,
1024                                          local_worker: &StealableTaskQueue<O>,
1025                                          local_worker_id: usize)
1026    -> Option<Arc<AsyncTask<StealableTaskPool<O>, O>>> {
1027    let task = local_worker
1028        .internal
1029        .pop();
1030    if task.is_some() {
1031        //如果工作者有内部任务,则立即返回
1032        pool
1033            .internal_consume
1034            .fetch_add(1, Ordering::Relaxed);
1035        task
1036    } else {
1037        //工作者的内部任务队列为空,则随机从其它工作者的内部任务队列中窃取任务
1038        let mut gen = thread_rng();
1039        let mut worker_stealers: Vec<&FIFOStealer<Arc<AsyncTask<StealableTaskPool<O>, O>>>> = pool
1040            .internal_stealers
1041            .iter()
1042            .enumerate()
1043            .filter_map(|(index, other)| {
1044                if index != local_worker_id {
1045                    Some(other)
1046                } else {
1047                    //忽略本地工作者
1048                    None
1049                }
1050            })
1051            .collect();
1052
1053        let remaining_len = local_worker.remaining_internal_capacity();
1054        loop {
1055            //随机窃取其它工作者的任务队列
1056            if worker_stealers.len() == 0 {
1057                //所有其它工作者的任务队列都为空,则返回空
1058                break;
1059            }
1060
1061            let index = gen.gen_range(0..worker_stealers.len());
1062            let worker_stealer = worker_stealers.swap_remove(index);
1063
1064            match worker_stealer.steal_and_pop(&local_worker.internal,
1065                                               |count| {
1066                                                   let stealable_len = count / 2;
1067                                                   if stealable_len <= remaining_len {
1068                                                       //当前工作者内部任务队列的剩余容量足够,则窃取指定的其它工作者的内部任务队列中一半的任务
1069                                                       if stealable_len == 0 {
1070                                                           1
1071                                                       } else {
1072                                                           stealable_len
1073                                                       }
1074                                                   } else {
1075                                                       //当前工作者内部任务队列的剩余容量不足够,则从指定的其它工作者的内部任务队列中窃取当前工作者内部任务队列剩余容量的任务
1076                                                       remaining_len
1077                                                   }
1078                                               }) {
1079                Err(StealError::Empty) => {
1080                    //指定的其它工作者的内部任务队列中没有可窃取的任务,则继续窃取下一个其它工作者的内部任务队列
1081                    continue;
1082                },
1083                Err(StealError::Busy) => {
1084                    //需要重试窃取指定的其它工作者的内部任务队列中的任务
1085                    continue;
1086                },
1087                Ok((task, _)) => {
1088                    //从从已窃取到的其它工作者内部任务中获取到首个任务,并立即返回
1089                    pool.internal_consume.fetch_add(1, Ordering::Relaxed);
1090                    return Some(task);
1091                },
1092            }
1093        }
1094
1095        None
1096    }
1097}
1098
1099// 尝试弹出外部任务队列的任务
1100#[inline]
1101fn try_pop_external<O: Default + 'static>(pool: &StealableTaskPool<O>,
1102                                          local_worker: &StealableTaskQueue<O>,
1103                                          local_worker_id: usize)
1104    -> Option<Arc<AsyncTask<StealableTaskPool<O>, O>>> {
1105    let task = local_worker
1106        .external
1107        .pop();
1108    if task.is_some() {
1109        //如果工作者有外部任务,则立即返回
1110        pool
1111            .external_consume
1112            .fetch_add(1, Ordering::Relaxed);
1113        task
1114    } else {
1115        //工作者的外部任务队列为空,则从公共任务池中弹出任务
1116        let task = try_pop_public(pool, local_worker);
1117        if task.is_some() {
1118            //如果公共任务池有外部任务,则立即返回
1119            pool
1120                .external_consume
1121                .fetch_add(1, Ordering::Relaxed);
1122            task
1123        } else {
1124            //公共任务池为空,则随机从其它工作者的外部任务队列中窃取任务
1125            let mut gen = thread_rng();
1126            let mut worker_stealers: Vec<&Stealer<Arc<AsyncTask<StealableTaskPool<O>, O>>>> = pool
1127                .external_stealers
1128                .iter()
1129                .enumerate()
1130                .filter_map(|(index, other)| {
1131                    if index != local_worker_id {
1132                        Some(other)
1133                    } else {
1134                        //忽略当前工作者
1135                        None
1136                    }
1137                })
1138                .collect();
1139
1140            loop {
1141                //随机窃取其它工作者的任务队列
1142                if worker_stealers.len() == 0 {
1143                    //所有其它工作者的外部任务队列都为空,则返回空
1144                    break;
1145                }
1146
1147                let index = gen.gen_range(0..worker_stealers.len());
1148                let worker_stealer = worker_stealers.swap_remove(index);
1149
1150                match worker_stealer.steal_batch_and_pop(&local_worker.external) {
1151                    Steal::Success(task) => {
1152                        //从从已窃取到的其它工作者外部任务中获取到首个任务,并立即返回
1153                        pool.external_consume.fetch_add(1, Ordering::Relaxed);
1154                        return Some(task);
1155                    },
1156                    Steal::Retry => {
1157                        //需要重试窃取指定的其它工作者的外部任务队列中的任务
1158                        continue;
1159                    },
1160                    Steal::Empty => {
1161                        //指定的其它工作者的外部任务队列中没有可窃取的任务,则继续窃取下一个其它工作者的外部任务队列
1162                        continue;
1163                    },
1164                }
1165            }
1166
1167            None
1168        }
1169    }
1170}
1171
1172// 尝试弹出公共任务池的任务
1173#[inline]
1174fn try_pop_public<O: Default + 'static>(pool: &StealableTaskPool<O>,
1175                                        local_worker: &StealableTaskQueue<O>)
1176    -> Option<Arc<AsyncTask<StealableTaskPool<O>, O>>> {
1177    loop {
1178        match pool.public.steal_batch_and_pop(&local_worker.external) {
1179            Steal::Empty => {
1180                //当前公共任务池没有任务
1181                return None;
1182            },
1183            Steal::Retry => {
1184                //需要重试窃取公共任务池的任务
1185                continue;
1186            },
1187            Steal::Success(task) => {
1188                //从已窃取到的公共任务中获取到首个任务,并立即返回
1189                pool.external_consume.fetch_add(1, Ordering::Relaxed);
1190                return Some(task);
1191            },
1192        }
1193    }
1194}
1195
1196impl<O: Default + 'static> AsyncTaskPoolExt<O> for StealableTaskPool<O> {
1197    #[inline]
1198    fn set_waits(&mut self, waits: Arc<ArrayQueue<Arc<(AtomicBool, Mutex<()>, Condvar)>>>) {
1199        self.waits = Some(waits);
1200    }
1201
1202    #[inline]
1203    fn get_waits(&self) -> Option<&Arc<ArrayQueue<Arc<(AtomicBool, Mutex<()>, Condvar)>>>> {
1204        self.waits.as_ref()
1205    }
1206
1207    #[inline]
1208    fn worker_len(&self) -> usize {
1209        self.workers.len()
1210    }
1211
1212    /// clone 当前 exact pool/worker 的休眠唤醒器。
1213    ///
1214    /// 合法 owner 返回 `Some(Arc<...>)`;wrong-pool 或未绑定线程在 worker lookup 前 panic。
1215    /// O(1),按既有语义执行一次 `Arc::clone`,owner guard 自身不 clone/分配/锁/阻塞,也不会
1216    /// notify 或产生错误唤醒。查询本身幂等但增加引用计数;调用方负责释放返回 Arc。线程、
1217    /// 内存和异步安全,不执行用户代码或 V8/FFI。
1218    #[inline]
1219    fn clone_thread_waker(&self) -> Option<Arc<(AtomicBool, Mutex<()>, Condvar)>> {
1220        let id = current_multi_thread_worker_context().owner_worker_id(self);
1221        if let Some(worker) = self.workers.get(id) {
1222            return Some(worker.thread_waker.clone());
1223        }
1224
1225        None
1226    }
1227}
1228
1229impl<O: Default + 'static> StealableTaskPool<O> {
1230    /// 使用平台默认 worker slot 数构建可窃取任务池。
1231    ///
1232    /// 非 wasm32 使用物理核数的两倍,wasm32 使用 1;内部 queue capacity、初始权重参数和刷新
1233    /// interval 保持既有默认值。返回值尚未启动线程或绑定 runtime。
1234    ///
1235    /// 时间/空间复杂度 O(W),会分配 W 组 queue/stealer/waker,因此不是纯函数且非幂等;不
1236    /// 阻塞、不执行 I/O 或 future。owner、安全和错误边界见 [`StealableTaskPool`] 类型文档。
1237    pub fn new() -> Self {
1238        #[cfg(not(target_arch = "wasm32"))]
1239            let size = num_cpus::get_physical() * 2; //默认最大工作者任务池数量是当前cpu物理核的2倍
1240        #[cfg(target_arch = "wasm32")]
1241            let size = 1; //默认最大工作者任务池数量是1
1242        StealableTaskPool::with(size,
1243                                0x8000,
1244                                [1, 1],
1245                                3000)
1246    }
1247
1248    /// 构建指定 worker slot、internal queue capacity、权重参数和刷新间隔的任务池。
1249    ///
1250    /// # 参数
1251    ///
1252    /// - `worker_size`:worker slot 数,必须大于 0;为 0 时保持旧行为并 panic。
1253    /// - `internal_queue_capacity`:每个 worker internal FIFO 的初始容量;允许 0,由底层 queue
1254    ///   按其既有规则归一化。值只在构建期消费。
1255    /// - `weights`:保留的两类队列权重配置。当前版本保持历史存储/调度行为,不新增“所有
1256    ///   selector 以该值初始化”或“全局收敛”保证;调用方不得据此假设硬公平性。
1257    /// - `interval`:近似流量统计刷新间隔,单位 ms,必须大于 0;为 0 时 panic。
1258    ///
1259    /// # 返回与副作用
1260    ///
1261    /// 返回未绑定 runtime 的独立 pool,持有 W 组 worker queue/stealer/waker 和一个 pool-wide
1262    /// 原子刷新时间。函数会分配内存,不创建线程、不执行用户 future、不进行 I/O。
1263    ///
1264    /// # 性能、幂等与安全
1265    ///
1266    /// 构建时间/空间 O(W);不是纯函数且每次创建不同资源,非幂等。运行时 owner 限制、panic
1267    /// 边界、线程/异步/内存安全和跨 runtime fallback 见类型文档。专项入口为
1268    /// `tests/stealable_task_pool_concurrency.rs`。
1269    pub fn with(worker_size: usize,
1270                internal_queue_capacity: usize,
1271                weights: [u8; 2],
1272                interval: usize) -> Self {
1273        if worker_size == 0 {
1274            //工作者任务池数量无效,则立即抛出异常
1275            panic!(
1276                "Create WorkerTaskPool failed, worker size: {}, reason: invalid worker size",
1277                worker_size
1278            );
1279        }
1280        if interval == 0 {
1281            panic!(
1282                "Create WorkerTaskPool failed, interval: {}, reason: invalid interval",
1283                worker_size
1284            );
1285        }
1286
1287        let public = Injector::new();
1288        let mut workers = Vec::with_capacity(worker_size);
1289        let mut internal_stealers = Vec::with_capacity(worker_size);
1290        let mut external_stealers = Vec::with_capacity(worker_size);
1291        for _ in 0..worker_size {
1292            //初始化指定初始作者任务池数量的工作者任务池和窃取者
1293            let thread_waker = Arc::new((AtomicBool::new(false), Mutex::new(()), Condvar::new()));
1294            let (worker,
1295                internal_stealer,
1296                external_stealer) =
1297                StealableTaskQueue::new(internal_queue_capacity,
1298                                        thread_waker);
1299            workers.push(worker);
1300            internal_stealers.push(internal_stealer);
1301            external_stealers.push(external_stealer);
1302        }
1303        let internal_consume = AtomicUsize::new(0);
1304        let internal_produce = AtomicUsize::new(0);
1305        let internal_traffic_statistics = AtomicUsize::new(0);
1306        let external_consume = AtomicUsize::new(0);
1307        let external_produce = AtomicUsize::new(0);
1308        let external_traffic_statistics = AtomicUsize::new(0);
1309        let clock = Clock::new();
1310        let last_time = AtomicCell::new(clock.recent());
1311
1312        StealableTaskPool {
1313            public,
1314            workers,
1315            internal_stealers,
1316            external_stealers,
1317            internal_consume,
1318            internal_produce,
1319            internal_traffic_statistics,
1320            external_consume,
1321            external_produce,
1322            external_traffic_statistics,
1323            weights,
1324            clock,
1325            interval,
1326            last_time,
1327            waits: None,
1328        }
1329    }
1330}
1331
1332///
1333/// 异步多线程任务运行时,支持运行时线程伸缩
1334///
1335pub struct MultiTaskRuntime<
1336    O: Default + 'static = (),
1337    P: AsyncTaskPoolExt<O> + AsyncTaskPool<O> = StealableTaskPool<O>,
1338>(
1339    Arc<(
1340        usize,                                                  //运行时唯一id
1341        Arc<P>,                                                 //异步任务池
1342        Option<
1343            Vec<(
1344                Sender<(usize, AsyncTimingTask<P, O>)>,
1345                Arc<AsyncTaskTimerByNotCancel<P, O>>,
1346            )>,
1347        >,                                                      //休眠的异步任务生产者和本地定时器
1348        AtomicUsize,                                            //定时任务计数器
1349        Arc<ArrayQueue<Arc<(AtomicBool, Mutex<()>, Condvar)>>>, //待唤醒的工作者唤醒器队列
1350        AtomicUsize,                                            //定时器生产计数
1351        AtomicUsize,                                            //定时器消费计数
1352    )>,
1353);
1354
1355unsafe impl<O: Default + 'static, P: AsyncTaskPoolExt<O> + AsyncTaskPool<O>> Send
1356    for MultiTaskRuntime<O, P>
1357{
1358}
1359unsafe impl<O: Default + 'static, P: AsyncTaskPoolExt<O> + AsyncTaskPool<O>> Sync
1360    for MultiTaskRuntime<O, P>
1361{
1362}
1363
1364impl<O: Default + 'static, P: AsyncTaskPoolExt<O> + AsyncTaskPool<O>> Clone
1365    for MultiTaskRuntime<O, P>
1366{
1367    fn clone(&self) -> Self {
1368        MultiTaskRuntime(self.0.clone())
1369    }
1370}
1371
1372impl<O: Default + 'static, P: AsyncTaskPoolExt<O> + AsyncTaskPool<O, Pool = P>> AsyncRuntime<O>
1373    for MultiTaskRuntime<O, P>
1374{
1375    type Pool = P;
1376
1377    /// 共享运行时内部任务池
1378    fn shared_pool(&self) -> Arc<Self::Pool> {
1379        (self.0).1.clone()
1380    }
1381
1382    /// 获取当前异步运行时的唯一id
1383    fn get_id(&self) -> usize {
1384        (self.0).0
1385    }
1386
1387    /// 获取当前异步运行时待处理任务数量
1388    fn wait_len(&self) -> usize {
1389        (self.0)
1390            .5
1391            .load(Ordering::Relaxed)
1392            .checked_sub((self.0).6.load(Ordering::Relaxed))
1393            .unwrap_or(0)
1394    }
1395
1396    /// 获取当前异步运行时任务数量
1397    fn len(&self) -> usize {
1398        (self.0).1.len()
1399    }
1400
1401    /// 分配异步任务的唯一id
1402    fn alloc<R: 'static>(&self) -> TaskId {
1403        TaskId(UnsafeCell::new((TaskHandle::<R>::default().into_raw() as u128) << 64 | self.get_id() as u128 & 0xffffffffffffffff))
1404    }
1405
1406    /// 派发一个指定的异步任务到异步运行时
1407    fn spawn<F>(&self, future: F) -> Result<TaskId>
1408    where
1409        F: Future<Output = O> + Send + 'static,
1410    {
1411        let task_id = self.alloc::<F::Output>();
1412        if let Err(e) = self.spawn_by_id(task_id.clone(), future) {
1413            return Err(e);
1414        }
1415
1416        Ok(task_id)
1417    }
1418
1419    /// 派发一个异步任务到本地异步运行时,如果本地没有本异步运行时,则会派发到当前运行时中
1420    fn spawn_local<F>(&self, future: F) -> Result<TaskId>
1421        where
1422            F: Future<Output=O> + Send + 'static {
1423        let task_id = self.alloc::<F::Output>();
1424        if let Err(e) = self.spawn_local_by_id(task_id.clone(), future) {
1425            return Err(e);
1426        }
1427
1428        Ok(task_id)
1429    }
1430
1431    /// 派发一个指定优先级的异步任务到异步运行时
1432    fn spawn_priority<F>(&self, priority: usize, future: F) -> Result<TaskId>
1433        where
1434            F: Future<Output=O> + Send + 'static {
1435        let task_id = self.alloc::<F::Output>();
1436        if let Err(e) = self.spawn_priority_by_id(task_id.clone(), priority, future) {
1437            return Err(e);
1438        }
1439
1440        Ok(task_id)
1441    }
1442
1443    /// 派发一个异步任务到异步运行时,并立即让出任务的当前运行
1444    fn spawn_yield<F>(&self, future: F) -> Result<TaskId>
1445        where
1446            F: Future<Output=O> + Send + 'static {
1447        let task_id = self.alloc::<F::Output>();
1448        if let Err(e) = self.spawn_yield_by_id(task_id.clone(), future) {
1449            return Err(e);
1450        }
1451
1452        Ok(task_id)
1453    }
1454
1455    /// 派发一个在指定时间后执行的异步任务到异步运行时,时间单位ms
1456    fn spawn_timing<F>(&self, future: F, time: usize) -> Result<TaskId>
1457    where
1458        F: Future<Output = O> + Send + 'static,
1459    {
1460        let task_id = self.alloc::<F::Output>();
1461        if let Err(e) = self.spawn_timing_by_id(task_id.clone(), future, time) {
1462            return Err(e);
1463        }
1464
1465        Ok(task_id)
1466    }
1467
1468    /// 派发一个指定任务唯一id的异步任务到异步运行时
1469    fn spawn_by_id<F>(&self, task_id: TaskId, future: F) -> Result<()>
1470        where
1471            F: Future<Output=O> + Send + 'static {
1472        let result = {
1473            (self.0).1.push(Arc::new(AsyncTask::new(
1474                task_id,
1475                (self.0).1.clone(),
1476                DEFAULT_MAX_LOW_PRIORITY_BOUNDED,
1477                Some(future.boxed()),
1478            )))
1479        };
1480
1481        let _ = wake_waiting_worker(&(self.0).4);
1482
1483        result
1484    }
1485
1486    fn spawn_local_by_id<F>(&self, task_id: TaskId, future: F) -> Result<()>
1487        where
1488            F: Future<Output=O> + Send + 'static {
1489        let should_wake = PI_ASYNC_THREAD_LOCAL_ID
1490            .try_with(|thread_id| unsafe { ((*thread_id.get()) >> 32) != self.get_id() })
1491            .unwrap_or(true);
1492        let result = (self.0).1.push_local(Arc::new(AsyncTask::new(
1493            task_id,
1494            (self.0).1.clone(),
1495            DEFAULT_HIGH_PRIORITY_BOUNDED,
1496            Some(future.boxed()),
1497        )));
1498
1499        if should_wake {
1500            let _ = wake_waiting_worker(&(self.0).4);
1501        }
1502
1503        result
1504    }
1505
1506    /// 派发一个指定任务唯一id和任务优先级的异步任务到异步运行时
1507    fn spawn_priority_by_id<F>(&self,
1508                               task_id: TaskId,
1509                               priority: usize,
1510                               future: F) -> Result<()>
1511        where
1512            F: Future<Output=O> + Send + 'static {
1513        let result = {
1514            (self.0).1.push_priority(priority, Arc::new(AsyncTask::new(
1515                task_id,
1516                (self.0).1.clone(),
1517                priority,
1518                Some(future.boxed()),
1519            )))
1520        };
1521
1522        let _ = wake_waiting_worker(&(self.0).4);
1523
1524        result
1525    }
1526
1527    /// 派发一个指定任务唯一id的异步任务到异步运行时,并立即让出任务的当前运行
1528    #[inline]
1529    fn spawn_yield_by_id<F>(&self, task_id: TaskId, future: F) -> Result<()>
1530        where
1531            F: Future<Output=O> + Send + 'static {
1532        self.spawn_priority_by_id(task_id,
1533                                  DEFAULT_HIGH_PRIORITY_BOUNDED,
1534                                  future)
1535    }
1536
1537    /// 派发一个指定任务唯一id和在指定时间后执行的异步任务到异步运行时,时间单位ms
1538    fn spawn_timing_by_id<F>(&self,
1539                             task_id: TaskId,
1540                             future: F,
1541                             time: usize) -> Result<()>
1542        where
1543            F: Future<Output=O> + Send + 'static {
1544        let rt = self.clone();
1545        self.spawn_by_id(task_id, async move {
1546            if let Some(timers) = &(rt.0).2 {
1547                //为定时器设置定时异步任务
1548                let id = (rt.0).1.get_thread_id() & 0xffffffff;
1549                let (_, timer) = &timers[id];
1550                timer.set_timer(
1551                    AsyncTimingTask::WaitRun(Arc::new(AsyncTask::new(
1552                        rt.alloc::<F::Output>(),
1553                        (rt.0).1.clone(),
1554                        DEFAULT_MAX_HIGH_PRIORITY_BOUNDED,
1555                        Some(future.boxed()),
1556                    ))),
1557                    time,
1558                );
1559
1560                (rt.0).5.fetch_add(1, Ordering::Relaxed);
1561            }
1562
1563            Default::default()
1564        })
1565    }
1566
1567    /// 挂起指定唯一id的异步任务
1568    fn pending<Output: 'static>(&self, task_id: &TaskId, waker: Waker) -> Poll<Output> {
1569        task_id.set_waker::<Output>(waker);
1570        Poll::Pending
1571    }
1572
1573    /// 唤醒指定唯一id的异步任务
1574    fn wakeup<Output: 'static>(&self, task_id: &TaskId) {
1575        task_id.wakeup::<Output>();
1576    }
1577
1578    /// 挂起当前异步运行时的当前任务,并在指定的其它运行时上派发一个指定的异步任务,等待其它运行时上的异步任务完成后,唤醒当前运行时的当前任务,并返回其它运行时上的异步任务的值
1579    fn wait<V: Send + 'static>(&self) -> AsyncWait<V> {
1580        AsyncWait(self.wait_any(2))
1581    }
1582
1583    /// 挂起当前异步运行时的当前任务,并在多个其它运行时上执行多个其它任务,其中任意一个任务完成,则唤醒当前运行时的当前任务,并返回这个已完成任务的值,而其它未完成的任务的值将被忽略
1584    fn wait_any<V: Send + 'static>(&self, capacity: usize) -> AsyncWaitAny<V> {
1585        let (producor, consumer) = async_bounded(capacity);
1586
1587        AsyncWaitAny {
1588            capacity,
1589            producor,
1590            consumer,
1591        }
1592    }
1593
1594    /// 挂起当前异步运行时的当前任务,并在多个其它运行时上执行多个其它任务,任务返回后需要通过用户指定的检查回调进行检查,其中任意一个任务检查通过,则唤醒当前运行时的当前任务,并返回这个已完成任务的值,而其它未完成或未检查通过的任务的值将被忽略,如果所有任务都未检查通过,则强制唤醒当前运行时的当前任务
1595    fn wait_any_callback<V: Send + 'static>(&self, capacity: usize) -> AsyncWaitAnyCallback<V> {
1596        let (producor, consumer) = async_bounded(capacity);
1597
1598        AsyncWaitAnyCallback {
1599            capacity,
1600            producor,
1601            consumer,
1602        }
1603    }
1604
1605    /// 构建用于派发多个异步任务到指定运行时的映射归并,需要指定映射归并的容量
1606    fn map_reduce<V: Send + 'static>(&self, capacity: usize) -> AsyncMapReduce<V> {
1607        let (producor, consumer) = async_bounded(capacity);
1608
1609        AsyncMapReduce {
1610            count: 0,
1611            capacity,
1612            producor,
1613            consumer,
1614        }
1615    }
1616
1617    /// 挂起当前异步运行时的当前任务,等待指定的时间后唤醒当前任务
1618    fn timeout(&self, timeout: usize) -> BoxFuture<'static, ()> {
1619        let rt = self.clone();
1620
1621        if let Some(timers) = &(self.0).2 {
1622            //有本地定时器,则异步等待指定时间
1623            match PI_ASYNC_THREAD_LOCAL_ID.try_with(move |thread_id| {
1624                //将休眠的异步任务投递到当前派发线程的定时器内
1625                let thread_id = unsafe { *thread_id.get() };
1626                let index = thread_id & 0xffffffff;
1627                if index > timers.len() {
1628                    //当前线程还未初始化运行时的线程id,说明当前线程不是当前多线程运行时的所属线程
1629                    TimerTaskProducor::Foreign(timers[(self.0).3.load(Ordering::Relaxed) % timers.len()].0.clone())
1630                } else {
1631                    TimerTaskProducor::Local(timers[index].1.clone())
1632                }
1633            }) {
1634                Err(_) => {
1635                    panic!("Multi thread runtime timeout failed, reason: local thread id not match")
1636                }
1637                Ok(producor) => match producor {
1638                    TimerTaskProducor::Local(timer) => {
1639                        LocalAsyncWaitTimeout::new(rt, timer, timeout).boxed()
1640                    },
1641                    TimerTaskProducor::Foreign(producor) => {
1642                        AsyncWaitTimeout::new(rt, producor, timeout).boxed()
1643                    },
1644                },
1645            }
1646        } else {
1647            //没有本地定时器,则同步休眠指定时间
1648            async move {
1649                thread::sleep(Duration::from_millis(timeout as u64));
1650            }
1651            .boxed()
1652        }
1653    }
1654
1655    /// 立即让出当前任务的执行
1656    fn yield_now(&self) -> BoxFuture<'static, ()> {
1657        async move {
1658            YieldNow(false).await;
1659        }.boxed()
1660    }
1661
1662    /// 生成一个异步管道,输入指定流,输入流的每个值通过过滤器生成输出流的值
1663    fn pipeline<S, SO, F, FO>(&self, input: S, mut filter: F) -> BoxStream<'static, FO>
1664    where
1665        S: Stream<Item = SO> + Send + 'static,
1666        SO: Send + 'static,
1667        F: FnMut(SO) -> AsyncPipelineResult<FO> + Send + 'static,
1668        FO: Send + 'static,
1669    {
1670        let output = stream! {
1671            for await value in input {
1672                match filter(value) {
1673                    AsyncPipelineResult::Disconnect => {
1674                        //立即中止管道
1675                        break;
1676                    },
1677                    AsyncPipelineResult::Filtered(result) => {
1678                        yield result;
1679                    },
1680                }
1681            }
1682        };
1683
1684        output.boxed()
1685    }
1686
1687    /// 关闭异步运行时,返回请求关闭是否成功
1688    fn close(&self) -> bool {
1689        false
1690    }
1691}
1692
1693impl<O: Default + 'static, P: AsyncTaskPoolExt<O> + AsyncTaskPool<O, Pool = P>> AsyncRuntimeExt<O>
1694    for MultiTaskRuntime<O, P>
1695{
1696    fn spawn_with_context<F, C>(&self, task_id: TaskId, future: F, context: C) -> Result<()>
1697    where
1698        F: Future<Output = O> + Send + 'static,
1699        C: 'static,
1700    {
1701        let task = Arc::new(AsyncTask::with_context(
1702            task_id,
1703            (self.0).1.clone(),
1704            DEFAULT_MAX_LOW_PRIORITY_BOUNDED,
1705            Some(future.boxed()),
1706            context,
1707        ));
1708        let result = (self.0).1.push(task);
1709
1710        let _ = wake_waiting_worker(&(self.0).4);
1711
1712        result
1713    }
1714
1715    fn spawn_timing_with_context<F, C>(
1716        &self,
1717        task_id: TaskId,
1718        future: F,
1719        context: C,
1720        time: usize,
1721    ) -> Result<()>
1722    where
1723        F: Future<Output = O> + Send + 'static,
1724        C: Send + 'static,
1725    {
1726        let rt = self.clone();
1727        self.spawn_by_id(task_id, async move {
1728            if let Some(timers) = &(rt.0).2 {
1729                //为定时器设置定时异步任务
1730                let id = (rt.0).1.get_thread_id() & 0xffffffff;
1731                let (_, timer) = &timers[id];
1732                timer.set_timer(
1733                    AsyncTimingTask::WaitRun(Arc::new(AsyncTask::with_context(
1734                        rt.alloc::<F::Output>(),
1735                        (rt.0).1.clone(),
1736                        DEFAULT_MAX_HIGH_PRIORITY_BOUNDED,
1737                        Some(future.boxed()),
1738                        context,
1739                    ))),
1740                    time,
1741                );
1742
1743                (rt.0).5.fetch_add(1, Ordering::Relaxed);
1744            }
1745
1746            Default::default()
1747        })
1748    }
1749
1750    fn block_on<F>(&self, future: F) -> Result<F::Output>
1751    where
1752        F: Future + Send + 'static,
1753        <F as Future>::Output: Default + Send + 'static,
1754    {
1755        //从本地线程获取当前异步运行时
1756        if let Some(local_rt) = local_async_runtime::<F::Output>() {
1757            //本地线程绑定了异步运行时
1758            if local_rt.get_id() == self.get_id() {
1759                //如果是相同运行时,则立即返回错误
1760                return Err(Error::new(
1761                    ErrorKind::WouldBlock,
1762                    format!("Block on failed, reason: would block"),
1763                ));
1764            }
1765        }
1766
1767        let (sender, receiver) = bounded(1);
1768        if let Err(e) = self.spawn(async move {
1769            //在指定运行时中执行,并返回结果
1770            let r = future.await;
1771            sender.send(r);
1772
1773            Default::default()
1774        }) {
1775            return Err(Error::new(
1776                ErrorKind::Other,
1777                format!("Block on failed, reason: {:?}", e),
1778            ));
1779        }
1780
1781        //同步阻塞等待异步任务返回
1782        match receiver.recv() {
1783            Err(e) => Err(Error::new(
1784                ErrorKind::Other,
1785                format!("Block on failed, reason: {:?}", e),
1786            )),
1787            Ok(result) => Ok(result),
1788        }
1789    }
1790}
1791
1792impl<O: Default + 'static, P: AsyncTaskPoolExt<O> + AsyncTaskPool<O, Pool = P>>
1793    MultiTaskRuntime<O, P>
1794{
1795    /// 获取当前运行时可新增的工作者数量
1796    pub fn idler_len(&self) -> usize {
1797        (self.0).1.idler_len()
1798    }
1799
1800    /// 获取当前运行时的工作者数量
1801    pub fn worker_len(&self) -> usize {
1802        (self.0).1.worker_len()
1803    }
1804
1805    /// 获取当前运行时缓冲区的任务数量,缓冲区的任务暂时没有分配给工作者
1806    pub fn buffer_len(&self) -> usize {
1807        (self.0).1.buffer_len()
1808    }
1809
1810    /// 获取当前多线程异步运行时的本地异步运行时
1811    pub fn to_local_runtime(&self) -> LocalAsyncRuntime<O> {
1812        LocalAsyncRuntime {
1813            inner: self.as_raw(),
1814            get_id_func: MultiTaskRuntime::<O, P>::get_id_raw,
1815            spawn_func: MultiTaskRuntime::<O, P>::spawn_raw,
1816            spawn_local_func: MultiTaskRuntime::<O, P>::spawn_local_raw,
1817            spawn_timing_func: MultiTaskRuntime::<O, P>::spawn_timing_raw,
1818            timeout_func: MultiTaskRuntime::<O, P>::timeout_raw,
1819        }
1820    }
1821
1822    /// 获取当前多线程异步运行时的指针
1823    #[inline]
1824    pub(crate) fn as_raw(&self) -> *const () {
1825        Arc::into_raw(self.0.clone()) as *const ()
1826    }
1827
1828    // 获取指定指针的单线程异步运行时
1829    #[inline]
1830    pub(crate) fn from_raw(raw: *const ()) -> Self {
1831        let inner = unsafe {
1832            Arc::from_raw(
1833                raw as *const (
1834                    usize,
1835                    Arc<P>,
1836                    Option<
1837                        Vec<(
1838                            Sender<(usize, AsyncTimingTask<P, O>)>,
1839                            Arc<AsyncTaskTimerByNotCancel<P, O>>,
1840                        )>,
1841                    >,
1842                    AtomicUsize,
1843                    Arc<ArrayQueue<Arc<(AtomicBool, Mutex<()>, Condvar)>>>,
1844                    AtomicUsize,
1845                    AtomicUsize,
1846                ),
1847            )
1848        };
1849        MultiTaskRuntime(inner)
1850    }
1851
1852    // 获取当前异步运行时的唯一id
1853    pub(crate) fn get_id_raw(raw: *const ()) -> usize {
1854        let rt = MultiTaskRuntime::<O, P>::from_raw(raw);
1855        let id = rt.get_id();
1856        Arc::into_raw(rt.0); //避免提前释放
1857        id
1858    }
1859
1860    // 派发一个指定的异步任务到异步运行时
1861    pub(crate) fn spawn_raw<F>(raw: *const (), future: F) -> Result<()>
1862    where
1863        F: Future<Output = O> + Send + 'static,
1864    {
1865        let rt = MultiTaskRuntime::<O, P>::from_raw(raw);
1866        let result = rt.spawn_by_id(rt.alloc::<F::Output>(), future);
1867        Arc::into_raw(rt.0); //避免提前释放
1868        result
1869    }
1870
1871    // 派发一个指定的异步任务到本地异步运行时
1872    pub(crate) fn spawn_local_raw<F>(raw: *const (), future: F) -> Result<()>
1873    where
1874        F: Future<Output = O> + Send + 'static,
1875    {
1876        let rt = MultiTaskRuntime::<O, P>::from_raw(raw);
1877        let result = rt.spawn_local_by_id(rt.alloc::<F::Output>(), future);
1878        Arc::into_raw(rt.0); //避免提前释放
1879        result
1880    }
1881
1882    // 定时派发一个指定的异步任务到异步运行时
1883    pub(crate) fn spawn_timing_raw(
1884        raw: *const (),
1885        future: BoxFuture<'static, O>,
1886        timeout: usize,
1887    ) -> Result<()> {
1888        let rt = MultiTaskRuntime::<O, P>::from_raw(raw);
1889        let result = rt.spawn_timing_by_id(rt.alloc::<O>(), future, timeout);
1890        Arc::into_raw(rt.0); //避免提前释放
1891        result
1892    }
1893
1894    // 挂起当前异步运行时的当前任务,等待指定的时间后唤醒当前任务
1895    pub(crate) fn timeout_raw(raw: *const (), timeout: usize) -> BoxFuture<'static, ()> {
1896        let rt = MultiTaskRuntime::<O, P>::from_raw(raw);
1897        let boxed = rt.timeout(timeout);
1898        Arc::into_raw(rt.0); //避免提前释放
1899        boxed
1900    }
1901}
1902
1903///
1904/// 异步多线程任务运行时构建器
1905///
1906pub struct MultiTaskRuntimeBuilder<
1907    O: Default + 'static = (),
1908    P: AsyncTaskPoolExt<O> + AsyncTaskPool<O> = StealableTaskPool<O>,
1909> {
1910    pool: P,                 //异步多线程任务运行时
1911    prefix: String,          //工作者线程名称前缀
1912    init: usize,             //初始工作者数量
1913    min: usize,              //最少工作者数量
1914    max: usize,              //最大工作者数量
1915    stack_size: usize,       //工作者线程栈大小
1916    timeout: u64,            //工作者空闲时最长休眠时间
1917    interval: Option<usize>, //工作者定时器间隔
1918    marker: PhantomData<O>,
1919}
1920
1921unsafe impl<O: Default + 'static, P: AsyncTaskPoolExt<O> + AsyncTaskPool<O>> Send
1922    for MultiTaskRuntimeBuilder<O, P>
1923{
1924}
1925unsafe impl<O: Default + 'static, P: AsyncTaskPoolExt<O> + AsyncTaskPool<O>> Sync
1926    for MultiTaskRuntimeBuilder<O, P>
1927{
1928}
1929
1930impl<O: Default + 'static> Default for MultiTaskRuntimeBuilder<O> {
1931    //默认构建可窃取可伸缩的多线程运行时
1932    fn default() -> Self {
1933        #[cfg(not(target_arch = "wasm32"))]
1934        let core_len = num_cpus::get(); //默认的工作者的数量为本机逻辑核数
1935        #[cfg(target_arch = "wasm32")]
1936        let core_len = 1; //默认的工作者的数量为1
1937        let pool = StealableTaskPool::with(core_len,
1938                                           65535,
1939                                           [1, 1],
1940                                           3000);
1941        MultiTaskRuntimeBuilder::new(pool)
1942            .thread_stack_size(2 * 1024 * 1024)
1943            .set_timer_interval(1)
1944    }
1945}
1946
1947impl<O: Default + 'static, P: AsyncTaskPoolExt<O> + AsyncTaskPool<O, Pool = P>>
1948    MultiTaskRuntimeBuilder<O, P>
1949{
1950    /// 构建指定任务池、线程名前缀、初始线程数量、最少线程数量、最大线程数量、线程栈大小、线程空闲时最长休眠时间和是否使用本地定时器的多线程任务池
1951    pub fn new(mut pool: P) -> Self {
1952        #[cfg(not(target_arch = "wasm32"))]
1953        let core_len = num_cpus::get(); //获取本机cpu逻辑核数
1954        #[cfg(target_arch = "wasm32")]
1955        let core_len = 1; //默认为1
1956
1957        MultiTaskRuntimeBuilder {
1958            pool,
1959            prefix: DEFAULT_WORKER_THREAD_PREFIX.to_string(),
1960            init: core_len,
1961            min: core_len,
1962            max: core_len,
1963            stack_size: DEFAULT_THREAD_STACK_SIZE,
1964            timeout: DEFAULT_WORKER_THREAD_SLEEP_TIME,
1965            interval: None,
1966            marker: PhantomData,
1967        }
1968    }
1969
1970    /// 设置工作者线程名称前缀
1971    pub fn thread_prefix(mut self, prefix: &str) -> Self {
1972        self.prefix = prefix.to_string();
1973        self
1974    }
1975
1976    /// 设置工作者线程栈大小
1977    pub fn thread_stack_size(mut self, stack_size: usize) -> Self {
1978        self.stack_size = stack_size;
1979        self
1980    }
1981
1982    /// 设置初始工作者数量
1983    pub fn init_worker_size(mut self, mut init: usize) -> Self {
1984        if init == 0 {
1985            //初始线程数量过小,则设置默认的初始线程数量
1986            init = DEFAULT_INIT_WORKER_SIZE;
1987        }
1988
1989        self.init = init;
1990        self
1991    }
1992
1993    /// 设置最小工作者数量和最大工作者数量
1994    pub fn set_worker_limit(mut self, mut min: usize, mut max: usize) -> Self {
1995        if self.init > max {
1996            //初始线程数量大于最大线程数量,则设置最大线程数量为初始线程数量
1997            max = self.init;
1998        }
1999
2000        if min == 0 || min > max {
2001            //最少线程数量无效,则设置最少线程数量为最大线程数量
2002            min = max;
2003        }
2004
2005        self.min = min;
2006        self.max = max;
2007        self
2008    }
2009
2010    /// 设置工作者空闲时最大休眠时长
2011    pub fn set_timeout(mut self, timeout: u64) -> Self {
2012        self.timeout = timeout;
2013        self
2014    }
2015
2016    /// 设置工作者定时器间隔
2017    pub fn set_timer_interval(mut self, interval: usize) -> Self {
2018        self.interval = Some(interval);
2019        self
2020    }
2021
2022    /// 构建并启动多线程异步运行时。
2023    ///
2024    /// 说明:
2025    /// - 该函数消费 builder,创建 runtime、定时器、waiting worker 队列,并启动初始
2026    ///   worker 线程。
2027    /// - 本轮保持公开 API 和启动流程不变,只在构建期增加 worker 数边界收敛,并确保
2028    ///   任务池保存 runtime 共享 waits 队列。
2029    ///
2030    /// 入参:
2031    /// - 使用 builder 中已经配置好的 pool、线程名前缀、栈大小、worker 数、sleep timeout
2032    ///   和 timer interval。
2033    ///
2034    /// 返回:
2035    /// - 已启动的 `MultiTaskRuntime<O, P>`。
2036    ///
2037    /// 边界条件:
2038    /// - 如果 pool 的 `worker_len()` 为 0,立即 panic;有效任务池不允许没有 worker slot。
2039    /// - 如果 `init/max` 大于 pool worker slot 数,会收敛到 `pool.worker_len()`。
2040    /// - 如果收敛后 `min > max`,会把 `min` 收敛到 `max`。
2041    /// - 上述收敛只避免内部 worker slot 越界,不改变已存在的公开方法签名。
2042    ///
2043    /// 性能:
2044    /// - 构建时间 O(W),空间 O(W),W 为最终 `max` worker 数。
2045    /// - 该函数不是任务调度热路径。
2046    ///
2047    /// 副作用:
2048    /// - 非纯函数,会分配 runtime 内部结构、注入 waits 队列、启动 worker 线程。
2049    /// - 不执行用户 future;worker 启动后由工作循环正常消费任务。
2050    ///
2051    /// 安全性:
2052    /// - 不引入新的 unsafe。
2053    /// - 线程安全依赖 `AsyncTaskPoolExt::set_waits` 在 pool 被放入 `Arc` 前完成,之后 waits
2054    ///   通过 `Arc<ArrayQueue<...>>` 在线程间共享。
2055    pub fn build(mut self) -> MultiTaskRuntime<O, P> {
2056        let pool_worker_len = self.pool.worker_len();
2057        if pool_worker_len == 0 {
2058            panic!("Build multi thread runtime failed, reason: worker pool is empty");
2059        }
2060        if self.init > pool_worker_len {
2061            self.init = pool_worker_len;
2062        }
2063        if self.max > pool_worker_len {
2064            self.max = pool_worker_len;
2065        }
2066        if self.min > self.max {
2067            self.min = self.max;
2068        }
2069
2070        //构建多线程任务运行时的本地定时器和定时异步任务生产者
2071        let interval = self.interval;
2072        let mut timers = if let Some(_) = interval {
2073            Some(Vec::with_capacity(self.max))
2074        } else {
2075            None
2076        };
2077        for _ in 0..self.max {
2078            //初始化指定的最大线程数量的本地定时器和定时异步任务生产者,定时器不会在关闭工作者时被移除
2079            if let Some(vec) = &mut timers {
2080                let timer = AsyncTaskTimerByNotCancel::new();
2081                let producor = timer.producor.clone();
2082                let timer = Arc::new(timer);
2083                vec.push((producor, timer));
2084            };
2085        }
2086
2087        //构建多线程任务运行时
2088        let rt_uid = alloc_rt_uid();
2089        let waits = Arc::new(ArrayQueue::new(self.max));
2090        let mut pool = self.pool;
2091        pool.set_waits(waits.clone()); //设置待唤醒的工作者唤醒器队列
2092        let pool = Arc::new(pool);
2093        let runtime = MultiTaskRuntime(Arc::new((
2094            rt_uid,
2095            pool,
2096            timers,
2097            AtomicUsize::new(0),
2098            waits,
2099            AtomicUsize::new(0),
2100            AtomicUsize::new(0),
2101        )));
2102
2103        //构建初始化线程数量的线程构建器
2104        let mut builders = Vec::with_capacity(self.init);
2105        for index in 0..self.init {
2106            let builder = Builder::new()
2107                .name(self.prefix.clone() + "-" + index.to_string().as_str())
2108                .stack_size(self.stack_size);
2109            builders.push(builder);
2110        }
2111
2112        //启动工作者线程
2113        let min = self.min;
2114        for index in 0..builders.len() {
2115            let builder = builders.remove(0);
2116            let runtime = runtime.clone();
2117            let timeout = self.timeout;
2118            let timer = if let Some(timers) = &(runtime.0).2 {
2119                let (_, timer) = &timers[index];
2120                Some(timer.clone())
2121            } else {
2122                None
2123            };
2124
2125            spawn_worker_thread(builder, index, runtime, min, timeout, interval, timer);
2126        }
2127
2128        runtime
2129    }
2130}
2131
2132/// 将当前 OS worker 绑定到既有 runtime thread id 和 exact task pool。
2133///
2134/// 参数:
2135/// - `thread_id`:由 runtime uid 和 worker index 组成的既有 packed id。
2136/// - `pool`:worker 即将驱动的 exact pool;只取稳定数据地址,不保存引用或增加引用计数。
2137///
2138/// 副作用与幂等:
2139/// - 非纯函数,会写当前线程的 `PI_ASYNC_THREAD_LOCAL_ID` 和模块私有 owner context。
2140/// - 对相同线程、相同 id/pool 重复调用是幂等的;本实现只在 worker 启动时调用一次。
2141///
2142/// 性能与阻塞:
2143/// - O(1) 时间、每线程 O(1) TLS 空间;无 heap allocation、Arc clone、锁、自旋、阻塞或 I/O。
2144/// - 不是任务 poll 热路径,只在 worker startup 执行。
2145///
2146/// 安全与错误边界:
2147/// - 写入旧 TLS 的 UnsafeCell 是安全的,因为 thread-local 实例只由当前 OS 线程访问。
2148/// - raw pool pointer 只在后续 owner guard 中比较,永不解引用;worker closure 持有 runtime Arc。
2149/// - TLS 已销毁时 panic;此时禁止启动/重绑 worker。无用户回调、V8/FFI 或跨 await 行为。
2150fn bind_multi_thread_worker_context<P>(thread_id: usize, pool: &P) {
2151    if let Err(e) = PI_ASYNC_THREAD_LOCAL_ID.try_with(|local_thread_id| unsafe {
2152        // SAFETY: this UnsafeCell belongs to the current OS thread's TLS instance. Worker startup
2153        // writes it before entering any task-pool operation, and no other thread can alias it.
2154        *local_thread_id.get() = thread_id;
2155    }) {
2156        panic!(
2157            "Multi thread runtime startup failed, thread id: {:?}, reason: {:?}",
2158            thread_id & MULTI_THREAD_WORKER_ID_MASK,
2159            e
2160        );
2161    }
2162
2163    let context = MultiThreadWorkerContext {
2164        thread_id,
2165        pool: pool as *const P as *const (),
2166    };
2167    if let Err(e) = PI_ASYNC_MULTI_THREAD_WORKER_CONTEXT.try_with(|current| {
2168        current.set(context);
2169    }) {
2170        panic!(
2171            "Bind multi-thread worker pool failed, thread id: {:?}, reason: {:?}",
2172            thread_id & MULTI_THREAD_WORKER_ID_MASK,
2173            e
2174        );
2175    }
2176}
2177
2178/// 创建一个 OS worker,绑定 runtime/pool 上下文,并进入 timer 或 non-timer 工作循环。
2179///
2180/// `builder`、`index`、`runtime`、worker limit、sleep timeout 和可选 timer 均由已校验的
2181/// `MultiTaskRuntimeBuilder::build` 传入。成功时函数只提交线程创建并立即返回;线程内部先绑定
2182/// TLS,再绑定 local runtime,最后进入原工作循环。线程创建失败保持旧行为,由 `spawn` 返回值
2183/// 被忽略;本轮不改变该既有错误语义。
2184///
2185/// 构建入口 O(1),每个 worker 固定 O(1) 额外 TLS;会创建线程和分配线程栈,但不是任务热路径。
2186/// 不在锁内绑定或执行 future,不增加阻塞/死锁/重入边界。`timer` 为 `Some` 时仍只归该 worker
2187/// 使用,本 helper 不共享或迁移 timer。公开 API、任务执行顺序和 V8/FFI 边界不变。
2188fn spawn_worker_thread<
2189    O: Default + 'static,
2190    P: AsyncTaskPoolExt<O> + AsyncTaskPool<O, Pool = P>,
2191>(
2192    builder: Builder,
2193    index: usize,
2194    runtime: MultiTaskRuntime<O, P>,
2195    min: usize,
2196    timeout: u64,
2197    interval: Option<usize>,
2198    timer: Option<Arc<AsyncTaskTimerByNotCancel<P, O>>>,
2199) {
2200    if let Some(timer) = timer {
2201        //设置了定时器
2202        let rt_uid = runtime.get_id();
2203        let _ = builder.spawn(move || {
2204            //设置线程本地唯一id并绑定exact pool owner上下文
2205            let thread_id = rt_uid << 32 | index & MULTI_THREAD_WORKER_ID_MASK;
2206            bind_multi_thread_worker_context(thread_id, (runtime.0).1.as_ref());
2207
2208            //绑定运行时到线程
2209            let runtime_copy = runtime.clone();
2210            match PI_ASYNC_LOCAL_THREAD_ASYNC_RUNTIME.try_with(move |rt| {
2211                let raw = Arc::into_raw(Arc::new(runtime_copy.to_local_runtime()))
2212                    as *mut LocalAsyncRuntime<O> as *mut ();
2213                rt.store(raw, Ordering::Relaxed);
2214            }) {
2215                Err(e) => {
2216                    panic!("Bind multi runtime to local thread failed, reason: {:?}", e);
2217                }
2218                Ok(_) => (),
2219            }
2220
2221            //执行有定时器的工作循环
2222            timer_work_loop(
2223                runtime,
2224                index,
2225                min,
2226                timeout,
2227                interval.unwrap() as u64,
2228                timer,
2229            );
2230        });
2231    } else {
2232        //未设置定时器
2233        let rt_uid = runtime.get_id();
2234        let _ = builder.spawn(move || {
2235            //设置线程本地唯一id并绑定exact pool owner上下文
2236            let thread_id = rt_uid << 32 | index & MULTI_THREAD_WORKER_ID_MASK;
2237            bind_multi_thread_worker_context(thread_id, (runtime.0).1.as_ref());
2238
2239            //绑定运行时到线程
2240            let runtime_copy = runtime.clone();
2241            match PI_ASYNC_LOCAL_THREAD_ASYNC_RUNTIME.try_with(move |rt| {
2242                let raw = Arc::into_raw(Arc::new(runtime_copy.to_local_runtime()))
2243                    as *mut LocalAsyncRuntime<O> as *mut ();
2244                rt.store(raw, Ordering::Relaxed);
2245            }) {
2246                Err(e) => {
2247                    panic!("Bind multi runtime to local thread failed, reason: {:?}", e);
2248                }
2249                Ok(_) => (),
2250            }
2251
2252            //执行无定时器的工作循环
2253            work_loop(runtime, index, min, timeout);
2254        });
2255    }
2256}
2257
2258/// worker 空闲等待的结果。
2259///
2260/// 说明:
2261/// - 该枚举只用于多线程运行时内部工作循环,不属于公开 API。
2262/// - 它把“休眠超时”“未进入休眠/被唤醒”“在休眠前二次检查直接拿到任务”三个结果
2263///   分开,避免工作循环用布尔值推断调度状态。
2264///
2265/// 业务边界:
2266/// - 不表达任务执行结果,也不表达 runtime 关闭状态。
2267/// - `Task` 只表示 worker 在进入 condvar wait 前从真实任务池取到了一个任务,调用方
2268///   必须立即走正常 `run_task` 路径。
2269///
2270/// 性能与安全:
2271/// - 纯数据枚举,本身无副作用、不分配、不阻塞。
2272/// - 持有 `Arc<AsyncTask<...>>` 的 `Task` 分支遵循原任务池所有权语义。
2273enum WorkerWaitResult<O: Default + 'static, P: AsyncTaskPoolExt<O> + AsyncTaskPool<O, Pool = P>> {
2274    TimedOut,
2275    NotSlept,
2276    Task(Arc<AsyncTask<P, O>>),
2277}
2278
2279/// 在 worker 空闲时注册可唤醒状态,并在必要时进入 condvar 等待。
2280///
2281/// 说明:
2282/// - 这是多线程运行时 worker sleep/wake 协议的唯一入口。
2283/// - 目标是保证外部线程在任务入队后只要存在 sleeping worker,就能即时唤醒一个 worker;
2284///   同时 worker 不会在“任务已入队但未被 notify”的状态下睡到 `sleep_timeout`。
2285/// - 该函数不改变任务执行语义,不创建/销毁任务,不修改公开 API。
2286///
2287/// 核心协议:
2288/// 1. 锁内注册:短暂持有当前 worker 的 `worker_waker` 锁,把唤醒器放入 waits 队列,
2289///    并在同一临界区发布 `is_sleep = true`。
2290/// 2. 锁外二次检查:释放 worker_waker 锁后检查真实任务池。如果已有任务,则取消
2291///    `is_sleep` 并直接返回任务或返回 `NotSlept`。
2292/// 3. 锁内等待:再次短暂持锁确认 `is_sleep` 仍为 true。若外部唤醒已把它置为 false,
2293///    直接返回 `NotSlept`;否则执行 `condvar.wait_for`。
2294///
2295/// 为什么这样设计:
2296/// - 注册和发布在同一把锁内连续完成,外部 wake 端弹出 waits 条目后会获取同一把锁,
2297///   因而不会把“已入队但尚未发布 true”的 worker 当成 stale,也不会漏唤醒。
2298/// - 任务队列 `try_pop` / `len` 放在 worker_waker 锁外,避免 worker_waker 临界区与
2299///   任务队列窃取、随机选择、统计更新等热路径逻辑重叠。
2300/// - `condvar.wait_for` 是唯一可能阻塞点;它只发生在确认队列无任务且 `is_sleep` 仍为
2301///   true 之后,并且 parking_lot 会在等待期间释放 mutex。
2302///
2303/// 参数:
2304/// - `runtime`:当前 worker 所属的多线程 runtime。
2305/// - `worker_waker`:当前 worker 独占使用的线程唤醒器。
2306/// - `sleep_timeout`:本次允许休眠的最长时长,单位 ms。定时器 worker 会传入计算后的
2307///   timer-aware timeout,普通 worker 会传入 builder 配置的 worker sleep timeout。
2308///
2309/// 返回:
2310/// - `TimedOut`:进入了 condvar wait,且本次由超时返回。调用方可增加连续休眠计数。
2311/// - `NotSlept`:没有进入有效休眠,或被 notify/取消后需要回到 poll loop 重新检查队列。
2312/// - `Task(task)`:休眠前二次检查直接取到任务,调用方应立即执行该任务。
2313///
2314/// 边界条件:
2315/// - waits 队列满时会释放当前 worker 锁,再清理 stale entry;若清理后仍无法注册,
2316///   返回 `NotSlept`,禁止无唤醒入口地休眠。
2317/// - 外部 wake 与 worker 二次检查竞态时,`is_sleep` 的 CAS/store 会收敛到最多一次
2318///   notify;额外的 `NotSlept` 只会让 worker 回到 poll loop,不会丢任务。
2319/// - sleep_timeout 为 0 时,`wait_for(0ms)` 会立即返回,不改变语义。
2320///
2321/// 性能:
2322/// - 快路径时间复杂度 O(1),空间复杂度 O(1)。
2323/// - waits 满且需要清理 stale 时最坏 O(W),W 为最大 worker 数;该慢路径只在注册失败
2324///   时触发,不在每次 wake 热路径上执行。
2325/// - 每个休眠周期最多 clone 一次 `worker_waker` Arc 用于队列登记;任务唤醒路径不额外
2326///   clone worker_waker。
2327///
2328/// 纯度与副作用:
2329/// - 非纯函数。会修改 waits 队列、当前 worker 的 `is_sleep` 状态,并可能从任务池取出
2330///   一个任务。
2331/// - 非幂等:每次调用代表一个新的 worker 空闲等待尝试。
2332///
2333/// 阻塞性:
2334/// - 除 `condvar.wait_for` 外不执行阻塞等待。
2335/// - 不在 worker_waker 锁内执行任务 poll、用户 future、I/O 或回调。
2336///
2337/// 安全性:
2338/// - 不引入新的 unsafe。
2339/// - 线程安全:依赖 `ArrayQueue`、`AtomicBool` 和 `Mutex/Condvar` 的组合协议。
2340/// - 内存安全:waits 中保存的是 worker_waker 的 Arc,生命周期由 runtime/worker 持有;
2341///   stale entry 被弹出后自然释放引用。
2342/// - 异步安全:只调度任务,不在锁内 poll future,不跨 await 持有锁。
2343/// - 运行时依赖:要求同一个 runtime 的所有 spawn/wake 路径在任务入队后调用
2344///   `wake_waiting_worker`。
2345#[inline]
2346fn worker_wait_for_task<O: Default + 'static, P: AsyncTaskPoolExt<O> + AsyncTaskPool<O, Pool = P>>(
2347    runtime: &MultiTaskRuntime<O, P>,
2348    worker_waker: &Arc<(AtomicBool, Mutex<()>, Condvar)>,
2349    sleep_timeout: u64,
2350) -> WorkerWaitResult<O, P> {
2351    let (is_sleep, lock, condvar) = &**worker_waker;
2352
2353    loop {
2354        let _locked = lock.lock();
2355        if is_sleep.load(Ordering::Acquire) {
2356            break;
2357        }
2358
2359        if register_waiting_worker(&(runtime.0).4, worker_waker) {
2360            is_sleep.store(true, Ordering::Release);
2361            break;
2362        }
2363
2364        drop(_locked);
2365        if prune_stale_waiting_workers(&(runtime.0).4) == 0 {
2366            return WorkerWaitResult::NotSlept;
2367        }
2368    }
2369
2370    if let Some(task) = (runtime.0).1.try_pop() {
2371        is_sleep.store(false, Ordering::Release);
2372        return WorkerWaitResult::Task(task);
2373    }
2374
2375    if runtime.len() > 0 {
2376        is_sleep.store(false, Ordering::Release);
2377        return WorkerWaitResult::NotSlept;
2378    }
2379
2380    let mut locked = lock.lock();
2381    if !is_sleep.load(Ordering::Acquire) {
2382        return WorkerWaitResult::NotSlept;
2383    }
2384
2385    let timed_out = condvar
2386        .wait_for(&mut locked, Duration::from_millis(sleep_timeout))
2387        .timed_out();
2388    is_sleep.store(false, Ordering::Release);
2389
2390    if timed_out {
2391        WorkerWaitResult::TimedOut
2392    } else {
2393        WorkerWaitResult::NotSlept
2394    }
2395}
2396
2397//线程工作循环
2398fn timer_work_loop<O: Default + 'static, P: AsyncTaskPoolExt<O> + AsyncTaskPool<O, Pool = P>>(
2399    runtime: MultiTaskRuntime<O, P>,
2400    index: usize,
2401    min: usize,
2402    sleep_timeout: u64,
2403    timer_interval: u64,
2404    timer: Arc<AsyncTaskTimerByNotCancel<P, O>>,
2405) {
2406    //初始化当前线程的线程id和线程活动状态
2407    let pool = (runtime.0).1.clone();
2408    let worker_waker = pool.clone_thread_waker().unwrap();
2409
2410    let mut sleep_count = 0; //连续休眠计数器
2411    let clock = Clock::new();
2412    loop {
2413        //设置新的定时异步任务,并唤醒已到期的定时异步任务
2414        let timer_run_millis = clock.recent(); //重置定时器运行时长
2415        let mut pop_len = 0;
2416        (runtime.0)
2417            .5
2418            .fetch_add(timer.consume(),
2419                       Ordering::Relaxed);
2420        loop {
2421            let current_time = timer.is_require_pop();
2422            if let Some(current_time) = current_time {
2423                //当前有到期的定时异步任务,则开始处理到期的所有定时异步任务
2424                loop {
2425                    let timed_out = timer.pop(current_time);
2426                    if let Some(timing_task) = timed_out {
2427                        match timing_task {
2428                            AsyncTimingTask::Pended(expired) => {
2429                                //唤醒休眠的异步任务,不需要立即在本工作者中执行,因为休眠的异步任务无法取消
2430                                runtime.wakeup::<O>(&expired);
2431                            }
2432                            AsyncTimingTask::WaitRun(expired) => {
2433                                //执行到期的定时异步任务,需要立即在本工作者中执行,因为定时异步任务可以取消
2434                                (runtime.0)
2435                                    .1
2436                                    .push_priority(DEFAULT_MAX_HIGH_PRIORITY_BOUNDED,
2437                                                   expired);
2438                                if let Some(task) = pool.try_pop() {
2439                                    sleep_count = 0; //重置连续休眠次数
2440                                    run_task(&runtime, task);
2441                                }
2442                            }
2443                            AsyncTimingTask::TimeoutWake(waiter) => {
2444                                //唤醒等待timeout到期的任务
2445                                waiter.fire();
2446                            }
2447                        }
2448                        pop_len += 1;
2449
2450                        if let Some(task) = pool.try_pop() {
2451                            //执行当前工作者任务池中的异步任务,避免定时异步任务占用当前工作者的所有工作时间
2452                            sleep_count = 0; //重置连续休眠次数
2453                            run_task(&runtime, task);
2454                        }
2455                    } else {
2456                        //当前所有的到期任务已处理完,则退出本次定时异步任务处理
2457                        break;
2458                    }
2459                }
2460            } else {
2461                //当前没有到期的定时异步任务,则退出本次定时异步任务处理
2462                break;
2463            }
2464        }
2465        (runtime.0)
2466            .6
2467            .fetch_add(pop_len,
2468                       Ordering::Relaxed);
2469
2470        //继续执行当前工作者任务池中的异步任务
2471        match pool.try_pop() {
2472            None => {
2473                if runtime.len() > 0 {
2474                    //确认当前还有任务需要处理,可能还没分配到当前工作者,则当前工作者继续工作
2475                    continue;
2476                }
2477
2478                //获取休眠的实际时长
2479                let diff_time = clock
2480                    .recent()
2481                    .duration_since(timer_run_millis)
2482                    .as_millis() as u64; //获取定时器运行时长
2483                let real_timeout = if timer.len() == 0 {
2484                    //当前定时器没有未到期的任务,则休眠指定时长
2485                    sleep_timeout
2486                } else {
2487                    //当前定时器还有未到期的任务,则计算需要休眠的时长
2488                    if diff_time >= timer_interval {
2489                        //定时器内部时间与当前时间差距过大,则忽略休眠,并继续工作
2490                        continue;
2491                    } else {
2492                        //定时器内部时间与当前时间差距不大,则休眠差值时间
2493                        timer_interval - diff_time
2494                    }
2495                };
2496
2497                //无任务,则准备休眠
2498                match worker_wait_for_task(&runtime, &worker_waker, real_timeout) {
2499                    WorkerWaitResult::TimedOut => {
2500                        //记录连续休眠次数,因为任务导致的唤醒不会计数
2501                        sleep_count += 1;
2502                    },
2503                    WorkerWaitResult::Task(task) => {
2504                        sleep_count = 0; //重置连续休眠次数
2505                        run_task(&runtime, task);
2506                    },
2507                    WorkerWaitResult::NotSlept => (),
2508                }
2509            }
2510            Some(task) => {
2511                //有任务,则执行
2512                sleep_count = 0; //重置连续休眠次数
2513                run_task(&runtime, task);
2514            }
2515        }
2516    }
2517
2518    //关闭当前工作者的任务池
2519    (runtime.0).1.close_worker();
2520    warn!(
2521        "Worker of runtime closed, runtime: {}, worker: {}, thread: {:?}",
2522        runtime.get_id(),
2523        index,
2524        thread::current()
2525    );
2526}
2527
2528//线程工作循环
2529fn work_loop<O: Default + 'static, P: AsyncTaskPoolExt<O> + AsyncTaskPool<O, Pool = P>>(
2530    runtime: MultiTaskRuntime<O, P>,
2531    index: usize,
2532    min: usize,
2533    sleep_timeout: u64,
2534) {
2535    //初始化当前线程的线程id和线程活动状态
2536    let pool = (runtime.0).1.clone();
2537    let worker_waker = pool.clone_thread_waker().unwrap();
2538
2539    let mut sleep_count = 0; //连续休眠计数器
2540    loop {
2541        match pool.try_pop() {
2542            None => {
2543                //无任务,则准备休眠
2544                if runtime.len() > 0 {
2545                    //确认当前还有任务需要处理,可能还没分配到当前工作者,则当前工作者继续工作
2546                    continue;
2547                }
2548
2549                match worker_wait_for_task(&runtime, &worker_waker, sleep_timeout) {
2550                    WorkerWaitResult::TimedOut => {
2551                        //记录连续休眠次数,因为任务导致的唤醒不会计数
2552                        sleep_count += 1;
2553                    },
2554                    WorkerWaitResult::Task(task) => {
2555                        sleep_count = 0; //重置连续休眠次数
2556                        run_task(&runtime, task);
2557                    },
2558                    WorkerWaitResult::NotSlept => (),
2559                }
2560            }
2561            Some(task) => {
2562                //有任务,则执行
2563                sleep_count = 0; //重置连续休眠次数
2564                run_task(&runtime, task);
2565            }
2566        }
2567    }
2568
2569    //关闭当前工作者的任务池
2570    (runtime.0).1.close_worker();
2571    warn!(
2572        "Worker of runtime closed, runtime: {}, worker: {}, thread: {:?}",
2573        runtime.get_id(),
2574        index,
2575        thread::current()
2576    );
2577}
2578
2579//执行异步任务
2580#[inline]
2581fn run_task<O: Default + 'static, P: AsyncTaskPoolExt<O> + AsyncTaskPool<O, Pool = P>>(
2582    runtime: &MultiTaskRuntime<O, P>,
2583    task: Arc<AsyncTask<P, O>>,
2584) {
2585    let waker = waker_ref(&task);
2586    let mut context = Context::from_waker(&*waker);
2587    if let Some(mut future) = task.get_inner() {
2588        if let Poll::Pending = future.as_mut().poll(&mut context) {
2589            //当前未准备好,则恢复异步任务,以保证异步服务后续访问异步任务和异步任务不被提前释放
2590            task.set_inner(Some(future));
2591        }
2592    } else {
2593        //当前异步任务在唤醒时还未被重置内部任务,则继续加入当前异步运行时队列,并等待下次被执行
2594        (runtime.0).1.push(task);
2595    }
2596}
2597
2598// 定时器任务生产者
2599enum TimerTaskProducor<
2600    O: Default + 'static = (),
2601    P: AsyncTaskPoolExt<O> + AsyncTaskPool<O> = StealableTaskPool<O>,
2602> {
2603    Local(Arc<AsyncTaskTimerByNotCancel<P, O>>),        //本地定时器任务生产者
2604    Foreign(Sender<(usize, AsyncTimingTask<P, O>)>),    //外部定时器任务生产者
2605}