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