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