millipede_storage_memory/
queue.rs1use 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
24pub 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 #[must_use]
56 pub fn new(name: impl Into<String>) -> Self {
57 Self::with_policy(name, MemoryQueuePolicy::Fifo)
58 }
59
60 #[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}