runledger_runtime/catalog/
registration.rs1use std::collections::BTreeMap;
2use std::sync::Arc;
3
4use runledger_core::jobs::{JobHandler, 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_handler<H>(self, handler: H) -> Result<Self, CatalogError>
39 where
40 H: JobHandler + 'static,
41 {
42 let handler_type = Self::validate_handler_job_type(&handler)?;
43 self.insert_handler(handler_type, Arc::new(handler))
44 }
45
46 #[must_use]
49 pub fn handler<H>(self, handler: H) -> Self
50 where
51 H: JobHandler + 'static,
52 {
53 self.try_handler(handler).unwrap_or_else(|error| {
54 panic!("invalid job catalog handler registration: {error}");
55 })
56 }
57
58 #[deprecated(
67 since = "0.11.0",
68 note = "use JobCatalog::try_handler(handler); the handler now supplies the catalog identity"
69 )]
70 pub fn try_job<H>(self, job_type: &'static str, handler: H) -> Result<Self, CatalogError>
71 where
72 H: JobHandler + 'static,
73 {
74 let declared = Self::validate_declared_job_type(job_type)?;
75 let handler_type = Self::validate_handler_job_type(&handler)?;
76 if declared != handler_type {
77 return Err(CatalogError::HandlerJobTypeMismatch {
78 declared: declared.as_str().to_owned(),
79 handler: handler_type.as_str().to_owned(),
80 });
81 }
82
83 self.insert_handler(handler_type, Arc::new(handler))
84 }
85
86 #[must_use]
92 #[deprecated(
93 since = "0.11.0",
94 note = "use JobCatalog::handler(handler); the handler now supplies the catalog identity"
95 )]
96 pub fn job<H>(self, job_type: &'static str, handler: H) -> Self
97 where
98 H: JobHandler + 'static,
99 {
100 #[allow(
101 deprecated,
102 reason = "compatibility wrapper delegates to legacy validation"
103 )]
104 self.try_job(job_type, handler).unwrap_or_else(|error| {
105 panic!("invalid job catalog registration for {job_type:?}: {error}");
106 })
107 }
108
109 pub fn try_handler_with_definition_overrides<H>(
116 self,
117 handler: H,
118 overrides: JobCatalogDefinitionOverrides,
119 ) -> Result<Self, CatalogError>
120 where
121 H: JobHandler + 'static,
122 {
123 let job_type = Self::validate_handler_job_type(&handler)?;
124 self.insert_handler(job_type, Arc::new(handler))?
125 .try_definition_overrides(job_type.as_str(), overrides)
126 }
127
128 #[must_use]
131 pub fn handler_with_definition_overrides<H>(
132 self,
133 handler: H,
134 overrides: JobCatalogDefinitionOverrides,
135 ) -> Self
136 where
137 H: JobHandler + 'static,
138 {
139 self.try_handler_with_definition_overrides(handler, overrides)
140 .unwrap_or_else(|error| {
141 panic!(
142 "invalid job catalog handler registration with definition overrides: {error}"
143 );
144 })
145 }
146
147 #[deprecated(
157 since = "0.11.0",
158 note = "use JobCatalog::try_handler_with_definition_overrides(handler, overrides); the handler now supplies the catalog identity"
159 )]
160 pub fn try_job_with_definition_overrides<H>(
161 self,
162 job_type: &'static str,
163 handler: H,
164 overrides: JobCatalogDefinitionOverrides,
165 ) -> Result<Self, CatalogError>
166 where
167 H: JobHandler + 'static,
168 {
169 #[allow(
170 deprecated,
171 reason = "compatibility wrapper preserves mismatch diagnostics"
172 )]
173 self.try_job(job_type, handler)?
174 .try_definition_overrides(job_type, overrides)
175 }
176
177 #[must_use]
184 #[deprecated(
185 since = "0.11.0",
186 note = "use JobCatalog::handler_with_definition_overrides(handler, overrides); the handler now supplies the catalog identity"
187 )]
188 pub fn job_with_definition_overrides<H>(
189 self,
190 job_type: &'static str,
191 handler: H,
192 overrides: JobCatalogDefinitionOverrides,
193 ) -> Self
194 where
195 H: JobHandler + 'static,
196 {
197 #[allow(deprecated, reason = "compatibility wrapper delegates to fallible legacy API")]
198 self.try_job_with_definition_overrides(job_type, handler, overrides)
199 .unwrap_or_else(|error| {
200 panic!(
201 "invalid job catalog registration with definition overrides for {job_type:?}: {error}"
202 );
203 })
204 }
205
206 pub fn try_definition_overrides(
211 mut self,
212 job_type: &str,
213 overrides: JobCatalogDefinitionOverrides,
214 ) -> Result<Self, CatalogError> {
215 let key = self.require_job_key(job_type)?;
216 overrides
217 .validate()
218 .map_err(|field| CatalogError::InvalidJobDefinitionValue {
219 job_type: job_type.to_owned(),
220 field,
221 })?;
222 self.jobs
223 .get_mut(&key)
224 .expect("job key validated")
225 .definition_overrides = overrides;
226 Ok(self)
227 }
228
229 #[must_use]
232 pub fn definition_overrides(
233 self,
234 job_type: &str,
235 overrides: JobCatalogDefinitionOverrides,
236 ) -> Self {
237 self.try_definition_overrides(job_type, overrides)
238 .unwrap_or_else(|error| {
239 panic!("invalid definition overrides for job type {job_type:?}: {error}");
240 })
241 }
242
243 pub fn try_retry_delay_override(
252 mut self,
253 job_type: &str,
254 failure_code: &'static str,
255 retry_delay_ms: i32,
256 ) -> Result<Self, CatalogError> {
257 let key = self.require_job_key(job_type)?;
258 Self::validate_failure_code(failure_code)?;
259 Self::validate_retry_delay(retry_delay_ms)?;
260 self.jobs
261 .get_mut(&key)
262 .expect("job key validated")
263 .retry_delay_overrides
264 .insert(failure_code, retry_delay_ms);
265 Ok(self)
266 }
267
268 #[must_use]
275 pub fn retry_delay_override(
276 self,
277 job_type: &str,
278 failure_code: &'static str,
279 retry_delay_ms: i32,
280 ) -> Self {
281 self.try_retry_delay_override(job_type, failure_code, retry_delay_ms)
282 .unwrap_or_else(|error| {
283 panic!(
284 "invalid retry delay override for job type {job_type:?}, failure code {failure_code:?}: {error}"
285 );
286 })
287 }
288
289 pub fn try_schedule(mut self, spec: CatalogJobScheduleSpec<'_>) -> Result<Self, CatalogError> {
299 spec.validate_shape()
300 .map_err(|field| CatalogError::InvalidScheduleSpec {
301 name: spec.name.to_owned(),
302 field,
303 })?;
304 self.require_catalog_enabled_job_type(spec.job_type)?;
305
306 if self.schedules.iter().any(|stored| stored.name == spec.name) {
307 return Err(CatalogError::DuplicateScheduleName {
308 name: spec.name.to_owned(),
309 });
310 }
311
312 self.schedules
313 .push(StoredCatalogJobScheduleSpec::from(&spec));
314 Ok(self)
315 }
316
317 #[must_use]
322 pub fn schedule(self, spec: CatalogJobScheduleSpec<'_>) -> Self {
323 self.try_schedule(spec).unwrap_or_else(|error| {
324 panic!("invalid job catalog schedule registration: {error}");
325 })
326 }
327
328 #[must_use]
333 pub fn to_registry(&self) -> JobRegistry {
334 let mut registry = JobRegistry::new();
335 for entry in self.jobs.values() {
336 registry.register_boxed_for_type(entry.job_type(), Arc::clone(&entry.handler));
337 for (failure_code, retry_delay_ms) in &entry.retry_delay_overrides {
338 registry.register_retry_delay_override(
339 entry.job_type(),
340 failure_code,
341 *retry_delay_ms,
342 );
343 }
344 }
345 registry
346 }
347
348 #[must_use]
350 pub fn contains(&self, job_type: JobType<'_>) -> bool {
351 self.jobs.contains_key(job_type.as_str())
352 }
353
354 pub fn require_job_type(&self, job_type: &str) -> Result<JobType<'static>, CatalogError> {
359 let key = self.require_job_key(job_type)?;
360 Ok(self.jobs.get(&key).expect("job key validated").job_type())
361 }
362
363 pub fn require_catalog_enabled_job_type(
373 &self,
374 job_type: &str,
375 ) -> Result<JobType<'static>, CatalogError> {
376 let key = self.require_job_key(job_type)?;
377 let entry = self.jobs.get(&key).expect("job key validated");
378 if !self.effective_defaults(entry).is_enabled {
379 return Err(CatalogError::DisabledJobType {
380 job_type: entry.job_type().as_str().to_owned(),
381 });
382 }
383 Ok(entry.job_type())
384 }
385
386 pub(super) fn validate_defaults(&self) -> Result<(), CatalogError> {
387 self.defaults
388 .validate()
389 .map_err(|field| CatalogError::InvalidDefinitionValue { field })?;
390
391 for entry in self.jobs.values() {
392 self.effective_defaults(entry).validate().map_err(|field| {
393 CatalogError::InvalidJobDefinitionValue {
394 job_type: entry.job_type().as_str().to_owned(),
395 field,
396 }
397 })?;
398 }
399
400 Ok(())
401 }
402
403 pub(super) fn validate_exact_sync_scope(
404 &self,
405 scope: &JobCatalogSyncScope,
406 ) -> Result<(), CatalogError> {
407 if self.jobs.is_empty() {
408 return Err(CatalogError::EmptyExactSyncCatalog);
409 }
410
411 for entry in self.jobs.values() {
412 if !scope.contains(entry.job_type()) {
413 return Err(CatalogError::JobTypeOutsideExactSyncScope {
414 job_type: entry.job_type().as_str().to_owned(),
415 });
416 }
417 }
418
419 Ok(())
420 }
421
422 pub(super) fn materialize_definition(
423 &self,
424 entry: &CatalogJob,
425 ) -> JobDefinitionUpsert<'static> {
426 let defaults = self.effective_defaults(entry);
427 JobDefinitionUpsert {
428 job_type: entry.job_type(),
429 version: defaults.version,
430 max_attempts: defaults.max_attempts,
431 default_timeout_seconds: defaults.default_timeout_seconds,
432 default_priority: defaults.default_priority,
433 is_enabled: defaults.is_enabled,
434 }
435 }
436
437 pub(super) fn effective_defaults(&self, entry: &CatalogJob) -> JobCatalogDefaults {
438 entry.definition_overrides.apply_to(self.defaults)
439 }
440
441 pub(super) fn require_job_key(&self, job_type: &str) -> Result<JobTypeName, CatalogError> {
442 let key = JobTypeName::new(job_type).map_err(|source| CatalogError::InvalidJobType {
443 job_type: job_type.to_owned(),
444 source,
445 })?;
446 if self.jobs.contains_key(&key) {
447 Ok(key)
448 } else {
449 Err(CatalogError::UnknownJobType {
450 job_type: job_type.to_owned(),
451 })
452 }
453 }
454
455 fn validate_declared_job_type(
456 job_type: &'static str,
457 ) -> Result<JobType<'static>, CatalogError> {
458 JobType::try_new(job_type).map_err(|source| CatalogError::InvalidJobType {
459 job_type: job_type.to_owned(),
460 source,
461 })
462 }
463
464 fn validate_handler_job_type<H: JobHandler + ?Sized>(
465 handler: &H,
466 ) -> Result<JobType<'static>, CatalogError> {
467 let handler_job_type = handler.job_type();
468 JobType::try_new(handler_job_type.as_str()).map_err(|source| {
469 CatalogError::InvalidHandlerJobType {
470 handler_job_type: handler_job_type.as_str().to_owned(),
471 source,
472 }
473 })
474 }
475
476 fn insert_handler(
477 mut self,
478 job_type: JobType<'static>,
479 handler: Arc<dyn JobHandler>,
480 ) -> Result<Self, CatalogError> {
481 let key = JobTypeName::new(job_type.as_str()).map_err(|source| {
482 CatalogError::InvalidHandlerJobType {
483 handler_job_type: job_type.as_str().to_owned(),
484 source,
485 }
486 })?;
487 if self.jobs.contains_key(&key) {
488 return Err(CatalogError::DuplicateJobType {
489 job_type: job_type.as_str().to_owned(),
490 });
491 }
492
493 self.jobs.insert(
494 key,
495 CatalogJob {
496 job_type,
497 handler,
498 definition_overrides: JobCatalogDefinitionOverrides::new(),
499 retry_delay_overrides: BTreeMap::new(),
500 },
501 );
502 Ok(self)
503 }
504
505 fn validate_failure_code(failure_code: &str) -> Result<(), CatalogError> {
506 if failure_code.trim().is_empty() {
507 Err(CatalogError::InvalidFailureCode)
508 } else {
509 Ok(())
510 }
511 }
512
513 fn validate_retry_delay(retry_delay_ms: i32) -> Result<(), CatalogError> {
514 if retry_delay_ms <= 0 {
515 Err(CatalogError::InvalidRetryDelay)
516 } else {
517 Ok(())
518 }
519 }
520}
521
522impl Default for JobCatalog {
523 fn default() -> Self {
524 Self::new()
525 }
526}