runledger_runtime/
registry.rs1use 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 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 #[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(®istry).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(®istry).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(®istry).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(®istry).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}