Skip to main content

solti_model/resource/
spec.rs

1//! # Task spec
2//!
3//! [`TaskSpec`] defines workload, slot, timeout, policies, and runner selection.
4
5use std::num::NonZeroU32;
6
7use serde::{Deserialize, Serialize};
8
9use crate::{
10    AdmissionPolicy, BackoffPolicy, LabelSelector, RestartPolicy, Slot, TaskWorkload, Timeout,
11    error::{ModelError, ModelResult},
12};
13
14/// Desired state for a task.
15///
16/// Use [`TaskSpec::builder`] to construct it.
17/// Resource constructors validate the completed spec.
18///
19/// ## Example
20///
21/// ```
22/// use solti_model::{EmbeddedSpec, RestartPolicy, TaskSpec, TaskWorkload};
23///
24/// let workload = TaskWorkload::Embedded(EmbeddedSpec::new("v1").unwrap());
25/// let spec = TaskSpec::builder("daily-cleanup", workload, 5_000u64)
26///     .restart(RestartPolicy::periodic(60_000))
27///     .build()
28///     .unwrap();
29///
30/// assert_eq!(spec.slot().as_str(), "daily-cleanup");
31/// assert_eq!(spec.timeout().as_millis(), 5_000);
32/// ```
33#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
34#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
35#[cfg_attr(feature = "schema", schemars(!try_from, deny_unknown_fields))]
36#[serde(rename_all = "camelCase")]
37#[serde(try_from = "raw::TaskSpecRaw")]
38pub struct TaskSpec {
39    slot: Slot,
40    workload: TaskWorkload,
41
42    timeout: Timeout,
43    restart: RestartPolicy,
44    backoff: BackoffPolicy,
45    admission: AdmissionPolicy,
46    #[serde(default, skip_serializing_if = "Option::is_none")]
47    max_retries: Option<NonZeroU32>,
48
49    #[serde(default, skip_serializing_if = "Option::is_none")]
50    runner_selector: Option<LabelSelector>,
51}
52
53impl TaskSpec {
54    /// Logical slot name for concurrency control.
55    #[inline]
56    pub fn slot(&self) -> &Slot {
57        &self.slot
58    }
59
60    /// Workload desired state.
61    #[inline]
62    pub fn workload(&self) -> &TaskWorkload {
63        &self.workload
64    }
65
66    /// Per-attempt timeout.
67    #[inline]
68    pub fn timeout(&self) -> Timeout {
69        self.timeout
70    }
71
72    /// Restart policy applied after completion or failure.
73    #[inline]
74    pub fn restart(&self) -> RestartPolicy {
75        self.restart
76    }
77
78    /// Backoff configuration between restart attempts.
79    #[inline]
80    pub fn backoff(&self) -> &BackoffPolicy {
81        &self.backoff
82    }
83
84    /// Admission policy for handling slot conflicts.
85    #[inline]
86    pub fn admission(&self) -> AdmissionPolicy {
87        self.admission
88    }
89
90    /// Maximum consecutive failure retries.
91    ///
92    /// `None` means unlimited.
93    /// Zero is not representable.
94    #[inline]
95    pub fn max_retries(&self) -> Option<NonZeroU32> {
96        self.max_retries
97    }
98
99    /// Label selector for runner routing, if present.
100    #[inline]
101    pub fn runner_selector(&self) -> Option<&LabelSelector> {
102        self.runner_selector.as_ref()
103    }
104}
105
106impl TaskSpec {
107    /// Creates a builder with the required fields.
108    ///
109    /// ```rust
110    /// use solti_model::{TaskSpec, TaskWorkload, SubprocessSpec, SubprocessMode, RestartPolicy};
111    ///
112    /// let spec = TaskSpec::builder(
113    ///     "my-slot",
114    ///     TaskWorkload::Subprocess(SubprocessSpec::new(
115    ///         SubprocessMode::Command {
116    ///             command: "echo".into(),
117    ///             args: vec!["hello".into()],
118    ///         },
119    ///         Default::default(),
120    ///         None,
121    ///         Default::default(),
122    ///     )),
123    ///     5_000u64,
124    /// )
125    /// .restart(RestartPolicy::OnFailure)
126    /// .build()
127    /// .expect("valid spec");
128    /// ```
129    pub fn builder(
130        slot: impl AsRef<str>,
131        workload: TaskWorkload,
132        timeout: impl Into<u64>,
133    ) -> TaskSpecBuilder {
134        TaskSpecBuilder::new(slot, workload, timeout)
135    }
136}
137
138impl TaskSpec {
139    /// Sets a runner selector on an existing spec.
140    ///
141    /// This method does not validate `sel`.
142    /// Call [`Self::validate`] before using the result outside a validated resource.
143    ///
144    /// ## Example
145    ///
146    /// ```
147    /// use solti_model::{EmbeddedSpec, Labels, LabelSelector, TaskSpec, TaskWorkload};
148    ///
149    /// let mut labels = Labels::new();
150    /// labels.insert("zone", "eu");
151    ///
152    /// let workload = TaskWorkload::Embedded(EmbeddedSpec::new("v1").unwrap());
153    /// let spec = TaskSpec::builder("build", workload, 1_000u64)
154    ///     .build()
155    ///     .unwrap()
156    ///     .with_runner_selector(LabelSelector::from_labels(labels));
157    ///
158    /// assert!(spec.runner_selector().is_some());
159    /// ```
160    #[inline]
161    pub fn with_runner_selector(mut self, sel: LabelSelector) -> Self {
162        self.runner_selector = Some(sel);
163        self
164    }
165
166    /// Sets the admission policy on an existing spec.
167    ///
168    /// ## Example
169    ///
170    /// ```
171    /// use solti_model::{AdmissionPolicy, EmbeddedSpec, TaskSpec, TaskWorkload};
172    ///
173    /// let workload = TaskWorkload::Embedded(EmbeddedSpec::new("v1").unwrap());
174    /// let spec = TaskSpec::builder("agent", workload, 1_000u64)
175    ///     .build()
176    ///     .unwrap()
177    ///     .with_admission(AdmissionPolicy::Replace);
178    ///
179    /// assert_eq!(spec.admission(), AdmissionPolicy::Replace);
180    /// ```
181    #[inline]
182    pub fn with_admission(mut self, admission: AdmissionPolicy) -> Self {
183        self.admission = admission;
184        self
185    }
186}
187
188impl TaskSpec {
189    /// Validates the complete spec.
190    ///
191    /// # Errors
192    ///
193    /// Returns [`ModelError::Invalid`] when the slot, workload, backoff, or runner selector is invalid.
194    ///
195    /// ## Example
196    ///
197    /// ```
198    /// use solti_model::{
199    ///     Flag, SubprocessMode, SubprocessSpec, TaskEnv, TaskSpec, TaskWorkload,
200    /// };
201    ///
202    /// let spec = TaskSpec::builder(
203    ///     "hello",
204    ///     TaskWorkload::Subprocess(SubprocessSpec::new(
205    ///         SubprocessMode::Command {
206    ///             command: "echo".into(),
207    ///             args: vec!["hello".into()],
208    ///         },
209    ///         TaskEnv::default(),
210    ///         None,
211    ///         Flag::enabled(),
212    ///     )),
213    ///     1_000u64,
214    /// )
215    /// .build()
216    /// .unwrap();
217    ///
218    /// spec.validate().unwrap();
219    /// ```
220    pub fn validate(&self) -> ModelResult<()> {
221        self.validate_structural()
222    }
223
224    /// Validates all structural fields.
225    fn validate_structural(&self) -> ModelResult<()> {
226        self.slot.validate_format()?;
227        self.workload.validate()?;
228        self.backoff.validate()?;
229        if let Some(ref sel) = self.runner_selector {
230            sel.validate()?;
231        }
232        Ok(())
233    }
234}
235
236/// Builder for [`TaskSpec`].
237///
238/// Required fields are set by [`TaskSpec::builder`].
239///
240/// Optional fields use these defaults:
241///
242/// - `backoff`: [`BackoffPolicy::default`] (full jitter, 1 second to 30 seconds, factor 2)
243/// - `admission`: [`AdmissionPolicy::DropIfRunning`]
244/// - `restart`: [`RestartPolicy::Never`]
245/// - `max_retries`: `None`
246/// - `runner_selector`: `None`
247///
248/// ## Example
249///
250/// ```
251/// use solti_model::{AdmissionPolicy, EmbeddedSpec, RestartPolicy, TaskSpec, TaskWorkload};
252///
253/// let workload = TaskWorkload::Embedded(EmbeddedSpec::new("v1").unwrap());
254/// let spec = TaskSpec::builder("service", workload, 5_000u64)
255///     .restart(RestartPolicy::always())
256///     .admission(AdmissionPolicy::Replace)
257///     .build()
258///     .unwrap();
259///
260/// assert_eq!(spec.restart(), RestartPolicy::always());
261/// assert_eq!(spec.admission(), AdmissionPolicy::Replace);
262/// ```
263pub struct TaskSpecBuilder {
264    runner_selector: Option<LabelSelector>,
265
266    workload: TaskWorkload,
267    slot: String,
268
269    backoff: BackoffPolicy,
270    restart: RestartPolicy,
271    timeout_ms: u64,
272    max_retries: Option<NonZeroU32>,
273
274    admission: AdmissionPolicy,
275}
276
277impl TaskSpecBuilder {
278    fn new(slot: impl AsRef<str>, workload: TaskWorkload, timeout: impl Into<u64>) -> Self {
279        Self {
280            runner_selector: None,
281
282            workload,
283            slot: slot.as_ref().to_owned(),
284
285            restart: RestartPolicy::default(),
286            backoff: BackoffPolicy::default(),
287            timeout_ms: timeout.into(),
288
289            admission: AdmissionPolicy::default(),
290            max_retries: None,
291        }
292    }
293
294    /// Sets the restart policy.
295    #[must_use]
296    pub fn restart(mut self, restart: RestartPolicy) -> Self {
297        self.restart = restart;
298        self
299    }
300
301    /// Sets the failure-retry budget.
302    ///
303    /// `None` means unlimited.
304    ///
305    /// ```rust
306    /// # use solti_model::{EmbeddedSpec, TaskSpec, TaskWorkload};
307    /// # use std::num::NonZeroU32;
308    /// let workload = TaskWorkload::Embedded(EmbeddedSpec::new("v1").unwrap());
309    /// let spec = TaskSpec::builder("s", workload, 1_000u64)
310    ///     .max_retries(NonZeroU32::new(3))
311    ///     .build()
312    ///     .expect("valid spec");
313    /// assert_eq!(spec.max_retries().map(NonZeroU32::get), Some(3));
314    /// ```
315    #[must_use]
316    pub fn max_retries(mut self, max_retries: impl Into<Option<NonZeroU32>>) -> Self {
317        self.max_retries = max_retries.into();
318        self
319    }
320
321    /// Sets the backoff policy.
322    #[must_use]
323    pub fn backoff(mut self, backoff: BackoffPolicy) -> Self {
324        self.backoff = backoff;
325        self
326    }
327
328    /// Sets the admission policy.
329    #[must_use]
330    pub fn admission(mut self, admission: AdmissionPolicy) -> Self {
331        self.admission = admission;
332        self
333    }
334
335    /// Sets the runner selector.
336    #[must_use]
337    pub fn runner_selector(mut self, sel: LabelSelector) -> Self {
338        self.runner_selector = Some(sel);
339        self
340    }
341
342    /// Builds and validates the spec.
343    ///
344    /// Runner availability is not checked by this crate.
345    ///
346    /// # Errors
347    ///
348    /// Returns [`ModelError::Invalid`] when any structural field is invalid.
349    ///
350    /// ## Example
351    ///
352    /// ```
353    /// use solti_model::{EmbeddedSpec, TaskSpec, TaskWorkload};
354    ///
355    /// let workload = TaskWorkload::Embedded(EmbeddedSpec::new("v1").unwrap());
356    /// let err = TaskSpec::builder("", workload, 1_000u64)
357    ///     .build()
358    ///     .unwrap_err();
359    ///
360    /// assert!(err.to_string().contains("slot"));
361    /// ```
362    pub fn build(self) -> ModelResult<TaskSpec> {
363        let spec = TaskSpec {
364            runner_selector: self.runner_selector,
365
366            workload: self.workload,
367            slot: Slot::new(self.slot)?,
368
369            restart: self.restart,
370            backoff: self.backoff,
371            timeout: Timeout::new(self.timeout_ms)?,
372
373            admission: self.admission,
374            max_retries: self.max_retries,
375        };
376        spec.validate_structural()?;
377        Ok(spec)
378    }
379}
380
381mod raw {
382    use super::*;
383
384    #[derive(Deserialize)]
385    #[serde(rename_all = "camelCase", deny_unknown_fields)]
386    pub(super) struct TaskSpecRaw {
387        slot: Slot,
388        workload: TaskWorkload,
389        timeout: Timeout,
390        restart: RestartPolicy,
391        backoff: BackoffPolicy,
392        admission: AdmissionPolicy,
393        #[serde(default)]
394        max_retries: Option<u32>,
395
396        #[serde(default)]
397        runner_selector: Option<LabelSelector>,
398    }
399
400    impl TryFrom<TaskSpecRaw> for TaskSpec {
401        type Error = ModelError;
402
403        fn try_from(r: TaskSpecRaw) -> Result<Self, Self::Error> {
404            let max_retries = match r.max_retries {
405                None => None,
406                Some(0) => {
407                    return Err(ModelError::Invalid(
408                        "maxRetries: 0 is not allowed; omit the field for an unlimited budget"
409                            .into(),
410                    ));
411                }
412                Some(n) => NonZeroU32::new(n),
413            };
414
415            let spec = Self {
416                runner_selector: r.runner_selector,
417
418                workload: r.workload,
419                slot: r.slot,
420
421                restart: r.restart,
422                backoff: r.backoff,
423                timeout: r.timeout,
424
425                admission: r.admission,
426                max_retries,
427            };
428            spec.validate_structural()?;
429            Ok(spec)
430        }
431    }
432}
433
434#[cfg(test)]
435mod tests {
436    use super::*;
437    use crate::{EmbeddedSpec, Flag, SubprocessMode, SubprocessSpec, TaskEnv};
438
439    fn embedded() -> TaskWorkload {
440        TaskWorkload::Embedded(EmbeddedSpec::new("test-v1").unwrap())
441    }
442
443    fn valid_spec() -> TaskSpec {
444        TaskSpec::builder(
445            "test",
446            TaskWorkload::Subprocess(SubprocessSpec {
447                mode: SubprocessMode::Command {
448                    command: "echo".into(),
449                    args: vec![],
450                },
451                env: TaskEnv::default(),
452                cwd: None,
453                fail_on_non_zero: Flag::enabled(),
454            }),
455            5_000u64,
456        )
457        .build()
458        .expect("test spec must be valid")
459    }
460
461    #[test]
462    fn builder_accepts_valid_specs_and_rejects_required_field_errors() {
463        valid_spec().validate().unwrap();
464
465        for (slot, timeout, field) in [("", 5_000_u64, "slot"), ("test", 0_u64, "timeout")] {
466            let error = TaskSpec::builder(slot, embedded(), timeout)
467                .build()
468                .unwrap_err();
469            assert!(error.to_string().contains(field), "got: {error}");
470        }
471    }
472
473    #[test]
474    fn embedded_workload_is_structurally_valid() {
475        let spec = TaskSpec::builder("test", embedded(), 5_000u64)
476            .build()
477            .expect("Embedded is structurally valid");
478        assert!(matches!(spec.workload(), TaskWorkload::Embedded(_)));
479        spec.validate().unwrap();
480    }
481
482    #[test]
483    fn builder_and_override_methods_expose_expected_values() {
484        let spec = TaskSpec::builder("my-slot", embedded(), 10_000u64)
485            .restart(RestartPolicy::OnFailure)
486            .admission(AdmissionPolicy::Replace)
487            .build()
488            .unwrap();
489
490        assert_eq!(spec.slot(), "my-slot");
491        assert_eq!(spec.timeout().as_millis(), 10_000);
492        assert_eq!(spec.restart(), RestartPolicy::OnFailure);
493        assert_eq!(spec.admission(), AdmissionPolicy::Replace);
494        assert_eq!(
495            valid_spec()
496                .with_admission(AdmissionPolicy::Replace)
497                .admission(),
498            AdmissionPolicy::Replace
499        );
500    }
501
502    #[test]
503    fn serde_roundtrip_and_unlimited_retry_shape_are_stable() {
504        let spec = valid_spec();
505        let json = serde_json::to_string(&spec).unwrap();
506        let back: TaskSpec = serde_json::from_str(&json).unwrap();
507        assert_eq!(back, spec);
508        let json = serde_json::to_value(valid_spec()).unwrap();
509        assert!(
510            json.get("maxRetries").is_none(),
511            "unlimited budget must serialize as an absent field"
512        );
513    }
514
515    #[test]
516    fn serde_validates_fields_and_rejects_unknown_fields() {
517        for (field, value, expected) in [
518            ("slot", serde_json::json!(""), "slot"),
519            ("timeout", serde_json::json!(0), "timeout"),
520            ("maxRetries", serde_json::json!(0), "maxRetries"),
521        ] {
522            let mut json = serde_json::to_value(valid_spec()).unwrap();
523            json[field] = value;
524            let error = serde_json::from_value::<TaskSpec>(json).unwrap_err();
525            assert!(error.to_string().contains(expected), "got: {error}");
526        }
527
528        let mut json = serde_json::to_value(valid_spec()).unwrap();
529        json["unexpected"] = serde_json::json!(true);
530        assert!(serde_json::from_value::<TaskSpec>(json).is_err());
531    }
532
533    #[test]
534    fn finite_retry_budget_roundtrips_through_json() {
535        let spec = valid_spec();
536        let mut json: serde_json::Value = serde_json::to_value(&spec).unwrap();
537        json["maxRetries"] = serde_json::json!(3);
538
539        let back: TaskSpec = serde_json::from_value(json).unwrap();
540        assert_eq!(back.max_retries().map(NonZeroU32::get), Some(3));
541    }
542}