Skip to main content

WorkQueue

Struct WorkQueue 

Source
pub struct WorkQueue<A: Address, T, WorkerRoute: DeliveryRoute<Protocol: Protocol<Addr = A, Msg = T>>, ReplyRoute: DeliveryRoute<Protocol: Protocol<Addr = A, Msg = WorkQueueOutcome<T>>>> { /* private fields */ }
Expand description

Bounded FIFO admission and worker-availability behavior.

Availability is a one-dispatch capability. It immediately consumes the oldest waiting value or joins a unique FIFO of available workers. Submission consumes the oldest worker or enters the bounded FIFO; at capacity it returns the owned value. Duplicate availability and withdrawal are idempotent. Initialization is empty, no actors are created, and the host never terminates by policy. FIFO selection and bounded admission are Bombay policy. Worker execution, mailbox admission, and physical backpressure are runtime responsibilities. Naming its protocol requires only typed routes; running the queue requires cloneable, comparable worker routes for state inspection and duplicate availability checks. No transition has a semantic panic condition.

Implementations§

Source§

impl<A, T, WorkerRoute, ReplyRoute> WorkQueue<A, T, WorkerRoute, ReplyRoute>
where A: Address, WorkerRoute: DeliveryRoute<Protocol: Protocol<Addr = A, Msg = T>> + Clone + PartialEq, ReplyRoute: DeliveryRoute<Protocol: Protocol<Addr = A, Msg = WorkQueueOutcome<T>>> + Clone,

Source

pub fn new(capacity: usize) -> Self

Construct an empty queue with explicit waiting capacity.

Source

pub fn state(&self) -> WorkQueueState<WorkerRoute>

Return complete observable queue state.

Trait Implementations§

Source§

impl<A, T, WorkerRoute, ReplyRoute> Behavior for WorkQueue<A, T, WorkerRoute, ReplyRoute>
where A: Address, WorkerRoute: DeliveryRoute<Protocol: Protocol<Addr = A, Msg = T>> + Clone + PartialEq, ReplyRoute: DeliveryRoute<Protocol: Protocol<Addr = A, Msg = WorkQueueOutcome<T>>> + Clone, WorkerRoute::Sends: SendsFor<User<A, WorkQueueMessage<T, WorkerRoute, ReplyRoute>>>, ReplyRoute::Sends: SendsFor<User<A, WorkQueueMessage<T, WorkerRoute, ReplyRoute>>>,

Source§

type Protocol = WorkQueue<A, T, WorkerRoute, ReplyRoute>

Stable public communication identity owned by this actor template. Read more
Source§

type Event = User<A, <<WorkQueue<A, T, WorkerRoute, ReplyRoute> as Behavior>::Protocol as Protocol>::Msg>

Source§

type Sends = WorkQueueSends<<WorkerRoute as DeliveryRoute>::Sends, <ReplyRoute as DeliveryRoute>::Sends>

Source§

type Ph = Never

Source§

type Error = Never

Source§

type Birth = NoBirths

Source§

fn transition( &mut self, _: ActiveTurn, event: Self::Event, ) -> BehaviorActed<Self>

Fold exactly one event into explicit actions and the next behavior. Read more
Source§

fn init( &mut self, _turn: InitializationTurn, ) -> Result<Actions<<Self::Protocol as Protocol>::Addr, Self::Ph, Self::Sends, Self::Birth>, Self::Error>
where Self: Sized,

Produce initialization actions before the first event is accepted. Read more
Source§

fn layer<L>(self, layer: L) -> <L as BehaviorLayer<Self>>::Output
where Self: Sized, L: BehaviorLayer<Self>,

Apply one statically dispatched construction layer. Read more
Source§

impl<A, T, WorkerRoute, ReplyRoute> BehaviorBase for WorkQueue<A, T, WorkerRoute, ReplyRoute>
where A: Address, WorkerRoute: DeliveryRoute<Protocol: Protocol<Addr = A, Msg = T>>, ReplyRoute: DeliveryRoute<Protocol: Protocol<Addr = A, Msg = WorkQueueOutcome<T>>>,

Source§

type Base = WorkQueue<A, T, WorkerRoute, ReplyRoute>

Source§

fn base(&self) -> &Self

Source§

impl<A, T, WorkerRoute, ReplyRoute> Protocol for WorkQueue<A, T, WorkerRoute, ReplyRoute>
where A: Address, WorkerRoute: DeliveryRoute<Protocol: Protocol<Addr = A, Msg = T>>, ReplyRoute: DeliveryRoute<Protocol: Protocol<Addr = A, Msg = WorkQueueOutcome<T>>>,

Source§

type Addr = A

Source§

type Msg = WorkQueueMessage<T, WorkerRoute, ReplyRoute>

Auto Trait Implementations§

§

