Skip to main content

wasmer_config/app/
job.rs

1use std::{borrow::Cow, fmt::Display, str::FromStr};
2
3use anyhow::anyhow;
4use serde::{Deserialize, Serialize, de::Error};
5
6use indexmap::IndexMap;
7
8use crate::package::PackageSource;
9
10use super::{AppConfigCapabilityMemoryV1, AppVolume, HttpRequest, pretty_duration::PrettyDuration};
11
12/// Job configuration.
13#[derive(
14    serde::Serialize, serde::Deserialize, schemars::JsonSchema, Clone, Debug, PartialEq, Eq,
15)]
16pub struct Job {
17    pub name: String,
18    pub trigger: JobTrigger,
19
20    #[serde(skip_serializing_if = "Option::is_none")]
21    pub timeout: Option<PrettyDuration>,
22
23    /// Don't start job if past the due time by this amount,
24    /// instead opting to wait for the next instance of it
25    /// to be triggered.
26    #[serde(skip_serializing_if = "Option::is_none")]
27    pub max_schedule_drift: Option<PrettyDuration>,
28
29    #[serde(skip_serializing_if = "Option::is_none")]
30    pub retries: Option<u32>,
31
32    /// Maximum percent of "jitter" to introduce between invocations.
33    ///
34    /// Value range: 0-100
35    ///
36    /// Jitter is used to spread out jobs over time.
37    /// The calculation works by multiplying the time between invocations
38    /// by a random amount, and taking the percentage of that random amount.
39    ///
40    /// See also [`Self::jitter_percent_min`] to set a minimum jitter.
41    #[serde(skip_serializing_if = "Option::is_none")]
42    pub jitter_percent_max: Option<u8>,
43
44    /// Minimum "jitter" to introduce between invocations.
45    ///
46    /// Value range: 0-100
47    ///
48    /// Jitter is used to spread out jobs over time.
49    /// The calculation works by multiplying the time between invocations
50    /// by a random amount, and taking the percentage of that random amount.
51    ///
52    /// If not specified while `jitter_percent_max` is, it will default to 10%.
53    ///
54    /// See also [`Self::jitter_percent_max`] to set a maximum jitter.
55    #[serde(skip_serializing_if = "Option::is_none")]
56    pub jitter_percent_min: Option<u8>,
57
58    pub action: JobAction,
59
60    /// Additional unknown fields.
61    ///
62    /// Exists for forward compatibility for newly added fields.
63    #[serde(flatten)]
64    pub other: IndexMap<String, serde_json::Value>,
65}
66
67// We need this wrapper struct to enable this formatting:
68// job:
69//   action:
70//     execute: ...
71#[derive(serde::Serialize, schemars::JsonSchema, Clone, Debug, PartialEq, Eq)]
72pub struct JobAction {
73    #[serde(flatten)]
74    pub action: JobActionCase,
75}
76
77impl From<JobActionCase> for JobAction {
78    fn from(action: JobActionCase) -> Self {
79        Self { action }
80    }
81}
82
83impl<'de> Deserialize<'de> for JobAction {
84    fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
85    where
86        D: serde::Deserializer<'de>,
87    {
88        #[derive(Deserialize)]
89        #[serde(deny_unknown_fields)]
90        struct Repr {
91            #[serde(default)]
92            fetch: Option<HttpRequest>,
93            #[serde(default)]
94            execute: Option<ExecutableJob>,
95        }
96
97        match Repr::deserialize(deserializer)? {
98            Repr {
99                fetch: Some(_),
100                execute: Some(_),
101            } => Err(D::Error::custom(
102                "job action must set exactly one of 'fetch' or 'execute', not both",
103            )),
104            Repr {
105                fetch: Some(fetch),
106                execute: None,
107            } => Ok(JobActionCase::Fetch(fetch).into()),
108            Repr {
109                fetch: None,
110                execute: Some(execute),
111            } => Ok(JobActionCase::Execute(execute).into()),
112            Repr {
113                fetch: None,
114                execute: None,
115            } => Err(D::Error::custom(
116                "job action must set one of 'fetch' or 'execute'",
117            )),
118        }
119    }
120}
121
122#[derive(
123    serde::Serialize, serde::Deserialize, schemars::JsonSchema, Clone, Debug, PartialEq, Eq,
124)]
125#[serde(rename_all = "lowercase")]
126pub enum JobActionCase {
127    Fetch(HttpRequest),
128    Execute(ExecutableJob),
129}
130
131#[derive(Clone, Debug, PartialEq, Eq)]
132pub struct CronExpression {
133    pub cron: saffron::parse::CronExpr,
134    // Keep the original string form around for serialization purposes.
135    pub parsed_from: String,
136}
137
138#[derive(Clone, Debug, PartialEq, Eq)]
139pub enum JobTrigger {
140    PreDeployment,
141    PostDeployment,
142    Cron(CronExpression),
143    Duration(PrettyDuration),
144}
145
146#[derive(
147    serde::Serialize, serde::Deserialize, schemars::JsonSchema, Clone, Debug, PartialEq, Eq,
148)]
149pub struct ExecutableJob {
150    /// The package that contains the command to run. Defaults to the app config's package.
151    #[serde(skip_serializing_if = "Option::is_none")]
152    pub package: Option<PackageSource>,
153
154    /// The command to run. Defaults to the package's entrypoint.
155    #[serde(skip_serializing_if = "Option::is_none")]
156    pub command: Option<String>,
157
158    /// CLI arguments passed to the runner.
159    /// Only applicable for runners that accept CLI arguments.
160    #[serde(skip_serializing_if = "Option::is_none")]
161    pub cli_args: Option<Vec<String>>,
162
163    /// Environment variables.
164    #[serde(default, skip_serializing_if = "Option::is_none")]
165    pub env: Option<IndexMap<String, String>>,
166
167    #[serde(skip_serializing_if = "Option::is_none")]
168    pub capabilities: Option<ExecutableJobCompatibilityMapV1>,
169
170    #[serde(skip_serializing_if = "Option::is_none")]
171    pub volumes: Option<Vec<AppVolume>>,
172}
173
174#[derive(
175    serde::Serialize, serde::Deserialize, schemars::JsonSchema, Clone, Debug, PartialEq, Eq,
176)]
177pub struct ExecutableJobCompatibilityMapV1 {
178    /// Instance memory settings.
179    #[serde(skip_serializing_if = "Option::is_none")]
180    pub memory: Option<AppConfigCapabilityMemoryV1>,
181
182    /// Additional unknown capabilities.
183    ///
184    /// This provides a small bit of forwards compatibility for newly added
185    /// capabilities.
186    #[serde(flatten)]
187    pub other: IndexMap<String, serde_json::Value>,
188}
189
190impl Serialize for JobTrigger {
191    fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
192    where
193        S: serde::Serializer,
194    {
195        self.to_string().serialize(serializer)
196    }
197}
198
199impl<'de> Deserialize<'de> for JobTrigger {
200    fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
201    where
202        D: serde::Deserializer<'de>,
203    {
204        let repr: Cow<'de, str> = Cow::deserialize(deserializer)?;
205        repr.parse().map_err(D::Error::custom)
206    }
207}
208
209impl Display for JobTrigger {
210    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
211        match self {
212            Self::PreDeployment => write!(f, "pre-deployment"),
213            Self::PostDeployment => write!(f, "post-deployment"),
214            Self::Cron(cron) => write!(f, "{}", cron.parsed_from),
215            Self::Duration(duration) => write!(f, "{duration}"),
216        }
217    }
218}
219
220impl FromStr for JobTrigger {
221    type Err = anyhow::Error;
222
223    fn from_str(s: &str) -> Result<Self, Self::Err> {
224        if s == "pre-deployment" {
225            Ok(Self::PreDeployment)
226        } else if s == "post-deployment" {
227            Ok(Self::PostDeployment)
228        } else {
229            match s.parse::<CronExpression>() {
230                Ok(expr) => Ok(Self::Cron(expr)),
231                _ => {
232                    if let Ok(duration) = s.parse::<PrettyDuration>() {
233                        Ok(Self::Duration(duration))
234                    } else {
235                        Err(anyhow!(
236                            "Invalid job trigger '{s}'. Must be 'pre-deployment', 'post-deployment', \
237                a valid cron expression such as '0 */5 * * *' or a duration such as '15m'.",
238                        ))
239                    }
240                }
241            }
242        }
243    }
244}
245
246impl FromStr for CronExpression {
247    type Err = Box<dyn std::error::Error + Send + Sync>;
248
249    fn from_str(s: &str) -> Result<Self, Self::Err> {
250        if let Some(predefined_sched) = s.strip_prefix('@') {
251            match predefined_sched {
252                "hourly" => Ok(Self {
253                    cron: "0 * * * *".parse().unwrap(),
254                    parsed_from: s.to_owned(),
255                }),
256                "daily" => Ok(Self {
257                    cron: "0 0 * * *".parse().unwrap(),
258                    parsed_from: s.to_owned(),
259                }),
260                "weekly" => Ok(Self {
261                    cron: "0 0 * * 1".parse().unwrap(),
262                    parsed_from: s.to_owned(),
263                }),
264                "monthly" => Ok(Self {
265                    cron: "0 0 1 * *".parse().unwrap(),
266                    parsed_from: s.to_owned(),
267                }),
268                "yearly" => Ok(Self {
269                    cron: "0 0 1 1 *".parse().unwrap(),
270                    parsed_from: s.to_owned(),
271                }),
272                _ => Err(format!("Invalid cron expression {s}").into()),
273            }
274        } else {
275            // Let's make sure the input string is valid...
276            match s.parse() {
277                Ok(expr) => Ok(Self {
278                    cron: expr,
279                    parsed_from: s.to_owned(),
280                }),
281                Err(_) => Err(format!("Invalid cron expression {s}").into()),
282            }
283        }
284    }
285}
286
287impl schemars::JsonSchema for JobTrigger {
288    fn schema_id() -> Cow<'static, str> {
289        Cow::Borrowed("JobTrigger")
290    }
291
292    fn json_schema(generator: &mut schemars::SchemaGenerator) -> schemars::Schema {
293        String::json_schema(generator)
294    }
295
296    fn schema_name() -> Cow<'static, str> {
297        Self::schema_id()
298    }
299}
300
301#[cfg(test)]
302mod tests {
303    use std::time::Duration;
304
305    use super::*;
306
307    #[test]
308    fn job_action_rejects_anything_but_exactly_one_action() {
309        for yaml in ["execute: {}\nfetch:\n  path: /", "{}", "sleep: {}"] {
310            assert!(
311                serde_yaml::from_str::<JobAction>(yaml).is_err(),
312                "accepted invalid action {yaml:?}"
313            );
314            let json: serde_json::Value = serde_yaml::from_str(yaml).unwrap();
315            assert!(serde_json::from_value::<JobAction>(json).is_err());
316        }
317    }
318
319    #[test]
320    fn job_action_serialization_roundtrip() {
321        for yaml in [
322            "execute:\n  command: php",
323            "fetch:\n  path: /\n  timeout: 30s",
324        ] {
325            let action: JobAction = serde_yaml::from_str(yaml).unwrap();
326            assert_eq!(serde_yaml::to_string(&action).unwrap().trim(), yaml);
327
328            let json: serde_json::Value = serde_yaml::from_str(yaml).unwrap();
329            assert_eq!(
330                serde_json::from_value::<JobAction>(json.clone()).unwrap(),
331                action
332            );
333            assert_eq!(serde_json::to_value(action).unwrap(), json);
334        }
335    }
336
337    #[test]
338    pub fn job_trigger_serialization_roundtrip() {
339        fn assert_roundtrip(
340            serialized: &str,
341            description: Option<&str>,
342            duration: Option<Duration>,
343        ) {
344            let parsed = serialized.parse::<JobTrigger>().unwrap();
345            assert_eq!(&parsed.to_string(), serialized);
346
347            if let JobTrigger::Cron(expr) = &parsed {
348                assert_eq!(
349                    &expr
350                        .cron
351                        .describe(saffron::parse::English::default())
352                        .to_string(),
353                    description.unwrap()
354                );
355            } else {
356                assert!(description.is_none());
357            }
358
359            if let JobTrigger::Duration(d) = &parsed {
360                assert_eq!(d.as_duration(), duration.unwrap());
361            } else {
362                assert!(duration.is_none());
363            }
364        }
365
366        assert_roundtrip("pre-deployment", None, None);
367        assert_roundtrip("post-deployment", None, None);
368
369        assert_roundtrip("@hourly", Some("Every hour"), None);
370        assert_roundtrip("@daily", Some("At 12:00 AM"), None);
371        assert_roundtrip("@weekly", Some("At 12:00 AM on Sunday"), None);
372        assert_roundtrip(
373            "@monthly",
374            Some("At 12:00 AM on the 1st of every month"),
375            None,
376        );
377        assert_roundtrip("@yearly", Some("At 12:00 AM on the 1st of January"), None);
378
379        // Note: the parsing code should keep the formatting of the source string.
380        // This is tested in assert_roundtrip.
381        assert_roundtrip(
382            "0/2 12 * JAN-APR 2",
383            Some(
384                "At every 2nd minute from 0 through 59 minutes past the hour, \
385                between 12:00 PM and 12:59 PM on Monday of January to April",
386            ),
387            None,
388        );
389
390        assert_roundtrip("10s", None, Some(Duration::from_secs(10)));
391        assert_roundtrip("15m", None, Some(Duration::from_secs(15 * 60)));
392        assert_roundtrip("20h", None, Some(Duration::from_secs(20 * 60 * 60)));
393        assert_roundtrip("2d", None, Some(Duration::from_secs(2 * 60 * 60 * 24)));
394    }
395
396    #[test]
397    pub fn job_serialization_roundtrip() {
398        fn parse_cron(expr: &str) -> CronExpression {
399            CronExpression {
400                cron: expr.parse().unwrap(),
401                parsed_from: expr.to_owned(),
402            }
403        }
404
405        let job = Job {
406            name: "my-job".to_owned(),
407            trigger: JobTrigger::Cron(parse_cron("0/2 12 * JAN-APR 2")),
408            timeout: Some("1m".parse().unwrap()),
409            max_schedule_drift: Some("2h".parse().unwrap()),
410            jitter_percent_max: None,
411            jitter_percent_min: None,
412            retries: None,
413            action: JobAction {
414                action: JobActionCase::Execute(super::ExecutableJob {
415                    package: Some(crate::package::PackageSource::Ident(
416                        crate::package::PackageIdent::Named(crate::package::NamedPackageIdent {
417                            registry: None,
418                            namespace: Some("ns".to_owned()),
419                            name: "pkg".to_owned(),
420                            tag: None,
421                        }),
422                    )),
423                    command: Some("cmd".to_owned()),
424                    cli_args: Some(vec!["arg-1".to_owned(), "arg-2".to_owned()]),
425                    env: Some([("VAR1".to_owned(), "Value".to_owned())].into()),
426                    capabilities: Some(super::ExecutableJobCompatibilityMapV1 {
427                        memory: Some(crate::app::AppConfigCapabilityMemoryV1 {
428                            limit: Some(bytesize::ByteSize::gib(1)),
429                        }),
430                        other: Default::default(),
431                    }),
432                    volumes: Some(vec![crate::app::AppVolume {
433                        name: "vol".to_owned(),
434                        mount: "/path/to/volume".to_owned(),
435                    }]),
436                }),
437            },
438            other: Default::default(),
439        };
440
441        let serialized = r#"
442name: my-job
443trigger: 0/2 12 * JAN-APR 2
444timeout: 1m
445max_schedule_drift: 2h
446action:
447  execute:
448    package: ns/pkg
449    command: cmd
450    cli_args:
451    - arg-1
452    - arg-2
453    env:
454      VAR1: Value
455    capabilities:
456      memory:
457        limit: 1.0 GiB
458    volumes:
459    - name: vol
460      mount: /path/to/volume"#;
461
462        assert_eq!(
463            serialized.trim(),
464            serde_yaml::to_string(&job).unwrap().trim()
465        );
466        assert_eq!(job, serde_yaml::from_str(serialized).unwrap());
467    }
468}