Skip to main content

lumen_server/service/
queue.rs

1use 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}