runledger_runtime/catalog/
registration.rs1use std::collections::BTreeMap;
2use std::sync::Arc;
3
4use runledger_core::jobs::{JobHandler, JobHandlerRegistry, JobType, JobTypeName};
5use runledger_postgres::jobs::JobDefinitionUpsert;
6
7use crate::registry::JobRegistry;
8
9use super::schedule_spec::{CatalogJobScheduleSpec, StoredCatalogJobScheduleSpec};
10use super::types::CatalogJob;
11use super::{
12 CatalogError, JobCatalog, JobCatalogDefaults, JobCatalogDefinitionOverrides,
13 JobCatalogSyncScope,
14};
15
16impl JobCatalog {
17 #[must_use]
19 pub fn new() -> Self {
20 Self {
21 defaults: JobCatalogDefaults::default(),
22 jobs: BTreeMap::new(),
23 schedules: Vec::new(),
24 }
25 }
26
27 #[must_use]
29 pub fn defaults(mut self, defaults: JobCatalogDefaults) -> Self {
30 self.defaults = defaults;
31 self
32 }
33
34 pub fn try_job<H>(mut self, job_type: &'static str, handler: H) -> Result<Self, CatalogError>
39 where
40 H: JobHandler + 'static,
41 {
42 let declared = Self::validate_declared_job_type(job_type)?;
43 let handler_type = Self::validate_handler_job_type(&handler)?;
44 if declared != handler_type {
45 return Err(CatalogError::HandlerJobTypeMismatch {
46 declared: declared.as_str().to_owned(),
47 handler: handler_type.as_str().to_owned(),
48 });
49 }
50
51 let key = JobTypeName::new(job_type).map_err(|source| CatalogError::InvalidJobType {
52 job_type: job_type.to_owned(),
53 source,
54 })?;
55 if self.jobs.contains_key(&key) {
56 return Err(CatalogError::DuplicateJobType {
57 job_type: job_type.to_owned(),
58 });
59 }
60
61 self.jobs.insert(
62 key,
63 CatalogJob {
64 job_type: declared,
65 handler: Arc::new(handler),
66 definition_overrides: JobCatalogDefinitionOverrides::new(),
67 retry_delay_overrides: BTreeMap::new(),
68 },
69 );
70 Ok(self)
71 }
72
73 #[must_use]
75 pub fn job<H>(self, job_type: &'static str, handler: H) -> Self
76 where
77 H: JobHandler + 'static,
78 {
79 self.try_job(job_type, handler).unwrap_or_else(|error| {
80 panic!("invalid job catalog registration for {job_type:?}: {error}");
81 })
82 }
83
84 pub fn try_job_with_definition_overrides<H>(
89 self,
90 job_type: &'static str,
91 handler: H,
92 overrides: JobCatalogDefinitionOverrides,
93 ) -> Result<Self, CatalogError>
94 where
95 H: JobHandler + 'static,
96 {
97 self.try_job(job_type, handler)?
98 .try_definition_overrides(job_type, overrides)
99 }
100
101 #[must_use]
104 pub fn job_with_definition_overrides<H>(
105 self,
106 job_type: &'static str,
107 handler: H,
108 overrides: JobCatalogDefinitionOverrides,
109 ) -> Self
110 where
111 H: JobHandler + 'static,
112 {
113 self.try_job_with_definition_overrides(job_type, handler, overrides)
114 .unwrap_or_else(|error| {
115 panic!(
116 "invalid job catalog registration with definition overrides for {job_type:?}: {error}"
117 );
118 })
119 }
120
121 pub fn try_definition_overrides(
126 mut self,
127 job_type: &str,
128 overrides: JobCatalogDefinitionOverrides,
129 ) -> Result<Self, CatalogError> {
130 let key = self.require_job_key(job_type)?;
131 overrides
132 .validate()
133 .map_err(|field| CatalogError::InvalidJobDefinitionValue {
134 job_type: job_type.to_owned(),
135 field,
136 })?;
137 self.jobs
138 .get_mut(&key)
139 .expect("job key validated")
140 .definition_overrides = overrides;
141 Ok(self)
142 }
143
144 #[must_use]
147 pub fn definition_overrides(
148 self,
149 job_type: &str,
150 overrides: JobCatalogDefinitionOverrides,
151 ) -> Self {
152 self.try_definition_overrides(job_type, overrides)
153 .unwrap_or_else(|error| {
154 panic!("invalid definition overrides for job type {job_type:?}: {error}");
155 })
156 }
157
158 pub fn try_retry_delay_override(
167 mut self,
168 job_type: &str,
169 failure_code: &'static str,
170 retry_delay_ms: i32,
171 ) -> Result<Self, CatalogError> {
172 let key = self.require_job_key(job_type)?;
173 Self::validate_failure_code(failure_code)?;
174 Self::validate_retry_delay(retry_delay_ms)?;
175 self.jobs
176 .get_mut(&key)
177 .expect("job key validated")
178 .retry_delay_overrides
179 .insert(failure_code, retry_delay_ms);
180 Ok(self)
181 }
182
183 #[must_use]
190 pub fn retry_delay_override(
191 self,
192 job_type: &str,
193 failure_code: &'static str,
194 retry_delay_ms: i32,
195 ) -> Self {
196 self.try_retry_delay_override(job_type, failure_code, retry_delay_ms)
197 .unwrap_or_else(|error| {
198 panic!(
199 "invalid retry delay override for job type {job_type:?}, failure code {failure_code:?}: {error}"
200 );
201 })
202 }
203
204 pub fn try_schedule(mut self, spec: CatalogJobScheduleSpec<'_>) -> Result<Self, CatalogError> {
214 spec.validate_shape()
215 .map_err(|field| CatalogError::InvalidScheduleSpec {
216 name: spec.name.to_owned(),
217 field,
218 })?;
219 self.require_catalog_enabled_job_type(spec.job_type)?;
220
221 if self.schedules.iter().any(|stored| stored.name == spec.name) {
222 return Err(CatalogError::DuplicateScheduleName {
223 name: spec.name.to_owned(),
224 });
225 }
226
227 self.schedules
228 .push(StoredCatalogJobScheduleSpec::from(&spec));
229 Ok(self)
230 }
231
232 #[must_use]
237 pub fn schedule(self, spec: CatalogJobScheduleSpec<'_>) -> Self {
238 self.try_schedule(spec).unwrap_or_else(|error| {
239 panic!("invalid job catalog schedule registration: {error}");
240 })
241 }
242
243 #[must_use]
248 pub fn to_registry(&self) -> JobRegistry {
249 let mut registry = JobRegistry::new();
250 for entry in self.jobs.values() {
251 registry.register_boxed(Arc::clone(&entry.handler));
252 for (failure_code, retry_delay_ms) in &entry.retry_delay_overrides {
253 registry.register_retry_delay_override(
254 entry.job_type,
255 failure_code,
256 *retry_delay_ms,
257 );
258 }
259 }
260 registry
261 }
262
263 #[must_use]
265 pub fn contains(&self, job_type: JobType<'_>) -> bool {
266 self.jobs.contains_key(job_type.as_str())
267 }
268
269 pub fn require_job_type(&self, job_type: &str) -> Result<JobType<'static>, CatalogError> {
274 let key = self.require_job_key(job_type)?;
275 Ok(self.jobs.get(&key).expect("job key validated").job_type)
276 }
277
278 pub fn require_catalog_enabled_job_type(
288 &self,
289 job_type: &str,
290 ) -> Result<JobType<'static>, CatalogError> {
291 let key = self.require_job_key(job_type)?;
292 let entry = self.jobs.get(&key).expect("job key validated");
293 if !self.effective_defaults(entry).is_enabled {
294 return Err(CatalogError::DisabledJobType {
295 job_type: entry.job_type.as_str().to_owned(),
296 });
297 }
298 Ok(entry.job_type)
299 }
300
301 pub(super) fn validate_defaults(&self) -> Result<(), CatalogError> {
302 self.defaults
303 .validate()
304 .map_err(|field| CatalogError::InvalidDefinitionValue { field })?;
305
306 for entry in self.jobs.values() {
307 self.effective_defaults(entry).validate().map_err(|field| {
308 CatalogError::InvalidJobDefinitionValue {
309 job_type: entry.job_type.as_str().to_owned(),
310 field,
311 }
312 })?;
313 }
314
315 Ok(())
316 }
317
318 pub(super) fn validate_exact_sync_scope(
319 &self,
320 scope: &JobCatalogSyncScope,
321 ) -> Result<(), CatalogError> {
322 if self.jobs.is_empty() {
323 return Err(CatalogError::EmptyExactSyncCatalog);
324 }
325
326 for entry in self.jobs.values() {
327 if !scope.contains(entry.job_type) {
328 return Err(CatalogError::JobTypeOutsideExactSyncScope {
329 job_type: entry.job_type.as_str().to_owned(),
330 });
331 }
332 }
333
334 Ok(())
335 }
336
337 pub(super) fn materialize_definition(
338 &self,
339 entry: &CatalogJob,
340 ) -> JobDefinitionUpsert<'static> {
341 let defaults = self.effective_defaults(entry);
342 JobDefinitionUpsert {
343 job_type: entry.job_type,
344 version: defaults.version,
345 max_attempts: defaults.max_attempts,
346 default_timeout_seconds: defaults.default_timeout_seconds,
347 default_priority: defaults.default_priority,
348 is_enabled: defaults.is_enabled,
349 }
350 }
351
352 pub(super) fn effective_defaults(&self, entry: &CatalogJob) -> JobCatalogDefaults {
353 entry.definition_overrides.apply_to(self.defaults)
354 }
355
356 pub(super) fn require_job_key(&self, job_type: &str) -> Result<JobTypeName, CatalogError> {
357 let key = JobTypeName::new(job_type).map_err(|source| CatalogError::InvalidJobType {
358 job_type: job_type.to_owned(),
359 source,
360 })?;
361 if self.jobs.contains_key(&key) {
362 Ok(key)
363 } else {
364 Err(CatalogError::UnknownJobType {
365 job_type: job_type.to_owned(),
366 })
367 }
368 }
369
370 fn validate_declared_job_type(
371 job_type: &'static str,
372 ) -> Result<JobType<'static>, CatalogError> {
373 JobType::try_new(job_type).map_err(|source| CatalogError::InvalidJobType {
374 job_type: job_type.to_owned(),
375 source,
376 })
377 }
378
379 fn validate_handler_job_type<H: JobHandler + ?Sized>(
380 handler: &H,
381 ) -> Result<JobType<'static>, CatalogError> {
382 let handler_job_type = handler.job_type();
383 JobType::try_new(handler_job_type.as_str()).map_err(|source| {
384 CatalogError::InvalidHandlerJobType {
385 handler_job_type: handler_job_type.as_str().to_owned(),
386 source,
387 }
388 })
389 }
390
391 fn validate_failure_code(failure_code: &str) -> Result<(), CatalogError> {
392 if failure_code.trim().is_empty() {
393 Err(CatalogError::InvalidFailureCode)
394 } else {
395 Ok(())
396 }
397 }
398
399 fn validate_retry_delay(retry_delay_ms: i32) -> Result<(), CatalogError> {
400 if retry_delay_ms <= 0 {
401 Err(CatalogError::InvalidRetryDelay)
402 } else {
403 Ok(())
404 }
405 }
406}
407
408impl Default for JobCatalog {
409 fn default() -> Self {
410 Self::new()
411 }
412}