Skip to main content

runledger_runtime/catalog/
registration.rs

1use 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    /// Creates an empty catalog with default definition values.
18    #[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    /// Replaces the definition defaults used by subsequent sync operations.
28    #[must_use]
29    pub fn defaults(mut self, defaults: JobCatalogDefaults) -> Self {
30        self.defaults = defaults;
31        self
32    }
33
34    /// Registers a handler after validating declared and handler job types match.
35    ///
36    /// # Errors
37    /// Returns [`CatalogError`] when job types are blank, mismatched, or duplicated.
38    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    /// Registers a handler, panicking when validation fails.
74    #[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    /// Registers a handler with job-specific definition overrides.
85    ///
86    /// # Errors
87    /// Returns [`CatalogError`] when job types are blank, mismatched, or duplicated.
88    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    /// Registers a handler with job-specific definition overrides, panicking
102    /// when validation fails.
103    #[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    /// Replaces the definition overrides for one registered catalog job.
122    ///
123    /// # Errors
124    /// Returns [`CatalogError::UnknownJobType`] when the job type is not registered.
125    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    /// Replaces the definition overrides for one registered catalog job,
145    /// panicking when validation fails.
146    #[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    /// Registers a policy retry-delay override for a catalog job type.
159    ///
160    /// A lower bound attached directly to a handler's
161    /// [`runledger_core::jobs::JobFailure`] may extend this delay but cannot
162    /// shorten it.
163    ///
164    /// # Errors
165    /// Returns [`CatalogError`] when the job type is unknown or override values are invalid.
166    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    /// Registers a policy retry-delay override, panicking when validation
184    /// fails.
185    ///
186    /// A lower bound attached directly to a handler's
187    /// [`runledger_core::jobs::JobFailure`] may extend this delay but cannot
188    /// shorten it.
189    #[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    /// Registers a catalog-owned schedule after validating its shape and job type.
205    ///
206    /// Registered schedules are used by [`Self::sync_schedules`] and
207    /// [`Self::sync_schedules_exact`]. The referenced job must already be
208    /// registered on this builder and effectively enabled.
209    ///
210    /// # Errors
211    /// Returns [`CatalogError`] when the schedule spec is invalid, the job type
212    /// is unknown or disabled, or the schedule name is already registered.
213    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    /// Registers a catalog-owned schedule, panicking when validation fails.
233    ///
234    /// Use [`Self::try_schedule`] when registration data is not static or should
235    /// be reported as a recoverable startup error.
236    #[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    /// Converts the catalog into a runtime [`JobRegistry`].
244    ///
245    /// Disabled catalog jobs still register handlers so workers can process
246    /// already-queued work and dead-letter hooks.
247    #[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    /// Returns whether the catalog has a registered job type.
264    #[must_use]
265    pub fn contains(&self, job_type: JobType<'_>) -> bool {
266        self.jobs.contains_key(job_type.as_str())
267    }
268
269    /// Returns a catalog job type when it is registered.
270    ///
271    /// # Errors
272    /// Returns [`CatalogError::UnknownJobType`] when the name is not in the catalog.
273    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    /// Returns a catalog job type when it is registered and catalog-enabled.
279    ///
280    /// This checks catalog configuration only. It does not read `job_definitions`;
281    /// operator-disabled database rows are enforced later by persistence APIs.
282    /// Job-specific definition overrides take precedence over the catalog
283    /// default enabled flag when present.
284    ///
285    /// # Errors
286    /// Returns [`CatalogError::UnknownJobType`] or [`CatalogError::DisabledJobType`].
287    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}