Skip to main content

ComputationalTaskPool

Struct ComputationalTaskPool 

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

计算型多线程任务池,适合 CPU 密集型应用,不支持运行时伸缩。

§使用方式

通过 ComputationalTaskPool::new 创建与 runtime worker slot 数匹配的 pool,再交给 MultiTaskRuntimeBuilder。外部线程可以使用 runtime 的 spawn/spawn_local;直接 try_poptry_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>

Source

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>

Source§

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<()>

将任务优先提交给所属 worker 的本地 queue,否则保持公共 queue fallback。

taskArc 所有权被 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<()>

按既有优先级语义提交任务,并保护最高优先级 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>>>

从当前 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>>>

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

空池返回空 iterator;wrong-pool 调用在首次 pop 前 panic。时间 O(N)、临时空间 O(N),会 按既有实现分配 Vec,其预分配容量来自近似 len()。非纯、非幂等;不阻塞、不 poll 用户 future,owner 契约内线程/内存/异步安全。该既有批量分配行为不属于本轮修改。

Source§

type Pool = ComputationalTaskPool<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 ComputationalTaskPool<O>

Source§

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 的释放由调用方负责,线程和内存安全。

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 ComputationalTaskPool<O>

Source§

fn default() -> Self

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

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

Source§

impl<O: Default + 'static> Sync for ComputationalTaskPool<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