millipede_core/storage/
queue.rs1use super::{StorageError, StorageResult};
4use crate::request::{Request, RequestId};
5use std::{fmt, time::Duration};
6
7#[derive(Debug, Clone, PartialEq, Eq, Hash)]
9pub struct LeaseId(u64);
10
11impl LeaseId {
12 pub fn new(raw: u64) -> Self {
14 Self(raw)
15 }
16
17 pub fn as_u64(&self) -> u64 {
19 self.0
20 }
21}
22
23impl fmt::Display for LeaseId {
24 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
25 self.0.fmt(formatter)
26 }
27}
28
29#[derive(Debug)]
35pub struct Lease {
36 pub request: Request,
38 pub lease_id: LeaseId,
40 pub expires_at: std::time::Instant,
42}
43
44#[derive(Debug, Clone, Default)]
46#[non_exhaustive]
47#[must_use = "add options do nothing unless passed to RequestQueue::add"]
48pub struct AddOptions {
49 pub forefront: bool,
51}
52
53#[derive(Debug, Clone)]
55#[non_exhaustive]
56#[must_use = "reclaim options do nothing unless passed to RequestQueue::reclaim"]
57pub struct ReclaimOptions {
58 pub forefront: bool,
60 pub increment_retry: bool,
62}
63
64impl Default for ReclaimOptions {
65 fn default() -> Self {
66 Self {
67 forefront: false,
68 increment_retry: true,
69 }
70 }
71}
72
73#[derive(Debug, Clone)]
75#[must_use = "queue insertion results report deduplication state"]
76pub struct ProcessedRequest {
77 pub request_id: RequestId,
79 pub unique_key: String,
81 pub was_already_present: bool,
83 pub was_already_handled: bool,
85}
86
87pub type QueueOpInfo = ProcessedRequest;
89
90#[derive(Debug, Clone)]
96#[non_exhaustive]
97pub enum RequestSource {
98 Request(Request),
100}
101
102impl From<Request> for RequestSource {
103 fn from(request: Request) -> Self {
104 Self::Request(request)
105 }
106}
107
108#[derive(Debug, Clone)]
110#[non_exhaustive]
111#[must_use = "batched insertion results report processed requests"]
112pub struct AddRequestsBatchedResult {
113 pub processed: Vec<ProcessedRequest>,
115}
116
117#[must_use = "batch handles must be awaited to observe completion"]
119pub struct BatchAddHandle {
120 pub added: Vec<ProcessedRequest>,
122 completion: Completion,
123}
124
125enum Completion {
126 Ready(AddRequestsBatchedResult),
127 Task(tokio::task::JoinHandle<StorageResult<AddRequestsBatchedResult>>),
128}
129
130impl fmt::Debug for BatchAddHandle {
131 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
132 formatter
133 .debug_struct("BatchAddHandle")
134 .field("added", &self.added)
135 .field("completion", &"<completion>")
136 .finish()
137 }
138}
139
140impl BatchAddHandle {
141 pub fn ready(added: Vec<ProcessedRequest>) -> Self {
143 Self {
144 completion: Completion::Ready(AddRequestsBatchedResult {
145 processed: added.clone(),
146 }),
147 added,
148 }
149 }
150
151 pub fn deferred(
153 added: Vec<ProcessedRequest>,
154 task: tokio::task::JoinHandle<StorageResult<AddRequestsBatchedResult>>,
155 ) -> Self {
156 Self {
157 added,
158 completion: Completion::Task(task),
159 }
160 }
161
162 pub(crate) fn notify_on_completion<F>(self, notify: F) -> Self
164 where
165 F: FnOnce() + Send + 'static,
166 {
167 let Self { added, completion } = self;
168 match completion {
169 Completion::Ready(result) => {
170 notify();
171 Self {
172 added,
173 completion: Completion::Ready(result),
174 }
175 }
176 Completion::Task(task) => Self {
177 added,
178 completion: Completion::Task(tokio::spawn(async move {
179 let result = task.await.map_err(|error| {
180 StorageError::Backend(anyhow::anyhow!("batch add task failed: {error}"))
181 });
182 notify();
183 result?
184 })),
185 },
186 }
187 }
188
189 pub async fn wait(self) -> StorageResult<AddRequestsBatchedResult> {
191 match self.completion {
192 Completion::Ready(result) => Ok(result),
193 Completion::Task(task) => task.await.map_err(|error| {
194 StorageError::Backend(anyhow::anyhow!("batch add task failed: {error}"))
195 })?,
196 }
197 }
198}
199
200#[async_trait::async_trait]
205pub trait RequestQueue: Send + Sync {
206 async fn add(&self, req: Request, opts: AddOptions) -> StorageResult<QueueOpInfo>;
208 async fn add_batch(
210 &self,
211 reqs: Vec<RequestSource>,
212 opts: AddOptions,
213 ) -> StorageResult<BatchAddHandle>;
214 async fn fetch_next(&self) -> StorageResult<Option<Lease>>;
216 async fn mark_handled(&self, lease: Lease) -> StorageResult<()>;
218 async fn reclaim(&self, lease: Lease, opts: ReclaimOptions) -> StorageResult<()>;
223 async fn renew(&self, lease_id: &LeaseId, extend_by: Duration) -> StorageResult<()>;
225 async fn abandon(&self, lease: Lease) -> StorageResult<()>;
230 async fn is_empty(&self) -> StorageResult<bool>;
232 async fn is_finished(&self) -> StorageResult<bool>;
234 async fn handled_count(&self) -> StorageResult<u64>;
236 async fn pending_count(&self) -> StorageResult<u64>;
238}