1use std::sync::Arc;
31use std::vec::IntoIter;
32use std::time::Duration;
33use std::future::Future;
34use std::cell::UnsafeCell;
35use std::marker::PhantomData;
36use std::io::{Error, ErrorKind, Result};
37use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
38use std::task::{Context, Poll, Waker};
39use std::thread::{self, Builder};
40
41use async_stream::stream;
42use crossbeam_channel::{bounded, Sender};
43use crossbeam_deque::{Injector, Steal, Stealer, Worker};
44use crossbeam_queue::{ArrayQueue, SegQueue};
45use st3::{StealError,
46 fifo::{Worker as FIFOWorker, Stealer as FIFOStealer}};
47use flume::bounded as async_bounded;
48use futures::{
49 future::{BoxFuture, FutureExt},
50 stream::{BoxStream, Stream, StreamExt},
51 task::waker_ref,
52 TryFuture,
53};
54use parking_lot::{Condvar, Mutex};
55use rand::{Rng, thread_rng};
56use num_cpus;
57use wrr::IWRRSelector;
58use quanta::{Clock, Instant as QInstant};
59use log::warn;
60
61use super::{
62 PI_ASYNC_LOCAL_THREAD_ASYNC_RUNTIME, PI_ASYNC_THREAD_LOCAL_ID, DEFAULT_MAX_HIGH_PRIORITY_BOUNDED, DEFAULT_HIGH_PRIORITY_BOUNDED, DEFAULT_MAX_LOW_PRIORITY_BOUNDED, alloc_rt_uid, local_async_runtime, AsyncMapReduce, AsyncPipelineResult, AsyncRuntime,
63 AsyncRuntimeExt, AsyncTask, AsyncTaskPool, AsyncTaskPoolExt, AsyncTaskTimerByNotCancel, AsyncTimingTask,
64 AsyncWait, AsyncWaitAny, AsyncWaitAnyCallback, AsyncWaitTimeout, LocalAsyncWaitTimeout, LocalAsyncRuntime, TaskId, TaskHandle, YieldNow, prune_stale_waiting_workers, register_waiting_worker, wake_waiting_worker
65};
66
67#[cfg(not(target_arch = "wasm32"))]
71const DEFAULT_INIT_WORKER_SIZE: usize = 2;
72#[cfg(target_arch = "wasm32")]
73const DEFAULT_INIT_WORKER_SIZE: usize = 1;
74
75const DEFAULT_WORKER_THREAD_PREFIX: &str = "Default-Multi-RT";
79
80const DEFAULT_THREAD_STACK_SIZE: usize = 1024 * 1024;
84
85const DEFAULT_WORKER_THREAD_SLEEP_TIME: u64 = 10;
89
90const DEFAULT_RUNTIME_SLEEP_TIME: u64 = 1000;
94
95const DEFAULT_MAX_WEIGHT: u8 = 254;
99
100const DEFAULT_MIN_WEIGHT: u8 = 1;
104
105struct ComputationalTaskQueue<O: Default + 'static> {
109 stack: Worker<Arc<AsyncTask<ComputationalTaskPool<O>, O>>>, queue: SegQueue<Arc<AsyncTask<ComputationalTaskPool<O>, O>>>, thread_waker: Arc<(AtomicBool, Mutex<()>, Condvar)>, }
113
114impl<O: Default + 'static> ComputationalTaskQueue<O> {
115 pub fn new(thread_waker: Arc<(AtomicBool, Mutex<()>, Condvar)>) -> Self {
117 let stack = Worker::new_lifo();
118 let queue = SegQueue::new();
119
120 ComputationalTaskQueue {
121 stack,
122 queue,
123 thread_waker,
124 }
125 }
126
127 pub fn len(&self) -> usize {
129 self.stack.len() + self.queue.len()
130 }
131}
132
133pub struct ComputationalTaskPool<O: Default + 'static> {
137 workers: Vec<ComputationalTaskQueue<O>>, waits: Option<Arc<ArrayQueue<Arc<(AtomicBool, Mutex<()>, Condvar)>>>>, consume_count: Arc<AtomicUsize>, produce_count: Arc<AtomicUsize>, }
142
143unsafe impl<O: Default + 'static> Send for ComputationalTaskPool<O> {}
144unsafe impl<O: Default + 'static> Sync for ComputationalTaskPool<O> {}
145
146impl<O: Default + 'static> Default for ComputationalTaskPool<O> {
147 fn default() -> Self {
148 #[cfg(not(target_arch = "wasm32"))]
149 let core_len = num_cpus::get(); #[cfg(target_arch = "wasm32")]
151 let core_len = 1; ComputationalTaskPool::new(core_len)
153 }
154}
155
156impl<O: Default + 'static> AsyncTaskPool<O> for ComputationalTaskPool<O> {
157 type Pool = ComputationalTaskPool<O>;
158
159 #[inline]
160 fn get_thread_id(&self) -> usize {
161 match PI_ASYNC_THREAD_LOCAL_ID.try_with(move |thread_id| unsafe { *thread_id.get() }) {
162 Err(e) => {
163 panic!(
165 "Get thread id failed, thread: {:?}, reason: {:?}",
166 thread::current(),
167 e
168 );
169 }
170 Ok(id) => id,
171 }
172 }
173
174 #[inline]
175 fn len(&self) -> usize {
176 if let Some(len) = self
177 .produce_count
178 .load(Ordering::Relaxed)
179 .checked_sub(self.consume_count.load(Ordering::Relaxed))
180 {
181 len
182 } else {
183 0
184 }
185 }
186
187 #[inline]
188 fn push(&self, task: Arc<AsyncTask<Self::Pool, O>>) -> Result<()> {
189 let index = self.produce_count.fetch_add(1, Ordering::Relaxed) % self.workers.len();
190 self.workers[index].queue.push(task);
191 Ok(())
192 }
193
194 #[inline]
195 fn push_local(&self, task: Arc<AsyncTask<Self::Pool, O>>) -> Result<()> {
196 let id = self.get_thread_id();
197 let rt_uid = task.owner();
198 if (id >> 32) == rt_uid {
199 let worker = &self.workers[id & 0xffffffff];
201 worker.queue.push(task);
202
203 self.produce_count.fetch_add(1, Ordering::Relaxed);
204 Ok(())
205 } else {
206 self.push(task)
208 }
209 }
210
211 #[inline]
212 fn push_priority(&self,
213 priority: usize,
214 task: Arc<AsyncTask<Self::Pool, O>>) -> Result<()> {
215 if priority >= DEFAULT_MAX_HIGH_PRIORITY_BOUNDED {
216 let id = self.get_thread_id();
218 let rt_uid = task.owner();
219 if (id >> 32) == rt_uid {
220 let worker = &self.workers[id & 0xffffffff];
221 worker.stack.push(task);
222
223 self.produce_count.fetch_add(1, Ordering::Relaxed);
224 Ok(())
225 } else {
226 self.push(task)
227 }
228 } else if priority >= DEFAULT_HIGH_PRIORITY_BOUNDED {
229 self.push_local(task)
231 } else {
232 self.push(task)
234 }
235 }
236
237 #[inline]
238 fn push_keep(&self, task: Arc<AsyncTask<Self::Pool, O>>) -> Result<()> {
239 self.push_priority(DEFAULT_HIGH_PRIORITY_BOUNDED, task)
240 }
241
242 #[inline]
243 fn try_pop(&self) -> Option<Arc<AsyncTask<Self::Pool, O>>> {
244 let id = self.get_thread_id() & 0xffffffff;
245 let worker = &self.workers[id];
246 let task = worker.stack.pop();
247 if task.is_some() {
248 self.consume_count.fetch_add(1, Ordering::Relaxed);
250 return task;
251 }
252
253 let task = worker.queue.pop();
254 if task.is_some() {
255 self.consume_count.fetch_add(1, Ordering::Relaxed);
256 }
257
258 task
259 }
260
261 #[inline]
262 fn try_pop_all(&self) -> IntoIter<Arc<AsyncTask<Self::Pool, O>>> {
263 let mut tasks = Vec::with_capacity(self.len());
264 while let Some(task) = self.try_pop() {
265 tasks.push(task);
266 }
267
268 tasks.into_iter()
269 }
270
271 #[inline]
272 fn get_thread_waker(&self) -> Option<&Arc<(AtomicBool, Mutex<()>, Condvar)>> {
273 None
275 }
276}
277
278impl<O: Default + 'static> AsyncTaskPoolExt<O> for ComputationalTaskPool<O> {
279 #[inline]
280 fn set_waits(&mut self, waits: Arc<ArrayQueue<Arc<(AtomicBool, Mutex<()>, Condvar)>>>) {
281 self.waits = Some(waits);
282 }
283
284 #[inline]
285 fn get_waits(&self) -> Option<&Arc<ArrayQueue<Arc<(AtomicBool, Mutex<()>, Condvar)>>>> {
286 self.waits.as_ref()
287 }
288
289 #[inline]
290 fn worker_len(&self) -> usize {
291 self.workers.len()
292 }
293
294 #[inline]
295 fn clone_thread_waker(&self) -> Option<Arc<(AtomicBool, Mutex<()>, Condvar)>> {
296 let worker = &self.workers[self.get_thread_id() & 0xffffffff];
297 Some(worker.thread_waker.clone())
298 }
299}
300
301impl<O: Default + 'static> ComputationalTaskPool<O> {
302 pub fn new(mut size: usize) -> Self {
304 if size < DEFAULT_INIT_WORKER_SIZE {
305 size = DEFAULT_INIT_WORKER_SIZE;
307 }
308
309 let mut workers = Vec::with_capacity(size);
310 for _ in 0..size {
311 let thread_waker = Arc::new((AtomicBool::new(false), Mutex::new(()), Condvar::new()));
312 let worker = ComputationalTaskQueue::new(thread_waker);
313 workers.push(worker);
314 }
315 let consume_count = Arc::new(AtomicUsize::new(0));
316 let produce_count = Arc::new(AtomicUsize::new(0));
317
318 ComputationalTaskPool {
319 workers,
320 waits: None,
321 consume_count,
322 produce_count,
323 }
324 }
325}
326
327struct StealableTaskQueue<O: Default + 'static> {
331 stack: UnsafeCell<Option<Arc<AsyncTask<StealableTaskPool<O>, O>>>>, internal: FIFOWorker<Arc<AsyncTask<StealableTaskPool<O>, O>>>, external: Worker<Arc<AsyncTask<StealableTaskPool<O>, O>>>, selector: UnsafeCell<IWRRSelector<2>>, thread_waker: Arc<(AtomicBool, Mutex<()>, Condvar)>, }
337
338impl<O: Default + 'static> StealableTaskQueue<O> {
339 pub fn new(
342 init_queue_capacity: usize,
343 thread_waker: Arc<(AtomicBool, Mutex<()>, Condvar)>,
344 ) -> (Self,
345 FIFOStealer<Arc<AsyncTask<StealableTaskPool<O>, O>>>,
346 Stealer<Arc<AsyncTask<StealableTaskPool<O>, O>>>) {
347 let stack = UnsafeCell::new(None);
348 let internal = FIFOWorker::new(init_queue_capacity);
349 let external = Worker::new_fifo();
350 let internal_stealer = internal.stealer();
351 let external_stealer = external.stealer();
352 let selector = UnsafeCell::new(IWRRSelector::new([2, 1]));
353
354 (
355 StealableTaskQueue {
356 stack,
357 internal,
358 external,
359 selector,
360 thread_waker,
361 },
362 internal_stealer,
363 external_stealer
364 )
365 }
366
367 pub const fn stack_capacity(&self) -> usize {
369 1
370 }
371
372 pub fn internal_capacity(&self) -> usize {
374 self.internal.capacity()
375 }
376
377 pub fn remaining_internal_capacity(&self) -> usize {
379 self.internal.spare_capacity()
380 }
381
382 #[inline]
384 pub fn stack_len(&self) -> usize {
385 unsafe {
386 if (&*self.stack.get()).is_some() {
387 1
388 } else {
389 0
390 }
391 }
392 }
393
394 pub fn internal_len(&self) -> usize {
396 self
397 .internal_capacity()
398 .checked_sub(self.remaining_internal_capacity())
399 .unwrap_or(0)
400 }
401
402 pub fn external_len(&self) -> usize {
404 self.external.len()
405 }
406}
407
408pub struct StealableTaskPool<O: Default + 'static> {
412 public: Injector<Arc<AsyncTask<StealableTaskPool<O>, O>>>, workers: Vec<StealableTaskQueue<O>>, internal_stealers: Vec<FIFOStealer<Arc<AsyncTask<StealableTaskPool<O>, O>>>>, external_stealers: Vec<Stealer<Arc<AsyncTask<StealableTaskPool<O>, O>>>>, internal_consume: AtomicUsize, internal_produce: AtomicUsize, internal_traffic_statistics: AtomicUsize, external_consume: AtomicUsize, external_produce: AtomicUsize, external_traffic_statistics: AtomicUsize, weights: [u8; 2], clock: Clock, interval: usize, last_time: UnsafeCell<QInstant>, waits: Option<Arc<ArrayQueue<Arc<(AtomicBool, Mutex<()>, Condvar)>>>>, }
428
429unsafe impl<O: Default + 'static> Send for StealableTaskPool<O> {}
430unsafe impl<O: Default + 'static> Sync for StealableTaskPool<O> {}
431
432impl<O: Default + 'static> Default for StealableTaskPool<O> {
433 fn default() -> Self {
434 StealableTaskPool::new()
435 }
436}
437
438impl<O: Default + 'static> AsyncTaskPool<O> for StealableTaskPool<O> {
439 type Pool = StealableTaskPool<O>;
440
441 #[inline]
442 fn get_thread_id(&self) -> usize {
443 match PI_ASYNC_THREAD_LOCAL_ID.try_with(move |thread_id| unsafe { *thread_id.get() }) {
444 Err(e) => {
445 panic!(
447 "Get thread id failed, thread: {:?}, reason: {:?}",
448 thread::current(),
449 e
450 );
451 }
452 Ok(id) => id,
453 }
454 }
455
456 #[inline]
457 fn len(&self) -> usize {
458 self.internal_produce
459 .load(Ordering::Relaxed)
460 .checked_sub(self.internal_consume.load(Ordering::Relaxed))
461 .unwrap_or(0)
462 +
463 self.external_produce
464 .load(Ordering::Relaxed)
465 .checked_sub(self.external_consume.load(Ordering::Relaxed))
466 .unwrap_or(0)
467 }
468
469 #[inline]
470 fn push(&self, task: Arc<AsyncTask<Self::Pool, O>>) -> Result<()> {
471 self.public.push(task);
472
473 self
474 .external_produce
475 .fetch_add(1, Ordering::Relaxed);
476 Ok(())
477 }
478
479 #[inline]
480 fn push_local(&self, task: Arc<AsyncTask<Self::Pool, O>>) -> Result<()> {
481 let id = self.get_thread_id();
482 let rt_uid = task.owner();
483 if (id >> 32) == rt_uid {
484 let worker = &self.workers[id & 0xffffffff];
486 if worker.remaining_internal_capacity() > 0 {
487 let _ = worker.internal.push(task);
489
490 self
491 .internal_produce
492 .fetch_add(1, Ordering::Relaxed);
493 Ok(())
494 } else {
495 self.push(task)
497 }
498 } else {
499 self.push(task)
501 }
502 }
503
504 #[inline]
505 fn push_priority(&self,
506 priority: usize,
507 task: Arc<AsyncTask<Self::Pool, O>>) -> Result<()> {
508 if priority >= DEFAULT_MAX_HIGH_PRIORITY_BOUNDED {
509 let id = self.get_thread_id();
511 let rt_uid = task.owner();
512 if (id >> 32) == rt_uid {
513 let worker = &self.workers[id & 0xffffffff];
515 if worker.stack_len() < 1 {
516 unsafe {
518 *worker.stack.get() = Some(task);
519 }
520 } else if worker.remaining_internal_capacity() > 0 {
521 let _ = worker.internal.push(task);
523 } else {
524 return self.push(task);
526 }
527
528 self
529 .internal_produce
530 .fetch_add(1, Ordering::Relaxed);
531 Ok(())
532 } else {
533 self.push(task)
535 }
536 } else if priority >= DEFAULT_HIGH_PRIORITY_BOUNDED {
537 self.push_local(task)
539 } else {
540 self.push(task)
542 }
543 }
544
545 #[inline]
546 fn push_keep(&self, task: Arc<AsyncTask<Self::Pool, O>>) -> Result<()> {
547 self.push_priority(DEFAULT_HIGH_PRIORITY_BOUNDED, task)
548 }
549
550 #[inline]
551 fn try_pop(&self) -> Option<Arc<AsyncTask<Self::Pool, O>>> {
552 let id = self.get_thread_id() & 0xffffffff;
553 let worker = &self.workers[id];
554 let task = unsafe { (&mut *worker
555 .stack
556 .get())
557 .take()
558 };
559 if task.is_some() {
560 return task;
562 }
563
564 try_pop_by_weight(self, worker, id)
566 }
567
568 #[inline]
569 fn try_pop_all(&self) -> IntoIter<Arc<AsyncTask<Self::Pool, O>>> {
570 let mut tasks = Vec::with_capacity(self.len());
571 while let Some(task) = self.try_pop() {
572 tasks.push(task);
573 }
574
575 tasks.into_iter()
576 }
577
578 #[inline]
579 fn get_thread_waker(&self) -> Option<&Arc<(AtomicBool, Mutex<()>, Condvar)>> {
580 None
582 }
583}
584
585const fn get_msb(n: usize) -> usize {
587 usize::BITS as usize - n.leading_zeros() as usize
588}
589
590fn try_pop_by_weight<O: Default + 'static>(pool: &StealableTaskPool<O>,
592 local_worker: &StealableTaskQueue<O>,
593 local_worker_id: usize)
594 -> Option<Arc<AsyncTask<StealableTaskPool<O>, O>>> {
595 unsafe {
596 let duration = pool
597 .clock
598 .recent()
599 .duration_since(*pool.last_time.get())
600 .as_millis() as usize;
601 if duration >= pool.interval {
602 let new_external_traffic_statistics = pool
604 .external_produce
605 .load(Ordering::Relaxed);
606 let new_internal_traffic_statistics = pool
607 .internal_produce
608 .load(Ordering::Relaxed);
609
610 let external_delta = if new_external_traffic_statistics == 0 {
612 1
614 } else {
615 new_external_traffic_statistics
617 .checked_sub(pool
618 .external_traffic_statistics
619 .load(Ordering::Relaxed))
620 .unwrap_or(1)
621 };
622 pool
623 .external_traffic_statistics
624 .store(new_external_traffic_statistics, Ordering::Relaxed); let internal_delta = if new_internal_traffic_statistics == 0 {
626 1
628 } else {
629 new_internal_traffic_statistics
631 .checked_sub(pool
632 .internal_traffic_statistics
633 .load(Ordering::Relaxed))
634 .unwrap_or(1)
635 };
636 pool
637 .internal_traffic_statistics
638 .store(new_internal_traffic_statistics, Ordering::Relaxed); let selector = &mut *local_worker.selector.get();
642 if external_delta > internal_delta {
643 let msb = get_msb(internal_delta);
645 let internal_weight
646 = (internal_delta >> msb.checked_sub(2).unwrap_or(0)).max(1);
647 let external_weight
648 = ((external_delta >> msb).min(DEFAULT_MAX_WEIGHT as usize)).max(1);
649
650 selector.change_weight(0, external_weight as u8);
651 selector.change_weight(1, internal_weight as u8);
652 } else if external_delta < internal_delta {
653 let msb = get_msb(external_delta);
655 let external_weight
656 = (external_delta >> msb.checked_sub(2).unwrap_or(0)).max(1);
657 let internal_weight
658 = ((internal_delta >> msb).min(DEFAULT_MAX_WEIGHT as usize)).max(1);
659
660 selector.change_weight(0, external_weight as u8);
661 selector.change_weight(1, internal_weight as u8);
662 } else {
663 selector.change_weight(0, 1);
665 selector.change_weight(1, 1);
666 }
667
668 *pool.last_time.get() = pool.clock.recent(); }
670
671 match (&mut *local_worker.selector.get()).select() {
673 0 => {
674 let task = try_pop_external(pool, local_worker, local_worker_id);
676 if task.is_some() {
677 task
678 } else {
679 try_pop_internal(pool, local_worker, local_worker_id)
681 }
682 },
683 _ => {
684 let task = try_pop_internal(pool, local_worker, local_worker_id);
686 if task.is_some() {
687 task
688 } else {
689 try_pop_external(pool, local_worker, local_worker_id)
691 }
692 },
693 }
694 }
695}
696
697#[inline]
699fn try_pop_internal<O: Default + 'static>(pool: &StealableTaskPool<O>,
700 local_worker: &StealableTaskQueue<O>,
701 local_worker_id: usize)
702 -> Option<Arc<AsyncTask<StealableTaskPool<O>, O>>> {
703 let task = local_worker
704 .internal
705 .pop();
706 if task.is_some() {
707 pool
709 .internal_consume
710 .fetch_add(1, Ordering::Relaxed);
711 task
712 } else {
713 let mut gen = thread_rng();
715 let mut worker_stealers: Vec<&FIFOStealer<Arc<AsyncTask<StealableTaskPool<O>, O>>>> = pool
716 .internal_stealers
717 .iter()
718 .enumerate()
719 .filter_map(|(index, other)| {
720 if index != local_worker_id {
721 Some(other)
722 } else {
723 None
725 }
726 })
727 .collect();
728
729 let remaining_len = local_worker.remaining_internal_capacity();
730 loop {
731 if worker_stealers.len() == 0 {
733 break;
735 }
736
737 let index = gen.gen_range(0..worker_stealers.len());
738 let worker_stealer = worker_stealers.swap_remove(index);
739
740 match worker_stealer.steal_and_pop(&local_worker.internal,
741 |count| {
742 let stealable_len = count / 2;
743 if stealable_len <= remaining_len {
744 if stealable_len == 0 {
746 1
747 } else {
748 stealable_len
749 }
750 } else {
751 remaining_len
753 }
754 }) {
755 Err(StealError::Empty) => {
756 continue;
758 },
759 Err(StealError::Busy) => {
760 continue;
762 },
763 Ok((task, _)) => {
764 pool.internal_consume.fetch_add(1, Ordering::Relaxed);
766 return Some(task);
767 },
768 }
769 }
770
771 None
772 }
773}
774
775#[inline]
777fn try_pop_external<O: Default + 'static>(pool: &StealableTaskPool<O>,
778 local_worker: &StealableTaskQueue<O>,
779 local_worker_id: usize)
780 -> Option<Arc<AsyncTask<StealableTaskPool<O>, O>>> {
781 let task = local_worker
782 .external
783 .pop();
784 if task.is_some() {
785 pool
787 .external_consume
788 .fetch_add(1, Ordering::Relaxed);
789 task
790 } else {
791 let task = try_pop_public(pool, local_worker);
793 if task.is_some() {
794 pool
796 .external_consume
797 .fetch_add(1, Ordering::Relaxed);
798 task
799 } else {
800 let mut gen = thread_rng();
802 let mut worker_stealers: Vec<&Stealer<Arc<AsyncTask<StealableTaskPool<O>, O>>>> = pool
803 .external_stealers
804 .iter()
805 .enumerate()
806 .filter_map(|(index, other)| {
807 if index != local_worker_id {
808 Some(other)
809 } else {
810 None
812 }
813 })
814 .collect();
815
816 loop {
817 if worker_stealers.len() == 0 {
819 break;
821 }
822
823 let index = gen.gen_range(0..worker_stealers.len());
824 let worker_stealer = worker_stealers.swap_remove(index);
825
826 match worker_stealer.steal_batch_and_pop(&local_worker.external) {
827 Steal::Success(task) => {
828 pool.external_consume.fetch_add(1, Ordering::Relaxed);
830 return Some(task);
831 },
832 Steal::Retry => {
833 continue;
835 },
836 Steal::Empty => {
837 continue;
839 },
840 }
841 }
842
843 None
844 }
845 }
846}
847
848#[inline]
850fn try_pop_public<O: Default + 'static>(pool: &StealableTaskPool<O>,
851 local_worker: &StealableTaskQueue<O>)
852 -> Option<Arc<AsyncTask<StealableTaskPool<O>, O>>> {
853 loop {
854 match pool.public.steal_batch_and_pop(&local_worker.external) {
855 Steal::Empty => {
856 return None;
858 },
859 Steal::Retry => {
860 continue;
862 },
863 Steal::Success(task) => {
864 pool.external_consume.fetch_add(1, Ordering::Relaxed);
866 return Some(task);
867 },
868 }
869 }
870}
871
872impl<O: Default + 'static> AsyncTaskPoolExt<O> for StealableTaskPool<O> {
873 #[inline]
874 fn set_waits(&mut self, waits: Arc<ArrayQueue<Arc<(AtomicBool, Mutex<()>, Condvar)>>>) {
875 self.waits = Some(waits);
876 }
877
878 #[inline]
879 fn get_waits(&self) -> Option<&Arc<ArrayQueue<Arc<(AtomicBool, Mutex<()>, Condvar)>>>> {
880 self.waits.as_ref()
881 }
882
883 #[inline]
884 fn worker_len(&self) -> usize {
885 self.workers.len()
886 }
887
888 #[inline]
889 fn clone_thread_waker(&self) -> Option<Arc<(AtomicBool, Mutex<()>, Condvar)>> {
890 if let Some(worker) = self.workers.get(self.get_thread_id() & 0xffffffff) {
891 return Some(worker.thread_waker.clone());
892 }
893
894 None
895 }
896}
897
898impl<O: Default + 'static> StealableTaskPool<O> {
899 pub fn new() -> Self {
901 #[cfg(not(target_arch = "wasm32"))]
902 let size = num_cpus::get_physical() * 2; #[cfg(target_arch = "wasm32")]
904 let size = 1; StealableTaskPool::with(size,
906 0x8000,
907 [1, 1],
908 3000)
909 }
910
911 pub fn with(worker_size: usize,
913 internal_queue_capacity: usize,
914 weights: [u8; 2],
915 interval: usize) -> Self {
916 if worker_size == 0 {
917 panic!(
919 "Create WorkerTaskPool failed, worker size: {}, reason: invalid worker size",
920 worker_size
921 );
922 }
923 if interval == 0 {
924 panic!(
925 "Create WorkerTaskPool failed, interval: {}, reason: invalid interval",
926 worker_size
927 );
928 }
929
930 let public = Injector::new();
931 let mut workers = Vec::with_capacity(worker_size);
932 let mut internal_stealers = Vec::with_capacity(worker_size);
933 let mut external_stealers = Vec::with_capacity(worker_size);
934 for _ in 0..worker_size {
935 let thread_waker = Arc::new((AtomicBool::new(false), Mutex::new(()), Condvar::new()));
937 let (worker,
938 internal_stealer,
939 external_stealer) =
940 StealableTaskQueue::new(internal_queue_capacity,
941 thread_waker);
942 workers.push(worker);
943 internal_stealers.push(internal_stealer);
944 external_stealers.push(external_stealer);
945 }
946 let internal_consume = AtomicUsize::new(0);
947 let internal_produce = AtomicUsize::new(0);
948 let internal_traffic_statistics = AtomicUsize::new(0);
949 let external_consume = AtomicUsize::new(0);
950 let external_produce = AtomicUsize::new(0);
951 let external_traffic_statistics = AtomicUsize::new(0);
952 let clock = Clock::new();
953 let last_time = UnsafeCell::new(clock.recent());
954
955 StealableTaskPool {
956 public,
957 workers,
958 internal_stealers,
959 external_stealers,
960 internal_consume,
961 internal_produce,
962 internal_traffic_statistics,
963 external_consume,
964 external_produce,
965 external_traffic_statistics,
966 weights,
967 clock,
968 interval,
969 last_time,
970 waits: None,
971 }
972 }
973}
974
975pub struct MultiTaskRuntime<
979 O: Default + 'static = (),
980 P: AsyncTaskPoolExt<O> + AsyncTaskPool<O> = StealableTaskPool<O>,
981>(
982 Arc<(
983 usize, Arc<P>, Option<
986 Vec<(
987 Sender<(usize, AsyncTimingTask<P, O>)>,
988 Arc<AsyncTaskTimerByNotCancel<P, O>>,
989 )>,
990 >, AtomicUsize, Arc<ArrayQueue<Arc<(AtomicBool, Mutex<()>, Condvar)>>>, AtomicUsize, AtomicUsize, )>,
996);
997
998unsafe impl<O: Default + 'static, P: AsyncTaskPoolExt<O> + AsyncTaskPool<O>> Send
999 for MultiTaskRuntime<O, P>
1000{
1001}
1002unsafe impl<O: Default + 'static, P: AsyncTaskPoolExt<O> + AsyncTaskPool<O>> Sync
1003 for MultiTaskRuntime<O, P>
1004{
1005}
1006
1007impl<O: Default + 'static, P: AsyncTaskPoolExt<O> + AsyncTaskPool<O>> Clone
1008 for MultiTaskRuntime<O, P>
1009{
1010 fn clone(&self) -> Self {
1011 MultiTaskRuntime(self.0.clone())
1012 }
1013}
1014
1015impl<O: Default + 'static, P: AsyncTaskPoolExt<O> + AsyncTaskPool<O, Pool = P>> AsyncRuntime<O>
1016 for MultiTaskRuntime<O, P>
1017{
1018 type Pool = P;
1019
1020 fn shared_pool(&self) -> Arc<Self::Pool> {
1022 (self.0).1.clone()
1023 }
1024
1025 fn get_id(&self) -> usize {
1027 (self.0).0
1028 }
1029
1030 fn wait_len(&self) -> usize {
1032 (self.0)
1033 .5
1034 .load(Ordering::Relaxed)
1035 .checked_sub((self.0).6.load(Ordering::Relaxed))
1036 .unwrap_or(0)
1037 }
1038
1039 fn len(&self) -> usize {
1041 (self.0).1.len()
1042 }
1043
1044 fn alloc<R: 'static>(&self) -> TaskId {
1046 TaskId(UnsafeCell::new((TaskHandle::<R>::default().into_raw() as u128) << 64 | self.get_id() as u128 & 0xffffffffffffffff))
1047 }
1048
1049 fn spawn<F>(&self, future: F) -> Result<TaskId>
1051 where
1052 F: Future<Output = O> + Send + 'static,
1053 {
1054 let task_id = self.alloc::<F::Output>();
1055 if let Err(e) = self.spawn_by_id(task_id.clone(), future) {
1056 return Err(e);
1057 }
1058
1059 Ok(task_id)
1060 }
1061
1062 fn spawn_local<F>(&self, future: F) -> Result<TaskId>
1064 where
1065 F: Future<Output=O> + Send + 'static {
1066 let task_id = self.alloc::<F::Output>();
1067 if let Err(e) = self.spawn_local_by_id(task_id.clone(), future) {
1068 return Err(e);
1069 }
1070
1071 Ok(task_id)
1072 }
1073
1074 fn spawn_priority<F>(&self, priority: usize, future: F) -> Result<TaskId>
1076 where
1077 F: Future<Output=O> + Send + 'static {
1078 let task_id = self.alloc::<F::Output>();
1079 if let Err(e) = self.spawn_priority_by_id(task_id.clone(), priority, future) {
1080 return Err(e);
1081 }
1082
1083 Ok(task_id)
1084 }
1085
1086 fn spawn_yield<F>(&self, future: F) -> Result<TaskId>
1088 where
1089 F: Future<Output=O> + Send + 'static {
1090 let task_id = self.alloc::<F::Output>();
1091 if let Err(e) = self.spawn_yield_by_id(task_id.clone(), future) {
1092 return Err(e);
1093 }
1094
1095 Ok(task_id)
1096 }
1097
1098 fn spawn_timing<F>(&self, future: F, time: usize) -> Result<TaskId>
1100 where
1101 F: Future<Output = O> + Send + 'static,
1102 {
1103 let task_id = self.alloc::<F::Output>();
1104 if let Err(e) = self.spawn_timing_by_id(task_id.clone(), future, time) {
1105 return Err(e);
1106 }
1107
1108 Ok(task_id)
1109 }
1110
1111 fn spawn_by_id<F>(&self, task_id: TaskId, future: F) -> Result<()>
1113 where
1114 F: Future<Output=O> + Send + 'static {
1115 let result = {
1116 (self.0).1.push(Arc::new(AsyncTask::new(
1117 task_id,
1118 (self.0).1.clone(),
1119 DEFAULT_MAX_LOW_PRIORITY_BOUNDED,
1120 Some(future.boxed()),
1121 )))
1122 };
1123
1124 let _ = wake_waiting_worker(&(self.0).4);
1125
1126 result
1127 }
1128
1129 fn spawn_local_by_id<F>(&self, task_id: TaskId, future: F) -> Result<()>
1130 where
1131 F: Future<Output=O> + Send + 'static {
1132 let should_wake = PI_ASYNC_THREAD_LOCAL_ID
1133 .try_with(|thread_id| unsafe { ((*thread_id.get()) >> 32) != self.get_id() })
1134 .unwrap_or(true);
1135 let result = (self.0).1.push_local(Arc::new(AsyncTask::new(
1136 task_id,
1137 (self.0).1.clone(),
1138 DEFAULT_HIGH_PRIORITY_BOUNDED,
1139 Some(future.boxed()),
1140 )));
1141
1142 if should_wake {
1143 let _ = wake_waiting_worker(&(self.0).4);
1144 }
1145
1146 result
1147 }
1148
1149 fn spawn_priority_by_id<F>(&self,
1151 task_id: TaskId,
1152 priority: usize,
1153 future: F) -> Result<()>
1154 where
1155 F: Future<Output=O> + Send + 'static {
1156 let result = {
1157 (self.0).1.push_priority(priority, Arc::new(AsyncTask::new(
1158 task_id,
1159 (self.0).1.clone(),
1160 priority,
1161 Some(future.boxed()),
1162 )))
1163 };
1164
1165 let _ = wake_waiting_worker(&(self.0).4);
1166
1167 result
1168 }
1169
1170 #[inline]
1172 fn spawn_yield_by_id<F>(&self, task_id: TaskId, future: F) -> Result<()>
1173 where
1174 F: Future<Output=O> + Send + 'static {
1175 self.spawn_priority_by_id(task_id,
1176 DEFAULT_HIGH_PRIORITY_BOUNDED,
1177 future)
1178 }
1179
1180 fn spawn_timing_by_id<F>(&self,
1182 task_id: TaskId,
1183 future: F,
1184 time: usize) -> Result<()>
1185 where
1186 F: Future<Output=O> + Send + 'static {
1187 let rt = self.clone();
1188 self.spawn_by_id(task_id, async move {
1189 if let Some(timers) = &(rt.0).2 {
1190 let id = (rt.0).1.get_thread_id() & 0xffffffff;
1192 let (_, timer) = &timers[id];
1193 timer.set_timer(
1194 AsyncTimingTask::WaitRun(Arc::new(AsyncTask::new(
1195 rt.alloc::<F::Output>(),
1196 (rt.0).1.clone(),
1197 DEFAULT_MAX_HIGH_PRIORITY_BOUNDED,
1198 Some(future.boxed()),
1199 ))),
1200 time,
1201 );
1202
1203 (rt.0).5.fetch_add(1, Ordering::Relaxed);
1204 }
1205
1206 Default::default()
1207 })
1208 }
1209
1210 fn pending<Output: 'static>(&self, task_id: &TaskId, waker: Waker) -> Poll<Output> {
1212 task_id.set_waker::<Output>(waker);
1213 Poll::Pending
1214 }
1215
1216 fn wakeup<Output: 'static>(&self, task_id: &TaskId) {
1218 task_id.wakeup::<Output>();
1219 }
1220
1221 fn wait<V: Send + 'static>(&self) -> AsyncWait<V> {
1223 AsyncWait(self.wait_any(2))
1224 }
1225
1226 fn wait_any<V: Send + 'static>(&self, capacity: usize) -> AsyncWaitAny<V> {
1228 let (producor, consumer) = async_bounded(capacity);
1229
1230 AsyncWaitAny {
1231 capacity,
1232 producor,
1233 consumer,
1234 }
1235 }
1236
1237 fn wait_any_callback<V: Send + 'static>(&self, capacity: usize) -> AsyncWaitAnyCallback<V> {
1239 let (producor, consumer) = async_bounded(capacity);
1240
1241 AsyncWaitAnyCallback {
1242 capacity,
1243 producor,
1244 consumer,
1245 }
1246 }
1247
1248 fn map_reduce<V: Send + 'static>(&self, capacity: usize) -> AsyncMapReduce<V> {
1250 let (producor, consumer) = async_bounded(capacity);
1251
1252 AsyncMapReduce {
1253 count: 0,
1254 capacity,
1255 producor,
1256 consumer,
1257 }
1258 }
1259
1260 fn timeout(&self, timeout: usize) -> BoxFuture<'static, ()> {
1262 let rt = self.clone();
1263
1264 if let Some(timers) = &(self.0).2 {
1265 match PI_ASYNC_THREAD_LOCAL_ID.try_with(move |thread_id| {
1267 let thread_id = unsafe { *thread_id.get() };
1269 let index = thread_id & 0xffffffff;
1270 if index > timers.len() {
1271 TimerTaskProducor::Foreign(timers[(self.0).3.load(Ordering::Relaxed) % timers.len()].0.clone())
1273 } else {
1274 TimerTaskProducor::Local(timers[index].1.clone())
1275 }
1276 }) {
1277 Err(_) => {
1278 panic!("Multi thread runtime timeout failed, reason: local thread id not match")
1279 }
1280 Ok(producor) => match producor {
1281 TimerTaskProducor::Local(timer) => {
1282 LocalAsyncWaitTimeout::new(rt, timer, timeout).boxed()
1283 },
1284 TimerTaskProducor::Foreign(producor) => {
1285 AsyncWaitTimeout::new(rt, producor, timeout).boxed()
1286 },
1287 },
1288 }
1289 } else {
1290 async move {
1292 thread::sleep(Duration::from_millis(timeout as u64));
1293 }
1294 .boxed()
1295 }
1296 }
1297
1298 fn yield_now(&self) -> BoxFuture<'static, ()> {
1300 async move {
1301 YieldNow(false).await;
1302 }.boxed()
1303 }
1304
1305 fn pipeline<S, SO, F, FO>(&self, input: S, mut filter: F) -> BoxStream<'static, FO>
1307 where
1308 S: Stream<Item = SO> + Send + 'static,
1309 SO: Send + 'static,
1310 F: FnMut(SO) -> AsyncPipelineResult<FO> + Send + 'static,
1311 FO: Send + 'static,
1312 {
1313 let output = stream! {
1314 for await value in input {
1315 match filter(value) {
1316 AsyncPipelineResult::Disconnect => {
1317 break;
1319 },
1320 AsyncPipelineResult::Filtered(result) => {
1321 yield result;
1322 },
1323 }
1324 }
1325 };
1326
1327 output.boxed()
1328 }
1329
1330 fn close(&self) -> bool {
1332 false
1333 }
1334}
1335
1336impl<O: Default + 'static, P: AsyncTaskPoolExt<O> + AsyncTaskPool<O, Pool = P>> AsyncRuntimeExt<O>
1337 for MultiTaskRuntime<O, P>
1338{
1339 fn spawn_with_context<F, C>(&self, task_id: TaskId, future: F, context: C) -> Result<()>
1340 where
1341 F: Future<Output = O> + Send + 'static,
1342 C: 'static,
1343 {
1344 let task = Arc::new(AsyncTask::with_context(
1345 task_id,
1346 (self.0).1.clone(),
1347 DEFAULT_MAX_LOW_PRIORITY_BOUNDED,
1348 Some(future.boxed()),
1349 context,
1350 ));
1351 let result = (self.0).1.push(task);
1352
1353 let _ = wake_waiting_worker(&(self.0).4);
1354
1355 result
1356 }
1357
1358 fn spawn_timing_with_context<F, C>(
1359 &self,
1360 task_id: TaskId,
1361 future: F,
1362 context: C,
1363 time: usize,
1364 ) -> Result<()>
1365 where
1366 F: Future<Output = O> + Send + 'static,
1367 C: Send + 'static,
1368 {
1369 let rt = self.clone();
1370 self.spawn_by_id(task_id, async move {
1371 if let Some(timers) = &(rt.0).2 {
1372 let id = (rt.0).1.get_thread_id() & 0xffffffff;
1374 let (_, timer) = &timers[id];
1375 timer.set_timer(
1376 AsyncTimingTask::WaitRun(Arc::new(AsyncTask::with_context(
1377 rt.alloc::<F::Output>(),
1378 (rt.0).1.clone(),
1379 DEFAULT_MAX_HIGH_PRIORITY_BOUNDED,
1380 Some(future.boxed()),
1381 context,
1382 ))),
1383 time,
1384 );
1385
1386 (rt.0).5.fetch_add(1, Ordering::Relaxed);
1387 }
1388
1389 Default::default()
1390 })
1391 }
1392
1393 fn block_on<F>(&self, future: F) -> Result<F::Output>
1394 where
1395 F: Future + Send + 'static,
1396 <F as Future>::Output: Default + Send + 'static,
1397 {
1398 if let Some(local_rt) = local_async_runtime::<F::Output>() {
1400 if local_rt.get_id() == self.get_id() {
1402 return Err(Error::new(
1404 ErrorKind::WouldBlock,
1405 format!("Block on failed, reason: would block"),
1406 ));
1407 }
1408 }
1409
1410 let (sender, receiver) = bounded(1);
1411 if let Err(e) = self.spawn(async move {
1412 let r = future.await;
1414 sender.send(r);
1415
1416 Default::default()
1417 }) {
1418 return Err(Error::new(
1419 ErrorKind::Other,
1420 format!("Block on failed, reason: {:?}", e),
1421 ));
1422 }
1423
1424 match receiver.recv() {
1426 Err(e) => Err(Error::new(
1427 ErrorKind::Other,
1428 format!("Block on failed, reason: {:?}", e),
1429 )),
1430 Ok(result) => Ok(result),
1431 }
1432 }
1433}
1434
1435impl<O: Default + 'static, P: AsyncTaskPoolExt<O> + AsyncTaskPool<O, Pool = P>>
1436 MultiTaskRuntime<O, P>
1437{
1438 pub fn idler_len(&self) -> usize {
1440 (self.0).1.idler_len()
1441 }
1442
1443 pub fn worker_len(&self) -> usize {
1445 (self.0).1.worker_len()
1446 }
1447
1448 pub fn buffer_len(&self) -> usize {
1450 (self.0).1.buffer_len()
1451 }
1452
1453 pub fn to_local_runtime(&self) -> LocalAsyncRuntime<O> {
1455 LocalAsyncRuntime {
1456 inner: self.as_raw(),
1457 get_id_func: MultiTaskRuntime::<O, P>::get_id_raw,
1458 spawn_func: MultiTaskRuntime::<O, P>::spawn_raw,
1459 spawn_local_func: MultiTaskRuntime::<O, P>::spawn_local_raw,
1460 spawn_timing_func: MultiTaskRuntime::<O, P>::spawn_timing_raw,
1461 timeout_func: MultiTaskRuntime::<O, P>::timeout_raw,
1462 }
1463 }
1464
1465 #[inline]
1467 pub(crate) fn as_raw(&self) -> *const () {
1468 Arc::into_raw(self.0.clone()) as *const ()
1469 }
1470
1471 #[inline]
1473 pub(crate) fn from_raw(raw: *const ()) -> Self {
1474 let inner = unsafe {
1475 Arc::from_raw(
1476 raw as *const (
1477 usize,
1478 Arc<P>,
1479 Option<
1480 Vec<(
1481 Sender<(usize, AsyncTimingTask<P, O>)>,
1482 Arc<AsyncTaskTimerByNotCancel<P, O>>,
1483 )>,
1484 >,
1485 AtomicUsize,
1486 Arc<ArrayQueue<Arc<(AtomicBool, Mutex<()>, Condvar)>>>,
1487 AtomicUsize,
1488 AtomicUsize,
1489 ),
1490 )
1491 };
1492 MultiTaskRuntime(inner)
1493 }
1494
1495 pub(crate) fn get_id_raw(raw: *const ()) -> usize {
1497 let rt = MultiTaskRuntime::<O, P>::from_raw(raw);
1498 let id = rt.get_id();
1499 Arc::into_raw(rt.0); id
1501 }
1502
1503 pub(crate) fn spawn_raw<F>(raw: *const (), future: F) -> Result<()>
1505 where
1506 F: Future<Output = O> + Send + 'static,
1507 {
1508 let rt = MultiTaskRuntime::<O, P>::from_raw(raw);
1509 let result = rt.spawn_by_id(rt.alloc::<F::Output>(), future);
1510 Arc::into_raw(rt.0); result
1512 }
1513
1514 pub(crate) fn spawn_local_raw<F>(raw: *const (), future: F) -> Result<()>
1516 where
1517 F: Future<Output = O> + Send + 'static,
1518 {
1519 let rt = MultiTaskRuntime::<O, P>::from_raw(raw);
1520 let result = rt.spawn_local_by_id(rt.alloc::<F::Output>(), future);
1521 Arc::into_raw(rt.0); result
1523 }
1524
1525 pub(crate) fn spawn_timing_raw(
1527 raw: *const (),
1528 future: BoxFuture<'static, O>,
1529 timeout: usize,
1530 ) -> Result<()> {
1531 let rt = MultiTaskRuntime::<O, P>::from_raw(raw);
1532 let result = rt.spawn_timing_by_id(rt.alloc::<O>(), future, timeout);
1533 Arc::into_raw(rt.0); result
1535 }
1536
1537 pub(crate) fn timeout_raw(raw: *const (), timeout: usize) -> BoxFuture<'static, ()> {
1539 let rt = MultiTaskRuntime::<O, P>::from_raw(raw);
1540 let boxed = rt.timeout(timeout);
1541 Arc::into_raw(rt.0); boxed
1543 }
1544}
1545
1546pub struct MultiTaskRuntimeBuilder<
1550 O: Default + 'static = (),
1551 P: AsyncTaskPoolExt<O> + AsyncTaskPool<O> = StealableTaskPool<O>,
1552> {
1553 pool: P, prefix: String, init: usize, min: usize, max: usize, stack_size: usize, timeout: u64, interval: Option<usize>, marker: PhantomData<O>,
1562}
1563
1564unsafe impl<O: Default + 'static, P: AsyncTaskPoolExt<O> + AsyncTaskPool<O>> Send
1565 for MultiTaskRuntimeBuilder<O, P>
1566{
1567}
1568unsafe impl<O: Default + 'static, P: AsyncTaskPoolExt<O> + AsyncTaskPool<O>> Sync
1569 for MultiTaskRuntimeBuilder<O, P>
1570{
1571}
1572
1573impl<O: Default + 'static> Default for MultiTaskRuntimeBuilder<O> {
1574 fn default() -> Self {
1576 #[cfg(not(target_arch = "wasm32"))]
1577 let core_len = num_cpus::get(); #[cfg(target_arch = "wasm32")]
1579 let core_len = 1; let pool = StealableTaskPool::with(core_len,
1581 65535,
1582 [1, 1],
1583 3000);
1584 MultiTaskRuntimeBuilder::new(pool)
1585 .thread_stack_size(2 * 1024 * 1024)
1586 .set_timer_interval(1)
1587 }
1588}
1589
1590impl<O: Default + 'static, P: AsyncTaskPoolExt<O> + AsyncTaskPool<O, Pool = P>>
1591 MultiTaskRuntimeBuilder<O, P>
1592{
1593 pub fn new(mut pool: P) -> Self {
1595 #[cfg(not(target_arch = "wasm32"))]
1596 let core_len = num_cpus::get(); #[cfg(target_arch = "wasm32")]
1598 let core_len = 1; MultiTaskRuntimeBuilder {
1601 pool,
1602 prefix: DEFAULT_WORKER_THREAD_PREFIX.to_string(),
1603 init: core_len,
1604 min: core_len,
1605 max: core_len,
1606 stack_size: DEFAULT_THREAD_STACK_SIZE,
1607 timeout: DEFAULT_WORKER_THREAD_SLEEP_TIME,
1608 interval: None,
1609 marker: PhantomData,
1610 }
1611 }
1612
1613 pub fn thread_prefix(mut self, prefix: &str) -> Self {
1615 self.prefix = prefix.to_string();
1616 self
1617 }
1618
1619 pub fn thread_stack_size(mut self, stack_size: usize) -> Self {
1621 self.stack_size = stack_size;
1622 self
1623 }
1624
1625 pub fn init_worker_size(mut self, mut init: usize) -> Self {
1627 if init == 0 {
1628 init = DEFAULT_INIT_WORKER_SIZE;
1630 }
1631
1632 self.init = init;
1633 self
1634 }
1635
1636 pub fn set_worker_limit(mut self, mut min: usize, mut max: usize) -> Self {
1638 if self.init > max {
1639 max = self.init;
1641 }
1642
1643 if min == 0 || min > max {
1644 min = max;
1646 }
1647
1648 self.min = min;
1649 self.max = max;
1650 self
1651 }
1652
1653 pub fn set_timeout(mut self, timeout: u64) -> Self {
1655 self.timeout = timeout;
1656 self
1657 }
1658
1659 pub fn set_timer_interval(mut self, interval: usize) -> Self {
1661 self.interval = Some(interval);
1662 self
1663 }
1664
1665 pub fn build(mut self) -> MultiTaskRuntime<O, P> {
1699 let pool_worker_len = self.pool.worker_len();
1700 if pool_worker_len == 0 {
1701 panic!("Build multi thread runtime failed, reason: worker pool is empty");
1702 }
1703 if self.init > pool_worker_len {
1704 self.init = pool_worker_len;
1705 }
1706 if self.max > pool_worker_len {
1707 self.max = pool_worker_len;
1708 }
1709 if self.min > self.max {
1710 self.min = self.max;
1711 }
1712
1713 let interval = self.interval;
1715 let mut timers = if let Some(_) = interval {
1716 Some(Vec::with_capacity(self.max))
1717 } else {
1718 None
1719 };
1720 for _ in 0..self.max {
1721 if let Some(vec) = &mut timers {
1723 let timer = AsyncTaskTimerByNotCancel::new();
1724 let producor = timer.producor.clone();
1725 let timer = Arc::new(timer);
1726 vec.push((producor, timer));
1727 };
1728 }
1729
1730 let rt_uid = alloc_rt_uid();
1732 let waits = Arc::new(ArrayQueue::new(self.max));
1733 let mut pool = self.pool;
1734 pool.set_waits(waits.clone()); let pool = Arc::new(pool);
1736 let runtime = MultiTaskRuntime(Arc::new((
1737 rt_uid,
1738 pool,
1739 timers,
1740 AtomicUsize::new(0),
1741 waits,
1742 AtomicUsize::new(0),
1743 AtomicUsize::new(0),
1744 )));
1745
1746 let mut builders = Vec::with_capacity(self.init);
1748 for index in 0..self.init {
1749 let builder = Builder::new()
1750 .name(self.prefix.clone() + "-" + index.to_string().as_str())
1751 .stack_size(self.stack_size);
1752 builders.push(builder);
1753 }
1754
1755 let min = self.min;
1757 for index in 0..builders.len() {
1758 let builder = builders.remove(0);
1759 let runtime = runtime.clone();
1760 let timeout = self.timeout;
1761 let timer = if let Some(timers) = &(runtime.0).2 {
1762 let (_, timer) = &timers[index];
1763 Some(timer.clone())
1764 } else {
1765 None
1766 };
1767
1768 spawn_worker_thread(builder, index, runtime, min, timeout, interval, timer);
1769 }
1770
1771 runtime
1772 }
1773}
1774
1775fn spawn_worker_thread<
1777 O: Default + 'static,
1778 P: AsyncTaskPoolExt<O> + AsyncTaskPool<O, Pool = P>,
1779>(
1780 builder: Builder,
1781 index: usize,
1782 runtime: MultiTaskRuntime<O, P>,
1783 min: usize,
1784 timeout: u64,
1785 interval: Option<usize>,
1786 timer: Option<Arc<AsyncTaskTimerByNotCancel<P, O>>>,
1787) {
1788 if let Some(timer) = timer {
1789 let rt_uid = runtime.get_id();
1791 let _ = builder.spawn(move || {
1792 if let Err(e) = PI_ASYNC_THREAD_LOCAL_ID.try_with(move |thread_id| unsafe {
1794 *thread_id.get() = rt_uid << 32 | index & 0xffffffff;
1795 }) {
1796 panic!(
1797 "Multi thread runtime startup failed, thread id: {:?}, reason: {:?}",
1798 index, e
1799 );
1800 }
1801
1802 let runtime_copy = runtime.clone();
1804 match PI_ASYNC_LOCAL_THREAD_ASYNC_RUNTIME.try_with(move |rt| {
1805 let raw = Arc::into_raw(Arc::new(runtime_copy.to_local_runtime()))
1806 as *mut LocalAsyncRuntime<O> as *mut ();
1807 rt.store(raw, Ordering::Relaxed);
1808 }) {
1809 Err(e) => {
1810 panic!("Bind multi runtime to local thread failed, reason: {:?}", e);
1811 }
1812 Ok(_) => (),
1813 }
1814
1815 timer_work_loop(
1817 runtime,
1818 index,
1819 min,
1820 timeout,
1821 interval.unwrap() as u64,
1822 timer,
1823 );
1824 });
1825 } else {
1826 let rt_uid = runtime.get_id();
1828 let _ = builder.spawn(move || {
1829 if let Err(e) = PI_ASYNC_THREAD_LOCAL_ID.try_with(move |thread_id| unsafe {
1831 *thread_id.get() = rt_uid << 32 | index & 0xffffffff;
1832 }) {
1833 panic!(
1834 "Multi thread runtime startup failed, thread id: {:?}, reason: {:?}",
1835 index, e
1836 );
1837 }
1838
1839 let runtime_copy = runtime.clone();
1841 match PI_ASYNC_LOCAL_THREAD_ASYNC_RUNTIME.try_with(move |rt| {
1842 let raw = Arc::into_raw(Arc::new(runtime_copy.to_local_runtime()))
1843 as *mut LocalAsyncRuntime<O> as *mut ();
1844 rt.store(raw, Ordering::Relaxed);
1845 }) {
1846 Err(e) => {
1847 panic!("Bind multi runtime to local thread failed, reason: {:?}", e);
1848 }
1849 Ok(_) => (),
1850 }
1851
1852 work_loop(runtime, index, min, timeout);
1854 });
1855 }
1856}
1857
1858enum WorkerWaitResult<O: Default + 'static, P: AsyncTaskPoolExt<O> + AsyncTaskPool<O, Pool = P>> {
1874 TimedOut,
1875 NotSlept,
1876 Task(Arc<AsyncTask<P, O>>),
1877}
1878
1879#[inline]
1946fn worker_wait_for_task<O: Default + 'static, P: AsyncTaskPoolExt<O> + AsyncTaskPool<O, Pool = P>>(
1947 runtime: &MultiTaskRuntime<O, P>,
1948 worker_waker: &Arc<(AtomicBool, Mutex<()>, Condvar)>,
1949 sleep_timeout: u64,
1950) -> WorkerWaitResult<O, P> {
1951 let (is_sleep, lock, condvar) = &**worker_waker;
1952
1953 loop {
1954 let _locked = lock.lock();
1955 if is_sleep.load(Ordering::Acquire) {
1956 break;
1957 }
1958
1959 if register_waiting_worker(&(runtime.0).4, worker_waker) {
1960 is_sleep.store(true, Ordering::Release);
1961 break;
1962 }
1963
1964 drop(_locked);
1965 if prune_stale_waiting_workers(&(runtime.0).4) == 0 {
1966 return WorkerWaitResult::NotSlept;
1967 }
1968 }
1969
1970 if let Some(task) = (runtime.0).1.try_pop() {
1971 is_sleep.store(false, Ordering::Release);
1972 return WorkerWaitResult::Task(task);
1973 }
1974
1975 if runtime.len() > 0 {
1976 is_sleep.store(false, Ordering::Release);
1977 return WorkerWaitResult::NotSlept;
1978 }
1979
1980 let mut locked = lock.lock();
1981 if !is_sleep.load(Ordering::Acquire) {
1982 return WorkerWaitResult::NotSlept;
1983 }
1984
1985 let timed_out = condvar
1986 .wait_for(&mut locked, Duration::from_millis(sleep_timeout))
1987 .timed_out();
1988 is_sleep.store(false, Ordering::Release);
1989
1990 if timed_out {
1991 WorkerWaitResult::TimedOut
1992 } else {
1993 WorkerWaitResult::NotSlept
1994 }
1995}
1996
1997fn timer_work_loop<O: Default + 'static, P: AsyncTaskPoolExt<O> + AsyncTaskPool<O, Pool = P>>(
1999 runtime: MultiTaskRuntime<O, P>,
2000 index: usize,
2001 min: usize,
2002 sleep_timeout: u64,
2003 timer_interval: u64,
2004 timer: Arc<AsyncTaskTimerByNotCancel<P, O>>,
2005) {
2006 let pool = (runtime.0).1.clone();
2008 let worker_waker = pool.clone_thread_waker().unwrap();
2009
2010 let mut sleep_count = 0; let clock = Clock::new();
2012 loop {
2013 let timer_run_millis = clock.recent(); let mut pop_len = 0;
2016 (runtime.0)
2017 .5
2018 .fetch_add(timer.consume(),
2019 Ordering::Relaxed);
2020 loop {
2021 let current_time = timer.is_require_pop();
2022 if let Some(current_time) = current_time {
2023 loop {
2025 let timed_out = timer.pop(current_time);
2026 if let Some(timing_task) = timed_out {
2027 match timing_task {
2028 AsyncTimingTask::Pended(expired) => {
2029 runtime.wakeup::<O>(&expired);
2031 }
2032 AsyncTimingTask::WaitRun(expired) => {
2033 (runtime.0)
2035 .1
2036 .push_priority(DEFAULT_MAX_HIGH_PRIORITY_BOUNDED,
2037 expired);
2038 if let Some(task) = pool.try_pop() {
2039 sleep_count = 0; run_task(&runtime, task);
2041 }
2042 }
2043 AsyncTimingTask::TimeoutWake(waiter) => {
2044 waiter.fire();
2046 }
2047 }
2048 pop_len += 1;
2049
2050 if let Some(task) = pool.try_pop() {
2051 sleep_count = 0; run_task(&runtime, task);
2054 }
2055 } else {
2056 break;
2058 }
2059 }
2060 } else {
2061 break;
2063 }
2064 }
2065 (runtime.0)
2066 .6
2067 .fetch_add(pop_len,
2068 Ordering::Relaxed);
2069
2070 match pool.try_pop() {
2072 None => {
2073 if runtime.len() > 0 {
2074 continue;
2076 }
2077
2078 let diff_time = clock
2080 .recent()
2081 .duration_since(timer_run_millis)
2082 .as_millis() as u64; let real_timeout = if timer.len() == 0 {
2084 sleep_timeout
2086 } else {
2087 if diff_time >= timer_interval {
2089 continue;
2091 } else {
2092 timer_interval - diff_time
2094 }
2095 };
2096
2097 match worker_wait_for_task(&runtime, &worker_waker, real_timeout) {
2099 WorkerWaitResult::TimedOut => {
2100 sleep_count += 1;
2102 },
2103 WorkerWaitResult::Task(task) => {
2104 sleep_count = 0; run_task(&runtime, task);
2106 },
2107 WorkerWaitResult::NotSlept => (),
2108 }
2109 }
2110 Some(task) => {
2111 sleep_count = 0; run_task(&runtime, task);
2114 }
2115 }
2116 }
2117
2118 (runtime.0).1.close_worker();
2120 warn!(
2121 "Worker of runtime closed, runtime: {}, worker: {}, thread: {:?}",
2122 runtime.get_id(),
2123 index,
2124 thread::current()
2125 );
2126}
2127
2128fn work_loop<O: Default + 'static, P: AsyncTaskPoolExt<O> + AsyncTaskPool<O, Pool = P>>(
2130 runtime: MultiTaskRuntime<O, P>,
2131 index: usize,
2132 min: usize,
2133 sleep_timeout: u64,
2134) {
2135 let pool = (runtime.0).1.clone();
2137 let worker_waker = pool.clone_thread_waker().unwrap();
2138
2139 let mut sleep_count = 0; loop {
2141 match pool.try_pop() {
2142 None => {
2143 if runtime.len() > 0 {
2145 continue;
2147 }
2148
2149 match worker_wait_for_task(&runtime, &worker_waker, sleep_timeout) {
2150 WorkerWaitResult::TimedOut => {
2151 sleep_count += 1;
2153 },
2154 WorkerWaitResult::Task(task) => {
2155 sleep_count = 0; run_task(&runtime, task);
2157 },
2158 WorkerWaitResult::NotSlept => (),
2159 }
2160 }
2161 Some(task) => {
2162 sleep_count = 0; run_task(&runtime, task);
2165 }
2166 }
2167 }
2168
2169 (runtime.0).1.close_worker();
2171 warn!(
2172 "Worker of runtime closed, runtime: {}, worker: {}, thread: {:?}",
2173 runtime.get_id(),
2174 index,
2175 thread::current()
2176 );
2177}
2178
2179#[inline]
2181fn run_task<O: Default + 'static, P: AsyncTaskPoolExt<O> + AsyncTaskPool<O, Pool = P>>(
2182 runtime: &MultiTaskRuntime<O, P>,
2183 task: Arc<AsyncTask<P, O>>,
2184) {
2185 let waker = waker_ref(&task);
2186 let mut context = Context::from_waker(&*waker);
2187 if let Some(mut future) = task.get_inner() {
2188 if let Poll::Pending = future.as_mut().poll(&mut context) {
2189 task.set_inner(Some(future));
2191 }
2192 } else {
2193 (runtime.0).1.push(task);
2195 }
2196}
2197
2198enum TimerTaskProducor<
2200 O: Default + 'static = (),
2201 P: AsyncTaskPoolExt<O> + AsyncTaskPool<O> = StealableTaskPool<O>,
2202> {
2203 Local(Arc<AsyncTaskTimerByNotCancel<P, O>>), Foreign(Sender<(usize, AsyncTimingTask<P, O>)>), }