Skip to main content

StealableTaskPool

Struct StealableTaskPool 

Source
pub struct StealableTaskPool<O: Default + 'static> { /* private fields */ }
Expand description

可窃取的混合多线程任务池。

§使用方式

pool 将 external/public 任务和 runtime worker 内产生的 internal/local 任务分开保存,并由 每个 worker 的 IWRRSelector 按既有 best-effort 权重选择顺序。典型用法:

use pi_async_rt::rt::{AsyncRuntime, multi_thread::{MultiTaskRuntimeBuilder, StealableTaskPool}};

let pool = StealableTaskPool::with(4, 4096, [1, 1], 3000);
let runtime = MultiTaskRuntimeBuilder::new(pool)
    .init_worker_size(4)
    .set_worker_limit(4, 4)
    .build();
runtime.spawn(async {}).unwrap();

§调度与业务边界

  • push 可从任意线程进入 public injector;所属 worker 的 push_local 优先进入 internal queue,跨 runtime 调用保持 public fallback。
  • try_pop/try_pop_all 和 worker waker 获取是 owner-only 操作,只允许拥有该 exact pool 的 worker 调用;错误 pool/外部线程会在访问 worker-local 状态前 panic。
  • pool-level 流量计数和刷新时间只用于近似调度启发式,不发布任务对象,也不承诺所有 worker selector 同时或最终应用同一权重。
  • weights 保留既有 API/存储语义;本版本不改变其历史行为,不在此处定义新的公平性保证。
  • timeout、取消、任务结果、worker sleep/wake 和 runtime 关闭由其它层负责。

§性能、纯度和副作用

  • local/public pop 快路径为 O(1);尝试其它 W-1 个 stealer 的最坏时间为 O(W)。空间 O(W + N),N 为排队任务数。
  • 每次 weighted pop 读取一次 AtomicCell<QInstant>;到刷新窗口时写入一次。x86_64 验收 目标要求该类型 lock-free;其它目标保证正确性但不承诺无内部锁。
  • owner 校验 O(1),一次 TLS Copy 读取和 raw pointer 比较;不分配、不 clone、不锁、不阻塞。
  • push/pop/刷新均非纯且非幂等,会修改队列、统计或当前 worker selector;不会在 pool 内 poll 用户 future、执行回调、I/O 或 V8/FFI。

§线程、内存与异步安全

pool-wide 状态使用并发容器/原子;stack、deque worker 和 selector 只由 owner worker 访问。raw pool pointer 仅比较不解引用,worker 闭包持有 runtime Arc 保证生命周期。没有锁 guard 跨 await,也不引入引用环。专项和 TSan 入口: tests/stealable_task_pool_concurrency.rs

Implementations§

Source§

impl<O: Default + 'static> StealableTaskPool<O>

Source

pub fn new() -> Self

使用平台默认 worker slot 数构建可窃取任务池。

非 wasm32 使用物理核数的两倍,wasm32 使用 1;内部 queue capacity、初始权重参数和刷新 interval 保持既有默认值。返回值尚未启动线程或绑定 runtime。

时间/空间复杂度 O(W),会分配 W 组 queue/stealer/waker,因此不是纯函数且非幂等;不 阻塞、不执行 I/O 或 future。owner、安全和错误边界见 StealableTaskPool 类型文档。

Source

pub fn with( worker_size: usize, internal_queue_capacity: usize, weights: [u8; 2], interval: usize, ) -> Self

构建指定 worker slot、internal queue capacity、权重参数和刷新间隔的任务池。

§参数
  • worker_size:worker slot 数,必须大于 0;为 0 时保持旧行为并 panic。
  • internal_queue_capacity:每个 worker internal FIFO 的初始容量;允许 0,由底层 queue 按其既有规则归一化。值只在构建期消费。
  • weights:保留的两类队列权重配置。当前版本保持历史存储/调度行为,不新增“所有 selector 以该值初始化”或“全局收敛”保证;调用方不得据此假设硬公平性。
  • interval:近似流量统计刷新间隔,单位 ms,必须大于 0;为 0 时 panic。
§返回与副作用

返回未绑定 runtime 的独立 pool,持有 W 组 worker queue/stealer/waker 和一个 pool-wide 原子刷新时间。函数会分配内存,不创建线程、不执行用户 future、不进行 I/O。

