pub struct PriorityQueue { /* private fields */ }Expand description
A priority queue for PromptRequest items that orders requests by
Priority and, within the same tier, by arrival time (FIFO).
§Capacity
PriorityQueue::new creates a queue with capacity 10 000. Use
PriorityQueue::with_capacity for custom sizing. When the queue is full,
PriorityQueue::push returns Err(QueueError::QueueFull).
§Deadline enforcement
Use PriorityQueue::pop_with_deadline_check to automatically skip
requests whose deadline has already passed. Each skipped request is
counted in expired_total (accessible via QueueStats) and a
tracing::warn! is emitted.
§Thread safety
PriorityQueue is Clone + Send + Sync. All clones share the same
underlying heap behind a Mutex.
§Examples
use std::collections::HashMap;
use tokio_prompt_orchestrator::{PromptRequest, SessionId};
use tokio_prompt_orchestrator::enhanced::{PriorityQueue, Priority};
let queue = PriorityQueue::new();
let req = PromptRequest {
session: SessionId::new("s1"),
request_id: "r1".to_string(),
input: "hello".to_string(),
meta: HashMap::new(),
deadline: None,
};
queue.push(Priority::High, req).await.unwrap();
if let Some((priority, request)) = queue.pop().await {
println!("{priority:?}: {}", request.input);
}Implementations§
Source§impl PriorityQueue
impl PriorityQueue
Sourcepub fn new() -> Self
pub fn new() -> Self
Create a PriorityQueue with the default capacity of 10 000 entries.
See also PriorityQueue::with_capacity.
Sourcepub fn with_capacity(max_size: usize) -> Self
pub fn with_capacity(max_size: usize) -> Self
Create a PriorityQueue with a specific maximum capacity.
A capacity of 0 causes every push to return
Err(QueueError::QueueFull) immediately.
Sourcepub async fn push(
&self,
priority: Priority,
request: PromptRequest,
) -> Result<(), QueueError>
pub async fn push( &self, priority: Priority, request: PromptRequest, ) -> Result<(), QueueError>
Push a request onto the queue with the given priority.
§Arguments
priority— Relative importance; seePriority.request— ThePromptRequestto enqueue.
§Errors
Returns Err(QueueError::QueueFull) if the queue has reached its
capacity. The request is not enqueued in that case.
Sourcepub async fn pop(&self) -> Option<(Priority, PromptRequest)>
pub async fn pop(&self) -> Option<(Priority, PromptRequest)>
Remove and return the highest-priority request.
Among requests with equal priority the one inserted first is returned (FIFO within a tier).
§Returns
Some((priority, request)) if the queue is non-empty, otherwise None.
Sourcepub async fn pop_with_deadline_check(&self) -> Option<(Priority, PromptRequest)>
pub async fn pop_with_deadline_check(&self) -> Option<(Priority, PromptRequest)>
Remove and return the highest-priority request whose deadline has not yet passed.
Requests are popped in priority order (highest first, then FIFO).
Any request whose deadline is Some(t) and t < Instant::now() is
silently discarded: the expired_total counter is incremented, a
tracing::warn! is emitted, and the next candidate is checked.
Requests with deadline = None are never expired.
§Returns
Some((priority, request)) for the first non-expired request found,
or None if the queue is empty or all remaining requests have expired.
§Panics
This function does not panic.
Sourcepub async fn stats(&self) -> QueueStats
pub async fn stats(&self) -> QueueStats
Get queue statistics by priority.
Reads from lock-free atomic counters — O(1), no mutex held.
Trait Implementations§
Source§impl Clone for PriorityQueue
impl Clone for PriorityQueue
Source§fn clone(&self) -> PriorityQueue
fn clone(&self) -> PriorityQueue
1.0.0 (const: unstable) · Source§fn clone_from(&mut self, source: &Self)
fn clone_from(&mut self, source: &Self)
source. Read moreAuto Trait Implementations§
impl !RefUnwindSafe for PriorityQueue
impl !UnwindSafe for PriorityQueue
impl Freeze for PriorityQueue
impl Send for PriorityQueue
impl Sync for PriorityQueue
impl Unpin for PriorityQueue
impl UnsafeUnpin for PriorityQueue
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> CloneToUninit for Twhere
T: Clone,
impl<T> CloneToUninit for Twhere
T: Clone,
§impl<T> FutureExt for T
impl<T> FutureExt for T
§fn with_context(self, otel_cx: Context) -> WithContext<Self>
fn with_context(self, otel_cx: Context) -> WithContext<Self>
§fn with_current_context(self) -> WithContext<Self>
fn with_current_context(self) -> WithContext<Self>
§impl<T> Instrument for T
impl<T> Instrument for T
§fn instrument(self, span: Span) -> Instrumented<Self>
fn instrument(self, span: Span) -> Instrumented<Self>
§fn in_current_span(self) -> Instrumented<Self>
fn in_current_span(self) -> Instrumented<Self>
Source§impl<T> IntoRequest<T> for T
impl<T> IntoRequest<T> for T
Source§fn into_request(self) -> Request<T>
fn into_request(self) -> Request<T>
T in a tonic::Request