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#[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 #[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 #[serde(skip_serializing_if = "Option::is_none")]
42 pub jitter_percent_max: Option<u8>,
43
44 #[serde(skip_serializing_if = "Option::is_none")]
56 pub jitter_percent_min: Option<u8>,
57
58 pub action: JobAction,
59
60 #[serde(flatten)]
64 pub other: IndexMap<String, serde_json::Value>,
65}
66
67#[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 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 #[serde(skip_serializing_if = "Option::is_none")]
152 pub package: Option<PackageSource>,
153
154 #[serde(skip_serializing_if = "Option::is_none")]
156 pub command: Option<String>,
157
158 #[serde(skip_serializing_if = "Option::is_none")]
161 pub cli_args: Option<Vec<String>>,
162
163 #[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 #[serde(skip_serializing_if = "Option::is_none")]
180 pub memory: Option<AppConfigCapabilityMemoryV1>,
181
182 #[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 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 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}