§性能、幂等与安全

构建时间/空间 O(W);不是纯函数且每次创建不同资源,非幂等。运行时 owner 限制、panic 边界、线程/异步/内存安全和跨 runtime fallback 见类型文档。专项入口为 tests/stealable_task_pool_concurrency.rs

Trait Implementations§

Source§

impl<O: Default + 'static> AsyncTaskPool<O> for StealableTaskPool<O>

Source§

fn get_thread_id(&self) -> usize

返回当前线程既有的 packed runtime/worker id。

worker 线程返回 runtime id 与 worker index 的组合值,未绑定线程保持既有 usize::MAX。本方法不执行新增 owner 校验;O(1)、纯只读、幂等、无分配/锁/阻塞, 只读取当前线程 TLS,调用目的、可见语义及线程/内存安全边界均保持不变。

Source§

fn push_local(&self, task: Arc<AsyncTask<Self::Pool, O>>) -> Result<()>

将任务提交到所属 worker 的 internal queue,无法使用本地路径时进入 public injector。

task 所有权被成功转移并返回 Ok(())。外部线程或跨 runtime 调用保持 public fallback; task owner 与当前 runtime 相同但 pool 身份错误时,在读取 internal queue 前 panic。正常 local/public 路径分摊 O(1)、零新增分配/clone/锁/阻塞;非纯、非幂等,不 poll 用户 future。合法调用域内线程安全、内存安全、异步安全,不触达 V8/FFI。

Source§

fn push_priority( &self, priority: usize, task: Arc<AsyncTask<Self::Pool, O>>, ) -> Result<()>

按既有优先级规则提交任务,并保护最高优先级 owner-only single-item stack。

priority 的阈值、internal/public fallback 和返回值保持不变。最高优先级本地分支先 校验 exact pool,错误上下文在 stack/queue 访问前 panic;其它 runtime 继续 public fallback。分摊 O(1),无新增分配/clone/锁/阻塞;非纯、非幂等,不执行任务、回调、 I/O 或 V8/FFI,owner 契约内线程/内存/异步安全。

Source§

fn try_pop(&self) -> Option<Arc<AsyncTask<Self::Pool, O>>>

从当前 exact pool 的 owner worker slot 按既有权重和 steal 顺序尝试取一个任务。

返回任务 Arc 或所有候选队列均空时的 None。只能由所属 worker 调用;wrong-pool 或 外部线程在索引、stack、selector 及 local deque 访问前 panic。local 快路径 O(1),尝试 W-1 个 stealer 最坏 O(W);原实现的空队列 steal 路径可能构造 O(W) 临时 Vec,本轮不 改该行为。新增 owner guard 零分配/clone/锁/阻塞。消费和统计更新使本方法非纯、非幂等; 不 poll 用户 future,owner 契约内线程/内存/异步安全。

Source§

fn try_pop_all(&self) -> IntoIter<Arc<AsyncTask<Self::Pool, O>>>

反复执行 owner-checked try_pop,返回当前可取得任务的 owned iterator。

空池返回空 iterator;wrong-pool 调用在首次 pop 前 panic。除每次 pop 的既有复杂度外, 汇总 N 个任务需 O(N) 时间和 O(N) Vec 空间;会分配但不阻塞、不执行任务。该操作非纯、 非幂等,owner 契约内线程/内存/异步安全;批量收集语义和分配行为未被本轮改变。

Source§

type Pool = StealableTaskPool<O>

Source§

fn len(&self) -> usize

获取当前异步任务池内任务数量
Source§

fn push(&self, task: Arc<AsyncTask<Self::Pool, O>>) -> Result<()>

将异步任务加入异步任务池
Source§

fn push_keep(&self, task: Arc<AsyncTask<Self::Pool, O>>) -> Result<()>

异步任务被唤醒时,将异步任务继续加入异步任务池
Source§

fn get_thread_waker(&self) -> Option<&Arc<(AtomicBool, Mutex<()>, Condvar)>>

获取本地线程的唤醒器
Source§

impl<O: Default + 'static> AsyncTaskPoolExt<O> for StealableTaskPool<O>

Source§

fn clone_thread_waker(&self) -> Option<Arc<(AtomicBool, Mutex<()>, Condvar)>>

