Skip to main content

runledger_runtime/
registry.rs

1use std::collections::HashMap;
2use std::sync::Arc;
3
4pub use runledger_core::jobs::JobHandler;
5use runledger_core::jobs::{JobHandlerRegistry, JobType};
6use thiserror::Error;
7
8#[derive(Debug, Clone, Error, Eq, PartialEq)]
9#[non_exhaustive]
10pub enum JobRegistryError {
11    #[error("job handler already registered for {job_type}")]
12    DuplicateJobType { job_type: JobType<'static> },
13}
14
15#[derive(Clone, Default)]
16pub struct JobRegistry {
17    handlers: HashMap<JobType<'static>, Arc<dyn JobHandler>>,
18    retry_delay_overrides: HashMap<JobType<'static>, HashMap<&'static str, i32>>,
19}
20
21impl JobRegistry {
22    #[must_use]
23    pub fn new() -> Self {
24        Self::default()
25    }
26
27    pub fn register<H>(&mut self, handler: H)
28    where
29        H: JobHandler + 'static,
30    {
31        let handler: Arc<dyn JobHandler> = Arc::new(handler);
32        let job_type = handler.job_type();
33        self.register_boxed_for_type(job_type, handler);
34    }
35
36    pub fn try_register<H>(&mut self, handler: H) -> Result<(), JobRegistryError>
37    where
38        H: JobHandler + 'static,
39    {
40        self.try_register_boxed(Arc::new(handler))
41    }
42
43    pub fn try_register_boxed(
44        &mut self,
45        handler: Arc<dyn JobHandler>,
46    ) -> Result<(), JobRegistryError> {
47        let job_type = handler.job_type();
48        if self.handlers.contains_key(job_type.as_str()) {
49            return Err(JobRegistryError::DuplicateJobType { job_type });
50        }
51
52        self.register_boxed_for_type(job_type, handler);
53        Ok(())
54    }
55
56    pub(crate) fn register_boxed_for_type(
57        &mut self,
58        job_type: JobType<'static>,
59        handler: Arc<dyn JobHandler>,
60    ) {
61        self.handlers.insert(job_type, handler);
62    }
63
64    /// Registers the policy retry delay for one job type and failure code.
65    ///
66    /// A lower bound attached directly to [`runledger_core::jobs::JobFailure`]
67    /// may extend this delay but can never shorten it. The worker uses its
68    /// exponential backoff when no matching override exists.
69    ///
70    /// # Panics
71    ///
72    /// Panics when `retry_delay_ms` is not positive.
73    pub fn register_retry_delay_override(
74        &mut self,
75        job_type: JobType<'static>,
76        failure_code: &'static str,
77        retry_delay_ms: i32,
78    ) {
79        assert!(retry_delay_ms > 0, "retry delay override must be positive");
80
81        self.retry_delay_overrides
82            .entry(job_type)
83            .or_default()
84            .insert(failure_code, retry_delay_ms);
85    }
86
87    #[must_use]
88    pub fn get(&self, job_type: JobType<'_>) -> Option<Arc<dyn JobHandler>> {
89        self.handlers.get(job_type.as_str()).cloned()
90    }
91
92    /// Returns the configured policy retry delay for an exact job type and
93    /// failure code.
94    ///
95    /// Handler retry timing is a lower bound that may extend this policy delay
96    /// but cannot shorten it.
97    #[must_use]
98    pub fn retry_delay_override(&self, job_type: JobType<'_>, failure_code: &str) -> Option<i32> {
99        self.retry_delay_overrides
100            .get(job_type.as_str())
101            .and_then(|overrides| overrides.get(failure_code).copied())
102    }
103
104    #[must_use]
105    pub(crate) fn registered_static_types(&self) -> Vec<JobType<'static>> {
106        let mut keys: Vec<JobType<'static>> = self.handlers.keys().copied().collect();
107        keys.sort_unstable();
108        keys
109    }
110
111    #[must_use]
112    pub fn registered_types(&self) -> Vec<JobType<'_>> {
113        self.registered_static_types()
114    }
115}
116
117impl JobHandlerRegistry for JobRegistry {
118    fn register_boxed(&mut self, handler: Arc<dyn JobHandler>) {
119        let job_type = handler.job_type();
120        self.register_boxed_for_type(job_type, handler);
121    }
122}
123
124#[cfg(test)]
125mod tests {
126    use std::sync::Arc;
127
128    use async_trait::async_trait;
129    use runledger_core::jobs::{
130        JobCompletion, JobContext, JobFailure, JobHandlerRegistry, JobType,
131    };
132    use serde_json::{Value, json};
133    use uuid::Uuid;
134
135    use super::{JobHandler, JobRegistry, JobRegistryError};
136
137    struct ExampleHandler;
138
139    #[async_trait]
140    impl JobHandler for ExampleHandler {
141        fn job_type(&self) -> JobType<'static> {
142            JobType::new("jobs.example")
143        }
144
145        async fn execute(
146            &self,
147            _context: JobContext,
148            _payload: Value,
149        ) -> Result<JobCompletion, JobFailure> {
150            Ok(JobCompletion::success())
151        }
152    }
153
154    struct OutputHandler(&'static str);
155
156    #[async_trait]
157    impl JobHandler for OutputHandler {
158        fn job_type(&self) -> JobType<'static> {
159            JobType::new("jobs.example")
160        }
161
162        async fn execute(
163            &self,
164            _context: JobContext,
165            _payload: Value,
166        ) -> Result<JobCompletion, JobFailure> {
167            Ok(JobCompletion::with_output(json!(self.0)))
168        }
169    }
170
171    fn test_context() -> JobContext {
172        JobContext {
173            job_id: Uuid::now_v7(),
174            run_number: 1,
175            attempt: 1,
176            organization_id: None,
177            worker_id: "registry-test-worker".to_string(),
178            checkpoint: None,
179        }
180    }
181
182    async fn registered_output(registry: &JobRegistry) -> Value {
183        registry
184            .get(JobType::new("jobs.example"))
185            .expect("handler exists")
186            .execute(test_context(), json!({}))
187            .await
188            .expect("registered handler should execute")
189            .output()
190            .cloned()
191            .expect("handler should return output")
192    }
193
194    #[tokio::test]
195    async fn registered_handler_executes_successfully_via_trait_object() {
196        let mut registry = JobRegistry::new();
197        registry.register(ExampleHandler);
198        let handler = registry
199            .get(JobType::new("jobs.example"))
200            .expect("handler exists");
201
202        handler
203            .execute(test_context(), json!({}))
204            .await
205            .expect("registered handler should execute");
206    }
207
208    #[test]
209    fn registered_types_returns_sorted_job_types() {
210        let mut registry = JobRegistry::new();
211        registry.register(ExampleHandler);
212
213        assert_eq!(
214            registry.registered_types(),
215            vec![JobType::new("jobs.example")]
216        );
217    }
218
219    #[tokio::test]
220    async fn try_register_rejects_duplicate_job_type_and_keeps_first_handler() {
221        let mut registry = JobRegistry::new();
222        registry
223            .try_register(OutputHandler("first"))
224            .expect("first handler registration should succeed");
225
226        assert_eq!(
227            registry.try_register(OutputHandler("second")),
228            Err(JobRegistryError::DuplicateJobType {
229                job_type: JobType::new("jobs.example"),
230            })
231        );
232        assert_eq!(registered_output(&registry).await, json!("first"));
233    }
234
235    #[tokio::test]
236    async fn try_register_boxed_rejects_duplicate_job_type_and_keeps_first_handler() {
237        let mut registry = JobRegistry::new();
238        registry
239            .try_register_boxed(Arc::new(OutputHandler("first")))
240            .expect("first handler registration should succeed");
241
242        assert_eq!(
243            registry.try_register_boxed(Arc::new(OutputHandler("second"))),
244            Err(JobRegistryError::DuplicateJobType {
245                job_type: JobType::new("jobs.example"),
246            })
247        );
248        assert_eq!(registered_output(&registry).await, json!("first"));
249    }
250
251    #[tokio::test]
252    async fn register_overwrites_duplicate_job_type() {
253        let mut registry = JobRegistry::new();
254        registry.register(OutputHandler("first"));
255        registry.register(OutputHandler("second"));
256
257        assert_eq!(registered_output(&registry).await, json!("second"));
258    }
259
260    #[tokio::test]
261    async fn trait_register_boxed_overwrites_duplicate_job_type() {
262        let mut registry = JobRegistry::new();
263        JobHandlerRegistry::register_boxed(&mut registry, Arc::new(OutputHandler("first")));
264        JobHandlerRegistry::register_boxed(&mut registry, Arc::new(OutputHandler("second")));
265
266        assert_eq!(registered_output(&registry).await, json!("second"));
267    }
268
269    #[test]
270    fn retry_delay_override_matches_job_type_and_failure_code() {
271        let mut registry = JobRegistry::new();
272        registry.register_retry_delay_override(
273            JobType::new("jobs.example"),
274            "job.example.wait",
275            42,
276        );
277
278        assert_eq!(
279            registry.retry_delay_override(JobType::new("jobs.example"), "job.example.wait"),
280            Some(42)
281        );
282        assert_eq!(
283            registry.retry_delay_override(JobType::new("jobs.other"), "job.example.wait"),
284            None
285        );
286        assert_eq!(
287            registry.retry_delay_override(JobType::new("jobs.example"), "job.example.other"),
288            None
289        );
290    }
291
292    #[test]
293    #[should_panic(expected = "retry delay override must be positive")]
294    fn retry_delay_override_rejects_zero_delay() {
295        let mut registry = JobRegistry::new();
296        registry.register_retry_delay_override(JobType::new("jobs.example"), "job.example.wait", 0);
297    }
298}