lumen_server/service/
queue.rs1use std::{
2 collections::{HashMap, VecDeque},
3 sync::{
4 Mutex,
5 atomic::{AtomicU64, Ordering},
6 },
7 time::{Duration, Instant},
8};
9
10use super::{
11 BoxFuture, RenderJob, RenderJobId, RenderLease, RenderLeaseId, RenderQueue, ServiceError,
12 ServiceResult, WorkerId,
13};
14
15#[derive(Debug, Default)]
16pub struct InMemoryRenderQueue {
17 next_lease: AtomicU64,
18 inner: Mutex<InMemoryRenderQueueState>,
19}
20
21#[derive(Debug, Default)]
22struct InMemoryRenderQueueState {
23 pending: VecDeque<RenderJob>,
24 leased: HashMap<RenderLeaseId, RenderLease>,
25}
26
27impl InMemoryRenderQueue {
28 pub fn new() -> Self {
29 Self::default()
30 }
31}
32
33impl RenderQueue for InMemoryRenderQueue {
34 fn enqueue<'a>(&'a self, job: RenderJob) -> BoxFuture<'a, ServiceResult<RenderJobId>> {
35 Box::pin(async move {
36 let id = job.id.clone();
37 self.inner
38 .lock()
39 .map_err(lock_error)?
40 .pending
41 .push_back(job);
42 Ok(id)
43 })
44 }
45
46 fn lease<'a>(
47 &'a self,
48 _worker_id: &'a WorkerId,
49 ttl: Duration,
50 ) -> BoxFuture<'a, ServiceResult<Option<RenderLease>>> {
51 Box::pin(async move {
52 let Some(job) = self.inner.lock().map_err(lock_error)?.pending.pop_front() else {
53 return Ok(None);
54 };
55 let lease = RenderLease {
56 id: RenderLeaseId(format!(
57 "lease-{}",
58 self.next_lease.fetch_add(1, Ordering::Relaxed)
59 )),
60 job,
61 leased_until: Instant::now() + ttl,
62 };
63 self.inner
64 .lock()
65 .map_err(lock_error)?
66 .leased
67 .insert(lease.id.clone(), lease.clone());
68 Ok(Some(lease))
69 })
70 }
71
72 fn ack<'a>(&'a self, lease_id: RenderLeaseId) -> BoxFuture<'a, ServiceResult<()>> {
73 Box::pin(async move {
74 self.inner
75 .lock()
76 .map_err(lock_error)?
77 .leased
78 .remove(&lease_id);
79 Ok(())
80 })
81 }
82
83 fn nack<'a>(
84 &'a self,
85 lease_id: RenderLeaseId,
86 _reason: ServiceError,
87 ) -> BoxFuture<'a, ServiceResult<()>> {
88 Box::pin(async move {
89 let mut inner = self.inner.lock().map_err(lock_error)?;
90 if let Some(lease) = inner.leased.remove(&lease_id) {
91 inner.pending.push_back(lease.job);
92 }
93 Ok(())
94 })
95 }
96
97 fn heartbeat<'a>(
98 &'a self,
99 lease_id: &'a RenderLeaseId,
100 ttl: Duration,
101 ) -> BoxFuture<'a, ServiceResult<()>> {
102 Box::pin(async move {
103 if let Some(lease) = self
104 .inner
105 .lock()
106 .map_err(lock_error)?
107 .leased
108 .get_mut(lease_id)
109 {
110 lease.leased_until = Instant::now() + ttl;
111 }
112 Ok(())
113 })
114 }
115}
116
117fn lock_error<T>(_err: std::sync::PoisonError<T>) -> ServiceError {
118 ServiceError {
119 code: "queue_lock_poisoned",
120 message: "render queue lock was poisoned".to_string(),
121 retryable: true,
122 }
123}
124
125#[cfg(test)]
126mod tests {
127 use super::{InMemoryRenderQueue, RenderJob, RenderQueue, WorkerId};
128
129 #[tokio::test]
130 async fn in_memory_queue_requeues_nacked_jobs() {
131 let queue = InMemoryRenderQueue::new();
132 queue
133 .enqueue(RenderJob::new("job-1", serde_json::json!({})))
134 .await
135 .expect("enqueue");
136
137 let worker = WorkerId("worker-1".to_string());
138 let lease = queue
139 .lease(&worker, std::time::Duration::from_secs(30))
140 .await
141 .expect("lease")
142 .expect("lease present");
143 assert_eq!(lease.job.id.0, "job-1");
144
145 queue
146 .nack(
147 lease.id,
148 super::ServiceError {
149 code: "test_failure",
150 message: "test failure".to_string(),
151 retryable: true,
152 },
153 )
154 .await
155 .expect("nack");
156
157 let lease = queue
158 .lease(&worker, std::time::Duration::from_secs(30))
159 .await
160 .expect("lease")
161 .expect("lease present");
162 assert_eq!(lease.job.id.0, "job-1");
163 }
164}