clone 当前 exact pool/worker 的休眠唤醒器。

合法 owner 返回 Some(Arc<...>);wrong-pool 或未绑定线程在 worker lookup 前 panic。 O(1),按既有语义执行一次 Arc::clone,owner guard 自身不 clone/分配/锁/阻塞,也不会 notify 或产生错误唤醒。查询本身幂等但增加引用计数;调用方负责释放返回 Arc。线程、 内存和异步安全,不执行用户代码或 V8/FFI。

Source§

fn set_waits( &mut self, waits: Arc<ArrayQueue<Arc<(AtomicBool, Mutex<()>, Condvar)>>>, )

设置待唤醒的工作者唤醒器队列
Source§

fn get_waits( &self, ) -> Option<&Arc<ArrayQueue<Arc<(AtomicBool, Mutex<()>, Condvar)>>>>

获取待唤醒的工作者唤醒器队列
Source§

fn worker_len(&self) -> usize

获取工作者的数量
Source§

fn idler_len(&self) -> usize

获取空闲的工作者的数量,这个数量大于0,表示可以新开线程来运行可分派的工作者
Source§

fn spawn_worker(&self) -> Option<usize>

分派一个空闲的工作者
Source§

fn buffer_len(&self) -> usize

获取缓冲区的任务数量,缓冲区任务是未分配给工作者的任务
Source§

fn set_thread_waker( &mut self, _thread_waker: Arc<(AtomicBool, Mutex<()>, Condvar)>, )

设置当前绑定本地线程的唤醒器
Source§

fn close_worker(&self)

关闭当前工作者
Source§

impl<O: Default + 'static> Default for StealableTaskPool<O>

Source§

fn default() -> Self

Returns the “default value” for a type. Read more
Source§

impl<O: Default + 'static> Send for StealableTaskPool<O>

Source§

impl<O: Default + 'static> Sync for StealableTaskPool<O>

Auto Trait Implementations§

Blanket Implementations§

Source§

impl<T> Any for T
where T: 'static + ?Sized,

Source§

fn type_id(&self) -> TypeId

Gets the TypeId of self. Read more
Source§

impl<T> Borrow<T> for T
where T: ?Sized,

Source§

fn borrow(&self) -> &T

Immutably borrows from an owned value. Read more
Source§

impl<T> BorrowMut<T> for T
where T: ?Sized,

Source§

fn borrow_mut(&mut self) -> &mut T

Mutably borrows from an owned value. Read more
Source§

impl<T> From<T> for T

Source§

fn from(t: T) -> T

Returns the argument unchanged.

Source§

impl<T, U> Into<U> for T
where U: From<T>,

Source§

fn into(self) -> U

Calls U::from(self).

That is, this conversion is whatever the implementation of From<T> for U chooses to do.

Source§

impl<T> Pointable for T

Source§

const ALIGN: usize

The alignment of pointer.
Source§

type Init = T

The type for initializers.
Source§

unsafe fn init(init: <T as Pointable>::Init) -> usize

Initializes a with the given initializer. Read more
Source§

unsafe fn deref<'a>(ptr: usize) -> &'a T

Dereferences the given pointer. Read more
Source§

unsafe fn deref_mut<'a>(ptr: usize) -> &'a mut T

Mutably dereferences the given pointer. Read more
Source§

unsafe fn drop(ptr: usize)

Drops the object pointed to by the given pointer. Read more
Source§

impl<T> ThreadSend for T
where T: Send,

Source§

impl<T> ThreadSync for T
where T: Sync + Send,

Source§

impl<T, U> TryFrom<U> for T
where U: Into<T>,

Source§

type Error = Infallible

The type returned in the event of a conversion error.
Source§

fn try_from(value: U) -> Result<T, <T as TryFrom<U>>::Error>

Performs the conversion.
Source§

impl<T, U> TryInto<U> for T
where U: TryFrom<T>,

Source§

type Error = <U as TryFrom<T>>::Error

The type returned in the event of a conversion error.
Source§

fn try_into(self) -> Result<U, <U as TryFrom<T>>::Error>

Performs the conversion.
Source§

impl<V, T> VZip<V> for T
where V: MultiLane<T>,

Source§

fn vzip(self) -> V