pub struct TaskScheduler { /* private fields */ }Expand description
Shared actor handle. One instance belongs to each Agent.
Implementations§
Source§impl TaskScheduler
impl TaskScheduler
Sourcepub fn new(config: TaskSchedulerConfig) -> Result<Self, TaskSchedulerError>
pub fn new(config: TaskSchedulerConfig) -> Result<Self, TaskSchedulerError>
Start a scheduler on the current Tokio runtime.
Sourcepub async fn acquire(
&self,
priority: TaskPriority,
label: impl Into<String>,
cancellation: &CancellationToken,
) -> Result<TaskLease, TaskSchedulerError>
pub async fn acquire( &self, priority: TaskPriority, label: impl Into<String>, cancellation: &CancellationToken, ) -> Result<TaskLease, TaskSchedulerError>
Wait until this task owns one global execution slot.
Sourcepub async fn acquire_with_identity(
&self,
priority: TaskPriority,
label: impl Into<String>,
identity: Option<ExecutionIdentityV1>,
cancellation: &CancellationToken,
) -> Result<TaskLease, TaskSchedulerError>
pub async fn acquire_with_identity( &self, priority: TaskPriority, label: impl Into<String>, identity: Option<ExecutionIdentityV1>, cancellation: &CancellationToken, ) -> Result<TaskLease, TaskSchedulerError>
Wait until this task owns one global execution slot and carry its semantic execution identity through the admission boundary.
The identity is optional for backwards compatibility with callers that only need capacity. When present it is validated before anything is queued, and the resulting lease retains it for tracing and downstream adapters.
Sourcepub async fn acquire_quota(
&self,
priority: TaskPriority,
label: impl Into<String>,
quota: &TaskSchedulerQuota,
cancellation: &CancellationToken,
) -> Result<TaskLease, TaskSchedulerError>
pub async fn acquire_quota( &self, priority: TaskPriority, label: impl Into<String>, quota: &TaskSchedulerQuota, cancellation: &CancellationToken, ) -> Result<TaskLease, TaskSchedulerError>
Wait for one or more quota dimensions without consuming another global execution slot.
This is used at leaf resource boundaries (for example, one model generation inside an already-admitted session run). The request still enters the same priority queue and actor as global admissions, so a provider limit cannot be bypassed with a local semaphore and a max-active=1 session does not deadlock while a nested model call waits.
Sourcepub async fn acquire_quotas(
&self,
priority: TaskPriority,
label: impl Into<String>,
quotas: &[TaskSchedulerQuota],
cancellation: &CancellationToken,
) -> Result<TaskLease, TaskSchedulerError>
pub async fn acquire_quotas( &self, priority: TaskPriority, label: impl Into<String>, quotas: &[TaskSchedulerQuota], cancellation: &CancellationToken, ) -> Result<TaskLease, TaskSchedulerError>
Multi-dimensional quota-only counterpart of Self::acquire_quota.
Sourcepub async fn acquire_with_quota(
&self,
priority: TaskPriority,
label: impl Into<String>,
quota: &TaskSchedulerQuota,
identity: Option<ExecutionIdentityV1>,
cancellation: &CancellationToken,
) -> Result<TaskLease, TaskSchedulerError>
pub async fn acquire_with_quota( &self, priority: TaskPriority, label: impl Into<String>, quota: &TaskSchedulerQuota, identity: Option<ExecutionIdentityV1>, cancellation: &CancellationToken, ) -> Result<TaskLease, TaskSchedulerError>
Wait until this task owns a global execution slot subject to a capacity quota. The quota reservation is made in the same scheduler actor as global admission, so a caller cannot bypass it by creating a fresh local semaphore or executor handle.
Sourcepub async fn acquire_with_quotas(
&self,
priority: TaskPriority,
label: impl Into<String>,
quotas: &[TaskSchedulerQuota],
identity: Option<ExecutionIdentityV1>,
cancellation: &CancellationToken,
) -> Result<TaskLease, TaskSchedulerError>
pub async fn acquire_with_quotas( &self, priority: TaskPriority, label: impl Into<String>, quotas: &[TaskSchedulerQuota], identity: Option<ExecutionIdentityV1>, cancellation: &CancellationToken, ) -> Result<TaskLease, TaskSchedulerError>
Wait until this task owns a global execution slot subject to multiple immutable quota dimensions. All dimensions are evaluated by the same scheduler actor, so a caller cannot bypass one limit by splitting the request across independent local gates.
Sourcepub async fn quota_snapshot(
&self,
quota: &TaskSchedulerQuota,
) -> Result<TaskSchedulerQuotaSnapshot, TaskSchedulerError>
pub async fn quota_snapshot( &self, quota: &TaskSchedulerQuota, ) -> Result<TaskSchedulerQuotaSnapshot, TaskSchedulerError>
Return the live occupancy projection for one owner quota.
Sourcepub async fn quota_health(
&self,
quota: &TaskSchedulerQuota,
) -> Result<TaskSchedulerQuotaHealthSnapshot, TaskSchedulerError>
pub async fn quota_health( &self, quota: &TaskSchedulerQuota, ) -> Result<TaskSchedulerQuotaHealthSnapshot, TaskSchedulerError>
Return bounded cumulative health for one quota identity.
The scheduler keeps a fixed number of recent idle quota epochs so a
host can inspect a completed provider generation without turning the
actor into an unbounded metrics store. A descriptor that has never been
admitted returns observed = false and zero counters.
Sourcepub async fn stats(&self) -> Result<TaskSchedulerStats, TaskSchedulerError>
pub async fn stats(&self) -> Result<TaskSchedulerStats, TaskSchedulerError>
Return a consistent actor-owned occupancy snapshot.
Sourcepub async fn health(
&self,
) -> Result<TaskSchedulerHealthSnapshot, TaskSchedulerError>
pub async fn health( &self, ) -> Result<TaskSchedulerHealthSnapshot, TaskSchedulerError>
Return occupancy plus bounded cumulative admission/fairness counters.
This is intentionally a separate method from Self::stats so the
long-lived counters can be added without changing the established
occupancy wire shape consumed by older SDKs.
Trait Implementations§
Auto Trait Implementations§
impl !Freeze for TaskScheduler
impl RefUnwindSafe for TaskScheduler
impl Send for TaskScheduler
impl Sync for TaskScheduler
impl Unpin for TaskScheduler
impl UnsafeUnpin for TaskScheduler
impl UnwindSafe for TaskScheduler
Blanket Implementations§
Source§impl<T> BorrowMut<T> for Twhere
T: ?Sized,
impl<T> BorrowMut<T> for Twhere
T: ?Sized,
Source§fn borrow_mut(&mut self) -> &mut T
fn borrow_mut(&mut self) -> &mut T
impl<ST, DT> CastableFrom<ST, Initialized, Initialized> for DT
impl<ST, DT> CastableFrom<ST, Uninit, Uninit> for DT
Source§impl<T> FutureExt for T
impl<T> FutureExt for T
Source§fn with_context(self, otel_cx: Context) -> WithContext<Self> ⓘ
fn with_context(self, otel_cx: Context) -> WithContext<Self> ⓘ
Source§fn with_current_context(self) -> WithContext<Self> ⓘ
fn with_current_context(self) -> WithContext<Self> ⓘ
Source§impl<T> Instrument for T
impl<T> Instrument for T
Source§fn instrument(self, span: Span) -> Instrumented<Self> ⓘ
fn instrument(self, span: Span) -> Instrumented<Self> ⓘ
Source§fn in_current_span(self) -> Instrumented<Self> ⓘ
fn in_current_span(self) -> Instrumented<Self> ⓘ
Source§impl<T> IntoEither for T
impl<T> IntoEither for T
Source§fn into_either(self, into_left: bool) -> Either<Self, Self> ⓘ
fn into_either(self, into_left: bool) -> Either<Self, Self> ⓘ
self into a Left variant of Either<Self, Self>
if into_left is true.
Converts self into a Right variant of Either<Self, Self>
otherwise. Read moreSource§fn into_either_with<F>(self, into_left: F) -> Either<Self, Self> ⓘ
fn into_either_with<F>(self, into_left: F) -> Either<Self, Self> ⓘ
self into a Left variant of Either<Self, Self>
if into_left(&self) returns true.
Converts self into a Right variant of Either<Self, Self>
otherwise. Read more