Skip to main content

lumen_server/service/
worker.rs

1use std::{sync::Arc, time::Duration};
2
3use super::{
4    ArtifactStore, ArtifactWrite, ProgressSink, RenderExecutor, RenderQueue, ServiceError,
5    ServiceResult, WorkerId,
6};
7
8pub struct RenderService<Q, E, A, P> {
9    pub queue: Arc<Q>,
10    pub executor: Arc<E>,
11    pub artifacts: Arc<A>,
12    pub progress: Arc<P>,
13}
14
15impl<Q, E, A, P> RenderService<Q, E, A, P>
16where
17    Q: RenderQueue,
18    E: RenderExecutor,
19    A: ArtifactStore,
20    P: ProgressSink,
21{
22    pub fn new(queue: Q, executor: E, artifacts: A, progress: P) -> Self {
23        Self {
24            queue: Arc::new(queue),
25            executor: Arc::new(executor),
26            artifacts: Arc::new(artifacts),
27            progress: Arc::new(progress),
28        }
29    }
30
31    pub async fn process_next(
32        &self,
33        worker_id: &WorkerId,
34        lease_ttl: Duration,
35    ) -> ServiceResult<ProcessNextOutcome> {
36        let Some(lease) = self.queue.lease(worker_id, lease_ttl).await? else {
37            return Ok(ProcessNextOutcome::NoJob);
38        };
39        let lease_id = lease.id.clone();
40        let job_id = lease.job.id.clone();
41
42        match self
43            .executor
44            .execute(lease.job, self.progress.as_ref())
45            .await
46        {
47            Ok(output) => {
48                let artifact = match self
49                    .artifacts
50                    .put(ArtifactWrite {
51                        job_id,
52                        bytes: output.bytes,
53                        content_type: output.content_type,
54                    })
55                    .await
56                {
57                    Ok(artifact) => artifact,
58                    Err(error) => {
59                        self.queue.nack(lease_id, error.clone()).await?;
60                        return Err(error);
61                    }
62                };
63                self.queue.ack(lease_id).await?;
64                Ok(ProcessNextOutcome::Completed {
65                    artifact_uri: artifact.uri,
66                })
67            }
68            Err(error) => {
69                self.queue.nack(lease_id, error.clone()).await?;
70                Err(error)
71            }
72        }
73    }
74}
75
76#[derive(Debug, Clone, PartialEq, Eq)]
77pub enum ProcessNextOutcome {
78    NoJob,
79    Completed { artifact_uri: String },
80}
81
82impl From<anyhow::Error> for ServiceError {
83    fn from(error: anyhow::Error) -> Self {
84        Self {
85            code: "service_error",
86            message: error.to_string(),
87            retryable: true,
88        }
89    }
90}