impl<A, T, WorkerRoute, ReplyRoute> Freeze for WorkQueue<A, T, WorkerRoute, ReplyRoute>

§

impl<A, T, WorkerRoute, ReplyRoute> RefUnwindSafe for WorkQueue<A, T, WorkerRoute, ReplyRoute>
where WorkerRoute: RefUnwindSafe, T: RefUnwindSafe, ReplyRoute: RefUnwindSafe,

§

impl<A, T, WorkerRoute, ReplyRoute> Send for WorkQueue<A, T, WorkerRoute, ReplyRoute>
where WorkerRoute: Send, T: Send, ReplyRoute: Send,

§

impl<A, T, WorkerRoute, ReplyRoute> Sync for WorkQueue<A, T, WorkerRoute, ReplyRoute>
where WorkerRoute: Sync, T: Sync, ReplyRoute: Sync,

§

impl<A, T, WorkerRoute, ReplyRoute> Unpin for WorkQueue<A, T, WorkerRoute, ReplyRoute>
where WorkerRoute: Unpin, T: Unpin, ReplyRoute: Unpin,

§

impl<A, T, WorkerRoute, ReplyRoute> UnsafeUnpin for WorkQueue<A, T, WorkerRoute, ReplyRoute>

§

impl<A, T, WorkerRoute, ReplyRoute> UnwindSafe for WorkQueue<A, T, WorkerRoute, ReplyRoute>
where WorkerRoute: UnwindSafe, T: UnwindSafe, ReplyRoute: UnwindSafe,

Blanket Implementations§

Source§

impl<B> Activate for B
where B: Behavior,

Source§

fn initialize(self) -> Result<Initialized<Self>, Self::Error>

Consume this definition, perform its one initialization fold, and return the active behavior together with the ordered initialization effects. Read more
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<B> BeginShutdownPhases for B
where B: Behavior, <<B as Behavior>::Birth as BirthMode>::Child: ChildOccurrenceProduct<AvailableChildren>,

Source§

impl<Node> BirthNodeAppend<Never> for Node
where Node: NonEmptyBirthNode,

Source§

type Output = Node

Closed child algebra containing the complete prefix followed by the complete appended tail.
Source§

fn append_prefix(self) -> <Node as BirthNodeAppend<Never>>::Output

Inject one child from the existing prefix without changing its structural occurrence.
Source§

fn append_tail(tail: Never) -> <Node as BirthNodeAppend<Never>>::Output

Inject one child from the appended tail after every prefix occurrence.
Source§

fn append_creations<A>( prefix: Creations<CreateChild<A, Self>>, tail: Creations<CreateChild<A, Tail>>, ) -> Creations<CreateChild<A, Self::Output>>
where A: Address,

Preserve and concatenate two ordered creation batches.
Source§

impl<Child, Tail> BirthNodeAppend<Tail> for Child
where Child: Behavior, Tail: NonEmptyBirthNode,

Source§

type Output = ChildChoice<Child, Tail>

Closed child algebra containing the complete prefix followed by the complete appended tail.
Source§

fn append_prefix(self) -> <Child as BirthNodeAppend<Tail>>::Output

Inject one child from the existing prefix without changing its structural occurrence.
Source§

fn append_tail(tail: Tail) -> <Child as BirthNodeAppend<Tail>>::Output

Inject one child from the appended tail after every prefix occurrence.
Source§

fn append_creations<A>( prefix: Creations<CreateChild<A, Self>>, tail: Creations<CreateChild<A, Tail>>, ) -> Creations<CreateChild<A, Self::Output>>
where A: Address,

Preserve and concatenate two ordered creation batches.
Source§

impl<B> BirthProtocols for B
where B: Behavior, <B as Behavior>::Birth: BirthModeProtocols,

Source§

type Protocols = BirthProtocol<<B as Behavior>::Protocol, <<B as Behavior>::Birth as BirthModeProtocols>::Protocols>

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<Node, Shape> ChildOccurrenceProduct<Shape> for Node
where Node: ChildOccurrenceProductAt<ChildHead, Shape>, Shape: ChildOccurrenceShape,

Source§

type Product = <Node as ChildOccurrenceProductAt<ChildHead, Shape>>::Product

Complete shape-owned representation of this closed birth node.
Source§

impl<A, Child, Host> DispatchBirth<A, Host> for Child
where A: Address, Child: CreateSelectedChild<A, ChildHead, Host>,

Source§

fn dispatch_birth( self, id: CreationId, route: <A as Address>::Nonce, kind: CreationKind, host: &mut Host, ) -> impl Future<Output = ItemSettlement<RoutedCreation<A, Child>, <Child as ChildCreationProduct<A, ChildHead>>::Result, CreationRejection, Never>> + Send

Select exactly one concrete child host while preserving creation data in every non-accepted settlement.
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, 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.