Skip to main content

isb_apps/jobs/
mod.rs

1//! Scheduled jobs: a command run on a cron schedule against an app or a
2//! stack service of an org.
3//!
4//! Two ways to run:
5//! - `exec` (default): in a running replica of the service, like
6//!   `isb exec` (the instance's environment, secrets included).
7//! - `run`: in a fresh one-off instance made from the service's deployed
8//!   image, environment and secrets (no published ports, no volumes, no
9//!   health check), deleted afterwards.
10//!
11//! Each run has a timeout and a record (status, exit code, duration, a
12//! bounded log) kept under the state directory, newest `keep`. A run that
13//! comes due while the previous one is still going is skipped (`concurrency:
14//! skip`, the default) or started anyway (`allow`). A finished run emits
15//! `job.succeeded` or `job.failed`.
16//!
17//! ```text
18//! <org root>/jobs/<name>/job.json
19//! <org root>/jobs/<name>/runs/<id>.json, <id>.log
20//! ```
21
22pub mod runs;
23pub mod scheduler;
24
25use std::collections::BTreeMap;
26use std::path::{Path, PathBuf};
27use std::sync::{Arc, Mutex};
28use std::time::Duration;
29
30use serde::{Deserialize, Serialize};
31use serde_json::{Value, json};
32
33use crate::app::Apps;
34use crate::cron::Schedule;
35use crate::error::{Error, Result};
36use crate::exec::{ExecEvent, ExecOptions};
37use crate::org::OrgId;
38use crate::sandbox::Sandbox;
39pub use runs::{Run, RunLog, RunStatus, RunStore, RunTrigger};
40pub use scheduler::{Entry, Scheduled, Scheduler};
41
42/// The longest a job may run.
43pub const MAX_TIMEOUT: Duration = Duration::from_secs(24 * 3600);
44
45/// What a job runs against.
46#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
47#[serde(deny_unknown_fields)]
48pub struct Target {
49    /// An app (its service in its project environment's stack).
50    #[serde(default, skip_serializing_if = "Option::is_none")]
51    pub app: Option<String>,
52    /// Or a stack and one of its services.
53    #[serde(default, skip_serializing_if = "Option::is_none")]
54    pub stack: Option<String>,
55    #[serde(default, skip_serializing_if = "Option::is_none")]
56    pub service: Option<String>,
57}
58
59#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)]
60#[serde(rename_all = "lowercase")]
61pub enum Mode {
62    /// In a running replica.
63    #[default]
64    Exec,
65    /// In a fresh one-off instance from the service's image.
66    Run,
67}
68
69#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)]
70#[serde(rename_all = "lowercase")]
71pub enum Concurrency {
72    /// A run that comes due while one is going is skipped (recorded).
73    #[default]
74    Skip,
75    /// Runs may overlap.
76    Allow,
77}
78
79fn yes() -> bool {
80    true
81}
82
83fn default_timeout() -> String {
84    "10m".into()
85}
86
87fn default_keep() -> u32 {
88    20
89}
90
91/// What a user sets on a job.
92#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
93#[serde(deny_unknown_fields)]
94pub struct JobSpec {
95    pub name: String,
96    /// Five cron fields or an alias (`@hourly`, `@daily`, ...).
97    pub schedule: String,
98    /// `UTC` (default) or a fixed offset such as `+02:00`.
99    #[serde(default, skip_serializing_if = "Option::is_none")]
100    pub timezone: Option<String>,
101    pub target: Target,
102    #[serde(default)]
103    pub mode: Mode,
104    /// argv (no shell unless you run one: `[sh, -c, ...]`).
105    pub command: Vec<String>,
106    #[serde(default = "default_timeout")]
107    pub timeout: String,
108    #[serde(default)]
109    pub concurrency: Concurrency,
110    /// Runs kept.
111    #[serde(default = "default_keep")]
112    pub keep: u32,
113    #[serde(default = "yes")]
114    pub enabled: bool,
115    #[serde(default, skip_serializing_if = "Option::is_none")]
116    pub user: Option<String>,
117    #[serde(default, skip_serializing_if = "Option::is_none")]
118    pub cwd: Option<String>,
119    /// Extra variables for the command.
120    #[serde(default, skip_serializing_if = "BTreeMap::is_empty")]
121    pub env: BTreeMap<String, String>,
122    /// How late a slot missed while the daemon was down may still run
123    /// (default 1h; `0s` never runs missed slots).
124    #[serde(default, skip_serializing_if = "Option::is_none")]
125    pub missed_grace: Option<String>,
126}
127
128/// A job as stored.
129#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
130pub struct Job {
131    pub spec: JobSpec,
132    pub created_at: u64,
133    pub updated_at: u64,
134    /// The last slot it fired for (Unix seconds; creation time before).
135    pub anchor: i64,
136}
137
138/// A job name: `[a-z0-9-]`, a path component and part of instance names.
139pub fn validate_name(what: &str, s: &str) -> Result<()> {
140    let ok = !s.is_empty()
141        && s.len() <= 30
142        && s.starts_with(|c: char| c.is_ascii_lowercase())
143        && !s.ends_with('-')
144        && s.chars()
145            .all(|c| c.is_ascii_lowercase() || c.is_ascii_digit() || c == '-');
146    if ok {
147        Ok(())
148    } else {
149        Err(Error::invalid(format!(
150            "{what} name {s:?}: up to 30 characters of [a-z0-9-], starting with a letter"
151        )))
152    }
153}
154
155/// A duration setting, bounded.
156pub fn parse_timeout(s: &str) -> Result<Duration> {
157    let d = crate::flex::parse_duration(s).map_err(Error::invalid)?;
158    if d.is_zero() || d > MAX_TIMEOUT {
159        return Err(Error::invalid(format!(
160            "timeout {s:?}: more than 0 and at most 24h"
161        )));
162    }
163    Ok(d)
164}
165
166/// A grace window in seconds (default [`scheduler::DEFAULT_GRACE`]).
167pub fn parse_grace(s: &Option<String>) -> Result<i64> {
168    match s {
169        None => Ok(scheduler::DEFAULT_GRACE),
170        Some(g) => {
171            let d = crate::flex::parse_duration(g).map_err(Error::invalid)?;
172            Ok(d.as_secs().min(31 * 86_400) as i64)
173        }
174    }
175}
176
177impl JobSpec {
178    pub fn schedule(&self) -> Result<Schedule> {
179        let off = crate::cron::parse_offset(self.timezone.as_deref().unwrap_or("UTC"))?;
180        Schedule::parse_in(&self.schedule, off)
181    }
182
183    pub fn validate(&self) -> Result<()> {
184        validate_name("job", &self.name)?;
185        self.schedule()?;
186        match (&self.target.app, &self.target.stack, &self.target.service) {
187            (Some(a), None, None) => crate::app::validate_app_name(a)?,
188            (None, Some(st), Some(_)) => crate::stack::validate_stack_name(st)?,
189            _ => {
190                return Err(Error::invalid(
191                    "target: {app: NAME}, or {stack: NAME, service: NAME}",
192                ));
193            }
194        }
195        if self.command.is_empty() || self.command[0].is_empty() {
196            return Err(Error::invalid("command: argv, at least the program"));
197        }
198        parse_timeout(&self.timeout)?;
199        parse_grace(&self.missed_grace)?;
200        if self.keep == 0 || self.keep > 1000 {
201            return Err(Error::invalid("keep: 1 to 1000 runs"));
202        }
203        Ok(())
204    }
205}
206
207/// Runs in progress, by (org, kind, name): the concurrency policy's view.
208#[derive(Default)]
209pub struct Running {
210    set: Mutex<BTreeMap<(OrgId, String, String), usize>>,
211}
212
213/// Held while a run goes; dropping it marks the run over.
214pub struct RunGuard {
215    r: Arc<Running>,
216    key: (OrgId, String, String),
217}
218
219impl Running {
220    /// Start a run of `kind`/`name`: `None` when one is going and `skip`.
221    pub fn enter(
222        self: &Arc<Self>,
223        org: &OrgId,
224        kind: &str,
225        name: &str,
226        skip: bool,
227    ) -> Option<RunGuard> {
228        let key = (org.clone(), kind.to_string(), name.to_string());
229        let mut s = self.set.lock().unwrap();
230        let n = s.entry(key.clone()).or_default();
231        if *n > 0 && skip {
232            return None;
233        }
234        *n += 1;
235        Some(RunGuard {
236            r: self.clone(),
237            key,
238        })
239    }
240
241    pub fn is_running(&self, org: &OrgId, kind: &str, name: &str) -> bool {
242        self.set
243            .lock()
244            .unwrap()
245            .get(&(org.clone(), kind.to_string(), name.to_string()))
246            .is_some_and(|n| *n > 0)
247    }
248}
249
250impl Drop for RunGuard {
251    fn drop(&mut self) {
252        let mut s = self.r.set.lock().unwrap();
253        if let Some(n) = s.get_mut(&self.key) {
254            *n = n.saturating_sub(1);
255            if *n == 0 {
256                s.remove(&self.key);
257            }
258        }
259    }
260}
261
262struct Inner {
263    state: PathBuf,
264    apps: Apps,
265    running: Arc<Running>,
266    edit: Mutex<()>,
267    scheduler: Mutex<Scheduler>,
268}
269
270/// Every org's jobs.
271#[derive(Clone)]
272pub struct Jobs {
273    inner: Arc<Inner>,
274}
275
276/// Every org with a state directory (the default org always).
277pub fn orgs(state: &Path) -> Vec<OrgId> {
278    let mut out = vec![OrgId::default_org()];
279    if let Ok(rd) = std::fs::read_dir(state.join("orgs")) {
280        for e in rd.flatten() {
281            if let Some(o) = e.file_name().to_str().and_then(|s| OrgId::new(s).ok()) {
282                if !o.is_default() {
283                    out.push(o);
284                }
285            }
286        }
287    }
288    out
289}
290
291/// The stack and service a target names, qualified by org.
292pub fn resolve_target(apps: &Apps, org: &OrgId, t: &Target) -> Result<(String, String)> {
293    match (&t.app, &t.stack, &t.service) {
294        (Some(a), _, _) => {
295            let app = apps.get(org, a)?;
296            Ok((crate::stack::qualified(org, &app.spec.stack()?), a.clone()))
297        }
298        (None, Some(st), Some(svc)) => Ok((crate::stack::qualified(org, st), svc.clone())),
299        _ => Err(Error::invalid("target: an app, or a stack and a service")),
300    }
301}
302
303/// The name of a running replica of a stack's service (the lowest slot).
304pub fn running_instance(
305    client: &crate::client::Client,
306    org: &OrgId,
307    stack: &str,
308    service: &str,
309) -> Result<String> {
310    let oc = crate::org::client(client, org);
311    let name = stack.rsplit('/').next().unwrap_or(stack);
312    let insts = crate::stack::controller::list_instances(&oc, name, Some(service))?;
313    insts
314        .iter()
315        .find(|i| i.is_running())
316        .map(|i| i.name.clone())
317        .ok_or_else(|| {
318            Error::invalid(format!(
319                "service {service} of stack {name} has no running replica"
320            ))
321        })
322}
323
324/// Run argv in `instance`, streaming output into `log`. Returns the exit
325/// code.
326pub fn exec_logged(
327    client: &crate::client::Client,
328    org: &OrgId,
329    instance: &str,
330    argv: &[String],
331    opts: ExecOptions,
332    log: &mut RunLog,
333) -> Result<i32> {
334    let oc = crate::org::client(client, org);
335    let sb = Sandbox::get(&oc, instance)?;
336    let timeout = opts.timeout;
337    let mut s = sb.exec_stream(argv.to_vec(), opts)?;
338    while let Some(ev) = s.next_event() {
339        match ev {
340            ExecEvent::Stdout(b) | ExecEvent::Stderr(b) => log.write(&b),
341        }
342    }
343    match s.wait() {
344        Err(Error::ExecTimeout { .. }) => Err(Error::invalid(format!(
345            "timed out after {} and was killed",
346            timeout.map(|t| format!("{t:?}")).unwrap_or_default()
347        ))),
348        r => r,
349    }
350}
351
352impl Jobs {
353    pub fn new(state: &Path, apps: Apps) -> Jobs {
354        let j = Jobs {
355            inner: Arc::new(Inner {
356                state: state.to_path_buf(),
357                apps,
358                running: Arc::default(),
359                edit: Mutex::new(()),
360                scheduler: Mutex::new(Scheduler::idle()),
361            }),
362        };
363        for org in orgs(state) {
364            for job in j.list(&org).unwrap_or_default() {
365                j.runs(&org, &job.spec.name).recover();
366            }
367        }
368        j
369    }
370
371    /// The scheduler to wake when schedules change.
372    pub fn set_scheduler(&self, s: Scheduler) {
373        *self.inner.scheduler.lock().unwrap() = s;
374    }
375
376    fn dir(&self, org: &OrgId) -> PathBuf {
377        crate::app::org_root(&self.inner.state, org).join("jobs")
378    }
379
380    fn job_path(&self, org: &OrgId, name: &str) -> PathBuf {
381        self.dir(org).join(name).join("job.json")
382    }
383
384    pub fn runs(&self, org: &OrgId, name: &str) -> RunStore {
385        RunStore::new(self.dir(org).join(name).join("runs"))
386    }
387
388    pub fn get(&self, org: &OrgId, name: &str) -> Result<Job> {
389        validate_name("job", name)?;
390        match std::fs::read(self.job_path(org, name)) {
391            Ok(b) => Ok(serde_json::from_slice(&b)?),
392            Err(e) if e.kind() == std::io::ErrorKind::NotFound => {
393                Err(Error::NotFound(format!("job {name} in org {org}")))
394            }
395            Err(e) => Err(e.into()),
396        }
397    }
398
399    pub fn list(&self, org: &OrgId) -> Result<Vec<Job>> {
400        let mut out = Vec::new();
401        let Ok(rd) = std::fs::read_dir(self.dir(org)) else {
402            return Ok(out);
403        };
404        for e in rd.flatten() {
405            let p = e.path().join("job.json");
406            if !p.is_file() {
407                continue;
408            }
409            match serde_json::from_slice::<Job>(&std::fs::read(&p)?) {
410                Ok(j) => out.push(j),
411                Err(e) => eprintln!("isb serve: skipping {}: {e}", p.display()),
412            }
413        }
414        out.sort_by(|a, b| a.spec.name.cmp(&b.spec.name));
415        Ok(out)
416    }
417
418    fn save(&self, org: &OrgId, j: &Job) -> Result<()> {
419        crate::app::write_atomic(
420            &self.job_path(org, &j.spec.name),
421            &serde_json::to_vec_pretty(j)?,
422        )
423    }
424
425    /// The target must exist now (a typo is better caught at create).
426    fn check_target(&self, org: &OrgId, spec: &JobSpec) -> Result<()> {
427        let (stack, service) = resolve_target(&self.inner.apps, org, &spec.target)?;
428        if spec.target.stack.is_some() {
429            self.inner
430                .apps
431                .controller()
432                .definition(&stack)?
433                .service(&service)?;
434        }
435        Ok(())
436    }
437
438    pub fn create(&self, org: &OrgId, spec: JobSpec) -> Result<Job> {
439        spec.validate()?;
440        self.check_target(org, &spec)?;
441        let _g = self.inner.edit.lock().unwrap();
442        if self.job_path(org, &spec.name).exists() {
443            return Err(Error::AlreadyExists(format!("job {}", spec.name)));
444        }
445        let now = crate::stack::now_secs();
446        let j = Job {
447            spec,
448            created_at: now,
449            updated_at: now,
450            anchor: now as i64,
451        };
452        self.save(org, &j)?;
453        self.inner.scheduler.lock().unwrap().wake();
454        Ok(j)
455    }
456
457    /// A JSON merge patch of the spec (the name is fixed). A changed
458    /// schedule counts from now.
459    pub fn update(&self, org: &OrgId, name: &str, patch: &Value) -> Result<Job> {
460        let _g = self.inner.edit.lock().unwrap();
461        let mut j = self.get(org, name)?;
462        let mut v = serde_json::to_value(&j.spec)?;
463        crate::app::merge_patch(&mut v, patch);
464        let spec: JobSpec =
465            serde_json::from_value(v).map_err(|e| Error::invalid(format!("job {name}: {e}")))?;
466        if spec.name != j.spec.name {
467            return Err(Error::invalid("a job's name is fixed"));
468        }
469        spec.validate()?;
470        self.check_target(org, &spec)?;
471        let now = crate::stack::now_secs();
472        if spec.schedule != j.spec.schedule
473            || spec.timezone != j.spec.timezone
474            || (spec.enabled && !j.spec.enabled)
475        {
476            j.anchor = now as i64;
477        }
478        j.spec = spec;
479        j.updated_at = now;
480        self.save(org, &j)?;
481        self.inner.scheduler.lock().unwrap().wake();
482        Ok(j)
483    }
484
485    pub fn delete(&self, org: &OrgId, name: &str) -> Result<()> {
486        let _g = self.inner.edit.lock().unwrap();
487        self.get(org, name)?;
488        if self.inner.running.is_running(org, "job", name) {
489            return Err(Error::invalid(format!(
490                "job {name} is running; delete it once that run finishes"
491            )));
492        }
493        std::fs::remove_dir_all(self.dir(org).join(name))?;
494        Ok(())
495    }
496
497    /// When the job next fires (Unix seconds), if enabled.
498    pub fn next_run(&self, j: &Job) -> Option<i64> {
499        if !j.spec.enabled {
500            return None;
501        }
502        j.spec.schedule().ok()?.next_after(j.anchor)
503    }
504
505    /// Run the job now, in the background. The run's record (or, when one
506    /// is going and the policy is skip, a refusal).
507    pub fn run_now(&self, org: &OrgId, name: &str, by: &str) -> Result<Run> {
508        let j = self.get(org, name)?;
509        self.start(org, j, RunTrigger::Manual, by, None)?
510            .ok_or_else(|| {
511                Error::invalid(format!("job {name} is still running (concurrency: skip)"))
512            })
513    }
514
515    /// Start a run on its own thread. `None`: skipped (recorded as such).
516    fn start(
517        &self,
518        org: &OrgId,
519        j: Job,
520        trigger: RunTrigger,
521        by: &str,
522        slot: Option<i64>,
523    ) -> Result<Option<Run>> {
524        let name = j.spec.name.clone();
525        let store = self.runs(org, &name);
526        let skip = j.spec.concurrency == Concurrency::Skip;
527        let Some(guard) = self.inner.running.enter(org, "job", &name, skip) else {
528            if trigger != RunTrigger::Manual {
529                let (mut r, mut log) =
530                    store.start("job", trigger, by, slot, j.spec.keep as usize)?;
531                r.error = Some("the previous run was still going (concurrency: skip)".into());
532                r.finish(RunStatus::Skipped);
533                store.finish(&mut r, &mut log)?;
534            }
535            return Ok(None);
536        };
537        let (r, log) = store.start("job", trigger, by, slot, j.spec.keep as usize)?;
538        let me = self.clone();
539        let org2 = org.clone();
540        let run = r.clone();
541        std::thread::spawn(move || {
542            let _guard = guard;
543            me.execute(&org2, &j, run, log);
544        });
545        Ok(Some(r))
546    }
547
548    fn execute(&self, org: &OrgId, j: &Job, mut r: Run, mut log: RunLog) {
549        let store = self.runs(org, &j.spec.name);
550        let res = self.attempt(org, j, &mut log);
551        let target = resolve_target(&self.inner.apps, org, &j.spec.target).ok();
552        let (stack, service) = target.unwrap_or_default();
553        let (kind, level, msg) = match res {
554            Ok(0) => {
555                r.exit_code = Some(0);
556                r.finish(RunStatus::Succeeded);
557                (
558                    "job.succeeded",
559                    "info",
560                    format!("job {}: run {} succeeded", j.spec.name, r.id),
561                )
562            }
563            Ok(code) => {
564                r.exit_code = Some(code);
565                r.error = Some(format!("exit code {code}"));
566                r.finish(RunStatus::Failed);
567                (
568                    "job.failed",
569                    "error",
570                    format!(
571                        "job {}: run {} failed with exit code {code}",
572                        j.spec.name, r.id
573                    ),
574                )
575            }
576            Err(e) => {
577                log.line(&format!("isb: {e}"));
578                r.error = Some(e.to_string());
579                r.finish(RunStatus::Failed);
580                (
581                    "job.failed",
582                    "error",
583                    format!("job {}: run {} failed: {e}", j.spec.name, r.id),
584                )
585            }
586        };
587        if let Err(e) = store.finish(&mut r, &mut log) {
588            eprintln!("isb serve: job {}: run {}: {e}", j.spec.name, r.id);
589        }
590        let stack = if stack.is_empty() {
591            crate::stack::qualified(org, "")
592        } else {
593            stack
594        };
595        self.inner
596            .apps
597            .controller()
598            .event(kind, level, &stack, &service, msg);
599    }
600
601    fn attempt(&self, org: &OrgId, j: &Job, log: &mut RunLog) -> Result<i32> {
602        let s = &j.spec;
603        let timeout = parse_timeout(&s.timeout)?;
604        let (stack, service) = resolve_target(&self.inner.apps, org, &s.target)?;
605        let mut opts = ExecOptions::default().timeout(timeout);
606        opts.user = s.user.clone();
607        opts.cwd = s.cwd.clone();
608        opts.env = s.env.clone();
609        let client = self.inner.apps.client().clone();
610        match s.mode {
611            Mode::Exec => {
612                let inst = running_instance(&client, org, &stack, &service)?;
613                log.line(&format!("isb: exec in {inst}: {}", s.command.join(" ")));
614                exec_logged(&client, org, &inst, &s.command, opts, log)
615            }
616            Mode::Run => {
617                let def = self.inner.apps.controller().definition(&stack)?;
618                let (mut spec, _) = one_off_spec(&def, &service, &s.name, timeout)?;
619                let keys: Vec<&str> = spec.env.secrets.values().map(String::as_str).collect();
620                let values = crate::stack::secrets::values(
621                    self.inner.apps.secrets(),
622                    org,
623                    &def.secrets,
624                    keys,
625                )?;
626                let secret_env = crate::supervise::secret_env(&spec, &values)?;
627                spec.env.secrets.clear();
628                spec.env.vars.extend(secret_env);
629                let name = spec.name.clone().unwrap_or_default();
630                log.line(&format!("isb: one-off instance {name} from {}", spec.image));
631                let oc = crate::org::client(&client, org);
632                let base = self.dir(org);
633                std::fs::create_dir_all(&base)?;
634                let made = Sandbox::connect_or_create_with_base(
635                    &oc,
636                    &spec,
637                    &Default::default(),
638                    &base,
639                    crate::sandbox::EnsureOptions {
640                        wait_ready: true,
641                        ..Default::default()
642                    },
643                    &mut |m| log.line(&format!("isb: {m}")),
644                );
645                let r = match made {
646                    Ok((sb, _)) => {
647                        if !spec.secrets.is_empty() {
648                            crate::supervise::push_secrets(&sb, &spec, &values)?;
649                        }
650                        log.line(&format!("isb: run: {}", s.command.join(" ")));
651                        exec_logged(&client, org, &name, &s.command, opts, log)
652                    }
653                    Err(e) => Err(e),
654                };
655                if let Err(e) = Sandbox::remove(&oc, &name, true) {
656                    if !e.is_not_found() {
657                        log.line(&format!("isb: removing {name}: {e}"));
658                    }
659                }
660                r
661            }
662        }
663    }
664}
665
666/// A one-off instance's spec from a deployed service: its image (as its
667/// instances get it), environment, secrets, resources and user; no ports,
668/// volumes, domains, health check or replicas. An OCI image's command is
669/// replaced by `sleep`, so the job's command runs as an exec beside it.
670pub fn one_off_spec(
671    def: &crate::stack::StackDef,
672    service: &str,
673    job: &str,
674    timeout: Duration,
675) -> Result<(crate::spec::SandboxSpec, bool)> {
676    let svc = def.service(service)?;
677    let image = def.instance_image(service, &svc.image);
678    let oci = crate::plan::ImageSource::parse(&image)?.is_oci();
679    let mut v = json!({
680        "image": image,
681        "container_name": format!("job-{job}-{}", crate::app::git::random_hex(3)),
682        "labels": {"isb.job": job},
683    });
684    let src = serde_json::to_value(svc)?;
685    for k in [
686        "environment",
687        "secrets",
688        "cpus",
689        "mem_limit",
690        "user",
691        "working_dir",
692        "type",
693    ] {
694        if let Some(x) = src.get(k) {
695            v[k] = x.clone();
696        }
697    }
698    if oci {
699        let secs = timeout.as_secs() + 300;
700        v["entrypoint"] = json!([]);
701        v["command"] = json!(["sleep", secs.to_string()]);
702    }
703    let spec: crate::spec::SandboxSpec = serde_json::from_value(v)
704        .map_err(|e| Error::invalid(format!("one-off instance for {service}: {e}")))?;
705    Ok((spec, oci))
706}
707
708impl Scheduled for Jobs {
709    fn entries(&self) -> Vec<Entry> {
710        let mut out = Vec::new();
711        for org in orgs(&self.inner.state) {
712            for j in self.list(&org).unwrap_or_default() {
713                if !j.spec.enabled {
714                    continue;
715                }
716                let (Ok(schedule), Ok(grace)) =
717                    (j.spec.schedule(), parse_grace(&j.spec.missed_grace))
718                else {
719                    continue;
720                };
721                out.push(Entry {
722                    org: org.clone(),
723                    name: j.spec.name.clone(),
724                    schedule,
725                    anchor: j.anchor,
726                    grace,
727                });
728            }
729        }
730        out
731    }
732
733    fn fire(&self, e: &Entry, slot: i64, late: bool) {
734        let j = {
735            let _g = self.inner.edit.lock().unwrap();
736            let Ok(mut j) = self.get(&e.org, &e.name) else {
737                return;
738            };
739            if j.anchor >= slot {
740                return;
741            }
742            j.anchor = slot;
743            if let Err(err) = self.save(&e.org, &j) {
744                eprintln!("isb serve: job {}: {err}", e.name);
745                return;
746            }
747            j
748        };
749        let trigger = if late {
750            RunTrigger::Missed
751        } else {
752            RunTrigger::Schedule
753        };
754        if let Err(err) = self.start(&e.org, j, trigger, "schedule", Some(slot)) {
755            eprintln!("isb serve: job {}: {err}", e.name);
756        }
757    }
758
759    fn advance(&self, e: &Entry, to: i64) {
760        let _g = self.inner.edit.lock().unwrap();
761        if let Ok(mut j) = self.get(&e.org, &e.name) {
762            j.anchor = to;
763            let _ = self.save(&e.org, &j);
764        }
765    }
766}
767
768#[cfg(test)]
769mod tests {
770    use super::*;
771
772    fn spec(v: Value) -> JobSpec {
773        serde_json::from_value(v).unwrap()
774    }
775
776    #[test]
777    fn spec_validation() {
778        let ok = spec(json!({
779            "name": "nightly", "schedule": "0 3 * * *",
780            "target": {"app": "web"}, "command": ["sh", "-c", "echo hi"],
781        }));
782        ok.validate().unwrap();
783        assert_eq!(ok.mode, Mode::Exec);
784        assert_eq!(ok.concurrency, Concurrency::Skip);
785        assert_eq!(ok.keep, 20);
786        for bad in [
787            json!({"schedule": "61 * * * *"}),
788            json!({"target": {"app": null}}),
789            json!({"target": {"app": "web", "stack": "s", "service": "x"}}),
790            json!({"target": {"app": null, "stack": "s"}}),
791            json!({"command": []}),
792            json!({"timeout": "0s"}),
793            json!({"timeout": "48h"}),
794            json!({"timezone": "Europe/Paris"}),
795            json!({"keep": 0}),
796            json!({"name": "Bad_Name"}),
797        ] {
798            let mut v = serde_json::to_value(&ok).unwrap();
799            crate::app::merge_patch(&mut v, &bad);
800            let s: JobSpec = serde_json::from_value(v).unwrap();
801            assert!(s.validate().is_err(), "{bad}");
802        }
803        assert!(serde_json::from_value::<JobSpec>(json!({"name": "x", "schedule": "@daily", "target": {"app": "a"}, "command": ["x"], "bogus": 1})).is_err());
804        let st = spec(json!({
805            "name": "x", "schedule": "@hourly", "timezone": "+02:00",
806            "target": {"stack": "shop-production", "service": "web"}, "command": ["true"],
807            "mode": "run", "concurrency": "allow",
808        }));
809        st.validate().unwrap();
810        assert_eq!(st.schedule().unwrap().as_str(), "@hourly");
811        assert_eq!(parse_grace(&None).unwrap(), 3600);
812        assert_eq!(parse_grace(&Some("0s".into())).unwrap(), 0);
813    }
814
815    #[test]
816    fn concurrency_policy() {
817        let r = Arc::new(Running::default());
818        let org = OrgId::default_org();
819        let g1 = r.enter(&org, "job", "a", true).expect("first run starts");
820        assert!(r.is_running(&org, "job", "a"));
821        assert!(
822            r.enter(&org, "job", "a", true).is_none(),
823            "skip while running"
824        );
825        assert!(
826            r.enter(&org, "job", "b", true).is_some(),
827            "other jobs are not held"
828        );
829        assert!(
830            r.enter(&org, "backup", "a", true).is_some(),
831            "other kinds neither"
832        );
833        {
834            let _g2 = r.enter(&org, "job", "a", false).expect("allow overlaps");
835        }
836        assert!(r.is_running(&org, "job", "a"), "the first still runs");
837        drop(g1);
838        assert!(!r.is_running(&org, "job", "a"));
839        assert!(r.enter(&org, "job", "a", true).is_some());
840    }
841
842    #[test]
843    fn one_off_from_a_deployed_service() {
844        let file: crate::spec::ComposeFile = serde_yaml_ng::from_str(
845            "services:\n  web:\n    image: docker:traefik/whoami@sha256:abc\n    environment: {A: '1', T: {secret: web.tok}}\n    ports: ['127.0.0.1:8080:80']\n    volumes: ['d:/data']\n    deploy: {replicas: 3}\n    healthcheck: {test: [CMD, /x]}\n    mem_limit: 256m\nvolumes: {d: {}}\nsecrets: {web.tok: {external: true, name: tok}}\n",
846        )
847        .unwrap();
848        let def = crate::stack::StackDef {
849            name: "shop-production".into(),
850            org: OrgId::default_org(),
851            file,
852            base_dir: "/".into(),
853            secrets: Default::default(),
854            force: Default::default(),
855            images: Default::default(),
856            deployed_at: 0,
857            deployed_by: String::new(),
858            previous: None,
859        };
860        let (s, oci) = one_off_spec(&def, "web", "nightly", Duration::from_secs(60)).unwrap();
861        assert!(oci);
862        assert!(s.name.as_deref().unwrap().starts_with("job-nightly-"));
863        assert_eq!(s.image, "docker:traefik/whoami@sha256:abc");
864        assert_eq!(s.env["A"], "1");
865        assert_eq!(s.env.secrets["T"], "web.tok");
866        assert!(s.ports.is_empty() && s.volumes.is_empty() && s.healthcheck.is_none());
867        assert_eq!(s.replicas(), 1);
868        assert_eq!(s.memory.as_deref(), Some("256m"));
869        assert_eq!(s.command.as_deref().unwrap()[0], "sleep");
870        assert_eq!(s.labels["isb.job"], "nightly");
871        assert!(!s.labels.contains_key("isb.stack"));
872        assert!(one_off_spec(&def, "nope", "j", Duration::from_secs(1)).is_err());
873    }
874}