Skip to main content

TaskScheduler

Struct TaskScheduler 

Source
pub struct TaskScheduler { /* private fields */ }
Expand description

Shared actor handle. One instance belongs to each Agent.

Implementations§

Source§

impl TaskScheduler

Source

pub fn new(config: TaskSchedulerConfig) -> Result<Self, TaskSchedulerError>

Start a scheduler on the current Tokio runtime.

Source

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.

Source

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.

Source

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.

Source

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.

Source

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.

Source

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.

Source

pub async fn quota_snapshot( &self, quota: &TaskSchedulerQuota, ) -> Result<TaskSchedulerQuotaSnapshot, TaskSchedulerError>

Return the live occupancy projection for one owner quota.

Source

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.

Source

pub async fn stats(&self) -> Result<TaskSchedulerStats, TaskSchedulerError>

Return a consistent actor-owned occupancy snapshot.

Source

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.

Source

pub async fn shutdown(&self)

Reject pending work and wait for already-admitted leases to finish.

Trait Implementations§

Source§

impl Debug for TaskScheduler

Source§

fn fmt(&self, f: &mut Formatter<'_>) -> Result

Formats the value using the given formatter. Read more

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<ST, DT> CastableFrom<ST, Initialized, Initialized> for DT
where ST: ?Sized, DT: ?Sized,

Source§

impl<ST, DT> CastableFrom<ST, Uninit, Uninit> for DT
where ST: ?Sized, DT: ?Sized,

Source§

impl<T> From<T> for T

Source§

fn from(t: T) -> T

Returns the argument unchanged.

Source§

impl<T> FutureExt for T

Source§

fn with_context(self, otel_cx: Context) -> WithContext<Self>

Attaches the provided Context to this type, returning a WithContext wrapper. Read more
Source§

fn with_current_context(self) -> WithContext<Self>

Attaches the current Context to this type, returning a WithContext wrapper. Read more
Source§

impl<T> Instrument for T

Source§

fn instrument(self, span: Span) -> Instrumented<Self>

Instruments this type with the provided Span, returning an Instrumented wrapper. Read more
Source§

fn in_current_span(self) -> Instrumented<Self>

Instruments this type with the current Span, returning an Instrumented wrapper. Read more
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> IntoEither for T

Source§

fn into_either(self, into_left: bool) -> Either<Self, Self>

Converts 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 more
Source§

fn into_either_with<F>(self, into_left: F) -> Either<Self, Self>
where F: FnOnce(&Self) -> bool,

Converts 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
Source§

impl<T> ParallelSend for T

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> PolicyExt for T
where T: ?Sized,

Source§

fn and<P, B, E>(self, other: P) -> And<T, P>
where T: Sized + Policy<B, E>, P: Policy<B, E>,

Create a new Policy that returns Action::Follow only if self and other return Action::Follow. Read more
Source§

fn or<P, B, E>(self, other: P) -> Or<T, P>
where T: Sized + Policy<B, E>, P: Policy<B, E>,

Create a new Policy that returns Action::Follow if either self or other returns Action::Follow. Read more
Source§

impl<T> Read<Exclusive, BecauseExclusive> for T
where T: ?Sized,

Source§

impl<T> Same for T

Source§

type Output = T

Should always be Self
Source§

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

Source§

type Error = !

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

fn try_from(value: U) -> Result<T, !>

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

Source§

impl<T> WithSubscriber for T

Source§

fn with_subscriber<S>(self, subscriber: S) -> WithDispatch<Self>
where S: Into<Dispatch>,

Attaches the provided Subscriber to this type, returning a WithDispatch wrapper. Read more
Source§

fn with_current_subscriber(self) -> WithDispatch<Self>

Attaches the current default Subscriber to this type, returning a WithDispatch wrapper. Read more