pub struct RealtimeSessionScheduler<M, S, R, C, O>where
M: SemanticStateTransaction,
M::Error: 'static,
S: Clone,
R: Clone,
C: Completion,
O: TransitionOutput,{ /* private fields */ }Expand description
Singular fair scheduler for one exact selected realtime model.
Implementations§
Source§impl<M, S, R, C, O> RealtimeSessionScheduler<M, S, R, C, O>where
M: SemanticStateTransaction,
M::Error: 'static,
S: Clone,
R: Clone,
C: Completion,
O: TransitionOutput,
impl<M, S, R, C, O> RealtimeSessionScheduler<M, S, R, C, O>where
M: SemanticStateTransaction,
M::Error: 'static,
S: Clone,
R: Clone,
C: Completion,
O: TransitionOutput,
Sourcepub const fn model_identity(&self) -> &RealtimeModelSessionIdentity
pub const fn model_identity(&self) -> &RealtimeModelSessionIdentity
Exact selected model identity shared by every admitted request.
Sourcepub fn new(
model: RealtimeModelSessionIdentity,
limits: SchedulerLimits,
) -> Result<Self, SchedulerError>
pub fn new( model: RealtimeModelSessionIdentity, limits: SchedulerLimits, ) -> Result<Self, SchedulerError>
Creates one scheduler for both single-request and concurrent production use.
Sourcepub fn register(
&mut self,
request: RequestId,
generation: RealtimeGenerationState<M, S, R, C>,
) -> Result<RealtimeSessionIncarnation, RealtimeSessionError>
pub fn register( &mut self, request: RequestId, generation: RealtimeGenerationState<M, S, R, C>, ) -> Result<RealtimeSessionIncarnation, RealtimeSessionError>
Registers a new canonical session with a fresh monotonic incarnation.
Sourcepub fn resume(
&mut self,
request: RequestId,
released: ReleasedRealtimeSession<M, S, R, C>,
) -> Result<(), RealtimeSessionResumeError<M, S, R, C>>
pub fn resume( &mut self, request: RequestId, released: ReleasedRealtimeSession<M, S, R, C>, ) -> Result<(), RealtimeSessionResumeError<M, S, R, C>>
Resumes released state only under the exact selected model identity.
Sourcepub fn enqueue(
&mut self,
request: RequestId,
frame: RealtimeInputFrame,
) -> Result<WorkId, SchedulerError>
pub fn enqueue( &mut self, request: RequestId, frame: RealtimeInputFrame, ) -> Result<WorkId, SchedulerError>
Enqueues one portable frame on the singular fair path.
Sourcepub fn enqueue_with_deadline(
&mut self,
request: RequestId,
frame: RealtimeInputFrame,
deadline: Option<Instant>,
) -> Result<WorkId, SchedulerError>
pub fn enqueue_with_deadline( &mut self, request: RequestId, frame: RealtimeInputFrame, deadline: Option<Instant>, ) -> Result<WorkId, SchedulerError>
Enqueues one portable frame with an absolute deadline.
Sourcepub fn enqueue_batch(
&mut self,
request: RequestId,
frames: Vec<RealtimeInputFrame>,
) -> Result<Vec<WorkId>, SchedulerError>
pub fn enqueue_batch( &mut self, request: RequestId, frames: Vec<RealtimeInputFrame>, ) -> Result<Vec<WorkId>, SchedulerError>
Atomically enqueues ordered frames on one request.
Sourcepub fn run_local_turn<E>(
&mut self,
now: Instant,
execute: impl FnMut(WorkId, &RealtimeInputFrame, &mut RealtimeSessionBranch<M::Branch, S, R, C>) -> Result<O, E>,
) -> Result<SchedulerProgress<RealtimeInputFrame, O>, SchedulerError>where
E: Error,
pub fn run_local_turn<E>(
&mut self,
now: Instant,
execute: impl FnMut(WorkId, &RealtimeInputFrame, &mut RealtimeSessionBranch<M::Branch, S, R, C>) -> Result<O, E>,
) -> Result<SchedulerProgress<RealtimeInputFrame, O>, SchedulerError>where
E: Error,
Runs one fair local turn using an injected family/backend-independent submission closure.
Sourcepub fn run_local_bounded<E>(
&mut self,
now: Instant,
maximum_frames: usize,
execute: impl FnMut(WorkId, &RealtimeInputFrame, &mut RealtimeSessionBranch<M::Branch, S, R, C>) -> Result<O, E>,
) -> Result<SchedulerProgress<RealtimeInputFrame, O>, SchedulerError>where
E: Error,
pub fn run_local_bounded<E>(
&mut self,
now: Instant,
maximum_frames: usize,
execute: impl FnMut(WorkId, &RealtimeInputFrame, &mut RealtimeSessionBranch<M::Branch, S, R, C>) -> Result<O, E>,
) -> Result<SchedulerProgress<RealtimeInputFrame, O>, SchedulerError>where
E: Error,
Runs one local turn while admitting at most maximum_frames new transitions.
Sourcepub fn run_distributed_turn<T, E>(
&mut self,
protocol: u64,
transport: &T,
now: Instant,
execute: impl FnMut(WorkId, &RealtimeInputFrame, &mut RealtimeSessionBranch<M::Branch, S, R, C>) -> Result<O, E>,
) -> Result<SchedulerProgress<RealtimeInputFrame, O>, SchedulerError>where
T: BoundedConsensusTransport,
<T::Completion as Completion>::Error: Display,
E: Error,
O: DistributedTransitionOutput,
pub fn run_distributed_turn<T, E>(
&mut self,
protocol: u64,
transport: &T,
now: Instant,
execute: impl FnMut(WorkId, &RealtimeInputFrame, &mut RealtimeSessionBranch<M::Branch, S, R, C>) -> Result<O, E>,
) -> Result<SchedulerProgress<RealtimeInputFrame, O>, SchedulerError>where
T: BoundedConsensusTransport,
<T::Completion as Completion>::Error: Display,
E: Error,
O: DistributedTransitionOutput,
Runs one fair turn with mandatory topology-wide schedule and completion consensus.
Sourcepub fn replace_sampling<E>(
&mut self,
request: RequestId,
sampling: RealtimeSampling,
realize: impl FnOnce(RealtimeSampling) -> Result<(Vec<S>, Option<R>), E>,
) -> Result<(), RealtimeSamplingUpdateError<E, M::Error, C::Error>>
pub fn replace_sampling<E>( &mut self, request: RequestId, sampling: RealtimeSampling, realize: impl FnOnce(RealtimeSampling) -> Result<(Vec<S>, Option<R>), E>, ) -> Result<(), RealtimeSamplingUpdateError<E, M::Error, C::Error>>
Atomically replaces sampling only when the request owns no queued or branched work.
Sourcepub fn cancel(&mut self, request: RequestId) -> Result<(), SchedulerError>
pub fn cancel(&mut self, request: RequestId) -> Result<(), SchedulerError>
Cancels queued/prepared/submitted work using core scheduler semantics.
Sourcepub fn finish(&mut self, request: RequestId) -> Result<(), SchedulerError>
pub fn finish(&mut self, request: RequestId) -> Result<(), SchedulerError>
Marks a request finished using core scheduler semantics.
Sourcepub fn release(
&mut self,
request: RequestId,
) -> Result<ReleasedRealtimeSession<M, S, R, C>, SchedulerError>
pub fn release( &mut self, request: RequestId, ) -> Result<ReleasedRealtimeSession<M, S, R, C>, SchedulerError>
Releases an idle canonical session for exact later resumption.
Sourcepub fn request_status(&self, request: RequestId) -> Option<RequestStatus>
pub fn request_status(&self, request: RequestId) -> Option<RequestStatus>
Active or terminal request status.
Sourcepub fn queued_for_request(&self, request: RequestId) -> usize
pub fn queued_for_request(&self, request: RequestId) -> usize
Number of portable frames still queued for one active request.
Sourcepub fn forget_terminal(
&mut self,
request: RequestId,
) -> Result<RequestStatus, SchedulerError>
pub fn forget_terminal( &mut self, request: RequestId, ) -> Result<RequestStatus, SchedulerError>
Removes a terminal identity so the caller may explicitly reuse it.
Sourcepub fn request_state(
&self,
request: RequestId,
) -> Option<&RealtimeSessionState<M, S, R, C>>
pub fn request_state( &self, request: RequestId, ) -> Option<&RealtimeSessionState<M, S, R, C>>
Immutable canonical session state, when active.
Sourcepub fn report(&self) -> SchedulerReport
pub fn report(&self) -> SchedulerReport
Current scheduler telemetry.
Sourcepub fn capabilities(&self) -> SchedulerCapabilities
pub fn capabilities(&self) -> SchedulerCapabilities
Configured and observed scheduler capabilities.