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