Skip to main content

millipede_storage_memory/
queue.rs

1//! In-process request queue implementation.
2//!
3//! Lease expiry is a documented no-op in this single-process backend: an expired lease is never
4//! reassigned. Marking handled, reclaiming, renewing, and abandoning still enforce the full lease
5//! contract so distributed backends can implement real expiry without an API break.
6
7use crate::policy::{Frontier, MemoryQueuePolicy};
8use millipede_core::{
9    request::{Request, RequestId},
10    storage::{
11        AddOptions, BatchAddHandle, Lease, LeaseId, ProcessedRequest, QueueOpInfo, ReclaimOptions,
12        RequestQueue, RequestSource, StorageError, StorageResult,
13    },
14};
15use std::{
16    collections::{HashMap, HashSet},
17    fmt,
18    sync::Mutex,
19    time::{Duration, Instant},
20};
21
22const LEASE_TTL: Duration = Duration::from_secs(180);
23
24/// A thread-safe, lease-based request queue stored entirely in memory.
25pub struct MemoryRequestQueue {
26    name: String,
27    state: Mutex<QueueState>,
28}
29
30struct QueueState {
31    dedup: HashMap<String, RequestId>,
32    handled: HashSet<String>,
33    pending: Frontier,
34    in_flight: HashMap<u64, LeasedRequest>,
35    handled_count: u64,
36    next_lease_id: u64,
37}
38
39struct LeasedRequest {
40    request: Request,
41    expires_at: Instant,
42}
43
44impl fmt::Debug for MemoryRequestQueue {
45    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
46        formatter
47            .debug_struct("MemoryRequestQueue")
48            .field("name", &self.name)
49            .finish_non_exhaustive()
50    }
51}
52
53impl MemoryRequestQueue {
54    /// Creates an empty FIFO queue with the given name.
55    #[must_use]
56    pub fn new(name: impl Into<String>) -> Self {
57        Self::with_policy(name, MemoryQueuePolicy::Fifo)
58    }
59
60    /// Creates an empty queue using the selected frontier policy.
61    #[must_use]
62    pub fn with_policy(name: impl Into<String>, policy: MemoryQueuePolicy) -> Self {
63        Self {
64            name: name.into(),
65            state: Mutex::new(QueueState {
66                dedup: HashMap::new(),
67                handled: HashSet::new(),
68                pending: Frontier::new(policy),
69                in_flight: HashMap::new(),
70                handled_count: 0,
71                next_lease_id: 0,
72            }),
73        }
74    }
75
76    fn add_locked(state: &mut QueueState, request: Request, opts: &AddOptions) -> QueueOpInfo {
77        let unique_key = request.unique_key.clone();
78        if let Some(request_id) = state.dedup.get(&unique_key) {
79            return ProcessedRequest {
80                request_id: request_id.clone(),
81                unique_key: unique_key.clone(),
82                was_already_present: true,
83                was_already_handled: state.handled.contains(&unique_key),
84            };
85        }
86
87        let request_id = request.id.clone();
88        state.dedup.insert(unique_key.clone(), request_id.clone());
89        if opts.forefront {
90            state.pending.push_front(request);
91        } else {
92            state.pending.push_back(request);
93        }
94        ProcessedRequest {
95            request_id,
96            unique_key,
97            was_already_present: false,
98            was_already_handled: false,
99        }
100    }
101
102    fn remove_lease(state: &mut QueueState, lease_id: &LeaseId) -> StorageResult<LeasedRequest> {
103        state
104            .in_flight
105            .remove(&lease_id.as_u64())
106            .ok_or_else(|| StorageError::LeaseNotFound {
107                lease_id: lease_id.clone(),
108            })
109    }
110}
111
112#[async_trait::async_trait]
113impl RequestQueue for MemoryRequestQueue {
114    async fn add(&self, request: Request, opts: AddOptions) -> StorageResult<QueueOpInfo> {
115        let mut state = self.state.lock().expect("request queue mutex poisoned");
116        Ok(Self::add_locked(&mut state, request, &opts))
117    }
118
119    async fn add_batch(
120        &self,
121        requests: Vec<RequestSource>,
122        opts: AddOptions,
123    ) -> StorageResult<BatchAddHandle> {
124        let mut state = self.state.lock().expect("request queue mutex poisoned");
125        let mut infos = Vec::with_capacity(requests.len());
126        for source in requests {
127            let request = match source {
128                RequestSource::Request(request) => request,
129                _ => return Err(StorageError::Unsupported("memory request source")),
130            };
131            infos.push(Self::add_locked(&mut state, request, &opts));
132        }
133        Ok(BatchAddHandle::ready(infos))
134    }
135
136    async fn fetch_next(&self) -> StorageResult<Option<Lease>> {
137        let mut state = self.state.lock().expect("request queue mutex poisoned");
138        let Some(request) = state.pending.pop_front() else {
139            return Ok(None);
140        };
141        let raw_lease_id = state.next_lease_id;
142        state.next_lease_id += 1;
143        let lease_id = LeaseId::new(raw_lease_id);
144        let expires_at = Instant::now() + LEASE_TTL;
145        state.in_flight.insert(
146            raw_lease_id,
147            LeasedRequest {
148                request: request.clone(),
149                expires_at,
150            },
151        );
152        Ok(Some(Lease {
153            request,
154            lease_id,
155            expires_at,
156        }))
157    }
158
159    async fn mark_handled(&self, lease: Lease) -> StorageResult<()> {
160        let mut state = self.state.lock().expect("request queue mutex poisoned");
161        let leased = Self::remove_lease(&mut state, &lease.lease_id)?;
162        state.handled_count += 1;
163        state.handled.insert(leased.request.unique_key);
164        Ok(())
165    }
166
167    async fn reclaim(&self, lease: Lease, opts: ReclaimOptions) -> StorageResult<()> {
168        let mut state = self.state.lock().expect("request queue mutex poisoned");
169        Self::remove_lease(&mut state, &lease.lease_id)?;
170        let mut request = lease.request;
171        if opts.increment_retry {
172            request.retry_count += 1;
173        }
174        if opts.forefront {
175            state.pending.push_front(request);
176        } else {
177            state.pending.push_back(request);
178        }
179        Ok(())
180    }
181
182    async fn renew(&self, lease_id: &LeaseId, extend_by: Duration) -> StorageResult<()> {
183        let mut state = self.state.lock().expect("request queue mutex poisoned");
184        let leased = state.in_flight.get_mut(&lease_id.as_u64()).ok_or_else(|| {
185            StorageError::LeaseNotFound {
186                lease_id: lease_id.clone(),
187            }
188        })?;
189        leased.expires_at += extend_by;
190        Ok(())
191    }
192
193    async fn abandon(&self, lease: Lease) -> StorageResult<()> {
194        let mut state = self.state.lock().expect("request queue mutex poisoned");
195        Self::remove_lease(&mut state, &lease.lease_id)?;
196        state.pending.push_front(lease.request);
197        Ok(())
198    }
199
200    async fn is_empty(&self) -> StorageResult<bool> {
201        let state = self.state.lock().expect("request queue mutex poisoned");
202        Ok(state.pending.is_empty())
203    }
204
205    async fn is_finished(&self) -> StorageResult<bool> {
206        let state = self.state.lock().expect("request queue mutex poisoned");
207        Ok(state.pending.is_empty() && state.in_flight.is_empty())
208    }
209
210    async fn handled_count(&self) -> StorageResult<u64> {
211        let state = self.state.lock().expect("request queue mutex poisoned");
212        Ok(state.handled_count)
213    }
214
215    async fn pending_count(&self) -> StorageResult<u64> {
216        let state = self.state.lock().expect("request queue mutex poisoned");
217        Ok(state.pending.len() as u64)
218    }
219}