Skip to main content

pi_async_rt/rt/
mod.rs

1//! # 提供了通用的异步运行时
2//!
3
4use std::thread;
5use std::pin::Pin;
6use std::sync::Arc;
7use std::ptr::null_mut;
8use std::vec::IntoIter;
9use std::future::Future;
10use std::panic::set_hook;
11use std::any::{Any, TypeId};
12use std::marker::PhantomData;
13use std::ops::{Deref, DerefMut};
14use std::cell::{RefCell, UnsafeCell};
15use std::task::{Waker, Context, Poll};
16use std::time::{Duration, SystemTime};
17use std::io::{Error, Result, ErrorKind};
18use std::alloc::{Layout, set_alloc_error_hook};
19use std::fmt::{Debug, Formatter, Result as FmtResult};
20use std::sync::atomic::{AtomicBool, AtomicU8, AtomicUsize, AtomicPtr, Ordering};
21
22pub mod single_thread;
23pub mod multi_thread;
24pub mod worker_thread;
25pub mod serial;
26pub mod serial_local_thread;
27pub mod serial_single_thread;
28pub mod serial_worker_thread;
29pub mod serial_local_compatible_wasm_runtime;
30
31use libc;
32use futures::{future::{FutureExt, BoxFuture},
33              stream::{Stream, BoxStream},
34              task::{ArcWake, AtomicWaker}};
35use parking_lot::{Mutex, Condvar};
36use crossbeam_channel::{Sender, Receiver, unbounded};
37use crossbeam_queue::ArrayQueue;
38use crossbeam_utils::atomic::AtomicCell;
39use flume::{Sender as AsyncSender, Receiver as AsyncReceiver};
40use num_cpus;
41use backtrace::Backtrace;
42use slotmap::{Key, KeyData};
43use quanta::{Clock, Upkeep, Handle, Instant as QInstant};
44
45use pi_hash::XHashMap;
46use pi_cancel_timer::Timer;
47use pi_timer::Timer as NotCancelTimer;
48
49use single_thread::SingleTaskRuntime;
50use worker_thread::{WorkerTaskRunner, WorkerRuntime};
51use multi_thread::{MultiTaskRuntimeBuilder, MultiTaskRuntime, StealableTaskPool};
52
53use crate::lock::spin;
54
55/*
56* 本地线程绑定的异步运行时
57*/
58thread_local! {
59    static PI_ASYNC_LOCAL_THREAD_ASYNC_RUNTIME: AtomicPtr<()> = AtomicPtr::new(null_mut());
60    static PI_ASYNC_LOCAL_THREAD_ASYNC_RUNTIME_DICT: UnsafeCell<XHashMap<TypeId, Box<dyn Any + 'static>>> = UnsafeCell::new(XHashMap::default());
61}
62
63/*
64* 本地线程唯一id
65*/
66thread_local! {
67    static PI_ASYNC_THREAD_LOCAL_ID: UnsafeCell<usize> = UnsafeCell::new(usize::MAX);
68}
69
70/*
71* 默认的最高优先级边界
72*/
73const DEFAULT_MAX_HIGH_PRIORITY_BOUNDED: usize = 10;
74
75/*
76* 默认的高优先级边界
77*/
78const DEFAULT_HIGH_PRIORITY_BOUNDED: usize = 5;
79
80/*
81* 默认的最低优先级
82*/
83const DEFAULT_MAX_LOW_PRIORITY_BOUNDED: usize = 0;
84
85/*
86* 异步运行时唯一id生成器
87*/
88static RUNTIME_UID_GEN: AtomicUsize = AtomicUsize::new(1);
89
90/*
91* 全局时间状态
92*/
93static GLOBAL_TIME_LOOP_STATUS: AtomicBool = AtomicBool::new(false);
94
95///
96/// 启动全局时间循环,成功则返回句柄,释放句柄将关闭全局时间循环,失败表示已启动,则返回空
97/// 更新间隔时长为毫秒
98///
99pub fn startup_global_time_loop(interval: u64) -> Option<GlobalTimeLoopHandle> {
100    if let Err(_) = GLOBAL_TIME_LOOP_STATUS.compare_exchange(false,
101                                                             true,
102                                                             Ordering::AcqRel,
103                                                             Ordering::Relaxed) {
104        //已启动
105        None
106    } else {
107        //未启动
108        let timer = Upkeep::new_with_clock(Duration::from_millis(interval), Clock::new());
109        let handle = timer.start().unwrap();
110        let clock = Clock::new();
111        let _now = clock.recent();
112
113        Some(GlobalTimeLoopHandle(handle))
114    }
115}
116
117///
118/// 全局时间循环句柄
119///
120pub struct GlobalTimeLoopHandle(Handle);
121
122impl Drop for GlobalTimeLoopHandle {
123    fn drop(&mut self) {
124        GLOBAL_TIME_LOOP_STATUS.store(false, Ordering::Release);
125    }
126}
127
128///
129/// 分配异步运行时唯一id
130///
131pub fn alloc_rt_uid() -> usize {
132    RUNTIME_UID_GEN.fetch_add(1, Ordering::Relaxed)
133}
134
135///
136/// 异步任务唯一id
137///
138pub struct TaskId(UnsafeCell<u128>);
139
140impl Debug for TaskId {
141    fn fmt(&self, f: &mut Formatter) -> FmtResult {
142        write!(f, "TaskId[inner = {}]", unsafe { *self.0.get() })
143    }
144}
145
146impl Clone for TaskId {
147    fn clone(&self) -> Self {
148        unsafe {
149            TaskId(UnsafeCell::new(*self.0.get()))
150        }
151    }
152}
153
154impl TaskId {
155    /// 线程安全的判断异步任务唯一id对应的异步任务的唤醒器是否存在
156    #[inline]
157    pub fn exist_waker<R: 'static>(&self) -> bool {
158        unsafe {
159            let handle = unsafe { TaskHandle::<R>::from_raw((*self.0.get() >> 64) as *const ()) };
160            let inner = &*handle.0;
161            let r = if let Some(waker) = inner.0.swap(None) {
162                inner.0.swap(Some(waker));
163                true
164            } else {
165                false
166            };
167
168            //避免提前释放
169            handle.into_raw();
170
171            r
172        }
173    }
174
175    /// 线程安全的唤醒异步任务唯一id对应的异步任务
176    #[inline]
177    pub fn wakeup<R: 'static>(&self) {
178        unsafe {
179            let handle = unsafe { TaskHandle::<R>::from_raw((*self.0.get() >> 64) as *const ()) };
180            let inner = &*handle.0;
181            if let Some(waker) = inner.0.swap(None) {
182                //当前异步任务的唤醒器存在,则唤醒
183                waker.wake();
184            }
185
186            //避免提前释放
187            handle.into_raw();
188        }
189    }
190
191    /// 线程安全的为异步任务唯一id对应的异步任务设置唤醒器
192    #[inline]
193    pub fn set_waker<R: 'static>(&self, waker: Waker) -> Option<Waker> {
194        unsafe {
195            let handle = unsafe { TaskHandle::<R>::from_raw((*self.0.get() >> 64) as *const ()) };
196            let inner = &*handle.0;
197            let r = inner.0.swap(Some(waker));
198
199            //避免提前释放
200            handle.into_raw();
201
202            r
203        }
204    }
205
206    /// 线程安全的获取异步任务唯一id对应的异步任务的返回值
207    #[inline]
208    pub fn result<R: 'static>(&self) -> Option<R> {
209        unsafe {
210            let handle = unsafe { TaskHandle::<R>::from_raw((*self.0.get() >> 64) as *const ()) };
211            let inner = &*handle.0;
212            let r = inner.1.swap(None);
213
214            //避免提前释放
215            handle.into_raw();
216
217            r
218        }
219    }
220
221    /// 线程安全的为异步任务唯一id对应的异步任务设置返回值
222    #[inline]
223    pub fn set_result<R: 'static>(&self, result: R) -> Option<R> {
224        unsafe {
225            let handle = unsafe { TaskHandle::<R>::from_raw((*self.0.get() >> 64) as *const ()) };
226            let inner = &*handle.0;
227            let r = inner.1.swap(Some(result));
228
229            //避免提前释放
230            handle.into_raw();
231
232            r
233        }
234    }
235}
236
237// 异步任务句柄
238pub(crate) struct TaskHandle<R: 'static>(Box<(
239    AtomicCell<Option<Waker>>,  //任务唤醒器
240    AtomicCell<Option<R>>,      //任务返回值
241)>);
242
243impl<R: 'static> Default for TaskHandle<R> {
244    fn default() -> Self {
245        TaskHandle(Box::new((AtomicCell::new(None), AtomicCell::new(None))))
246    }
247}
248
249impl<R: 'static> TaskHandle<R> {
250    /// 将祼指针转换为异步任务句柄
251    pub unsafe fn from_raw(raw: *const ()) -> TaskHandle<R> {
252        let inner
253            = Box::from_raw(raw as *const (AtomicCell<Option<Waker>>, AtomicCell<Option<R>>) as *mut (AtomicCell<Option<Waker>>, AtomicCell<Option<R>>));
254        TaskHandle(inner)
255    }
256
257    /// 将异步任务句柄转换为祼指针
258    pub fn into_raw(self) -> *const () {
259        Box::into_raw(self.0)
260            as *mut (AtomicCell<Option<Waker>>, AtomicCell<Option<R>>)
261            as *const (AtomicCell<Option<Waker>>, AtomicCell<Option<R>>)
262            as *const ()
263    }
264}
265
266/// timeout专用等待句柄
267pub(crate) struct TimeoutWaiter {
268    fired: AtomicBool,
269    waker: AtomicWaker,
270}
271
272impl TimeoutWaiter {
273    #[inline]
274    pub fn new() -> Self {
275        TimeoutWaiter {
276            fired: AtomicBool::new(false),
277            waker: AtomicWaker::new(),
278        }
279    }
280
281    #[inline]
282    pub fn is_fired(&self) -> bool {
283        self.fired.load(Ordering::Acquire)
284    }
285
286    #[inline]
287    pub fn register(&self, waker: &Waker) {
288        self.waker.register(waker);
289    }
290
291    #[inline]
292    pub fn fire(&self) {
293        if !self.fired.swap(true, Ordering::AcqRel) {
294            self.waker.wake();
295        }
296    }
297
298    #[inline]
299    pub fn clear_waker(&self) {
300        let _ = self.waker.take();
301    }
302}
303
304/// 唤醒一个已经从等待队列中取出的工作者线程唤醒器。
305///
306/// 说明:
307/// - 该 helper 只用于“调用方已经确认这个 `worker_waker` 来自等待队列”的场景。
308/// - 它会先获取 `worker_waker` 内部的互斥锁,再检查并切换 `is_sleep`。
309/// - 锁内检查是为了覆盖 worker 进入休眠时的发布窗口:worker 先把唤醒器放入
310///   waits 队列,再在同一把锁保护下发布 `is_sleep = true`,外部唤醒端必须等这个
311///   发布动作完成后再判断是否 notify。
312///
313/// 参数:
314/// - `worker_waker`:工作者线程的 `(is_sleep, lock, condvar)` 三元组。
315///
316/// 返回:
317/// - `true`:本次成功把 `is_sleep` 从 `true` 切为 `false`,并调用了 `notify_one()`。
318/// - `false`:该唤醒器已经失效、已被其它唤醒者消费,或 worker 已经自行取消休眠。
319///
320/// 边界与业务范围:
321/// - 不创建任务、不修改任务队列,不负责判断 runtime 中是否已有任务。
322/// - 只唤醒一个已注册的 worker,不广播,不循环 notify,不负责 worker 选择策略。
323/// - 如果 worker 在被 notify 前已经通过二次检查取到任务并取消休眠,本函数会返回
324///   `false`,这是正确的无操作。
325///
326/// 性能:
327/// - 时间复杂度 O(1),空间复杂度 O(1),不分配内存,不 clone。
328/// - 可能短暂获取 parking_lot mutex;不在 poll future 的内部持锁等待,也不会执行
329///   condvar wait。
330///
331/// 纯度与副作用:
332/// - 非纯函数。副作用是原子状态切换和一次条件变量通知。
333/// - 对同一个已注册唤醒器重复调用是幂等收敛的:最多一次调用能从 `true` 切到
334///   `false` 并 notify。
335///
336/// 安全性:
337/// - 不使用 unsafe。
338/// - 线程安全:依赖 `AtomicBool` 的 Acquire/AcqRel 可见性和 `Mutex` 对休眠发布窗口
339///   的互斥保护。
340/// - 异步安全:不会阻塞 executor worker 的异步任务 poll;只在外部唤醒或 spawn
341///   入队后的线程级唤醒路径上短暂执行。
342/// - 运行时依赖:要求传入的 `worker_waker` 与对应 worker 的 condvar wait 使用同一把
343///   lock。
344#[inline]
345pub(crate) fn wake_registered_thread_waker(worker_waker: &Arc<(AtomicBool, Mutex<()>, Condvar)>) -> bool {
346    let (is_sleep, lock, condvar) = &**worker_waker;
347    let _locked = lock.lock();
348    if is_sleep
349        .compare_exchange(true, false, Ordering::AcqRel, Ordering::Acquire)
350        .is_ok()
351    {
352        condvar.notify_one();
353        return true;
354    }
355
356    false
357}
358
359/// 快速唤醒单个线程唤醒器。
360///
361/// 说明:
362/// - 该 helper 用于没有 waits 队列的单 worker / 单线程运行时唤醒路径。
363/// - 与 `wake_registered_thread_waker` 相比,它先用一次 Acquire load 做快速过滤;
364///   当 `is_sleep == false` 时不获取锁。
365///
366/// 使用指导:
367/// - waits 队列中弹出的条目必须使用 `wake_registered_thread_waker`,因为队列条目可能
368///   处于“已入队但尚未发布 `is_sleep = true`”的临界窗口。
369/// - 直接持有线程唤醒器、且没有队列发布窗口时,可以使用本函数。
370///
371/// 参数与返回:
372/// - 参数同 `wake_registered_thread_waker`。
373/// - 返回 `true` 表示实际 notify 了一次;返回 `false` 表示无需唤醒或已被消费。
374///
375/// 性能与副作用:
376/// - 常见无休眠路径 O(1) 且无锁;需要唤醒时 O(1) 并短暂持锁。
377/// - 不分配内存,不 clone,不广播,不会形成唤醒风暴。
378///
379/// 安全性:
380/// - 不使用 unsafe。
381/// - 线程安全、内存安全;依赖同一 `worker_waker` 被 worker wait 和 wake 端共享。
382#[inline]
383pub(crate) fn wake_thread_waker(worker_waker: &Arc<(AtomicBool, Mutex<()>, Condvar)>) -> bool {
384    if !worker_waker.0.load(Ordering::Acquire) {
385        return false;
386    }
387
388    wake_registered_thread_waker(worker_waker)
389}
390
391/// 从多线程运行时的等待队列中唤醒一个可唤醒 worker。
392///
393/// 说明:
394/// - waits 是一个有界队列,队列项是 worker 注册的线程唤醒器。
395/// - 本函数每次最多成功唤醒一个 worker;遇到已经失效的陈旧项会丢弃并继续扫描。
396/// - 该行为用于避免漏唤醒,同时避免对所有 worker 广播造成唤醒风暴。
397///
398/// 使用指导:
399/// - 任务被外部线程入队后调用,例如 `spawn_by_id`、`spawn_priority_by_id` 和
400///   `AsyncTask::wake_by_ref`。
401/// - 调用方不应在持有任务队列内部锁时调用;当前任务池入队 API 本身不暴露需要
402///   调用方持有的锁。
403///
404/// 参数:
405/// - `waits`:当前 runtime 共享的 sleeping worker 等待队列。
406///
407/// 返回:
408/// - `true`:成功唤醒了一个仍处于休眠发布状态的 worker。
409/// - `false`:队列为空,或扫描到的条目均已失效。
410///
411/// 边界条件:
412/// - 队列容量等于 runtime 最大 worker 数。扫描上限固定为 `capacity`,不会无限循环。
413/// - 并发唤醒同一个队列项时,只有一个调用者能 CAS 成功并 notify。
414/// - 陈旧项来自 worker timeout 或二次检查取消休眠;丢弃它们不会丢任务,因为任务
415///   已经在任务队列中,或 worker 已经自行继续轮询。
416///
417/// 性能:
418/// - 最坏时间复杂度 O(W),W 为 waits 容量,即最大 worker 数;常见路径接近 O(1)。
419/// - 空间复杂度 O(1),不分配内存,不 clone。
420/// - 每次调用最多一次 notify,避免唤醒风暴。
421///
422/// 纯度与副作用:
423/// - 非纯函数。会从 waits 队列弹出条目,可能切换 worker 休眠状态并 notify。
424/// - 对同一批陈旧项重复调用是幂等收敛的:陈旧项会被逐步清理。
425///
426/// 安全性:
427/// - 不使用 unsafe。
428/// - 线程安全:ArrayQueue 提供并发队列安全;worker 状态由原子和 mutex 保护。
429/// - 异步安全:不会执行 condvar wait,不会阻塞当前异步任务,只在调度唤醒路径短暂
430///   执行。
431#[inline]
432pub(crate) fn wake_waiting_worker(
433    waits: &ArrayQueue<Arc<(AtomicBool, Mutex<()>, Condvar)>>,
434) -> bool {
435    let scan_len = waits.capacity();
436    for _ in 0..scan_len {
437        match waits.pop() {
438            Some(worker_waker) => {
439                if wake_registered_thread_waker(&worker_waker) {
440                    return true;
441                }
442            },
443            None => {
444                return false;
445            },
446        }
447    }
448
449    false
450}
451
452/// 清理 waits 队列中的陈旧 worker 唤醒器。
453///
454/// 说明:
455/// - worker 可能因为 sleep timeout、二次检查发现任务、或被其它唤醒者消费而把
456///   `is_sleep` 清为 `false`,但旧队列项仍留在有界 waits 队列中。
457/// - 本函数用于注册新 sleep 前释放这些陈旧槽位,避免 waits 被 stale entry 填满后
458///   worker 只能忙等。
459///
460/// 使用指导:
461/// - 只在 `register_waiting_worker` 遇到队列满时调用。
462/// - 不作为常规 wake 路径使用,避免在热唤醒路径上做额外扫描。
463///
464/// 参数与返回:
465/// - `waits`:当前 runtime 的 waiting worker 队列。
466/// - 返回实际移除的陈旧项数量。
467///
468/// 边界条件:
469/// - 函数只扫描调用开始时观察到的 `waits.len()` 个条目,不无限循环。
470/// - 每个弹出的条目都会先短暂获取该 worker 的锁再判断 `is_sleep`,避免把“已入队但
471///   尚未发布 `is_sleep = true`”的注册窗口误判为 stale。
472/// - 仍为 `true` 的 live 条目会放回队列;如果并发竞争导致放回失败,则立即尝试唤醒
473///   该 live worker,避免丢失一个真实 sleeping worker 的唤醒入口。
474///
475/// 性能:
476/// - 最坏时间复杂度 O(N),N 为调用开始时的队列长度,N <= worker 上限。
477/// - 空间复杂度 O(1),不分配内存;live 条目放回时复用已弹出的 Arc,不额外 clone。
478///
479/// 纯度与副作用:
480/// - 非纯函数。会重排 waits 队列中的 live 条目,移除 stale 条目,极端竞争下可能
481///   notify 一个 live worker。
482/// - 幂等:重复调用会逐步收敛到没有 stale entry。
483///
484/// 安全性:
485/// - 不使用 unsafe。
486/// - 线程安全;依赖 ArrayQueue、AtomicBool 和 worker_waker mutex。
487#[inline]
488pub(crate) fn prune_stale_waiting_workers(
489    waits: &ArrayQueue<Arc<(AtomicBool, Mutex<()>, Condvar)>>,
490) -> usize {
491    let scan_len = waits.len();
492    let mut pruned = 0;
493
494    for _ in 0..scan_len {
495        let Some(worker_waker) = waits.pop() else {
496            break;
497        };
498
499        let is_live = {
500            let _locked = worker_waker.1.lock();
501            worker_waker.0.load(Ordering::Acquire)
502        };
503
504        if is_live {
505            match waits.push(worker_waker) {
506                Ok(()) => (),
507                Err(worker_waker) => {
508                    let _ = wake_registered_thread_waker(&worker_waker);
509                },
510            }
511        } else {
512            pruned += 1;
513        }
514    }
515
516    pruned
517}
518
519/// 将当前 worker 注册为可被外部任务入队唤醒的候选 worker。
520///
521/// 说明:
522/// - 本函数只负责把 `worker_waker` 放入 waits 队列,不负责把 `is_sleep` 置为 true。
523/// - 调用方必须在持有 `worker_waker` 内部 mutex 的情况下调用成功快路径,并在成功入队
524///   后、同一把锁释放前发布 `is_sleep = true`。这样外部唤醒端弹出队列项后会在锁上等待
525///   发布完成,不会出现“队列中已有条目但状态尚未可唤醒”的漏唤醒窗口。
526/// - 队列满时调用方应释放当前 worker 锁后再调用 `prune_stale_waiting_workers`,避免当前
527///   worker 锁与其它 worker 唤醒锁形成嵌套临界区。
528///
529/// 参数:
530/// - `waits`:当前 runtime 的 waiting worker 队列。
531/// - `worker_waker`:当前 worker 的线程唤醒器。
532///
533/// 返回:
534/// - `true`:注册成功;调用方可以继续发布 `is_sleep = true` 并进入二次检查/等待。
535/// - `false`:队列满;调用方不得进入 condvar wait,应先释放锁并尝试清理 stale,或
536///   继续 poll loop,避免无唤醒入口地睡眠。
537///
538/// 边界条件:
539/// - 若队列全是 live worker,返回 `false` 是允许的:已有其它 worker 可被唤醒,当前
540///   worker 继续循环即可。
541///
542/// 性能:
543/// - 成功快路径 O(1),一次 Arc clone 用于把 worker 唤醒器登记到队列。
544/// - 队列满时 O(1) 返回,不在该 helper 内扫描队列。
545///
546/// 纯度与副作用:
547/// - 非纯函数。会向 waits 入队。
548/// - 非幂等:重复成功调用会重复登记同一个 worker,因此必须由调用方的 `is_sleep`
549///   状态保证同一 worker 同一休眠周期只注册一次。
550///
551/// 安全性:
552/// - 不使用 unsafe。
553/// - 线程安全;要求调用方遵守“持锁注册,锁内发布 true”的协议。
554#[inline]
555pub(crate) fn register_waiting_worker(
556    waits: &ArrayQueue<Arc<(AtomicBool, Mutex<()>, Condvar)>>,
557    worker_waker: &Arc<(AtomicBool, Mutex<()>, Condvar)>,
558) -> bool {
559    waits.push(worker_waker.clone()).is_ok()
560}
561
562#[cfg(test)]
563mod timeout_waiter_tests {
564    use super::{
565        prune_stale_waiting_workers, register_waiting_worker, wake_thread_waker,
566        wake_waiting_worker, TimeoutWaiter,
567    };
568    use crossbeam_queue::ArrayQueue;
569    use futures::task::{waker_ref, ArcWake};
570    use parking_lot::{Condvar, Mutex};
571    use std::sync::{
572        atomic::{AtomicBool, AtomicUsize, Ordering},
573        Arc,
574    };
575
576    struct WakeCounter(AtomicUsize);
577
578    impl ArcWake for WakeCounter {
579        fn wake_by_ref(arc_self: &Arc<Self>) {
580            arc_self.0.fetch_add(1, Ordering::SeqCst);
581        }
582    }
583
584    #[test]
585    fn test_timeout_waiter_fire_wakes_once() {
586        let waiter = TimeoutWaiter::new();
587        let counter = Arc::new(WakeCounter(AtomicUsize::new(0)));
588        let waker = waker_ref(&counter);
589
590        waiter.register(&waker);
591        waiter.fire();
592        waiter.fire();
593
594        assert!(waiter.is_fired());
595        assert_eq!(counter.0.load(Ordering::SeqCst), 1);
596    }
597
598    #[test]
599    fn test_timeout_waiter_clear_waker_before_fire() {
600        let waiter = TimeoutWaiter::new();
601        let counter = Arc::new(WakeCounter(AtomicUsize::new(0)));
602        let waker = waker_ref(&counter);
603
604        waiter.register(&waker);
605        waiter.clear_waker();
606        waiter.fire();
607
608        assert!(waiter.is_fired());
609        assert_eq!(counter.0.load(Ordering::SeqCst), 0);
610    }
611
612    #[test]
613    fn test_timeout_waiter_replaces_waker() {
614        let waiter = TimeoutWaiter::new();
615        let old_counter = Arc::new(WakeCounter(AtomicUsize::new(0)));
616        let new_counter = Arc::new(WakeCounter(AtomicUsize::new(0)));
617        let old_waker = waker_ref(&old_counter);
618        let new_waker = waker_ref(&new_counter);
619
620        waiter.register(&old_waker);
621        waiter.register(&new_waker);
622        waiter.fire();
623
624        assert!(waiter.is_fired());
625        assert_eq!(old_counter.0.load(Ordering::SeqCst), 0);
626        assert_eq!(new_counter.0.load(Ordering::SeqCst), 1);
627    }
628
629    #[test]
630    fn test_worker_waker_wakes_once() {
631        let worker_waker = Arc::new((AtomicBool::new(true), Mutex::new(()), Condvar::new()));
632
633        assert!(wake_thread_waker(&worker_waker));
634        assert!(!worker_waker.0.load(Ordering::SeqCst));
635        assert!(!wake_thread_waker(&worker_waker));
636    }
637
638    #[test]
639    fn test_worker_waker_wait_queue_skips_stale_and_wakes_one_sleeping_worker() {
640        let waits = ArrayQueue::new(4);
641        let stale = Arc::new((AtomicBool::new(false), Mutex::new(()), Condvar::new()));
642        let sleeping = Arc::new((AtomicBool::new(true), Mutex::new(()), Condvar::new()));
643
644        waits.push(stale).unwrap();
645        waits.push(sleeping.clone()).unwrap();
646
647        assert!(wake_waiting_worker(&waits));
648        assert!(!sleeping.0.load(Ordering::SeqCst));
649        assert!(!wake_waiting_worker(&waits));
650    }
651
652    #[test]
653    fn test_worker_waker_register_and_prune_stale_waiter() {
654        let waits = ArrayQueue::new(1);
655        let stale = Arc::new((AtomicBool::new(false), Mutex::new(()), Condvar::new()));
656        let current = Arc::new((AtomicBool::new(false), Mutex::new(()), Condvar::new()));
657
658        waits.push(stale).unwrap();
659        assert!(!register_waiting_worker(&waits, &current));
660        assert_eq!(prune_stale_waiting_workers(&waits), 1);
661        assert!(register_waiting_worker(&waits, &current));
662    }
663
664    #[test]
665    fn test_worker_waker_prune_keeps_live_waiter() {
666        let waits = ArrayQueue::new(1);
667        let sleeping = Arc::new((AtomicBool::new(true), Mutex::new(()), Condvar::new()));
668
669        waits.push(sleeping.clone()).unwrap();
670        assert_eq!(prune_stale_waiting_workers(&waits), 0);
671        assert!(wake_waiting_worker(&waits));
672        assert!(!sleeping.0.load(Ordering::SeqCst));
673    }
674}
675
676/*
677* 运行时托管的 AsyncTask 调度状态。
678*
679* MANAGED 区分由本库单线程/多线程运行时驱动的任务,以及通过公开
680* get_inner/set_inner 接口交给外部手工驱动的任务。SCHEDULED 表示一个尚未履行的
681* 轮询义务,RUNNING 表示独占轮询所有权,COMPLETED 表示终态。
682*/
683const ASYNC_TASK_STATE_SCHEDULED: u8 = 0b0000_0001;
684const ASYNC_TASK_STATE_RUNNING: u8 = 0b0000_0010;
685const ASYNC_TASK_STATE_COMPLETED: u8 = 0b0000_0100;
686const ASYNC_TASK_STATE_MANAGED: u8 = 0b1000_0000;
687const ASYNC_TASK_STATE_INITIAL: u8 =
688    ASYNC_TASK_STATE_MANAGED | ASYNC_TASK_STATE_SCHEDULED;
689
690/// 运行时尝试认领一个已出队任务进行轮询时的结果。
691///
692/// 只有匹配的运行时驱动可以推进托管任务,因此本枚举仅在库内可见。`Legacy` 保留公开
693/// 手工驱动路径,`Managed` 授予独占轮询所有权,`Discard` 表示陈旧或已完成的队列
694/// 引用。生成本结果的时间复杂度为 O(1),不阻塞、不分配、线程安全,也不会访问或
695/// 轮询用户 Future。
696pub(crate) enum AsyncTaskPollClaim {
697    Legacy,
698    Managed,
699    Discard,
700}
701
702enum AsyncTaskWakeAction {
703    LegacyEnqueue,
704    ManagedEnqueue,
705    Coalesced,
706}
707
708/// 单次托管 `Future::poll` 的异常安全所有者。
709///
710/// 本守卫不捕获或压制异常。若驱动在记录 `Pending` 或 `Ready` 前发生栈展开,
711/// `Drop` 会发布已完成终态,防止仍被保留或陈旧的唤醒器让任务永久停留在 `RUNNING`。
712/// 它不持锁、不分配、不阻塞、不唤醒工作线程,也不调用用户代码。创建、正常完成和
713/// 栈展开清理均为 O(1)。它会推进任务状态,因此不是纯函数,也不幂等。
714pub(crate) struct AsyncTaskPollGuard<
715    'a,
716    P: AsyncTaskPoolExt<O> + AsyncTaskPool<O, Pool = P>,
717    O: Default + 'static = (),
718> {
719    task:  &'a AsyncTask<P, O>,
720    armed: bool,
721}
722
723impl<
724    'a,
725    P: AsyncTaskPoolExt<O> + AsyncTaskPool<O, Pool = P>,
726    O: Default + 'static,
727> AsyncTaskPollGuard<'a, P, O> {
728    /// 在 `try_begin_runtime_poll` 返回 `Managed` 后创建已启用的守卫。
729    #[inline]
730    pub(crate) fn new(task: &'a AsyncTask<P, O>) -> Self {
731        AsyncTaskPollGuard {
732            task,
733            armed: true,
734        }
735    }
736
737    /// 在 Future 已恢复到任务槽位后完成一次返回 `Pending` 的轮询。
738    ///
739    /// 仅当轮询期间发生过唤醒、必须生成一个后续队列项时返回 `true`。调用方必须在
740    /// 释放其任务 `Arc` 前完成入队。时间复杂度 O(1),无锁、无分配且不阻塞。
741    #[inline]
742    pub(crate) fn finish_pending(mut self) -> bool {
743        let should_enqueue = self.task.finish_runtime_poll_pending();
744        self.armed = false;
745        should_enqueue
746    }
747
748    /// 为成功完成的 Future 发布终态。
749    ///
750    /// 调用后,迟到或陈旧的唤醒均为空操作。时间复杂度 O(1),不分配、不阻塞;
751    /// 不访问 Future、任务池、工作线程锁、回调、I/O 或 FFI。
752    #[inline]
753    pub(crate) fn finish_ready(mut self) {
754        self.task.finish_runtime_poll_ready();
755        self.armed = false;
756    }
757}
758
759impl<
760    'a,
761    P: AsyncTaskPoolExt<O> + AsyncTaskPool<O, Pool = P>,
762    O: Default + 'static,
763> Drop for AsyncTaskPollGuard<'a, P, O> {
764    fn drop(&mut self) {
765        if self.armed {
766            self.task.finish_runtime_poll_ready();
767        }
768    }
769}
770
771/// 异步任务及其运行时调度/生命周期状态。
772///
773/// # 运行时契约
774///
775/// 运行时托管任务是一次性的:其 Future 被持续轮询,首次返回 `Ready` 后永久完成。
776/// 运行时任务通过私有原子状态合并重复唤醒、排除并发轮询,并拒绝迟到唤醒。公开
777/// `get_inner`/`set_inner` 选择旧手工驱动契约,包括显式的低层替换/复用能力,因此
778/// 既有自定义驱动无需新增特征方法。
779///
780/// 构造函数只创建一个逻辑上的首次轮询义务,并不执行物理入队。运行时或自定义任务池
781/// 必须通过既有 `push*` 接口把新任务准确入队一次。`Waker` 用于在首次提交后安排后续
782/// 轮询,不能替代首次入队。
783///
784/// # 示例
785///
786/// ```
787/// use std::sync::Arc;
788/// use futures::FutureExt;
789/// use pi_async_rt::rt::{
790///     AsyncRuntime, AsyncTask, AsyncTaskPool,
791///     single_thread::SingleTaskRunner,
792/// };
793///
794/// let runner = SingleTaskRunner::<()>::default();
795/// let runtime = runner.startup().unwrap();
796/// let task = Arc::new(AsyncTask::new(
797///     runtime.alloc::<()>(),
798///     runtime.shared_pool(),
799///     0,
800///     Some(async {}.boxed()),
801/// ));
802/// runtime.shared_pool().push(task).unwrap();
803/// runner.run_once().unwrap();
804/// ```
805///
806/// `future` 锁只覆盖取出或恢复装箱的 Future;绝不会跨越 `Future::poll`、任务池访问、
807/// 工作线程通知、用户回调/析构、I/O 或 FFI。运行时状态操作为 O(1)、无分配且无锁,
808/// 但比较并交换循环不具备无等待性,在竞争唤醒/轮询状态持续推进时可能重试。任务不是
809/// `repr(C)`,不提供稳定的 FFI/Rust 布局 ABI。
810///
811/// # 布局与内存
812///
813/// 在已验收的 x86_64 目标上,加入内联调度状态后,
814/// `AsyncTask<StealableTaskPool<()>, ()>` 为 96 字节,之前的实现为 80 字节。
815/// 一字节状态跨过了该特化的 16 字节对齐边界,因此实际内联增量是 16 字节,而不是
816/// 一字节。百万个同时存活的任务会增加 16,000,000 字节(约 15.26 MiB)任务本体
817/// 存储。这是并发存活/保留成本,不会按历史累计执行过的任务数增长。
818///
819/// 状态不会新增独立堆分配。`Arc` 管理信息、装箱的 Future、任务句柄、context 负载和
820/// 队列存储仍是独立的既有成本,不包含在 96 字节本体内。分配器尺寸分级和队列
821/// 容量会使实际 RSS 与逻辑本体增量不同。若要把本体恢复到 80 字节,需要重新组织
822/// context/state 表示;该优化会触及 V8 敏感的所有权表示,必须单独设计、审查和
823/// 验证下游,因此本轮明确延期。
824///
825/// # 安全性
826///
827/// 托管轮询所有权由 `state` 同步,Future 所有权由 `future` 同步。既有 `TaskId` 和
828/// context 安全要求不变。调用方不得并发手工轮询同一任务,也不得从任务池窃取
829/// 运行时所有的任务。
830pub struct AsyncTask<
831    P: AsyncTaskPoolExt<O> + AsyncTaskPool<O>,
832    O: Default + 'static = (),
833> {
834    uid:        TaskId,                                 //任务唯一标识
835    future:     Mutex<Option<BoxFuture<'static, O>>>,   //异步任务
836    pool:       Arc<P>,                                 //异步任务池
837    priority:   usize,                                  //异步任务优先级
838    context:    Option<UnsafeCell<Box<dyn Any>>>,       //异步任务上下文
839    state:      AtomicU8,                               //内联调度状态;x86_64 上使当前特化由 80B 对齐至 96B
840}
841
842impl<
843    P: AsyncTaskPoolExt<O> + AsyncTaskPool<O>,
844    O: Default + 'static,
845> Drop for AsyncTask<P, O> {
846    fn drop(&mut self) {
847        let _ = unsafe { TaskHandle::<O>::from_raw((*self.uid.0.get() >> 64) as usize as *const ()) };
848    }
849}
850
851unsafe impl<
852    P: AsyncTaskPoolExt<O> + AsyncTaskPool<O>,
853    O: Default + 'static,
854> Send for AsyncTask<P, O> {}
855unsafe impl<
856    P: AsyncTaskPoolExt<O> + AsyncTaskPool<O>,
857    O: Default + 'static,
858> Sync for AsyncTask<P, O> {}
859
860impl<
861    P: AsyncTaskPoolExt<O> + AsyncTaskPool<O, Pool = P>,
862    O: Default + 'static,
863> ArcWake for AsyncTask<P, O> {
864    /// 发布一个可运行义务,并在需要时调度任务。
865    ///
866    /// 托管任务在已入队或运行时会合并重复唤醒。使空闲任务转为已调度状态的唤醒
867    /// 准确执行一次 `Arc` 克隆、一次 `push_keep`,并且最多通知一个工作线程。运行中
868    /// 的任务只记录延期调度;当前轮询所有者在恢复返回 `Pending` 的 Future 后入队。
869    /// 对已完成任务的唤醒是空操作。
870    ///
871    /// 无竞争时的状态处理为 O(1) 且不新增分配;`push_keep` 继续服从具体任务池既有的
872    /// 时间、容量和扩容成本。该路径不访问 Future 互斥锁,不阻塞等待,不执行用户回调、
873    /// I/O 或 FFI。无锁比较并交换循环不具备无等待性,在竞争状态持续推进时可能重试。
874    /// 与修改前相同,为保证运行时活性,任务池的 `push_keep` 必须能够接受可运行任务。
875    fn wake_by_ref(arc_self: &Arc<Self>) {
876        let notify_on_push_error = match arc_self.prepare_wake() {
877            AsyncTaskWakeAction::Coalesced => return,
878            AsyncTaskWakeAction::LegacyEnqueue => true,
879            AsyncTaskWakeAction::ManagedEnqueue => false,
880        };
881
882        let pool = arc_self.get_pool();
883        let pushed = pool.push_keep(arc_self.clone()).is_ok();
884        if pushed || notify_on_push_error {
885            notify_runtime_task_pool(pool);
886        }
887    }
888}
889
890/// 在可运行任务已物理入队后,最多通知一个工作线程。
891///
892/// 多线程任务池使用 `wake_waiting_worker` 实现的有界陈旧等待者扫描;直接/单线程
893/// 任务池使用既有受谓词保护的线程唤醒器。本辅助函数不入队也不轮询任务。直接
894/// 唤醒器路径为 O(1),既有 `waits` 注册表路径为有界 O(工作线程数量);除既有短谓词锁
895/// 外,不分配也不阻塞。托管路径必须只在成功入队后调用;兼容手工路径在
896/// `push_keep` 返回错误时仍会调用,以保持修复前的通知语义。
897#[inline]
898pub(crate) fn notify_runtime_task_pool<
899    P: AsyncTaskPoolExt<O> + AsyncTaskPool<O, Pool = P>,
900    O: Default + 'static,
901>(pool: &P) {
902    if let Some(waits) = pool.get_waits() {
903        let _ = wake_waiting_worker(waits);
904    } else if let Some(thread_waker) = pool.get_thread_waker() {
905        let _ = wake_thread_waker(thread_waker);
906    }
907}
908
909/// 把托管任务轮询期间观察到的一个延期唤醒入队。
910///
911/// 调用方必须先恢复 `Pending` Future,并完成 `RUNNING -> SCHEDULED` 转换。这里保留
912/// 一次 `Arc` 克隆,因为公开任务池特征会消费队列参数,并且出错时不能返还所有权;
913/// 其成本与修复前第一次唤醒相同,同时消除了全部重复唤醒克隆。入队成功后最多
914/// 通知一个工作线程。复杂度为 O(1) 加所选任务池声明的队列复杂度;不会进入 Future
915/// 锁、用户代码、I/O 或 FFI。
916#[inline]
917pub(crate) fn requeue_runtime_task<
918    P: AsyncTaskPoolExt<O> + AsyncTaskPool<O, Pool = P>,
919    O: Default + 'static,
920>(pool: &P, task: &Arc<AsyncTask<P, O>>) {
921    if pool.push_keep(task.clone()).is_ok() {
922        notify_runtime_task_pool(pool);
923    }
924}
925
926impl<
927    P: AsyncTaskPoolExt<O> + AsyncTaskPool<O, Pool = P>,
928    O: Default + 'static,
929> AsyncTask<P, O> {
930    /// 构造一个一次性异步任务。
931    ///
932    /// 任务初始带有一个托管轮询义务。若任务进入本库运行时,私有驱动会保证唤醒
933    /// 合并和独占轮询。调用公开 `get_inner` 或 `set_inner` 会选择旧手工驱动模式,
934    /// 但不改变这两个方法的签名或值语义。
935    ///
936    /// 构造时间复杂度为 O(1),除参数已经拥有的值外不新增分配;不执行队列操作、
937    /// 唤醒、轮询、加锁、I/O 或回调。每个值拥有一个 TaskId,因此构造不是幂等操作。
938    pub fn new(uid: TaskId,
939               pool: Arc<P>,
940               priority: usize,
941               future: Option<BoxFuture<'static, O>>) -> AsyncTask<P, O> {
942        AsyncTask {
943            uid,
944            future: Mutex::new(future),
945            pool,
946            priority,
947            context: None,
948            state: AtomicU8::new(ASYNC_TASK_STATE_INITIAL),
949        }
950    }
951
952    /// 构造一个带调用方 context 的一次性任务。
953    ///
954    /// 调度和唤醒语义与 `AsyncTask::new` 相同。本函数为 `context` 执行一次 `Box`
955    /// 分配,因此构造的时间和空间复杂度均为 O(1)。它不入队、不轮询、不唤醒、
956    /// 不获取运行时锁、不执行用户代码、不执行 I/O,也不接触 FFI。保存的 context
957    /// 继续服从既有所有者线程/context 访问契约。
958    pub fn with_context<C: 'static>(uid: TaskId,
959                                    pool: Arc<P>,
960                                    priority: usize,
961                                    future: Option<BoxFuture<'static, O>>,
962                                    context: C) -> AsyncTask<P, O> {
963        let any = Box::new(context);
964
965        AsyncTask {
966            uid,
967            future: Mutex::new(future),
968            pool,
969            priority,
970            context: Some(UnsafeCell::new(any)),
971            state: AtomicU8::new(ASYNC_TASK_STATE_INITIAL),
972        }
973    }
974
975    /// 构造一个绑定到 `runtime` 且携带调用方 context 的任务。
976    ///
977    /// 返回值随后由本库运行时驱动消费时属于托管任务,这也包括
978    /// `pi_v8::VmTaskPool` 等外部定时器适配器。公开 get/set 仍会选择旧手工驱动。
979    /// 构造为 O(1),执行既有 `runtime.alloc` TaskHandle 分配、一次 context `Box`
980    /// 分配和一次共享任务池 `Arc` 克隆。分配失败继续保持这些既有操作的进程级行为。
981    /// 本函数不入队、不轮询、不唤醒、不阻塞、不执行用户代码、I/O 或 FFI。
982    pub fn with_runtime_and_context<RT, C>(runtime: &RT,
983                                           priority: usize,
984                                           future: Option<BoxFuture<'static, O>>,
985                                           context: C) -> AsyncTask<P, O>
986        where RT: AsyncRuntime<O, Pool = P>,
987              C: Send + 'static {
988        let any = Box::new(context);
989
990        AsyncTask {
991            uid: runtime.alloc::<O>(),
992            future: Mutex::new(future),
993            pool: runtime.shared_pool(),
994            priority,
995            context: Some(UnsafeCell::new(any)),
996            state: AtomicU8::new(ASYNC_TASK_STATE_INITIAL),
997        }
998    }
999
1000    /// 判断一次唤醒是否需要一个物理队列项。
1001    ///
1002    /// 即使 `next == current`,成功的比较并交换也使用 `AcqRel`。该同值读改写会发布每次被合并
1003    /// 唤醒之前的写入;后续轮询认领在读取 Future 关联共享状态前获取最新原子
1004    /// 修改。弱比较并交换可能伪失败,因此必须使用重试循环。本操作无锁但不具备无等待性;
1005    /// 不涉及互斥锁、分配、克隆、队列或用户代码。
1006    #[inline]
1007    fn prepare_wake(&self) -> AsyncTaskWakeAction {
1008        let mut current = self.state.load(Ordering::Acquire);
1009        loop {
1010            if current & ASYNC_TASK_STATE_MANAGED == 0 {
1011                return AsyncTaskWakeAction::LegacyEnqueue;
1012            }
1013            if current & ASYNC_TASK_STATE_COMPLETED != 0 {
1014                return AsyncTaskWakeAction::Coalesced;
1015            }
1016
1017            let next = current | ASYNC_TASK_STATE_SCHEDULED;
1018            match self.state.compare_exchange_weak(
1019                current,
1020                next,
1021                Ordering::AcqRel,
1022                Ordering::Acquire,
1023            ) {
1024                Ok(_) => {
1025                    if current & (ASYNC_TASK_STATE_SCHEDULED | ASYNC_TASK_STATE_RUNNING) != 0 {
1026                        return AsyncTaskWakeAction::Coalesced;
1027                    }
1028                    return AsyncTaskWakeAction::ManagedEnqueue;
1029                },
1030                Err(actual) => current = actual,
1031            }
1032        }
1033    }
1034
1035    /// 认领一个尚未履行的运行时轮询义务。
1036    ///
1037    /// `Managed` 通过原子清除 `SCHEDULED` 并设置 `RUNNING` 授予独占所有权。
1038    /// `Discard` 表示物理队列项陈旧、重复或已完成,必须在不接触 Future 的情况下
1039    /// 丢弃。`Legacy` 委托给保持不变的公开取出/轮询/恢复行为。
1040    ///
1041    /// 无竞争路径为 O(1),不分配、无锁且不阻塞。它不具备无等待性:弱比较并交换循环
1042    /// 可能因竞争或伪失败而重试。`AcqRel` 与唤醒发布同步;不接触用户代码、队列、
1043    /// 定时器、工作线程锁、I/O 或 FFI。
1044    #[inline]
1045    pub(crate) fn try_begin_runtime_poll(&self) -> AsyncTaskPollClaim {
1046        let mut current = self.state.load(Ordering::Acquire);
1047        loop {
1048            if current & ASYNC_TASK_STATE_MANAGED == 0 {
1049                return AsyncTaskPollClaim::Legacy;
1050            }
1051            if current & ASYNC_TASK_STATE_COMPLETED != 0
1052                || current & ASYNC_TASK_STATE_SCHEDULED == 0
1053                || current & ASYNC_TASK_STATE_RUNNING != 0
1054            {
1055                return AsyncTaskPollClaim::Discard;
1056            }
1057
1058            let next =
1059                (current & !ASYNC_TASK_STATE_SCHEDULED) | ASYNC_TASK_STATE_RUNNING;
1060            match self.state.compare_exchange_weak(
1061                current,
1062                next,
1063                Ordering::AcqRel,
1064                Ordering::Acquire,
1065            ) {
1066                Ok(_) => return AsyncTaskPollClaim::Managed,
1067                Err(actual) => current = actual,
1068            }
1069        }
1070    }
1071
1072    /// 在成功认领托管轮询后取出 Future。
1073    ///
1074    /// 只有本库运行时驱动可以调用本方法。互斥锁临界区只包含 `Option::take`,绝不与
1075    /// `Future::poll`、唤醒、任务池访问、工作线程通知或用户代码重叠。时间复杂度
1076    /// O(1),不分配。
1077    #[inline]
1078    pub(crate) fn take_inner_for_runtime_poll(&self) -> Option<BoxFuture<'static, O>> {
1079        self.future.lock().take()
1080    }
1081
1082    /// 恢复一个返回 `Pending` 的托管 Future。
1083    ///
1084    /// 必须在清除 `RUNNING` 前恢复,否则后续工作线程可能认领延期调度却观察到 Future
1085    /// 缺失。短互斥锁临界区只包含 `Option::replace`。若非法并发手工访问导致槽位中
1086    /// 已有值,该值只会在解锁后析构。时间复杂度 O(1),不新增分配,不在锁内执行回调,
1087    /// 也不执行队列操作、I/O 或 FFI。
1088    #[inline]
1089    pub(crate) fn restore_inner_after_runtime_poll(
1090        &self,
1091        inner: BoxFuture<'static, O>,
1092    ) {
1093        let replaced = {
1094            let mut future = self.future.lock();
1095            future.replace(inner)
1096        };
1097        drop(replaced);
1098    }
1099
1100    /// 完成一次托管 `Pending` 状态转换。
1101    ///
1102    /// 返回任务运行期间是否有唤醒设置了 `SCHEDULED`;调用方随后准确创建一个队列项。
1103    /// Future 恢复和互斥锁解锁先行发生于本次 `AcqRel` 转换;下一次认领通过
1104    /// `Acquire` 获取前述写入。
1105    #[inline]
1106    fn finish_runtime_poll_pending(&self) -> bool {
1107        let mut current = self.state.load(Ordering::Acquire);
1108        loop {
1109            if current & ASYNC_TASK_STATE_MANAGED == 0
1110                || current & ASYNC_TASK_STATE_COMPLETED != 0
1111                || current & ASYNC_TASK_STATE_RUNNING == 0
1112            {
1113                return false;
1114            }
1115
1116            let next = current & !ASYNC_TASK_STATE_RUNNING;
1117            match self.state.compare_exchange_weak(
1118                current,
1119                next,
1120                Ordering::AcqRel,
1121                Ordering::Acquire,
1122            ) {
1123                Ok(_) => return current & ASYNC_TASK_STATE_SCHEDULED != 0,
1124                Err(actual) => current = actual,
1125            }
1126        }
1127    }
1128
1129    /// 发布托管任务终态。
1130    ///
1131    /// 原子存储会有意同时清除 `RUNNING` 和并发延期的 `SCHEDULED` 位。竞争唤醒要么先于
1132    /// 本次 `Release` 存储发布并被完成状态吸收,要么观察到 `COMPLETED` 后成为空操作。
1133    /// 本操作为 O(1),不分配、无锁且不阻塞。
1134    #[inline]
1135    fn finish_runtime_poll_ready(&self) {
1136        self.state.store(
1137            ASYNC_TASK_STATE_MANAGED | ASYNC_TASK_STATE_COMPLETED,
1138            Ordering::Release,
1139        );
1140    }
1141
1142    /// 选择既有公开手工驱动行为。
1143    ///
1144    /// 历史上,公开 get/set 调用方拥有取出/轮询/恢复协议,无法调用新的私有完成
1145    /// 钩子。因此暴露 Future 前会把空闲或已入队的托管任务原子切换为兼容手工模式。已在
1146    /// 正在轮询的运行时任务不会降级。`allow_completed` 只供公开 `set_inner` 使用,以
1147    /// 保留显式低层任务复用。
1148    ///
1149    /// 无竞争时为 O(1),不分配、无锁且不阻塞;竞争状态转换下不具备无等待性。不访问
1150    /// 互斥锁、任务池或用户代码。
1151    #[inline]
1152    fn select_legacy_manual_driver(&self, allow_completed: bool) {
1153        let mut current = self.state.load(Ordering::Acquire);
1154        loop {
1155            if current & ASYNC_TASK_STATE_MANAGED == 0
1156                || current & ASYNC_TASK_STATE_RUNNING != 0
1157                || (!allow_completed && current & ASYNC_TASK_STATE_COMPLETED != 0)
1158            {
1159                return;
1160            }
1161
1162            match self.state.compare_exchange_weak(
1163                current,
1164                0,
1165                Ordering::AcqRel,
1166                Ordering::Acquire,
1167            ) {
1168                Ok(_) => return,
1169                Err(actual) => current = actual,
1170            }
1171        }
1172    }
1173
1174    /// 检查是否允许唤醒
1175    pub fn is_enable_wakeup(&self) -> bool {
1176        self.uid.exist_waker::<O>()
1177    }
1178
1179    /// 为外部/手工任务驱动取出装箱的 Future。
1180    ///
1181    /// 本方法在取出 Future 前选择兼容手工调度,因此既有自定义驱动无需新增特征
1182    /// 方法即可保留原有唤醒后调用 `push_keep` 的行为。本方法不是运行时驱动入口。
1183    ///
1184    /// 当其它驱动拥有 Future 或任务已完成时返回 `None`。时间复杂度 O(1),不分配;
1185    /// 只获取任务的短 Future 互斥锁,可能与其它 get/set 短暂竞争。不得与本库运行时驱动
1186    /// 并发调用,也不得用于并发轮询同一任务。本操作非纯、非幂等,并转移 Future
1187    /// 所有权。
1188    pub fn get_inner(&self) -> Option<BoxFuture<'static, O>> {
1189        self.select_legacy_manual_driver(false);
1190        self.future.lock().take()
1191    }
1192
1193    /// 为外部/手工任务驱动替换装箱的 Future。
1194    ///
1195    /// 本方法选择兼容手工调度,也允许显式复用已完成的低层任务,从而保留旧公开 get/set
1196    /// 能力。运行时所有的 `Pending` 恢复使用私有辅助函数,仍保持托管。复用已完成任务
1197    /// 前,手工驱动必须回收全部旧唤醒器;在保留的兼容手工契约下,旧唤醒器否则可能
1198    /// 把替换后的 Future 入队。
1199    ///
1200    /// 时间复杂度 O(1),除调用方拥有的装箱值外不分配;只获取短 Future 互斥锁,不轮询、
1201    /// 不入队也不唤醒。本操作非纯且非幂等。被替换的 Future 会先移出互斥锁临界区,
1202    /// 再在调用线程析构,因此其析构函数不能在持有 Future 锁时重入本任务。
1203    pub fn set_inner(&self, inner: Option<BoxFuture<'static, O>>) {
1204        self.select_legacy_manual_driver(true);
1205        let replaced = {
1206            let mut future = self.future.lock();
1207            std::mem::replace(&mut *future, inner)
1208        };
1209        drop(replaced);
1210    }
1211
1212    /// 获取任务的所有者
1213    #[inline]
1214    pub fn owner(&self) -> usize {
1215        unsafe {
1216            *self.uid.0.get() as usize
1217        }
1218    }
1219
1220    /// 获取异步任务优先级
1221    #[inline]
1222    pub fn priority(&self) -> usize {
1223        self.priority
1224    }
1225
1226    //判断异步任务是否有上下文
1227    pub fn exist_context(&self) -> bool {
1228        self.context.is_some()
1229    }
1230
1231    //获取异步任务上下文的只读引用
1232    pub fn get_context<C: Send + 'static>(&self) -> Option<&C> {
1233        if let Some(context) = &self.context {
1234            //存在上下文
1235            let any = unsafe { &*context.get() };
1236            return <dyn Any>::downcast_ref::<C>(&**any);
1237        }
1238
1239        None
1240    }
1241
1242    //获取异步任务上下文的可写引用
1243    pub fn get_context_mut<C: Send + 'static>(&self) -> Option<&mut C> {
1244        if let Some(context) = &self.context {
1245            //存在上下文
1246            let any = unsafe { &mut *context.get() };
1247            return <dyn Any>::downcast_mut::<C>(&mut **any);
1248        }
1249
1250        None
1251    }
1252
1253    //设置异步任务上下文,返回上一个异步任务上下文
1254    pub fn set_context<C: Send + 'static>(&self, new: C) {
1255        if let Some(context) = &self.context {
1256            //存在上一个上下文,则释放上一个上下文
1257            let _ = unsafe { &*context.get() };
1258
1259            //设置新的上下文
1260            let any: Box<dyn Any + 'static> = Box::new(new);
1261            unsafe { *context.get() = any; }
1262        }
1263    }
1264
1265    //获取异步任务的任务池
1266    pub fn get_pool(&self) -> &P {
1267        self.pool.as_ref()
1268    }
1269}
1270
1271///
1272/// 异步任务池
1273///
1274pub trait AsyncTaskPool<O: Default + 'static = ()>: Default + Send + Sync + 'static {
1275    type Pool: AsyncTaskPoolExt<O> + AsyncTaskPool<O>;
1276
1277    /// 获取绑定的线程唯一id
1278    fn get_thread_id(&self) -> usize;
1279
1280    /// 获取当前异步任务池内任务数量
1281    fn len(&self) -> usize;
1282
1283    /// 将异步任务加入异步任务池
1284    fn push(&self, task: Arc<AsyncTask<Self::Pool, O>>) -> Result<()>;
1285
1286    /// 将异步任务加入本地异步任务池
1287    fn push_local(&self, task: Arc<AsyncTask<Self::Pool, O>>) -> Result<()>;
1288
1289    /// 将指定了优先级的异步任务加入任务池
1290    fn push_priority(&self,
1291                     priority: usize,
1292                     task: Arc<AsyncTask<Self::Pool, O>>) -> Result<()>;
1293
1294    /// 将一次唤醒或托管延期唤醒产生的一个可运行任务入队。
1295    ///
1296    /// 每次成功调用准确拥有一个物理队列项。托管 `AsyncTask` 状态会在调用本方法前
1297    /// 合并重复唤醒;自定义任务池不得自行增加轮询或重复队列项。返回 `Err` 表示
1298    /// 运行时无法保证该次唤醒的进度,因为本特征会消费任务 `Arc`,不能返还所有权。
1299    ///
1300    /// 供运行时使用的实现必须线程安全,不得与 `Future::poll` 重入,并应在唤醒热路径
1301    /// 上保持不阻塞/O(1)。实现不得调用用户代码或通知多个工作线程;工作线程通知由
1302    /// 运行时在成功入队后执行。
1303    fn push_keep(&self, task: Arc<AsyncTask<Self::Pool, O>>) -> Result<()>;
1304
1305    /// 尝试从异步任务池中弹出一个异步任务
1306    fn try_pop(&self) -> Option<Arc<AsyncTask<Self::Pool, O>>>;
1307
1308    /// 尝试从异步任务池中弹出所有异步任务
1309    fn try_pop_all(&self) -> IntoIter<Arc<AsyncTask<Self::Pool, O>>>;
1310
1311    /// 获取本地线程的唤醒器
1312    fn get_thread_waker(&self) -> Option<&Arc<(AtomicBool, Mutex<()>, Condvar)>>;
1313}
1314
1315///
1316/// 异步任务池扩展
1317///
1318pub trait AsyncTaskPoolExt<O: Default + 'static = ()>: Send + Sync + 'static {
1319    /// 设置待唤醒的工作者唤醒器队列
1320    fn set_waits(&mut self,
1321                 _waits: Arc<ArrayQueue<Arc<(AtomicBool, Mutex<()>, Condvar)>>>) {}
1322
1323    /// 获取待唤醒的工作者唤醒器队列
1324    fn get_waits(&self) -> Option<&Arc<ArrayQueue<Arc<(AtomicBool, Mutex<()>, Condvar)>>>> {
1325        //默认没有待唤醒的工作者唤醒器队列
1326        None
1327    }
1328
1329    /// 获取空闲的工作者的数量,这个数量大于0,表示可以新开线程来运行可分派的工作者
1330    fn idler_len(&self) -> usize {
1331        //默认不分派
1332        0
1333    }
1334
1335    /// 分派一个空闲的工作者
1336    fn spawn_worker(&self) -> Option<usize> {
1337        //默认不分派
1338        None
1339    }
1340
1341    /// 获取工作者的数量
1342    fn worker_len(&self) -> usize {
1343        //默认工作者数量和本机逻辑核数相同
1344        #[cfg(not(target_arch = "wasm32"))]
1345        return num_cpus::get();
1346        #[cfg(target_arch = "wasm32")]
1347        return 1;
1348    }
1349
1350    /// 获取缓冲区的任务数量,缓冲区任务是未分配给工作者的任务
1351    fn buffer_len(&self) -> usize {
1352        //默认没有缓冲区
1353        0
1354    }
1355
1356    /// 设置当前绑定本地线程的唤醒器
1357    fn set_thread_waker(&mut self, _thread_waker: Arc<(AtomicBool, Mutex<()>, Condvar)>) {
1358        //默认不设置
1359    }
1360
1361    /// 复制当前绑定本地线程的唤醒器
1362    fn clone_thread_waker(&self) -> Option<Arc<(AtomicBool, Mutex<()>, Condvar)>> {
1363        //默认不复制
1364        None
1365    }
1366
1367    /// 关闭当前工作者
1368    fn close_worker(&self) {
1369        //默认不允许关闭工作者
1370    }
1371}
1372
1373///
1374/// 异步运行时
1375///
1376pub trait AsyncRuntime<O: Default + 'static = ()>: Clone + Send + Sync + 'static {
1377    type Pool: AsyncTaskPoolExt<O> + AsyncTaskPool<O, Pool = Self::Pool>;
1378
1379    /// 共享运行时内部任务池
1380    fn shared_pool(&self) -> Arc<Self::Pool>;
1381
1382    /// 获取当前异步运行时的唯一id
1383    fn get_id(&self) -> usize;
1384
1385    /// 获取当前异步运行时待处理任务数量
1386    fn wait_len(&self) -> usize;
1387
1388    /// 获取当前异步运行时任务数量
1389    fn len(&self) -> usize;
1390
1391    /// 分配异步任务的唯一id
1392    fn alloc<R: 'static>(&self) -> TaskId;
1393
1394    /// 派发一个指定的异步任务到异步运行时
1395    fn spawn<F>(&self, future: F) -> Result<TaskId>
1396        where F: Future<Output = O> + Send + 'static;
1397
1398    /// 派发一个异步任务到本地异步运行时,如果本地没有本异步运行时,则会派发到当前运行时中
1399    fn spawn_local<F>(&self, future: F) -> Result<TaskId>
1400        where F: Future<Output = O> + Send + 'static;
1401
1402    /// 派发一个指定优先级的异步任务到异步运行时
1403    fn spawn_priority<F>(&self, priority: usize, future: F) -> Result<TaskId>
1404        where F: Future<Output = O> + Send + 'static;
1405
1406    /// 派发一个异步任务到异步运行时,并立即让出任务的当前运行
1407    fn spawn_yield<F>(&self, future: F) -> Result<TaskId>
1408        where F: Future<Output = O> + Send + 'static;
1409
1410    /// 派发一个在指定时间后执行的异步任务到异步运行时,时间单位ms
1411    fn spawn_timing<F>(&self, future: F, time: usize) -> Result<TaskId>
1412        where F: Future<Output = O> + Send + 'static;
1413
1414    /// 派发一个指定任务唯一id的异步任务到异步运行时
1415    fn spawn_by_id<F>(&self, task_id: TaskId, future: F) -> Result<()>
1416        where F: Future<Output = O> + Send + 'static;
1417
1418    /// 派发一个指定任务唯一id的异步任务到本地异步运行时,如果本地没有本异步运行时,则会派发到当前运行时中
1419    fn spawn_local_by_id<F>(&self, task_id: TaskId, future: F) -> Result<()>
1420        where F: Future<Output = O> + Send + 'static;
1421
1422    /// 派发一个指定任务唯一id和任务优先级的异步任务到异步运行时
1423    fn spawn_priority_by_id<F>(&self,
1424                               task_id: TaskId,
1425                               priority: usize,
1426                               future: F) -> Result<()>
1427        where F: Future<Output = O> + Send + 'static;
1428
1429    /// 派发一个指定任务唯一id的异步任务到异步运行时,并立即让出任务的当前运行
1430    fn spawn_yield_by_id<F>(&self, task_id: TaskId, future: F) -> Result<()>
1431        where F: Future<Output = O> + Send + 'static;
1432
1433    /// 派发一个指定任务唯一id和在指定时间后执行的异步任务到异步运行时,时间单位ms
1434    fn spawn_timing_by_id<F>(&self,
1435                             task_id: TaskId,
1436                             future: F,
1437                             time: usize) -> Result<()>
1438        where F: Future<Output = O> + Send + 'static;
1439
1440    /// 挂起指定唯一id的异步任务
1441    fn pending<Output: 'static>(&self, task_id: &TaskId, waker: Waker) -> Poll<Output>;
1442
1443    /// 唤醒指定唯一id的异步任务
1444    fn wakeup<Output: 'static>(&self, task_id: &TaskId);
1445
1446    /// 挂起当前异步运行时的当前任务,并在指定的其它运行时上派发一个指定的异步任务,等待其它运行时上的异步任务完成后,唤醒当前运行时的当前任务,并返回其它运行时上的异步任务的值
1447    fn wait<V: Send + 'static>(&self) -> AsyncWait<V>;
1448
1449    /// 挂起当前异步运行时的当前任务,并在多个其它运行时上执行多个其它任务,其中任意一个任务完成,则唤醒当前运行时的当前任务,并返回这个已完成任务的值,而其它未完成的任务的值将被忽略
1450    fn wait_any<V: Send + 'static>(&self, capacity: usize) -> AsyncWaitAny<V>;
1451
1452    /// 挂起当前异步运行时的当前任务,并在多个其它运行时上执行多个其它任务,任务返回后需要通过用户指定的检查回调进行检查,其中任意一个任务检查通过,则唤醒当前运行时的当前任务,并返回这个已完成任务的值,而其它未完成或未检查通过的任务的值将被忽略,如果所有任务都未检查通过,则强制唤醒当前运行时的当前任务
1453    fn wait_any_callback<V: Send + 'static>(&self, capacity: usize) -> AsyncWaitAnyCallback<V>;
1454
1455    /// 构建用于派发多个异步任务到指定运行时的映射归并,需要指定映射归并的容量
1456    fn map_reduce<V: Send + 'static>(&self, capacity: usize) -> AsyncMapReduce<V>;
1457
1458    /// 挂起当前异步运行时的当前任务,等待指定的时间后唤醒当前任务
1459    fn timeout(&self, timeout: usize) -> BoxFuture<'static, ()>;
1460
1461    /// 立即让出当前任务的执行
1462    fn yield_now(&self) -> BoxFuture<'static, ()>;
1463
1464    /// 生成一个异步管道,输入指定流,输入流的每个值通过过滤器生成输出流的值
1465    fn pipeline<S, SO, F, FO>(&self, input: S, filter: F) -> BoxStream<'static, FO>
1466        where S: Stream<Item = SO> + Send + 'static,
1467              SO: Send + 'static,
1468              F: FnMut(SO) -> AsyncPipelineResult<FO> + Send + 'static,
1469              FO: Send + 'static;
1470
1471    /// 关闭异步运行时,返回请求关闭是否成功
1472    fn close(&self) -> bool;
1473}
1474
1475///
1476/// 异步运行时扩展
1477///
1478pub trait AsyncRuntimeExt<O: Default + 'static = ()> {
1479    /// 派发一个指定的异步任务到异步运行时,并指定异步任务的初始化上下文
1480    fn spawn_with_context<F, C>(&self,
1481                                task_id: TaskId,
1482                                future: F,
1483                                context: C) -> Result<()>
1484        where F: Future<Output = O> + Send + 'static,
1485              C: 'static;
1486
1487    /// 派发一个在指定时间后执行的异步任务到异步运行时,并指定异步任务的初始化上下文,时间单位ms
1488    fn spawn_timing_with_context<F, C>(&self,
1489                                       task_id: TaskId,
1490                                       future: F,
1491                                       context: C,
1492                                       time: usize) -> Result<()>
1493        where F: Future<Output = O> + Send + 'static,
1494              C: Send + 'static;
1495
1496    /// 立即创建一个指定任务池的异步运行时,并执行指定的异步任务,阻塞当前线程,等待异步任务完成后返回
1497    fn block_on<F>(&self, future: F) -> Result<F::Output>
1498        where F: Future + Send + 'static,
1499              <F as Future>::Output: Default + Send + 'static;
1500}
1501
1502///
1503/// 异步运行时构建器
1504///
1505pub struct AsyncRuntimeBuilder<O: Default + 'static = ()>(PhantomData<O>);
1506
1507impl<O: Default + 'static> AsyncRuntimeBuilder<O> {
1508    /// 构建默认的工作者异步运行时
1509    pub fn default_worker_thread(worker_name: Option<&str>,
1510                                 worker_stack_size: Option<usize>,
1511                                 worker_sleep_timeout: Option<u64>,
1512                                 worker_loop_interval: Option<Option<u64>>) -> WorkerRuntime<O> {
1513        let runner = WorkerTaskRunner::default();
1514
1515        let thread_name = if let Some(name) = worker_name {
1516            name
1517        } else {
1518            //默认的线程名称
1519            "Default-Single-Worker"
1520        };
1521        let thread_stack_size = if let Some(size) = worker_stack_size {
1522            size
1523        } else {
1524            //默认的线程堆栈大小
1525            2 * 1024 * 1024
1526        };
1527        let sleep_timeout = if let Some(timeout) = worker_sleep_timeout {
1528            timeout
1529        } else {
1530            //默认的线程休眠时长
1531            1
1532        };
1533        let loop_interval = if let Some(interval) = worker_loop_interval {
1534            interval
1535        } else {
1536            //默认的线程循环间隔时长
1537            None
1538        };
1539
1540        //创建线程并在线程中执行异步运行时
1541        let clock = Clock::new();
1542        let runner_copy = runner.clone();
1543        let rt_copy = runner.get_runtime();
1544        let rt = runner.startup(
1545            thread_name,
1546            thread_stack_size,
1547            sleep_timeout,
1548            loop_interval,
1549            move || {
1550                let last = clock.recent();
1551                match runner_copy.run_once() {
1552                    Err(e) => {
1553                        panic!("Run runner failed, reason: {:?}", e);
1554                    },
1555                    Ok(len) => {
1556                        (len == 0,
1557                         clock
1558                             .recent()
1559                             .duration_since(last))
1560                    },
1561                }
1562            },
1563            move || {
1564                rt_copy.wait_len() + rt_copy.len()
1565            },
1566        );
1567
1568        rt
1569    }
1570
1571    /// 构建自定义的工作者异步运行时
1572    pub fn custom_worker_thread<P, F0, F1>(pool: P,
1573                                           worker_handle: Arc<AtomicBool>,
1574                                           worker_condvar: Arc<(AtomicBool, Mutex<()>, Condvar)>,
1575                                           thread_name: &str,
1576                                           thread_stack_size: usize,
1577                                           sleep_timeout: u64,
1578                                           loop_interval: Option<u64>,
1579                                           loop_func: F0,
1580                                           get_queue_len: F1) -> WorkerRuntime<O, P>
1581        where P: AsyncTaskPoolExt<O> + AsyncTaskPool<O, Pool = P>,
1582              F0: Fn() -> (bool, Duration) + Send + 'static,
1583              F1: Fn() -> usize + Send + 'static {
1584        let runner = WorkerTaskRunner::new(pool,
1585                                           worker_handle,
1586                                           worker_condvar);
1587
1588        //创建线程并在线程中执行异步运行时
1589        let rt_copy = runner.get_runtime();
1590        let rt = runner.startup(
1591            thread_name,
1592            thread_stack_size,
1593            sleep_timeout,
1594            loop_interval,
1595            loop_func,
1596            move || {
1597                rt_copy.wait_len() + get_queue_len()
1598            },
1599        );
1600
1601        rt
1602    }
1603
1604    /// 构建默认的多线程异步运行时。
1605    ///
1606    /// 说明:
1607    /// - 这是 `AsyncRuntimeBuilder` 对外提供的默认多线程 runtime 构建入口。
1608    /// - 本轮保持函数签名、返回类型和既有启动语义不变。
1609    /// - 当调用方显式传入大于 0 的 `worker_size` 时,会同时创建相同 worker slot 数量
1610    ///   的 `StealableTaskPool`,避免启动 worker 数量大于 pool 内实际 worker slot 时,
1611    ///   worker 线程中 `clone_thread_waker().unwrap()` panic。
1612    /// - `worker_size=Some(0)` 保留旧语义:仍通过 builder 的 `init_worker_size(0)` 和
1613    ///   `set_worker_limit(0, 0)` 兜底,不会把 0 直接传给 `StealableTaskPool::with`。
1614    ///
1615    /// 参数:
1616    /// - `worker_prefix`:worker 线程名前缀;`None` 使用 builder 默认值。
1617    /// - `worker_stack_size`:worker 栈大小;`None` 使用默认 2 MiB。
1618    /// - `worker_size`:固定 worker 数量;`None` 使用默认 builder 和默认 pool 尺寸;
1619    ///   `Some(0)` 使用 builder 原有的默认初始 worker 兜底语义。
1620    /// - `worker_sleep_timeout`:worker 空闲休眠最长时长,单位 ms;`None` 使用默认值。
1621    ///
1622    /// 返回:
1623    /// - 已启动的 `MultiTaskRuntime<O>`。
1624    ///
1625    /// 边界条件:
1626    /// - `worker_size=Some(size > 0)` 时,实际 worker 数和 pool worker slot 数保持一致。
1627    /// - `worker_size=Some(0)` 时不创建 0 worker pool,避免兼容性回归。
1628    /// - `worker_size=None` 时不改变原默认构建路径。
1629    ///
1630    /// 性能:
1631    /// - 构建时间 O(W),W 为 worker 数;空间 O(W)。
1632    /// - 该函数不是任务调度热路径。
1633    ///
1634    /// 副作用与安全性:
1635    /// - 非纯函数,会创建任务池、runtime 和 worker 线程。
1636    /// - 不阻塞等待 worker 完成;不执行用户 future。
1637    /// - 线程安全由 `MultiTaskRuntimeBuilder::build` 和底层任务池保证。
1638    pub fn default_multi_thread(worker_prefix: Option<&str>,
1639                                worker_stack_size: Option<usize>,
1640                                worker_size: Option<usize>,
1641                                worker_sleep_timeout: Option<u64>) -> MultiTaskRuntime<O> {
1642        let mut builder = if let Some(size) = worker_size.filter(|size| *size > 0) {
1643            let pool = StealableTaskPool::with(size,
1644                                               65535,
1645                                               [1, 1],
1646                                               3000);
1647            MultiTaskRuntimeBuilder::new(pool)
1648                .thread_stack_size(2 * 1024 * 1024)
1649                .set_timer_interval(1)
1650        } else {
1651            MultiTaskRuntimeBuilder::default()
1652        };
1653
1654        if let Some(size) = worker_size {
1655            builder = builder
1656                .init_worker_size(size)
1657                .set_worker_limit(size, size);
1658        }
1659        if let Some(thread_prefix) = worker_prefix {
1660            builder = builder.thread_prefix(thread_prefix);
1661        }
1662        if let Some(thread_stack_size) = worker_stack_size {
1663            builder = builder.thread_stack_size(thread_stack_size);
1664        }
1665        if let Some(sleep_timeout) = worker_sleep_timeout {
1666            builder = builder.set_timeout(sleep_timeout);
1667        }
1668
1669        builder.build()
1670    }
1671
1672    /// 构建自定义的多线程异步运行时
1673    pub fn custom_multi_thread<P>(pool: P,
1674                                  worker_prefix: &str,
1675                                  worker_stack_size: usize,
1676                                  worker_size: usize,
1677                                  worker_sleep_timeout: u64,
1678                                  worker_timer_interval: usize) -> MultiTaskRuntime<O, P>
1679        where P: AsyncTaskPoolExt<O> + AsyncTaskPool<O, Pool = P> {
1680        MultiTaskRuntimeBuilder::new(pool)
1681            .thread_prefix(worker_prefix)
1682            .thread_stack_size(worker_stack_size)
1683            .init_worker_size(worker_size)
1684            .set_worker_limit(worker_size, worker_size)
1685            .set_timeout(worker_sleep_timeout)
1686            .set_timer_interval(worker_timer_interval)
1687            .build()
1688    }
1689}
1690
1691/// 绑定指定异步运行时到本地线程
1692pub fn bind_local_thread<O: Default + 'static>(runtime: LocalAsyncRuntime<O>) {
1693    match PI_ASYNC_LOCAL_THREAD_ASYNC_RUNTIME.try_with(move |rt| {
1694        let raw = Arc::into_raw(Arc::new(runtime)) as *mut LocalAsyncRuntime<O> as *mut ();
1695        rt.store(raw, Ordering::Relaxed);
1696    }) {
1697        Err(e) => {
1698            panic!("Bind single runtime to local thread failed, reason: {:?}", e);
1699        },
1700        Ok(_) => (),
1701    }
1702}
1703
1704/// 从本地线程解绑单线程异步任务执行器
1705pub fn unbind_local_thread() {
1706    let _ = PI_ASYNC_LOCAL_THREAD_ASYNC_RUNTIME.try_with(move |rt| {
1707        rt.store(null_mut(), Ordering::Relaxed);
1708    });
1709}
1710
1711///
1712/// 本地线程绑定的异步运行时
1713///
1714pub struct LocalAsyncRuntime<O: Default + 'static> {
1715    inner:              *const (),                                                  //内部运行时指针
1716    get_id_func:        fn(*const ()) -> usize,                                     //获取本地运行时的id的函数
1717    spawn_func:         fn(*const (), BoxFuture<'static, O>) -> Result<()>,         //派发函数
1718    spawn_local_func:   fn(*const (), BoxFuture<'static, O>) -> Result<()>,         //本地派发函数
1719    spawn_timing_func:  fn(*const (), BoxFuture<'static, O>, usize) -> Result<()>,  //定时派发函数
1720    timeout_func:       fn(*const (), usize) -> BoxFuture<'static, ()>,             //超时函数
1721}
1722
1723unsafe impl<O: Default + 'static> Send for LocalAsyncRuntime<O> {}
1724unsafe impl<O: Default + 'static> Sync for LocalAsyncRuntime<O> {}
1725
1726impl<O: Default + 'static> LocalAsyncRuntime<O> {
1727    /// 创建本地线程绑定的异步运行时
1728    pub fn new(inner: *const (),
1729               get_id_func: fn(*const ()) -> usize,
1730               spawn_func: fn(*const (), BoxFuture<'static, O>) -> Result<()>,
1731               spawn_local_func: fn(*const (), BoxFuture<'static, O>) -> Result<()>,
1732               spawn_timing_func: fn(*const (), BoxFuture<'static, O>, usize) -> Result<()>,
1733               timeout_func: fn(*const (), usize) -> BoxFuture<'static, ()>) -> Self {
1734        LocalAsyncRuntime {
1735            inner,
1736            get_id_func,
1737            spawn_func,
1738            spawn_local_func,
1739            spawn_timing_func,
1740            timeout_func,
1741        }
1742    }
1743
1744    /// 获取本地运行时的id
1745    #[inline]
1746    pub fn get_id(&self) -> usize {
1747        (self.get_id_func)(self.inner)
1748    }
1749
1750    /// 派发一个指定的异步任务到异步运行时
1751    #[inline]
1752    pub fn spawn<F>(&self, future: F) -> Result<()>
1753        where F: Future<Output = O> + Send + 'static {
1754        (self.spawn_func)(self.inner, async move {
1755            future.await
1756        }.boxed())
1757    }
1758
1759    /// 派发一个指定的异步任务到本地线程绑定的异步运行时
1760    #[inline]
1761    pub fn spawn_local<F>(&self, future: F) -> Result<()>
1762    where F: Future<Output = O> + Send + 'static {
1763        (self.spawn_local_func)(self.inner, async move {
1764            future.await
1765        }.boxed())
1766    }
1767
1768    /// 定时派发一个指定的异步任务到本地线程绑定的异步运行时
1769    #[inline]
1770    pub fn sapwn_timing_func<F>(&self, future: F, timeout: usize) -> Result<()>
1771        where F: Future<Output = O> + Send + 'static {
1772        (self.spawn_timing_func)(self.inner,
1773                                 async move {
1774                                     future.await
1775                                 }.boxed(),
1776                                 timeout)
1777    }
1778
1779    /// 挂起本地线程绑定的异步运行时的当前任务,等待指定的时间后唤醒当前任务
1780    #[inline]
1781    pub fn timeout(&self, timeout: usize) -> BoxFuture<'static, ()> {
1782        (self.timeout_func)(self.inner, timeout)
1783    }
1784}
1785
1786///
1787/// 获取本地线程绑定的异步运行时
1788/// 注意:O如果与本地线程绑定的运行时的O不相同,则无法获取本地线程绑定的运行时
1789///
1790pub fn local_async_runtime<O: Default + 'static>() -> Option<Arc<LocalAsyncRuntime<O>>> {
1791    match PI_ASYNC_LOCAL_THREAD_ASYNC_RUNTIME.try_with(move |ptr| {
1792        let raw = ptr.load(Ordering::Relaxed) as *const LocalAsyncRuntime<O>;
1793        unsafe {
1794            if raw.is_null() {
1795                //本地线程未绑定异步运行时
1796                None
1797            } else {
1798                //本地线程已绑定异步运行时
1799                let shared: Arc<LocalAsyncRuntime<O>> = unsafe { Arc::from_raw(raw) };
1800                let result = shared.clone();
1801                Arc::into_raw(shared); //避免提前释放
1802                Some(result)
1803            }
1804        }
1805    }) {
1806        Err(_) => None, //本地线程没有绑定异步运行时
1807        Ok(rt) => rt,
1808    }
1809}
1810
1811///
1812/// 派发任务到本地线程绑定的异步运行时,如果本地线程没有异步运行时,则返回错误
1813/// 注意:F::Output如果与本地线程绑定的运行时的O不相同,则无法执行指定任务
1814///
1815pub fn spawn_local<O, F>(future: F) -> Result<()>
1816    where O: Default + 'static,
1817          F: Future<Output = O> + Send + 'static {
1818    if let Some(rt) = local_async_runtime::<O>() {
1819        rt.spawn(future)
1820    } else {
1821        Err(Error::new(ErrorKind::Other, format!("Spawn task to local thread failed, reason: runtime not exist")))
1822    }
1823}
1824
1825///
1826/// 从本地线程绑定的字典中获取指定类型的值的只读引用
1827///
1828pub fn get_local_dict<T: 'static>() -> Option<&'static T> {
1829    match PI_ASYNC_LOCAL_THREAD_ASYNC_RUNTIME_DICT.try_with(move |dict| {
1830        unsafe {
1831            if let Some(any) = (&*dict.get()).get(&TypeId::of::<T>()) {
1832                //指定类型的值存在
1833                <dyn Any>::downcast_ref::<T>(&**any)
1834            } else {
1835                //指定类型的值不存在
1836                None
1837            }
1838        }
1839    }) {
1840        Err(_) => {
1841            None
1842        },
1843        Ok(result) => {
1844            result
1845        }
1846    }
1847}
1848
1849///
1850/// 从本地线程绑定的字典中获取指定类型的值的可写引用
1851///
1852pub fn get_local_dict_mut<T: 'static>() -> Option<&'static mut T> {
1853    match PI_ASYNC_LOCAL_THREAD_ASYNC_RUNTIME_DICT.try_with(move |dict| {
1854        unsafe {
1855            if let Some(any) = (&mut *dict.get()).get_mut(&TypeId::of::<T>()) {
1856                //指定类型的值存在
1857                <dyn Any>::downcast_mut::<T>(&mut **any)
1858            } else {
1859                //指定类型的值不存在
1860                None
1861            }
1862        }
1863    }) {
1864        Err(_) => {
1865            None
1866        },
1867        Ok(result) => {
1868            result
1869        }
1870    }
1871}
1872
1873///
1874/// 在本地线程绑定的字典中设置指定类型的值,返回上一个设置的值
1875///
1876pub fn set_local_dict<T: 'static>(value: T) -> Option<T> {
1877    match PI_ASYNC_LOCAL_THREAD_ASYNC_RUNTIME_DICT.try_with(move |dict| {
1878        unsafe {
1879            let result = if let Some(any) = (&mut *dict.get()).remove(&TypeId::of::<T>()) {
1880                //指定类型的上一个值存在
1881                if let Ok(r) = any.downcast() {
1882                    //造型成功,则返回
1883                    Some(*r)
1884                } else {
1885                    None
1886                }
1887            } else {
1888                //指定类型的上一个值不存在
1889                None
1890            };
1891
1892            //设置指定类型的新值
1893            (&mut *dict.get()).insert(TypeId::of::<T>(), Box::new(value) as Box<dyn Any>);
1894
1895            result
1896        }
1897    }) {
1898        Err(_) => {
1899            None
1900        },
1901        Ok(result) => {
1902            result
1903        }
1904    }
1905}
1906
1907///
1908/// 在本地线程绑定的字典中移除指定类型的值,并返回移除的值
1909///
1910pub fn remove_local_dict<T: 'static>() -> Option<T> {
1911    match PI_ASYNC_LOCAL_THREAD_ASYNC_RUNTIME_DICT.try_with(move |dict| {
1912        unsafe {
1913            if let Some(any) = (&mut *dict.get()).remove(&TypeId::of::<T>()) {
1914                //指定类型的上一个值存在
1915                if let Ok(r) = any.downcast() {
1916                    //造型成功,则返回
1917                    Some(*r)
1918                } else {
1919                    None
1920                }
1921            } else {
1922                //指定类型的上一个值不存在
1923                None
1924            }
1925        }
1926    }) {
1927        Err(_) => {
1928            None
1929        },
1930        Ok(result) => {
1931            result
1932        }
1933    }
1934}
1935
1936///
1937/// 清空本地线程绑定的字典
1938///
1939pub fn clear_local_dict() -> Result<()> {
1940    match PI_ASYNC_LOCAL_THREAD_ASYNC_RUNTIME_DICT.try_with(move |dict| {
1941        unsafe {
1942            (&mut *dict.get()).clear();
1943        }
1944    }) {
1945        Err(e) => {
1946            Err(Error::new(ErrorKind::Other, format!("Clear local dict failed, reason: {:?}", e)))
1947        },
1948        Ok(_) => {
1949            Ok(())
1950        }
1951    }
1952}
1953
1954const ASYNC_VALUE_EMPTY: u8 = 0;
1955const ASYNC_VALUE_WAITING: u8 = 1;
1956const ASYNC_VALUE_SETTING: u8 = 2;
1957const ASYNC_VALUE_READY: u8 = 3;
1958const ASYNC_VALUE_TAKING: u8 = 4;
1959const ASYNC_VALUE_CONSUMED: u8 = 5;
1960
1961/// 同步非阻塞的异步值,只允许被同步非阻塞设置一次值。
1962///
1963/// 说明:
1964/// - `AsyncValue` 是一个 single-shot future,`set()` 成功一次后,等待方可通过
1965///   `Future::poll` 取出该值。
1966/// - 本类型允许 pending 后重复 poll;重复 poll 会更新最新 waker,并继续返回
1967///   `Poll::Pending`。
1968/// - `set()` 保持旧语义:首次设置成功,后续设置静默失败且不会覆盖已设置的值。
1969/// - 当前 API 不表达 sender/receiver 拆分、关闭或取消语义;never set 的 future 会继续
1970///   pending。
1971/// - poll after ready 属于调用方违反 Future 契约,可能 panic。
1972///
1973/// 性能:
1974/// - `poll` 和 `set` 均为 O(1),只在 CAS 竞争或极短取值窗口内有限重试。
1975/// - 不执行阻塞等待,不持有互斥锁,不在热路径分配队列节点。
1976///
1977/// 安全性:
1978/// - 内部 value 使用 `UnsafeCell<Option<V>>` 保存;只有成功进入 `SETTING` 的 setter
1979///   可以写入,只有成功进入 `TAKING` 的 receiver 可以取出。
1980/// - 状态转换使用原子 Acquire/Release/AcqRel 保证跨线程可见性。
1981pub struct AsyncValue<V: Send + 'static>(Arc<InnerAsyncValue<V>>);
1982
1983unsafe impl<V: Send + 'static> Send for AsyncValue<V> {}
1984unsafe impl<V: Send + 'static> Sync for AsyncValue<V> {}
1985
1986impl<V: Send + 'static> Clone for AsyncValue<V> {
1987    fn clone(&self) -> Self {
1988        AsyncValue(self.0.clone())
1989    }
1990}
1991
1992impl<V: Send + 'static> Debug for AsyncValue<V> {
1993    fn fmt(&self, f: &mut Formatter<'_>) -> FmtResult {
1994        write!(f,
1995               "AsyncValue[status = {}]",
1996               self.0.status.load(Ordering::Acquire))
1997    }
1998}
1999
2000impl<V: Send + 'static> Future for AsyncValue<V> {
2001    type Output = V;
2002
2003    fn poll(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output> {
2004        let mut spin_len = 1;
2005        loop {
2006            match self.0.status.load(Ordering::Acquire) {
2007                ASYNC_VALUE_EMPTY => {
2008                    self.0.waker.register(cx.waker());
2009                    match self.0.status.compare_exchange(ASYNC_VALUE_EMPTY,
2010                                                         ASYNC_VALUE_WAITING,
2011                                                         Ordering::AcqRel,
2012                                                         Ordering::Acquire) {
2013                        Ok(_) => {
2014                            return Poll::Pending;
2015                        },
2016                        Err(ASYNC_VALUE_EMPTY) => {
2017                            continue;
2018                        },
2019                        Err(ASYNC_VALUE_WAITING) | Err(ASYNC_VALUE_SETTING) => {
2020                            return Poll::Pending;
2021                        },
2022                        Err(ASYNC_VALUE_READY) => {
2023                            continue;
2024                        },
2025                        Err(ASYNC_VALUE_TAKING) => {
2026                            spin_len = spin(spin_len);
2027                            continue;
2028                        },
2029                        Err(ASYNC_VALUE_CONSUMED) => {
2030                            panic!("AsyncValue polled after completion");
2031                        },
2032                        Err(_) => {
2033                            panic!("AsyncValue entered invalid state");
2034                        },
2035                    }
2036                },
2037                ASYNC_VALUE_WAITING | ASYNC_VALUE_SETTING => {
2038                    self.0.waker.register(cx.waker());
2039                    match self.0.status.load(Ordering::Acquire) {
2040                        ASYNC_VALUE_READY => {
2041                            continue;
2042                        },
2043                        ASYNC_VALUE_TAKING => {
2044                            spin_len = spin(spin_len);
2045                            continue;
2046                        },
2047                        ASYNC_VALUE_CONSUMED => {
2048                            panic!("AsyncValue polled after completion");
2049                        },
2050                        _ => {
2051                            return Poll::Pending;
2052                        },
2053                    }
2054                },
2055                ASYNC_VALUE_READY => {
2056                    match self.0.status.compare_exchange(ASYNC_VALUE_READY,
2057                                                         ASYNC_VALUE_TAKING,
2058                                                         Ordering::AcqRel,
2059                                                         Ordering::Acquire) {
2060                        Ok(_) => {
2061                            let value = unsafe { (*self.0.value.get()).take().unwrap() };
2062                            self.0.status.store(ASYNC_VALUE_CONSUMED, Ordering::Release);
2063                            return Poll::Ready(value);
2064                        },
2065                        Err(ASYNC_VALUE_TAKING) => {
2066                            spin_len = spin(spin_len);
2067                            continue;
2068                        },
2069                        Err(ASYNC_VALUE_CONSUMED) => {
2070                            panic!("AsyncValue polled after completion");
2071                        },
2072                        Err(_) => {
2073                            continue;
2074                        },
2075                    }
2076                },
2077                ASYNC_VALUE_TAKING => {
2078                    //其它 clone 已经获得取值权,等待其完成状态推进,避免同时访问 value。
2079                    spin_len = spin(spin_len);
2080                    continue;
2081                },
2082                ASYNC_VALUE_CONSUMED => {
2083                    panic!("AsyncValue polled after completion");
2084                },
2085                _ => {
2086                    panic!("AsyncValue entered invalid state");
2087                },
2088            }
2089        }
2090    }
2091}
2092
2093/*
2094* 同步非阻塞的异步值同步方法
2095*/
2096impl<V: Send + 'static> AsyncValue<V> {
2097    /// 构建异步值,默认值为未就绪
2098    pub fn new() -> Self {
2099        let inner = InnerAsyncValue {
2100            value: UnsafeCell::new(None),
2101            waker: AtomicWaker::new(),
2102            status: AtomicU8::new(ASYNC_VALUE_EMPTY),
2103        };
2104
2105        AsyncValue(Arc::new(inner))
2106    }
2107
2108    /// 判断异步值是否已完成设置
2109    pub fn is_complete(&self) -> bool {
2110        match self.0.status.load(Ordering::Acquire) {
2111            ASYNC_VALUE_READY | ASYNC_VALUE_TAKING | ASYNC_VALUE_CONSUMED => true,
2112            _ => false,
2113        }
2114    }
2115
2116    /// 设置异步值
2117    pub fn set(self, value: V) {
2118        let mut value = Some(value);
2119        loop {
2120            match self.0.status.load(Ordering::Acquire) {
2121                ASYNC_VALUE_EMPTY => {
2122                    match self.0.status.compare_exchange(ASYNC_VALUE_EMPTY,
2123                                                         ASYNC_VALUE_SETTING,
2124                                                         Ordering::AcqRel,
2125                                                         Ordering::Acquire) {
2126                        Ok(_) => {
2127                            unsafe { *self.0.value.get() = value.take(); }
2128                            self.0.status.store(ASYNC_VALUE_READY, Ordering::Release);
2129                            self.0.waker.wake();
2130                            return;
2131                        },
2132                        Err(_) => {
2133                            continue;
2134                        },
2135                    }
2136                },
2137                ASYNC_VALUE_WAITING => {
2138                    match self.0.status.compare_exchange(ASYNC_VALUE_WAITING,
2139                                                         ASYNC_VALUE_SETTING,
2140                                                         Ordering::AcqRel,
2141                                                         Ordering::Acquire) {
2142                        Ok(_) => {
2143                            unsafe { *self.0.value.get() = value.take(); }
2144                            self.0.status.store(ASYNC_VALUE_READY, Ordering::Release);
2145                            self.0.waker.wake();
2146                            return;
2147                        },
2148                        Err(_) => {
2149                            continue;
2150                        },
2151                    }
2152                },
2153                _ => {
2154                    //异步值正在设置、已设置或已消费,则保持旧语义:重复 set 静默失败。
2155                    return;
2156                }
2157            }
2158        }
2159    }
2160}
2161
2162// 同步非阻塞的内部异步值,只允许被同步非阻塞的设置一次值
2163pub struct InnerAsyncValue<V: Send + 'static> {
2164    value:  UnsafeCell<Option<V>>,  // 值,访问权由 status 的 SETTING/TAKING 状态独占保护。
2165    waker:  AtomicWaker,            // 最近一次 pending poll 注册的唤醒器。
2166    status: AtomicU8,               // 状态机,见 ASYNC_VALUE_* 常量。
2167}
2168
2169///
2170/// 异步非阻塞可变值的守护者
2171///
2172pub struct AsyncVariableGuard<'a, V: Send + 'static> {
2173    value:  &'a UnsafeCell<Option<V>>,      //值
2174    waker:  &'a UnsafeCell<Option<Waker>>,  //唤醒器
2175    status: &'a AtomicU8,                   //值状态
2176}
2177
2178unsafe impl<V: Send + 'static> Send for AsyncVariableGuard<'_, V> {}
2179
2180impl<V: Send + 'static> Drop for AsyncVariableGuard<'_, V> {
2181    fn drop(&mut self) {
2182        //当前异步可变值已锁定,则解除锁定
2183        //当前异步可变值的状态为2或6,表示当前异步可变值的唤醒器未就绪并已锁定,或当前异步可变值不需要唤醒并已完成所有修改
2184        //当前异步可变值的状态为3或7,表示当前异步可变值的唤醒器已就绪并已锁定,或当前异步可变值已唤醒并已完成所有修改
2185        self.status.fetch_sub(2, Ordering::Relaxed);
2186    }
2187}
2188
2189impl<V: Send + 'static> Deref for AsyncVariableGuard<'_, V> {
2190    type Target = Option<V>;
2191
2192    fn deref(&self) -> &Self::Target {
2193        unsafe {
2194            &*self.value.get()
2195        }
2196    }
2197}
2198
2199impl<V: Send + 'static> DerefMut for AsyncVariableGuard<'_, V> {
2200    fn deref_mut(&mut self) -> &mut Self::Target {
2201        unsafe {
2202            &mut *self.value.get()
2203        }
2204    }
2205}
2206
2207impl<V: Send + 'static> AsyncVariableGuard<'_, V> {
2208    /// 完成异步可变值的修改
2209    pub fn finish(self) {
2210        //设置异步可变值的状态为已完成修改
2211        if self.status.fetch_add(4, Ordering::Relaxed) == 3 {
2212            if let Some(waker) = unsafe { (&mut *self.waker.get()).take() } {
2213                //当前异步可变值需要唤醒,则立即唤醒异步可变值
2214                waker.wake();
2215            }
2216        }
2217    }
2218}
2219
2220///
2221/// 异步非阻塞可变值,在完成前允许被同步非阻塞的修改多次
2222///
2223pub struct AsyncVariable<V: Send + 'static>(Arc<InnerAsyncVariable<V>>);
2224
2225unsafe impl<V: Send + 'static> Send for AsyncVariable<V> {}
2226unsafe impl<V: Send + 'static> Sync for AsyncVariable<V> {}
2227
2228impl<V: Send + 'static> Clone for AsyncVariable<V> {
2229    fn clone(&self) -> Self {
2230        AsyncVariable(self.0.clone())
2231    }
2232}
2233
2234impl<V: Send + 'static> Future for AsyncVariable<V> {
2235    type Output = V;
2236
2237    fn poll(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output> {
2238        unsafe {
2239            *self.0.waker.get() = Some(cx.waker().clone()); //设置异步可变值的唤醒器准备就绪
2240        }
2241
2242        let mut spin_len = 1;
2243        loop {
2244            match self.0.status.compare_exchange(0,
2245                                                 1,
2246                                                 Ordering::Acquire,
2247                                                 Ordering::Relaxed) {
2248                Err(current) if current & 4 != 0 => {
2249                    //异步可变值已完成所有修改,则立即返回
2250                    unsafe {
2251                        let _ = (&mut *self.0.waker.get()).take(); //释放异步可变值的唤醒器
2252                        return Poll::Ready((&mut *(&self).0.value.get()).take().unwrap());
2253                    }
2254                },
2255                Err(_) => {
2256                    //还未完成值修改,则自旋等待
2257                    spin_len = spin(spin_len);
2258                },
2259                Ok(_) => {
2260                    //异步可变值已挂起
2261                    return Poll::Pending;
2262                },
2263            }
2264        }
2265    }
2266}
2267
2268impl<V: Send + 'static> AsyncVariable<V> {
2269    /// 构建异步可变值,默认值为未就绪
2270    pub fn new() -> Self {
2271        let inner = InnerAsyncVariable {
2272            value: UnsafeCell::new(None),
2273            waker: UnsafeCell::new(None),
2274            status: AtomicU8::new(0),
2275        };
2276
2277        AsyncVariable(Arc::new(inner))
2278    }
2279
2280    /// 判断异步可变值是否已完成设置
2281    pub fn is_complete(&self) -> bool {
2282        self
2283            .0
2284            .status
2285            .load(Ordering::Acquire) & 4 != 0
2286    }
2287
2288    /// 锁住待修改的异步可变值,并返回当前异步可变值的守护者,如果异步可变值已完成修改则返回空
2289    pub fn lock(&self) -> Option<AsyncVariableGuard<V>> {
2290        let mut spin_len = 1;
2291        loop {
2292            match self
2293                .0
2294                .status
2295                .compare_exchange(1,
2296                                  3,
2297                                  Ordering::Acquire,
2298                                  Ordering::Relaxed) {
2299                Err(0) => {
2300                    //异步可变值还未就绪,则自旋等待
2301                    match self
2302                        .0
2303                        .status
2304                        .compare_exchange(0,
2305                                          2,
2306                                          Ordering::Acquire,
2307                                          Ordering::Relaxed) {
2308                        Err(1) => {
2309                            //异步可变值已就绪,则继续尝试获取锁
2310                            continue;
2311                        },
2312                        Err(2) => {
2313                            //异步可变值的唤醒器未就绪且已锁,但未获取到锁,则自旋等待
2314                            spin_len = spin(spin_len);
2315                        },
2316                        Err(3) => {
2317                            //异步可变值的唤醒器已就绪且已锁,但未获取到锁,则自旋等待
2318                            spin_len = spin(spin_len);
2319                        },
2320                        Err(_) => {
2321                            //已完成,则返回空
2322                            return None;
2323                        },
2324                        Ok(_) => {
2325                            //异步可变值的唤醒器未就绪且获取到锁,则返回异步可变值的守护者
2326                            let guard = AsyncVariableGuard {
2327                                value: &self.0.value,
2328                                waker: &self.0.waker,
2329                                status: &self.0.status,
2330                            };
2331
2332                            return Some(guard)
2333                        },
2334                    }
2335                },
2336                Err(2) => {
2337                    //异步可变值的唤醒器未就绪且已锁,但未获取到锁,则自旋等待
2338                    spin_len = spin(spin_len);
2339                },
2340                Err(3) => {
2341                    //异步可变值的唤醒器已就绪且已锁,但未获取到锁,则自旋等待
2342                    spin_len = spin(spin_len);
2343                },
2344                Err(_) => {
2345                    //已完成,则返回空
2346                    return None;
2347                }
2348                Ok(_) => {
2349                    //异步可变值的唤醒器已就绪且获取到锁,则返回异步可变值的守护者
2350                    let guard = AsyncVariableGuard {
2351                        value: &self.0.value,
2352                        waker: &self.0.waker,
2353                        status: &self.0.status,
2354                    };
2355
2356                    return Some(guard)
2357                },
2358            }
2359        }
2360    }
2361}
2362
2363// 内部异步非阻塞可变值,在完成前允许被同步非阻塞的修改多次
2364pub struct InnerAsyncVariable<V: Send + 'static> {
2365    value:  UnsafeCell<Option<V>>,      //值
2366    waker:  UnsafeCell<Option<Waker>>,  //唤醒器
2367    status: AtomicU8,                   //状态
2368}
2369
2370///
2371/// 等待异步任务运行的结果
2372///
2373pub struct AsyncWaitResult<V: Send + 'static>(pub Arc<RefCell<Option<Result<V>>>>);
2374
2375unsafe impl<V: Send + 'static> Send for AsyncWaitResult<V> {}
2376unsafe impl<V: Send + 'static> Sync for AsyncWaitResult<V> {}
2377
2378impl<V: Send + 'static> Clone for AsyncWaitResult<V> {
2379    fn clone(&self) -> Self {
2380        AsyncWaitResult(self.0.clone())
2381    }
2382}
2383
2384///
2385/// 等待异步任务运行的结果集
2386///
2387pub struct AsyncWaitResults<V: Send + 'static>(pub Arc<RefCell<Option<Vec<Result<V>>>>>);
2388
2389unsafe impl<V: Send + 'static> Send for AsyncWaitResults<V> {}
2390unsafe impl<V: Send + 'static> Sync for AsyncWaitResults<V> {}
2391
2392impl<V: Send + 'static> Clone for AsyncWaitResults<V> {
2393    fn clone(&self) -> Self {
2394        AsyncWaitResults(self.0.clone())
2395    }
2396}
2397
2398///
2399/// 异步定时器任务
2400///
2401pub enum AsyncTimingTask<
2402    P: AsyncTaskPoolExt<O> + AsyncTaskPool<O>,
2403    O: Default + 'static = (),
2404> {
2405    Pended(TaskId),                     //已挂起的定时任务
2406    WaitRun(Arc<AsyncTask<P, O>>),      //等待执行的定时任务
2407    TimeoutWake(Arc<TimeoutWaiter>),    //等待timeout到期的唤醒句柄
2408}
2409
2410///
2411/// 异步任务本地定时器
2412///
2413pub struct AsyncTaskTimer<
2414    P: AsyncTaskPoolExt<O> + AsyncTaskPool<O>,
2415    O: Default + 'static = (),
2416> {
2417    producor:   Sender<(usize, AsyncTimingTask<P, O>)>,                     //定时任务生产者
2418    consumer:   Receiver<(usize, AsyncTimingTask<P, O>)>,                   //定时任务消费者
2419    timer:      Arc<RefCell<Timer<AsyncTimingTask<P, O>, 1000, 60, 3>>>,    //定时器
2420    clock:      Clock,                                                      //定时器时钟
2421    now:        QInstant,                                                   //当前时间
2422}
2423
2424unsafe impl<
2425    P: AsyncTaskPoolExt<O> + AsyncTaskPool<O>,
2426    O: Default + 'static,
2427> Send for AsyncTaskTimer<P, O> {}
2428unsafe impl<
2429    P: AsyncTaskPoolExt<O> + AsyncTaskPool<O>,
2430    O: Default + 'static,
2431> Sync for AsyncTaskTimer<P, O> {}
2432
2433impl<
2434    P: AsyncTaskPoolExt<O> + AsyncTaskPool<O>,
2435    O: Default + 'static,
2436> AsyncTaskTimer<P, O> {
2437    /// 构建异步任务本地定时器
2438    pub fn new() -> Self {
2439        let (producor, consumer) = unbounded();
2440        let clock = Clock::new();
2441        let now = clock.recent();
2442
2443        AsyncTaskTimer {
2444            producor,
2445            consumer,
2446            timer: Arc::new(RefCell::new(Timer::<AsyncTimingTask<P, O>, 1000, 60, 3>::default())),
2447            clock,
2448            now,
2449        }
2450    }
2451
2452    /// 获取定时任务生产者
2453    #[inline]
2454    pub fn get_producor(&self) -> &Sender<(usize, AsyncTimingTask<P, O>)> {
2455        &self.producor
2456    }
2457
2458    /// 获取剩余未到期的定时器任务数量
2459    #[inline]
2460    pub fn len(&self) -> usize {
2461        let timer = self.timer.as_ref().borrow();
2462        timer.add_count() - timer.remove_count()
2463    }
2464
2465    /// 设置定时器
2466    pub fn set_timer(&self, task: AsyncTimingTask<P, O>, timeout: usize) -> usize {
2467        let current_time = self
2468            .clock
2469            .recent()
2470            .duration_since(self.now)
2471            .as_millis() as u64;
2472        self
2473            .timer
2474            .borrow_mut()
2475            .push_time(current_time + timeout as u64, task)
2476            .data()
2477            .as_ffi() as usize
2478    }
2479
2480    /// 取消定时器
2481    pub fn cancel_timer(&self, timer_ref: usize) -> Option<AsyncTimingTask<P, O>> {
2482        if let Some(item) = self
2483            .timer
2484            .borrow_mut()
2485            .cancel(KeyData::from_ffi(timer_ref as u64).into()) {
2486            Some(item)
2487        } else {
2488            None
2489        }
2490    }
2491
2492    /// 消费所有定时任务,返回定时任务数量
2493    pub fn consume(&self) -> usize {
2494        let timer_tasks = self.consumer.try_iter().collect::<Vec<(usize, AsyncTimingTask<P, O>)>>();
2495        let len = timer_tasks.len();
2496        for (timeout, task) in timer_tasks {
2497            self.set_timer(task, timeout);
2498        }
2499
2500        len
2501    }
2502
2503    /// 判断当前时间是否有可以弹出的任务,如果有可以弹出的任务,则返回当前时间,否则返回空
2504    pub fn is_require_pop(&self) -> Option<u64> {
2505        let current_time = self
2506            .clock
2507            .recent()
2508            .duration_since(self.now)
2509            .as_millis() as u64;
2510        if self.timer.borrow_mut().is_ok(current_time) {
2511            Some(current_time)
2512        } else {
2513            None
2514        }
2515    }
2516
2517    /// 从定时器中弹出指定时间的一个到期任务
2518    pub fn pop(&self, current_time: u64) -> Option<(usize, AsyncTimingTask<P, O>)> {
2519        if let Some((key, item)) = self.timer.borrow_mut().pop_kv(current_time) {
2520            Some((key.data().as_ffi() as usize, item))
2521        } else {
2522            None
2523        }
2524    }
2525}
2526
2527///
2528/// 异步任务本地定时器,不支持取消定时任务
2529///
2530pub struct AsyncTaskTimerByNotCancel<
2531    P: AsyncTaskPoolExt<O> + AsyncTaskPool<O>,
2532    O: Default + 'static = (),
2533> {
2534    producor:   Sender<(usize, AsyncTimingTask<P, O>)>,                             //定时任务生产者
2535    consumer:   Receiver<(usize, AsyncTimingTask<P, O>)>,                           //定时任务消费者
2536    timer:      Arc<RefCell<NotCancelTimer<AsyncTimingTask<P, O>, 1000, 60, 3>>>,   //定时器
2537    clock:      Clock,                                                              //定时器时钟
2538    now:        QInstant,                                                           //当前时间
2539}
2540
2541unsafe impl<
2542    P: AsyncTaskPoolExt<O> + AsyncTaskPool<O>,
2543    O: Default + 'static,
2544> Send for AsyncTaskTimerByNotCancel<P, O> {}
2545unsafe impl<
2546    P: AsyncTaskPoolExt<O> + AsyncTaskPool<O>,
2547    O: Default + 'static,
2548> Sync for AsyncTaskTimerByNotCancel<P, O> {}
2549
2550impl<
2551    P: AsyncTaskPoolExt<O> + AsyncTaskPool<O>,
2552    O: Default + 'static,
2553> AsyncTaskTimerByNotCancel<P, O> {
2554    /// 构建异步任务本地定时器
2555    pub fn new() -> Self {
2556        let (producor, consumer) = unbounded();
2557        let clock = Clock::new();
2558        let now = clock.recent();
2559
2560        AsyncTaskTimerByNotCancel {
2561            producor,
2562            consumer,
2563            timer: Arc::new(RefCell::new(NotCancelTimer::<AsyncTimingTask<P, O>, 1000, 60, 3>::default())),
2564            clock,
2565            now,
2566        }
2567    }
2568
2569    /// 获取定时任务生产者
2570    #[inline]
2571    pub fn get_producor(&self) -> &Sender<(usize, AsyncTimingTask<P, O>)> {
2572        &self.producor
2573    }
2574
2575    /// 获取剩余未到期的定时器任务数量
2576    #[inline]
2577    pub fn len(&self) -> usize {
2578        let timer = self.timer.as_ref().borrow();
2579        timer.add_count() - timer.remove_count()
2580    }
2581
2582    /// 设置定时器
2583    pub fn set_timer(&self, task: AsyncTimingTask<P, O>, timeout: usize) {
2584        self
2585            .timer
2586            .borrow_mut()
2587            .push(timeout, task);
2588    }
2589
2590    /// 消费所有定时任务,返回定时任务数量
2591    pub fn consume(&self) -> usize {
2592        let timer_tasks = self.consumer.try_iter().collect::<Vec<(usize, AsyncTimingTask<P, O>)>>();
2593        let len = timer_tasks.len();
2594        for (timeout, task) in timer_tasks {
2595            self.set_timer(task, timeout);
2596        }
2597
2598        len
2599    }
2600
2601    /// 判断当前时间是否有可以弹出的任务,如果有可以弹出的任务,则返回当前时间,否则返回空
2602    pub fn is_require_pop(&self) -> Option<u64> {
2603        let current_time = self
2604            .clock
2605            .recent()
2606            .duration_since(self.now)
2607            .as_millis() as u64;
2608        if self.timer.borrow_mut().is_ok(current_time) {
2609            Some(current_time)
2610        } else {
2611            None
2612        }
2613    }
2614
2615    /// 从定时器中弹出指定时间的一个到期任务
2616    pub fn pop(&self, current_time: u64) -> Option<AsyncTimingTask<P, O>> {
2617        if let Some(item) = self.timer.borrow_mut().pop(current_time) {
2618            Some(item)
2619        } else {
2620            None
2621        }
2622    }
2623}
2624
2625///
2626/// 等待指定超时
2627///
2628pub struct AsyncWaitTimeout<
2629    RT: AsyncRuntime<O>,
2630    P: AsyncTaskPoolExt<O> + AsyncTaskPool<O>,
2631    O: Default + 'static = (),
2632> {
2633    rt:         RT,                                     //当前运行时
2634    producor:   Sender<(usize, AsyncTimingTask<P, O>)>, //超时请求生产者
2635    timeout:    usize,                                  //超时时长,单位ms
2636    registered: AtomicBool,                             //是否已注册到定时器
2637    waiter:     Arc<TimeoutWaiter>,                     //timeout专用等待句柄
2638}
2639
2640unsafe impl<
2641    RT: AsyncRuntime<O>,
2642    P: AsyncTaskPoolExt<O> + AsyncTaskPool<O>,
2643    O: Default + 'static,
2644> Send for AsyncWaitTimeout<RT, P, O> {}
2645unsafe impl<
2646    RT: AsyncRuntime<O>,
2647    P: AsyncTaskPoolExt<O> + AsyncTaskPool<O>,
2648    O: Default + 'static,
2649> Sync for AsyncWaitTimeout<RT, P, O> {}
2650
2651impl<
2652    RT: AsyncRuntime<O>,
2653    P: AsyncTaskPoolExt<O> + AsyncTaskPool<O>,
2654    O: Default + 'static,
2655> Future for AsyncWaitTimeout<RT, P, O> {
2656    type Output = ();
2657
2658    fn poll(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output> {
2659        if self.waiter.is_fired() {
2660            //已到期,则返回
2661            return Poll::Ready(());
2662        }
2663
2664        self.waiter.register(cx.waker());
2665
2666        if !self.registered.swap(true, Ordering::AcqRel) {
2667            //发送超时请求,并返回
2668            let _ = self
2669                .producor
2670                .send((self.timeout, AsyncTimingTask::TimeoutWake(self.waiter.clone())));
2671        }
2672
2673        if self.waiter.is_fired() {
2674            Poll::Ready(())
2675        } else {
2676            Poll::Pending
2677        }
2678    }
2679}
2680
2681impl<
2682    RT: AsyncRuntime<O>,
2683    P: AsyncTaskPoolExt<O> + AsyncTaskPool<O>,
2684    O: Default + 'static,
2685> Drop for AsyncWaitTimeout<RT, P, O> {
2686    fn drop(&mut self) {
2687        self.waiter.clear_waker();
2688    }
2689}
2690
2691impl<
2692    RT: AsyncRuntime<O>,
2693    P: AsyncTaskPoolExt<O> + AsyncTaskPool<O>,
2694    O: Default + 'static,
2695> AsyncWaitTimeout<RT, P, O> {
2696    /// 构建等待指定超时任务的方法
2697    pub fn new(rt: RT,
2698               producor: Sender<(usize, AsyncTimingTask<P, O>)>,
2699               timeout: usize) -> Self {
2700        AsyncWaitTimeout {
2701            rt,
2702            producor,
2703            timeout,
2704            registered: AtomicBool::new(false), //设置初始值
2705            waiter: Arc::new(TimeoutWaiter::new()),
2706        }
2707    }
2708}
2709
2710///
2711/// 本地等待指定超时
2712///
2713pub struct LocalAsyncWaitTimeout<
2714    RT: AsyncRuntime<O>,
2715    P: AsyncTaskPoolExt<O> + AsyncTaskPool<O>,
2716    O: Default + 'static = (),
2717> {
2718    rt:         RT,                                     //当前运行时
2719    timer:      Arc<AsyncTaskTimerByNotCancel<P, O>>,   //定时器
2720    timeout:    usize,                                  //超时时长,单位ms
2721    registered: AtomicBool,                             //是否已注册到定时器
2722    waiter:     Arc<TimeoutWaiter>,                     //timeout专用等待句柄
2723}
2724
2725unsafe impl<
2726    RT: AsyncRuntime<O>,
2727    P: AsyncTaskPoolExt<O> + AsyncTaskPool<O>,
2728    O: Default + 'static,
2729> Send for LocalAsyncWaitTimeout<RT, P, O> {}
2730unsafe impl<
2731    RT: AsyncRuntime<O>,
2732    P: AsyncTaskPoolExt<O> + AsyncTaskPool<O>,
2733    O: Default + 'static,
2734> Sync for LocalAsyncWaitTimeout<RT, P, O> {}
2735
2736impl<
2737    RT: AsyncRuntime<O>,
2738    P: AsyncTaskPoolExt<O> + AsyncTaskPool<O>,
2739    O: Default + 'static,
2740> Future for LocalAsyncWaitTimeout<RT, P, O> {
2741    type Output = ();
2742
2743    fn poll(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output> {
2744        if self.waiter.is_fired() {
2745            //已到期,则返回
2746            return Poll::Ready(());
2747        }
2748
2749        self.waiter.register(cx.waker());
2750
2751        if !self.registered.swap(true, Ordering::AcqRel) {
2752            //设置本地超时请求,并返回
2753            self
2754                .timer
2755                .set_timer(AsyncTimingTask::TimeoutWake(self.waiter.clone()),
2756                           self.timeout);
2757        }
2758
2759        if self.waiter.is_fired() {
2760            Poll::Ready(())
2761        } else {
2762            Poll::Pending
2763        }
2764    }
2765}
2766
2767impl<
2768    RT: AsyncRuntime<O>,
2769    P: AsyncTaskPoolExt<O> + AsyncTaskPool<O>,
2770    O: Default + 'static,
2771> Drop for LocalAsyncWaitTimeout<RT, P, O> {
2772    fn drop(&mut self) {
2773        self.waiter.clear_waker();
2774    }
2775}
2776
2777impl<
2778    RT: AsyncRuntime<O>,
2779    P: AsyncTaskPoolExt<O> + AsyncTaskPool<O>,
2780    O: Default + 'static,
2781> LocalAsyncWaitTimeout<RT, P, O> {
2782    /// 构建等待指定超时任务的方法
2783    pub fn new(rt: RT,
2784               timer: Arc<AsyncTaskTimerByNotCancel<P, O>>,
2785               timeout: usize) -> Self {
2786        LocalAsyncWaitTimeout {
2787            rt,
2788            timer,
2789            timeout,
2790            registered: AtomicBool::new(false), //设置初始值
2791            waiter: Arc::new(TimeoutWaiter::new()),
2792        }
2793    }
2794}
2795
2796///
2797/// 等待异步任务执行完成
2798///
2799pub struct AsyncWait<V: Send + 'static>(AsyncWaitAny<V>);
2800
2801unsafe impl<V: Send + 'static> Send for AsyncWait<V> {}
2802unsafe impl<V: Send + 'static> Sync for AsyncWait<V> {}
2803
2804/*
2805* 等待异步任务执行完成同步方法
2806*/
2807impl<V: Send + 'static> AsyncWait<V> {
2808    /// 派发指定超时时间的指定任务到指定的运行时,并返回派发是否成功
2809    pub fn spawn<RT, O, F>(&self,
2810                           rt: RT,
2811                           timeout: Option<usize>,
2812                           future: F) -> Result<()>
2813        where RT: AsyncRuntime<O>,
2814              O: Default + 'static,
2815              F: Future<Output = Result<V>> + Send + 'static {
2816        self.0.spawn(rt.clone(), future)?;
2817
2818        if let Some(timeout) = timeout {
2819            //设置了超时时间
2820            let rt_copy = rt.clone();
2821            self.0.spawn(rt, async move {
2822                rt_copy.timeout(timeout).await;
2823
2824                //返回超时错误
2825                Err(Error::new(ErrorKind::TimedOut, format!("Time out")))
2826            })
2827        } else {
2828            //未设置超时时间
2829            Ok(())
2830        }
2831    }
2832
2833    /// 派发指定超时时间的指定任务到本地运行时,并返回派发是否成功
2834    pub fn spawn_local<O, F>(&self,
2835                             timeout: Option<usize>,
2836                             future: F) -> Result<()>
2837        where O: Default + 'static,
2838              F: Future<Output = Result<V>> + Send + 'static {
2839        if let Some(rt) = local_async_runtime::<O>() {
2840            //当前线程有绑定运行时
2841            self.0.spawn_local(future)?;
2842
2843            if let Some(timeout) = timeout {
2844                //设置了超时时间
2845                let rt_copy = rt.clone();
2846                self.0.spawn_local(async move {
2847                    rt_copy.timeout(timeout).await;
2848
2849                    //返回超时错误
2850                    Err(Error::new(ErrorKind::TimedOut, format!("Time out")))
2851                })
2852            } else {
2853                //未设置超时时间
2854                Ok(())
2855            }
2856        } else {
2857            //当前线程未绑定运行时
2858            Err(Error::new(ErrorKind::Other, format!("Spawn wait task failed, reason: local async runtime not exist")))
2859        }
2860    }
2861}
2862
2863/*
2864* 等待异步任务执行完成异步方法
2865*/
2866impl<V: Send + 'static> AsyncWait<V> {
2867    /// 异步等待已派发任务的结果
2868    pub async fn wait_result(self) -> Result<V> {
2869        self.0.wait_result().await
2870    }
2871}
2872
2873///
2874/// 等待任意异步任务执行完成
2875///
2876pub struct AsyncWaitAny<V: Send + 'static> {
2877    capacity:       usize,                      //派发任务的容量
2878    producor:       AsyncSender<Result<V>>,     //异步返回值生成器
2879    consumer:       AsyncReceiver<Result<V>>,   //异步返回值接收器
2880}
2881
2882unsafe impl<V: Send + 'static> Send for AsyncWaitAny<V> {}
2883unsafe impl<V: Send + 'static> Sync for AsyncWaitAny<V> {}
2884
2885/*
2886* 等待任意异步任务执行完成同步方法
2887*/
2888impl<V: Send + 'static> AsyncWaitAny<V> {
2889    /// 派发指定任务到指定的运行时,并返回派发是否成功
2890    pub fn spawn<RT, O, F>(&self,
2891                           rt: RT,
2892                           future: F) -> Result<()>
2893        where RT: AsyncRuntime<O>,
2894              O: Default + 'static,
2895              F: Future<Output = Result<V>> + Send + 'static {
2896        let producor = self.producor.clone();
2897        rt.spawn_by_id(rt.alloc::<O>(), async move {
2898            let value = future.await;
2899            producor.into_send_async(value).await;
2900
2901            //返回异步任务的默认值
2902            Default::default()
2903        })
2904    }
2905
2906    /// 派发指定任务到本地运行时,并返回派发是否成功
2907    pub fn spawn_local<F>(&self,
2908                          future: F) -> Result<()>
2909        where F: Future<Output = Result<V>> + Send + 'static {
2910        if let Some(rt) = local_async_runtime() {
2911            //本地线程有绑定运行时
2912            let producor = self.producor.clone();
2913            rt.spawn(async move {
2914                let value = future.await;
2915                producor.into_send_async(value).await;
2916            })
2917        } else {
2918            //本地线程未绑定运行时
2919            Err(Error::new(ErrorKind::Other, format!("Spawn wait any task failed, reason: local async runtime not exist")))
2920        }
2921    }
2922}
2923
2924/*
2925* 等待任意异步任务执行完成异步方法
2926*/
2927impl<V: Send + 'static> AsyncWaitAny<V> {
2928    /// 异步等待任意已派发任务的结果
2929    pub async fn wait_result(self) -> Result<V> {
2930        match self.consumer.recv_async().await {
2931            Err(e) => {
2932                //接收错误,则立即返回
2933                Err(Error::new(ErrorKind::Other, format!("Wait any result failed, reason: {:?}", e)))
2934            },
2935            Ok(result) => {
2936                //接收成功,则立即返回
2937                result
2938            },
2939        }
2940    }
2941}
2942
2943///
2944/// 等待任意异步任务执行完成
2945///
2946pub struct AsyncWaitAnyCallback<V: Send + 'static> {
2947    capacity:   usize,                      //派发任务的容量
2948    producor:   AsyncSender<Result<V>>,     //异步返回值生成器
2949    consumer:   AsyncReceiver<Result<V>>,   //异步返回值接收器
2950}
2951
2952unsafe impl<V: Send + 'static> Send for AsyncWaitAnyCallback<V> {}
2953unsafe impl<V: Send + 'static> Sync for AsyncWaitAnyCallback<V> {}
2954
2955/*
2956* 等待任意异步任务执行完成同步方法
2957*/
2958impl<V: Send + 'static> AsyncWaitAnyCallback<V> {
2959    /// 派发指定任务到指定的运行时,并返回派发是否成功
2960    pub fn spawn<RT, O, F>(&self,
2961                           rt: RT,
2962                           future: F) -> Result<()>
2963        where RT: AsyncRuntime<O>,
2964              O: Default + 'static,
2965              F: Future<Output = Result<V>> + Send + 'static {
2966        let producor = self.producor.clone();
2967        rt.spawn_by_id(rt.alloc::<O>(), async move {
2968            let value = future.await;
2969            producor.into_send_async(value).await;
2970
2971            //返回异步任务的默认值
2972            Default::default()
2973        })
2974    }
2975
2976    /// 派发指定任务到本地运行时,并返回派发是否成功
2977    pub fn spawn_local<F>(&self,
2978                          future: F) -> Result<()>
2979        where F: Future<Output = Result<V>> + Send + 'static {
2980        if let Some(rt) = local_async_runtime() {
2981            //当前线程有绑定运行时
2982            let producor = self.producor.clone();
2983            rt.spawn(async move {
2984                let value = future.await;
2985                producor.into_send_async(value).await;
2986            })
2987        } else {
2988            //当前线程未绑定运行时
2989            Err(Error::new(ErrorKind::Other, format!("Spawn wait any task failed by callback, reason: current async runtime not exist")))
2990        }
2991    }
2992}
2993
2994/*
2995* 等待任意异步任务执行完成异步方法
2996*/
2997impl<V: Send + 'static> AsyncWaitAnyCallback<V> {
2998    /// 异步等待满足用户回调需求的已派发任务的结果
2999    pub async fn wait_result(mut self,
3000                             callback: impl Fn(&Result<V>) -> bool + Send + Sync + 'static) -> Result<V> {
3001        let checker = create_checker(self.capacity, callback);
3002        loop {
3003            match self.consumer.recv_async().await {
3004                Err(e) => {
3005                    //接收错误,则立即返回
3006                    return Err(Error::new(ErrorKind::Other, format!("Wait any result failed by callback, reason: {:?}", e)));
3007                },
3008                Ok(result) => {
3009                    //接收成功,则检查是否立即返回
3010                    if checker(&result) {
3011                        //检查通过,则立即唤醒等待的任务,否则等待其它任务唤醒
3012                        return result;
3013                    }
3014                },
3015            }
3016        }
3017    }
3018}
3019
3020// 根据用户提供的回调,生成检查器
3021fn create_checker<V, F>(len: usize,
3022                        callback: F) -> Arc<dyn Fn(&Result<V>) -> bool + Send + Sync + 'static>
3023    where V: Send + 'static,
3024          F: Fn(&Result<V>) -> bool + Send + Sync + 'static {
3025    let mut check_counter = AtomicUsize::new(len); //初始化检查计数器
3026    Arc::new(move |result| {
3027        if check_counter.fetch_sub(1, Ordering::SeqCst) == 1 {
3028            //最后一个任务的检查,则忽略用户回调,并立即返回成功
3029            true
3030        } else {
3031            //不是最后一个任务的检查,则调用用户回调,并根据用户回调确定是否成功
3032            callback(result)
3033        }
3034    })
3035}
3036
3037///
3038/// 异步映射归并
3039///
3040pub struct AsyncMapReduce<V: Send + 'static> {
3041    count:          usize,                              //派发的任务数量
3042    capacity:       usize,                              //派发任务的容量
3043    producor:       AsyncSender<(usize, Result<V>)>,    //异步返回值生成器
3044    consumer:       AsyncReceiver<(usize, Result<V>)>,  //异步返回值接收器
3045}
3046
3047unsafe impl<V: Send + 'static> Send for AsyncMapReduce<V> {}
3048
3049/*
3050* 异步映射归并同步方法
3051*/
3052impl<V: Send + 'static> AsyncMapReduce<V> {
3053    /// 映射指定任务到指定的运行时,并返回任务序号
3054    pub fn map<RT, O, F>(&mut self, rt: RT, future: F) -> Result<usize>
3055        where RT: AsyncRuntime<O>,
3056              O: Default + 'static,
3057              F: Future<Output = Result<V>> + Send + 'static {
3058        if self.count >= self.capacity {
3059            //已派发任务已达可派发任务的限制,则返回错误
3060            return Err(Error::new(ErrorKind::Other, format!("Map task to runtime failed, capacity: {}, reason: out of capacity", self.capacity)));
3061        }
3062
3063        let index = self.count;
3064        let producor = self.producor.clone();
3065        rt.spawn_by_id(rt.alloc::<O>(), async move {
3066            let value = future.await;
3067            producor.into_send_async((index, value)).await;
3068
3069            //返回异步任务的默认值
3070            Default::default()
3071        })?;
3072
3073        self.count += 1; //派发任务成功,则计数
3074        Ok(index)
3075    }
3076}
3077
3078/*
3079* 异步映射归并异步方法
3080*/
3081impl<V: Send + 'static> AsyncMapReduce<V> {
3082    /// 归并所有派发的任务
3083    pub async fn reduce(self, order: bool) -> Result<Vec<Result<V>>> {
3084        let mut count = self.count;
3085        let mut results = Vec::with_capacity(count);
3086        while count > 0 {
3087            match self.consumer.recv_async().await {
3088                Err(e) => {
3089                    //接收错误,则立即返回
3090                    return Err(Error::new(ErrorKind::Other, format!("Reduce result failed, reason: {:?}", e)));
3091                },
3092                Ok((index, result)) => {
3093                    //接收成功,则继续
3094                    results.push((index, result));
3095                    count -= 1;
3096                },
3097            }
3098        }
3099
3100        if order {
3101            //需要对结果集进行排序
3102            results.sort_by_key(|(key, _value)| {
3103                key.clone()
3104            });
3105        }
3106        let (_, values) = results
3107            .into_iter()
3108            .unzip::<usize, Result<V>, Vec<usize>, Vec<Result<V>>>();
3109
3110        Ok(values)
3111    }
3112}
3113
3114///
3115/// 异步管道过滤器结果
3116///
3117pub enum AsyncPipelineResult<O: 'static> {
3118    Disconnect,     //关闭管道
3119    Filtered(O),    //过滤后的值
3120}
3121
3122///
3123/// 派发一个工作线程
3124/// 返回线程的句柄,可以通过句柄关闭线程
3125/// 线程在没有任务可以执行时会休眠,当派发任务或唤醒任务时会自动唤醒线程
3126///
3127pub fn spawn_worker_thread<F0, F1>(thread_name: &str,
3128                                   thread_stack_size: usize,
3129                                   thread_handler: Arc<AtomicBool>,
3130                                   thread_waker: Arc<(AtomicBool, Mutex<()>, Condvar)>, //用于唤醒运行时所在线程的条件变量
3131                                   sleep_timeout: u64,                                  //休眠超时时长,单位毫秒
3132                                   loop_interval: Option<u64>,                          //工作者线程循环的间隔时长,None为无间隔,单位毫秒
3133                                   loop_func: F0,
3134                                   get_queue_len: F1) -> Arc<AtomicBool>
3135    where F0: Fn() -> (bool, Duration) + Send + 'static,
3136          F1: Fn() -> usize + Send + 'static {
3137    let thread_status_copy = thread_handler.clone();
3138
3139    thread::Builder::new()
3140        .name(thread_name.to_string())
3141        .stack_size(thread_stack_size).spawn(move || {
3142        let mut sleep_count = 0;
3143
3144        while thread_handler.load(Ordering::Relaxed) {
3145            let (is_no_task, run_time) = loop_func();
3146
3147            if is_no_task {
3148                //当前没有任务
3149                if sleep_count > 1 {
3150                    //当前没有任务连续达到2次,则休眠线程
3151                    sleep_count = 0; //重置休眠计数
3152                    let (is_sleep, lock, condvar) = &*thread_waker;
3153                    if get_queue_len() > 0 {
3154                        //当前有任务,则继续工作
3155                        continue;
3156                    }
3157
3158                    {
3159                        let _locked = lock.lock();
3160                        if !is_sleep.load(Ordering::Acquire) {
3161                            //发布休眠状态,外部唤醒端会在同一把锁内确认后再notify
3162                            is_sleep.store(true, Ordering::Release);
3163                        }
3164                    }
3165
3166                    if get_queue_len() > 0 {
3167                        //发布休眠后再次检查任务,避免外部唤醒落在发布窗口内
3168                        is_sleep.store(false, Ordering::Release);
3169                        continue;
3170                    }
3171
3172                    let mut locked = lock.lock();
3173                    if is_sleep.load(Ordering::Acquire) {
3174                        let _ = condvar.wait_for(
3175                            &mut locked,
3176                            Duration::from_millis(sleep_timeout),
3177                        );
3178                    }
3179                    is_sleep.store(false, Ordering::Release);
3180
3181                    continue; //唤醒后立即尝试执行任务
3182                }
3183
3184                sleep_count += 1; //休眠计数
3185                if let Some(interval) = &loop_interval {
3186                    //设置了循环间隔时长
3187                    if let Some(remaining_interval) = Duration::from_millis(*interval).checked_sub(run_time){
3188                        //本次运行少于循环间隔,则休眠剩余的循环间隔,并继续执行任务
3189                        thread::sleep(remaining_interval);
3190                    }
3191                }
3192            } else {
3193                //当前有任务
3194                sleep_count = 0; //重置休眠计数
3195                if let Some(interval) = &loop_interval {
3196                    //设置了循环间隔时长
3197                    if let Some(remaining_interval) = Duration::from_millis(*interval).checked_sub(run_time){
3198                        //本次运行少于循环间隔,则休眠剩余的循环间隔,并继续执行任务
3199                        thread::sleep(remaining_interval);
3200                    }
3201                }
3202            }
3203        }
3204    });
3205
3206    thread_status_copy
3207}
3208
3209/// 唤醒工作者所在线程,如果线程当前正在运行,则忽略
3210pub fn wakeup_worker_thread<O: Default + 'static, P: AsyncTaskPoolExt<O> + AsyncTaskPool<O, Pool = P>>(worker_waker: &Arc<(AtomicBool, Mutex<()>, Condvar)>, rt: &SingleTaskRuntime<O, P>) {
3211    //检查工作者所在线程是否需要唤醒
3212    if worker_waker.0.load(Ordering::Relaxed) && rt.len() > 0 {
3213        let _ = wake_thread_waker(worker_waker);
3214    }
3215}
3216
3217/// 注册全局异常处理器,会替换当前全局异常处理器
3218pub fn register_global_panic_handler<Handler>(handler: Handler)
3219    where Handler: Fn(thread::Thread, String, Option<String>, Option<(String, u32, u32)>) -> Option<i32> + Send + Sync + 'static {
3220    set_hook(Box::new(move |panic_info| {
3221        let thread_info = thread::current();
3222
3223        let payload = panic_info.payload();
3224        let payload_info = match payload.downcast_ref::<&str>() {
3225            None => {
3226                //不是String
3227                match payload.downcast_ref::<String>() {
3228                    None => {
3229                        //不是&'static str,则返回未知异常
3230                        "Unknow panic".to_string()
3231                    },
3232                    Some(info) => {
3233                        info.clone()
3234                    }
3235                }
3236            },
3237            Some(info) => {
3238                info.to_string()
3239            }
3240        };
3241
3242        let other_info = if let Some(arg) = panic_info.payload_as_str() {
3243            Some(arg.to_string())
3244        } else {
3245            None
3246        };
3247
3248        let location = if let Some(location) = panic_info.location() {
3249            Some((location.file().to_string(), location.line(), location.column()))
3250        } else {
3251            None
3252        };
3253
3254        if let Some(exit_code) = handler(thread_info, payload_info, other_info, location) {
3255            //需要关闭当前进程
3256            std::process::exit(exit_code);
3257        }
3258    }));
3259}
3260
3261/// 替换全局内存分配错误处理器
3262pub fn replace_global_alloc_error_handler() {
3263    set_alloc_error_hook(global_alloc_error_handle);
3264}
3265
3266fn global_alloc_error_handle(layout: Layout) {
3267    let bt = Backtrace::new();
3268    eprintln!("[UTC: {}][Thread: {}]Global memory allocation of {:?} bytes failed, stacktrace: \n{:?}",
3269              SystemTime::now().duration_since(SystemTime::UNIX_EPOCH).unwrap().as_millis(),
3270              thread::current().name().unwrap_or(""),
3271              layout.size(),
3272              bt);
3273}
3274
3275// 立即异步让出当前任务执行
3276pub(crate) struct YieldNow(bool);
3277
3278impl Future for YieldNow {
3279    type Output = ();
3280
3281    fn poll(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output> {
3282        if self.0 {
3283            Poll::Ready(())
3284        } else {
3285            self.0 = true;
3286            cx.waker().wake_by_ref();
3287            Poll::Pending
3288        }
3289    }
3290}