lumen_server/service/
worker.rs1use 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}