pub struct ComputationalTaskPool<O: Default + 'static> { /* private fields */ }Expand description
计算型多线程任务池,适合 CPU 密集型应用,不支持运行时伸缩。
§使用方式
通过 ComputationalTaskPool::new 创建与 runtime worker slot 数匹配的 pool,再交给
MultiTaskRuntimeBuilder。外部线程可以使用 runtime 的 spawn/spawn_local;直接
try_pop、try_pop_all 或取得当前 worker waker 只允许在拥有该 pool 的 worker 上执行。
use pi_async_rt::rt::{AsyncRuntime, multi_thread::{ComputationalTaskPool, MultiTaskRuntimeBuilder}};
let pool = ComputationalTaskPool::new(2);
let runtime = MultiTaskRuntimeBuilder::new(pool)
.init_worker_size(2)
.set_worker_limit(2, 2)
.build();
runtime.spawn(async {}).unwrap();§业务与错误边界
push允许任意线程调用;push_local在非所属 runtime 线程上保持公共队列 fallback。- owner-only pop/waker 操作在未绑定线程或其它 pool 的 worker 上调用会 panic,并保证在访问
worker-local
crossbeam_deque::Worker前失败。 - pool 不定义任务结果、取消、timeout 或 runtime 关闭语义。
§性能与安全
- push/pop 快路径为 O(1);总长度为原子计数近似值,空间为 O(W + N)。
- owner 校验一次 TLS 读取和一次指针比较,无分配、clone、锁、自旋或阻塞。
- 类型不是纯值容器:push/pop 会修改队列和统计,非幂等;不会在内部 poll 用户 future。
- pool 可在线程间共享;worker-local stack 只由通过 pool identity 校验的 owner worker 访问。
- 不持有跨 await guard,不触达 V8/FFI。专项入口:
tests/stealable_task_pool_concurrency.rs的 owner/cross-runtime 用例。
Implementations§
Source§impl<O: Default + 'static> ComputationalTaskPool<O>
impl<O: Default + 'static> ComputationalTaskPool<O>
Sourcepub fn new(size: usize) -> Self
pub fn new(size: usize) -> Self
构建指定 worker slot 数量的计算型多线程任务池。
§参数与返回
size:请求的 worker slot 数;小于默认初始 worker 数时保持旧行为,提升到默认值。- 返回独立 pool;尚未绑定 runtime/waits,也不会启动线程或执行 future。
§性能与副作用
时间和空间复杂度均为 O(W),W 为实际 slot 数。该函数会分配 worker/queue/waker,但不 阻塞、不执行 I/O、不创建线程。它不是纯函数;每次调用创建不同 pool,因此非幂等。
§安全与边界
返回值可安全移动并由 builder 在线程间共享。调用方必须让 builder 的最大 worker 数 不超过 slot 数;builder 会再次收敛该边界。owner-only API 的线程/pool 限制见类型文档。
Trait Implementations§
Source§impl<O: Default + 'static> AsyncTaskPool<O> for ComputationalTaskPool<O>
impl<O: Default + 'static> AsyncTaskPool<O> for ComputationalTaskPool<O>
Source§fn get_thread_id(&self) -> usize
fn get_thread_id(&self) -> usize
返回当前线程既有的 packed runtime/worker id。
worker 线程返回 runtime id 与 worker index 的组合值,未绑定线程保持既有
usize::MAX。本方法不校验 pool owner,调用目的、返回语义和 panic 边界均未改变。
O(1)、纯只读、幂等、无分配/锁/阻塞/唤醒;只读取当前线程 TLS,线程和内存安全。
Source§fn push_local(&self, task: Arc<AsyncTask<Self::Pool, O>>) -> Result<()>
fn push_local(&self, task: Arc<AsyncTask<Self::Pool, O>>) -> Result<()>
将任务优先提交给所属 worker 的本地 queue,否则保持公共 queue fallback。
task 的 Arc 所有权被 queue 接收;本实现当前成功返回 Ok(())。当 task owner 与当前
runtime 不同(含外部线程)时不触碰 owner-only 状态;owner 相同但 pool 身份错误时在
queue 访问前 panic。O(1)、零额外分配/clone/锁/阻塞,不 poll 用户 future。该操作非纯、
非幂等;合法调用域内线程安全、内存安全和异步安全,不触达 V8/FFI。
Source§fn push_priority(
&self,
priority: usize,
task: Arc<AsyncTask<Self::Pool, O>>,
) -> Result<()>
fn push_priority( &self, priority: usize, task: Arc<AsyncTask<Self::Pool, O>>, ) -> Result<()>
按既有优先级语义提交任务,并保护最高优先级 owner-only stack。
priority 的阈值和 fallback 顺序不变;task 被成功转移后返回 Ok(())。只有最高优先级
且 task 属于当前 runtime 时执行 exact-pool owner 校验,错误 pool 在 stack 访问前 panic。
O(1)、无新增分配/clone/锁/阻塞;非纯、非幂等,不执行任务或用户回调,线程/异步/内存
安全边界与 ComputationalTaskPool::push_local 相同。
Source§fn try_pop(&self) -> Option<Arc<AsyncTask<Self::Pool, O>>>
fn try_pop(&self) -> Option<Arc<AsyncTask<Self::Pool, O>>>
从当前 exact pool 的 owner worker slot 尝试取一个任务。
返回任务的 Arc 或空队列时的 None。只能由该 pool 绑定的 worker 调用;外部线程或
wrong-pool worker 会在索引及 local stack 访问前 panic。快路径 O(1)、零分配、零 clone、
无锁/自旋/阻塞;消费队列使其非纯、非幂等。本方法不 poll/drop 任务,不触达 V8/FFI,
在 owner 契约内线程、内存与异步安全。
Source§fn try_pop_all(&self) -> IntoIter<Arc<AsyncTask<Self::Pool, O>>>
fn try_pop_all(&self) -> IntoIter<Arc<AsyncTask<Self::Pool, O>>>
反复执行 owner-checked try_pop,返回当前可取得任务的 owned iterator。
空池返回空 iterator;wrong-pool 调用在首次 pop 前 panic。时间 O(N)、临时空间 O(N),会
按既有实现分配 Vec,其预分配容量来自近似 len()。非纯、非幂等;不阻塞、不 poll
用户 future,owner 契约内线程/内存/异步安全。该既有批量分配行为不属于本轮修改。
type Pool = ComputationalTaskPool<O>
Source§fn push_keep(&self, task: Arc<AsyncTask<Self::Pool, O>>) -> Result<()>
fn push_keep(&self, task: Arc<AsyncTask<Self::Pool, O>>) -> Result<()>
Source§fn get_thread_waker(&self) -> Option<&Arc<(AtomicBool, Mutex<()>, Condvar)>>
fn get_thread_waker(&self) -> Option<&Arc<(AtomicBool, Mutex<()>, Condvar)>>
Source§impl<O: Default + 'static> AsyncTaskPoolExt<O> for ComputationalTaskPool<O>
impl<O: Default + 'static> AsyncTaskPoolExt<O> for ComputationalTaskPool<O>
Source§fn clone_thread_waker(&self) -> Option<Arc<(AtomicBool, Mutex<()>, Condvar)>>
fn clone_thread_waker(&self) -> Option<Arc<(AtomicBool, Mutex<()>, Condvar)>>
clone 当前 exact pool/worker 的休眠唤醒器。
合法 owner 返回 Some(Arc<...>);wrong-pool 或未绑定线程在索引前 panic。O(1),会按
既有语义执行一次 Arc::clone,但 owner 校验自身不 clone/分配/加锁/阻塞。该只读查询
不唤醒线程,本身幂等但增加引用计数;返回 Arc 的释放由调用方负责,线程和内存安全。