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