use std::collections::{HashMap, HashSet, VecDeque};
use std::fs;
use std::io::{BufRead, BufReader, Write};
use std::os::unix::net::{UnixListener, UnixStream};
use std::path::{Path, PathBuf};
use std::process::{Command, Stdio};
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::{Arc, Mutex};
use std::thread;
use std::time::{Duration, Instant};
use anyhow::{bail, Context, Result};
use chrono::Local;
use serde::{Deserialize, Serialize};
use crate::config;
use crate::cron::{CronExpr, TickTime, TICK_MS};
use crate::inspect::{collect_scripts, ScriptEntry};
use crate::unifier_events::{self, Wakeup, WakeupMetrics};
use crate::{load_spec, resolve_preferred_spec};
const SERVICE_NAME: &str = "jan-cron";
const SOCKET_FILE: &str = "cron.sock";
const PID_FILE: &str = "cron.pid";
const DISABLED_FILE: &str = "cron-disabled.json";
const DEFAULT_MAX_CONCURRENT: usize = 32;
const DEFAULT_MAX_DEFERRED: usize = 256;
const RECENT_SPAWN_CAP: usize = 32;
const DEFERRED_PEEK: usize = 8;
static STOP_REQUESTED: AtomicBool = AtomicBool::new(false);
fn now_rfc3339() -> String {
Local::now().to_rfc3339_opts(chrono::SecondsFormat::Secs, true)
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum TriggerKind {
Cron,
Mailbox,
Event,
}
impl TriggerKind {
pub fn as_str(self) -> &'static str {
match self {
Self::Cron => "cron",
Self::Mailbox => "mailbox",
Self::Event => "event",
}
}
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct JobView {
pub leaf: String,
pub chain: Vec<String>,
pub trigger: TriggerKind,
#[serde(skip_serializing_if = "Option::is_none")]
pub system: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub wakeup_id: Option<String>,
pub started_at: String,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct SpawnRecord {
pub leaf: String,
pub chain: Vec<String>,
pub trigger: TriggerKind,
#[serde(skip_serializing_if = "Option::is_none")]
pub system: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub wakeup_id: Option<String>,
pub started_at: String,
#[serde(skip_serializing_if = "Option::is_none")]
pub finished_at: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub exit_code: Option<i32>,
#[serde(default)]
pub overlap_skip: bool,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
pub struct WakeupLeaf {
pub leaf: String,
pub chain: Vec<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub system: Option<String>,
pub has_cron: bool,
pub disabled: bool,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct StatusDetail {
pub running: Vec<JobView>,
pub deferred: Vec<JobView>,
pub recent: Vec<SpawnRecord>,
pub wakeups: Vec<WakeupLeaf>,
}
pub(crate) fn stop_requested() -> bool {
STOP_REQUESTED.load(Ordering::SeqCst)
}
pub fn disabled_path() -> PathBuf {
config::config_dir().join(DISABLED_FILE)
}
pub fn load_disabled() -> HashSet<String> {
let path = disabled_path();
let Ok(text) = fs::read_to_string(&path) else {
return HashSet::new();
};
serde_json::from_str::<Vec<String>>(&text)
.unwrap_or_default()
.into_iter()
.map(|s| s.trim().to_string())
.filter(|s| !s.is_empty())
.collect()
}
fn save_disabled(set: &HashSet<String>) -> Result<()> {
let dir = config::config_dir();
fs::create_dir_all(&dir).with_context(|| format!("create {}", dir.display()))?;
let mut names: Vec<_> = set.iter().cloned().collect();
names.sort();
let text = serde_json::to_string_pretty(&names).context("serialize disabled agents")?;
fs::write(disabled_path(), format!("{text}\n"))
.with_context(|| format!("write {}", disabled_path().display()))?;
Ok(())
}
pub fn is_disabled_name(name: &str) -> bool {
load_disabled().contains(name)
}
pub fn is_disabled_chain(chain: &[String]) -> bool {
chain.last().is_some_and(|name| is_disabled_name(name))
}
pub fn runtime_dir() -> PathBuf {
if let Ok(dir) = std::env::var("XDG_RUNTIME_DIR") {
let dir = dir.trim();
if !dir.is_empty() {
return PathBuf::from(dir).join("jan-cli");
}
}
config::config_dir().join("run")
}
pub fn socket_path() -> PathBuf {
runtime_dir().join(SOCKET_FILE)
}
pub fn pid_path() -> PathBuf {
runtime_dir().join(PID_FILE)
}
fn systemd_unit_path() -> PathBuf {
dirs::home_dir()
.unwrap_or_else(|| PathBuf::from("."))
.join(".config/systemd/user")
.join(format!("{SERVICE_NAME}.service"))
}
fn jan_bin() -> Result<PathBuf> {
let exe = std::env::current_exe().context("resolve jan binary path")?;
exe.canonicalize()
.with_context(|| format!("canonicalize {}", exe.display()))
}
fn write_pid_file() -> Result<()> {
let dir = runtime_dir();
fs::create_dir_all(&dir).with_context(|| format!("create {}", dir.display()))?;
fs::write(pid_path(), format!("{}\n", std::process::id()))
.with_context(|| format!("write {}", pid_path().display()))?;
Ok(())
}
fn remove_runtime_files() {
let _ = fs::remove_file(socket_path());
let _ = fs::remove_file(pid_path());
}
fn load_schedule_cache() -> Result<ScheduleCache> {
let (spec_path, identity) = resolve_preferred_spec()?;
let spec = load_spec(&spec_path).with_context(|| format!("load {}", spec_path.display()))?;
let scripts = collect_scripts(&spec);
let mut agents_by_name = HashMap::new();
for s in &scripts {
agents_by_name.insert(
s.name.clone(),
AgentMeta {
chain: s.chain.clone(),
system: s.system.clone(),
},
);
}
let mut schedules = Vec::new();
for s in scripts.into_iter().filter(|s| !s.cron.is_empty()) {
schedules.push(ParsedSchedule::try_from_entry(&s)?);
}
schedules.sort_by(|a, b| a.chain.cmp(&b.chain));
Ok(ScheduleCache {
schedules,
agents_by_name,
spec_dir: identity.spec_dir,
root_yaml: identity.root_yaml,
loaded_at: Instant::now(),
})
}
#[derive(Debug, Clone)]
struct AgentMeta {
chain: Vec<String>,
system: Option<String>,
}
#[derive(Debug, Clone)]
struct ParsedSchedule {
chain: Vec<String>,
system: Option<String>,
cron_raw: Vec<String>,
exprs: Vec<CronExpr>,
}
impl ParsedSchedule {
fn try_from_entry(s: &ScriptEntry) -> Result<Self> {
let mut exprs = Vec::with_capacity(s.cron.len());
for raw in &s.cron {
exprs.push(
CronExpr::parse(raw)
.with_context(|| format!("{}: invalid cron `{raw}`", s.chain.join(" ")))?,
);
}
Ok(Self {
chain: s.chain.clone(),
system: s.system.clone(),
cron_raw: s.cron.clone(),
exprs,
})
}
fn matches_tick(&self, now: &TickTime) -> bool {
self.exprs.iter().any(|e| e.matches_tick(now))
}
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
pub struct CachedCronEntry {
pub chain: Vec<String>,
pub cron: Vec<String>,
#[serde(default)]
pub disabled: bool,
}
impl CachedCronEntry {
pub fn chain_str(&self) -> String {
self.chain.join(" ")
}
}
#[derive(Debug)]
struct ScheduleCache {
schedules: Vec<ParsedSchedule>,
agents_by_name: HashMap<String, AgentMeta>,
spec_dir: String,
root_yaml: String,
loaded_at: Instant,
}
#[derive(Debug)]
struct DaemonState {
cache: ScheduleCache,
last_fired: HashMap<String, TickTime>,
wakeups: Arc<Mutex<VecDeque<Wakeup>>>,
events_connected: Arc<AtomicBool>,
wakeup_metrics: Arc<WakeupMetrics>,
jobs: Arc<Mutex<JobGate>>,
disabled: HashSet<String>,
}
#[derive(Debug, Clone)]
struct PendingSpawn {
chain: Vec<String>,
system: Option<String>,
extra_args: Vec<String>,
env: Vec<(String, String)>,
trigger: TriggerKind,
wakeup_id: Option<String>,
}
impl PendingSpawn {
fn leaf(&self) -> String {
self.chain
.last()
.cloned()
.unwrap_or_else(|| "unknown".into())
}
fn as_job_view(&self, started_at: &str) -> JobView {
JobView {
leaf: self.leaf(),
chain: self.chain.clone(),
trigger: self.trigger,
system: self.system.clone(),
wakeup_id: self.wakeup_id.clone(),
started_at: started_at.to_string(),
}
}
}
#[derive(Debug, Clone)]
struct RunningJob {
leaf: String,
chain: Vec<String>,
trigger: TriggerKind,
system: Option<String>,
wakeup_id: Option<String>,
started_at: String,
}
#[derive(Debug)]
struct JobGate {
max_concurrent: usize,
max_deferred: usize,
allow_overlap: bool,
running_total: usize,
running_by_leaf: HashMap<String, usize>,
running_jobs: Vec<RunningJob>,
deferred: VecDeque<PendingSpawn>,
recent: VecDeque<SpawnRecord>,
deferred_drops: u64,
skipped_overlap: u64,
spawned: u64,
}
#[derive(Debug, Clone, Copy)]
struct JobGateSnapshot {
running: usize,
max_concurrent: usize,
deferred: usize,
deferred_drops: u64,
skipped_overlap: u64,
spawned: u64,
allow_overlap: bool,
}
impl JobGate {
fn from_env() -> Self {
let max_concurrent = std::env::var("JAN_CRON_MAX_CONCURRENT")
.ok()
.and_then(|s| s.parse().ok())
.filter(|&n| n > 0)
.unwrap_or(DEFAULT_MAX_CONCURRENT);
let max_deferred = std::env::var("JAN_CRON_MAX_DEFERRED")
.ok()
.and_then(|s| s.parse().ok())
.filter(|&n| n > 0)
.unwrap_or(DEFAULT_MAX_DEFERRED);
let allow_overlap = matches!(
std::env::var("JAN_CRON_ALLOW_OVERLAP")
.unwrap_or_default()
.to_ascii_lowercase()
.as_str(),
"1" | "true" | "yes" | "on"
);
Self {
max_concurrent,
max_deferred,
allow_overlap,
running_total: 0,
running_by_leaf: HashMap::new(),
running_jobs: Vec::new(),
deferred: VecDeque::new(),
recent: VecDeque::new(),
deferred_drops: 0,
skipped_overlap: 0,
spawned: 0,
}
}
fn snapshot(&self) -> JobGateSnapshot {
JobGateSnapshot {
running: self.running_total,
max_concurrent: self.max_concurrent,
deferred: self.deferred.len(),
deferred_drops: self.deferred_drops,
skipped_overlap: self.skipped_overlap,
spawned: self.spawned,
allow_overlap: self.allow_overlap,
}
}
fn leaf_running(&self, leaf: &str) -> usize {
self.running_by_leaf.get(leaf).copied().unwrap_or(0)
}
fn push_recent(&mut self, record: SpawnRecord) {
if self.recent.len() >= RECENT_SPAWN_CAP {
self.recent.pop_front();
}
self.recent.push_back(record);
}
fn note_overlap_skip(&mut self, job: &PendingSpawn) {
self.skipped_overlap += 1;
let at = now_rfc3339();
self.push_recent(SpawnRecord {
leaf: job.leaf(),
chain: job.chain.clone(),
trigger: job.trigger,
system: job.system.clone(),
wakeup_id: job.wakeup_id.clone(),
started_at: at.clone(),
finished_at: Some(at),
exit_code: None,
overlap_skip: true,
});
}
fn note_started(&mut self, job: &PendingSpawn) -> String {
let started_at = now_rfc3339();
let leaf = job.leaf();
self.running_total += 1;
*self.running_by_leaf.entry(leaf.clone()).or_insert(0) += 1;
self.spawned += 1;
self.running_jobs.push(RunningJob {
leaf,
chain: job.chain.clone(),
trigger: job.trigger,
system: job.system.clone(),
wakeup_id: job.wakeup_id.clone(),
started_at: started_at.clone(),
});
started_at
}
fn note_finished(&mut self, leaf: &str, started_at: &str, exit_code: Option<i32>) {
self.running_total = self.running_total.saturating_sub(1);
if let Some(n) = self.running_by_leaf.get_mut(leaf) {
*n = n.saturating_sub(1);
if *n == 0 {
self.running_by_leaf.remove(leaf);
}
}
let mut matched: Option<RunningJob> = None;
if let Some(pos) = self
.running_jobs
.iter()
.position(|j| j.leaf == leaf && j.started_at == started_at)
{
matched = Some(self.running_jobs.remove(pos));
} else if let Some(pos) = self.running_jobs.iter().position(|j| j.leaf == leaf) {
matched = Some(self.running_jobs.remove(pos));
}
if let Some(job) = matched {
self.push_recent(SpawnRecord {
leaf: job.leaf,
chain: job.chain,
trigger: job.trigger,
system: job.system,
wakeup_id: job.wakeup_id,
started_at: job.started_at,
finished_at: Some(now_rfc3339()),
exit_code,
overlap_skip: false,
});
}
}
fn running_views(&self) -> Vec<JobView> {
self.running_jobs
.iter()
.map(|j| JobView {
leaf: j.leaf.clone(),
chain: j.chain.clone(),
trigger: j.trigger,
system: j.system.clone(),
wakeup_id: j.wakeup_id.clone(),
started_at: j.started_at.clone(),
})
.collect()
}
fn deferred_views(&self) -> Vec<JobView> {
let queued = now_rfc3339();
self.deferred
.iter()
.take(DEFERRED_PEEK)
.map(|j| j.as_job_view(&queued))
.collect()
}
fn recent_records(&self) -> Vec<SpawnRecord> {
self.recent.iter().rev().cloned().collect()
}
fn defer(&mut self, job: PendingSpawn, verbose: bool) {
if self.deferred.len() >= self.max_deferred {
self.deferred.pop_front();
self.deferred_drops += 1;
if verbose {
eprintln!("jan cron daemon: deferred queue full; dropped oldest pending job");
}
}
if verbose {
eprintln!(
"jan cron daemon: defer `{}` (running={}/{})",
job.chain.join(" "),
self.running_total,
self.max_concurrent
);
}
self.deferred.push_back(job);
}
}
impl DaemonState {
fn empty() -> Self {
Self {
cache: ScheduleCache {
schedules: Vec::new(),
agents_by_name: HashMap::new(),
spec_dir: String::new(),
root_yaml: String::new(),
loaded_at: Instant::now(),
},
last_fired: HashMap::new(),
wakeups: Arc::new(Mutex::new(VecDeque::new())),
events_connected: Arc::new(AtomicBool::new(false)),
wakeup_metrics: Arc::new(WakeupMetrics::default()),
jobs: Arc::new(Mutex::new(JobGate::from_env())),
disabled: load_disabled(),
}
}
fn refresh(&mut self) -> Result<usize> {
self.cache = load_schedule_cache()?;
self.last_fired.clear();
self.disabled = load_disabled();
Ok(self.cache.schedules.len())
}
fn agent_disabled(&self, name: &str) -> bool {
self.disabled.contains(name)
}
fn resolve_disable_target(&self, target: &str) -> Result<String> {
let target = target.trim();
if target.is_empty() {
bail!("agent name required");
}
let parts: Vec<&str> = target.split_whitespace().collect();
if parts.len() == 1 {
let name = parts[0];
if self.cache.agents_by_name.contains_key(name)
|| self
.cache
.schedules
.iter()
.any(|s| s.chain.last().is_some_and(|n| n == name))
|| self.disabled.contains(name)
{
return Ok(name.to_string());
}
bail!("unknown agent `{name}` (not a script leaf in the preferred tree)");
}
let chain: Vec<String> = parts.iter().map(|s| (*s).to_string()).collect();
if let Some(name) = chain.last() {
if self
.cache
.agents_by_name
.get(name)
.is_some_and(|m| m.chain == chain)
|| self.cache.schedules.iter().any(|s| s.chain == chain)
{
return Ok(name.clone());
}
if self.cache.agents_by_name.contains_key(name) {
return Ok(name.clone());
}
}
bail!("unknown agent `{}`", chain.join(" "))
}
fn disable_agent(&mut self, target: &str) -> Result<String> {
let name = self.resolve_disable_target(target)?;
if self.disabled.insert(name.clone()) {
save_disabled(&self.disabled)?;
}
Ok(name)
}
fn enable_agent(&mut self, target: &str) -> Result<String> {
let name = self.resolve_disable_target(target)?;
if self.disabled.remove(&name) {
save_disabled(&self.disabled)?;
}
Ok(name)
}
fn disabled_names(&self) -> Vec<String> {
let mut names: Vec<_> = self.disabled.iter().cloned().collect();
names.sort();
names
}
fn tick(&mut self, now: &TickTime, jan: &Path, verbose: bool) {
self.drain_deferred(jan, verbose);
for sched in &self.cache.schedules {
let chain_key = sched.chain.join(" ");
let leaf = sched.chain.last().map(String::as_str).unwrap_or("");
if self.agent_disabled(leaf) {
if verbose && sched.matches_tick(now) {
eprintln!("jan cron daemon: disabled skip `{chain_key}`");
}
continue;
}
if !sched.matches_tick(now) {
continue;
}
if self
.last_fired
.get(&chain_key)
.is_some_and(|prev| prev == now)
{
continue;
}
self.last_fired.insert(chain_key.clone(), *now);
self.enqueue_or_spawn(
jan,
PendingSpawn {
chain: sched.chain.clone(),
system: sched.system.clone(),
extra_args: Vec::new(),
env: sched
.system
.as_ref()
.map(|s| vec![("JAN_SYSTEM".into(), s.clone())])
.unwrap_or_default(),
trigger: TriggerKind::Cron,
wakeup_id: None,
},
verbose,
);
}
self.drain_unifier_wakeups(jan, verbose);
}
fn drain_deferred(&self, jan: &Path, verbose: bool) {
loop {
let job = {
let mut jobs = self.jobs.lock().unwrap();
let Some(front) = jobs.deferred.front() else {
break;
};
let leaf = front
.chain
.last()
.cloned()
.unwrap_or_else(|| "unknown".into());
if !jobs.allow_overlap && jobs.leaf_running(&leaf) > 0 {
break;
}
if jobs.running_total >= jobs.max_concurrent {
break;
}
jobs.deferred.pop_front()
};
let Some(job) = job else {
break;
};
if !spawn_tracked(
jan,
&job,
&self.jobs,
self.events_connected.load(Ordering::SeqCst),
verbose,
) {
break;
}
}
}
fn enqueue_or_spawn(&self, jan: &Path, job: PendingSpawn, verbose: bool) {
let leaf = job.leaf();
{
let mut jobs = self.jobs.lock().unwrap();
if !jobs.allow_overlap && jobs.leaf_running(&leaf) > 0 {
jobs.note_overlap_skip(&job);
if verbose {
eprintln!(
"jan cron daemon: skip overlap for `{leaf}` (already running; set JAN_CRON_ALLOW_OVERLAP=1 to allow)"
);
}
return;
}
if jobs.running_total >= jobs.max_concurrent {
jobs.defer(job, verbose);
return;
}
}
let _ = spawn_tracked(
jan,
&job,
&self.jobs,
self.events_connected.load(Ordering::SeqCst),
verbose,
);
}
fn drain_unifier_wakeups(&self, jan: &Path, verbose: bool) {
for wakeup in unifier_events::drain_wakeups(&self.wakeups) {
let Some(name) = wakeup.agent_name().map(str::to_string) else {
if verbose {
eprintln!("jan cron daemon: skip event wakeup without name ({wakeup:?})");
}
continue;
};
if self.agent_disabled(&name) {
if verbose {
eprintln!("jan cron daemon: disabled skip wakeup for `{name}`");
}
continue;
}
let Some(meta) = self.cache.agents_by_name.get(&name).cloned() else {
if verbose {
eprintln!("jan cron daemon: no script leaf named `{name}` for unifier wakeup");
}
continue;
};
let (extra_args, mut env, trigger, wakeup_id) = wakeup_spawn_meta(&wakeup);
if let Some(ref sys) = meta.system {
env.push(("JAN_SYSTEM".into(), sys.clone()));
}
self.enqueue_or_spawn(
jan,
PendingSpawn {
chain: meta.chain,
system: meta.system,
extra_args,
env,
trigger,
wakeup_id,
},
verbose,
);
}
}
fn status_line(&self, started: Instant) -> String {
let exprs: usize = self.cache.schedules.iter().map(|s| s.cron_raw.len()).sum();
let events = if self.events_connected.load(Ordering::SeqCst) {
"connected"
} else {
"disconnected"
};
let wm = self.wakeup_metrics.snapshot();
let jobs = self.jobs.lock().unwrap().snapshot();
let overlap = if jobs.allow_overlap { "allow" } else { "deny" };
format!(
"ok pid={} uptime_s={} schedules={} exprs={} agents={} disabled={} events={} wakeups={} drops={} tick_notices={} reconnects={} running={}/{} deferred={} deferred_drops={} overlap_skips={} overlap={} spawned={} cached=1 age_s={} spec={}/{} events_sock={}",
std::process::id(),
started.elapsed().as_secs(),
self.cache.schedules.len(),
exprs,
self.cache.agents_by_name.len(),
self.disabled.len(),
events,
unifier_events::wakeup_depth(&self.wakeups),
wm.dropped,
wm.tick_notices,
wm.reconnects,
jobs.running,
jobs.max_concurrent,
jobs.deferred,
jobs.deferred_drops,
jobs.skipped_overlap,
overlap,
jobs.spawned,
self.cache.loaded_at.elapsed().as_secs(),
self.cache.spec_dir,
self.cache.root_yaml,
unifier_events::events_socket_path().display()
)
}
fn status_detail(&self) -> StatusDetail {
let jobs = self.jobs.lock().unwrap();
StatusDetail {
running: jobs.running_views(),
deferred: jobs.deferred_views(),
recent: jobs.recent_records(),
wakeups: self.wakeup_leaves(),
}
}
fn wakeup_leaves(&self) -> Vec<WakeupLeaf> {
let cron_leaves: HashSet<String> = self
.cache
.schedules
.iter()
.filter_map(|s| s.chain.last().cloned())
.collect();
let mut leaves: Vec<WakeupLeaf> = self
.cache
.agents_by_name
.iter()
.map(|(leaf, meta)| WakeupLeaf {
leaf: leaf.clone(),
chain: meta.chain.clone(),
system: meta.system.clone(),
has_cron: cron_leaves.contains(leaf),
disabled: self.disabled.contains(leaf),
})
.collect();
leaves.sort_by(|a, b| a.leaf.cmp(&b.leaf));
leaves
}
fn cached_entries(&self) -> Vec<CachedCronEntry> {
self.cache
.schedules
.iter()
.map(|s| {
let leaf = s.chain.last().cloned().unwrap_or_default();
CachedCronEntry {
chain: s.chain.clone(),
cron: s.cron_raw.clone(),
disabled: self.disabled.contains(&leaf),
}
})
.collect()
}
}
fn wakeup_spawn_meta(
wakeup: &Wakeup,
) -> (Vec<String>, Vec<(String, String)>, TriggerKind, Option<String>) {
match wakeup {
Wakeup::Mailbox { id, from, to } => (
vec!["--message-id".into(), id.clone()],
vec![
("JAN_UNIFIER_KIND".into(), "mailbox".into()),
("JAN_UNIFIER_MESSAGE_ID".into(), id.clone()),
("JAN_UNIFIER_FROM".into(), from.clone()),
("JAN_UNIFIER_TO".into(), to.clone()),
],
TriggerKind::Mailbox,
Some(id.clone()),
),
Wakeup::Event { id, name } => {
let mut args = vec!["--event-id".into(), id.clone()];
let mut env = vec![
("JAN_UNIFIER_KIND".into(), "event".into()),
("JAN_UNIFIER_EVENT_ID".into(), id.clone()),
];
if let Some(n) = name {
args.push("--event-name".into());
args.push(n.clone());
env.push(("JAN_UNIFIER_EVENT_NAME".into(), n.clone()));
}
(
args,
env,
TriggerKind::Event,
Some(id.clone()),
)
}
}
}
fn cron_span_enabled(events_connected: bool) -> bool {
match std::env::var("JAN_CRON_SPAN")
.unwrap_or_default()
.to_ascii_lowercase()
.as_str()
{
"0" | "false" | "no" | "off" => false,
"1" | "true" | "yes" | "on" => true,
"" => events_connected,
_ => events_connected,
}
}
fn unifier_log_start(
leaf: &str,
trigger: TriggerKind,
wakeup_id: Option<&str>,
system: Option<&str>,
) -> Option<String> {
let mut cmd = Command::new("unifier");
cmd.args(["log", "start", leaf, "-f", &format!("trigger={}", trigger.as_str())]);
if let Some(id) = wakeup_id {
cmd.args(["-f", &format!("wakeup={id}")]);
}
if let Some(sys) = system {
cmd.args(["-f", &format!("system={sys}")]);
}
cmd.stdin(Stdio::null())
.stdout(Stdio::piped())
.stderr(Stdio::null());
let out = cmd.output().ok()?;
if !out.status.success() {
return None;
}
let id = String::from_utf8_lossy(&out.stdout).trim().to_string();
if id.is_empty() {
None
} else {
Some(id)
}
}
fn unifier_log_end(span_id: &str) {
let _ = Command::new("unifier")
.args(["log", "end", span_id])
.stdin(Stdio::null())
.stdout(Stdio::null())
.stderr(Stdio::null())
.status();
}
fn spawn_tracked(
jan: &Path,
job: &PendingSpawn,
jobs: &Arc<Mutex<JobGate>>,
events_connected: bool,
verbose: bool,
) -> bool {
let leaf = job.leaf();
let mut args = vec!["--no-log".to_string()];
args.extend(job.chain.iter().cloned());
args.push("run".into());
args.extend(job.extra_args.iter().cloned());
if verbose {
eprintln!(
"jan cron daemon: spawn `{} {}` trigger={}",
jan.display(),
args.join(" "),
job.trigger.as_str()
);
}
let mut cmd = Command::new(jan);
cmd.args(&args)
.stdin(Stdio::null())
.stdout(Stdio::null())
.stderr(Stdio::inherit());
for (k, v) in &job.env {
cmd.env(k, v);
}
if let Some(ref sys) = job.system {
cmd.env("JAN_SYSTEM", sys);
}
match cmd.spawn() {
Ok(mut child) => {
let started_at = {
let mut gate = jobs.lock().unwrap();
gate.note_started(job)
};
let span_id = if cron_span_enabled(events_connected) {
unifier_log_start(
&leaf,
job.trigger,
job.wakeup_id.as_deref(),
job.system.as_deref(),
)
} else {
None
};
let jobs = Arc::clone(jobs);
thread::spawn(move || {
let code = child.wait().ok().and_then(|s| s.code());
if let Some(ref span) = span_id {
unifier_log_end(span);
}
if let Ok(mut gate) = jobs.lock() {
gate.note_finished(&leaf, &started_at, code);
}
});
true
}
Err(e) => {
eprintln!(
"jan cron daemon: failed to spawn `{} {}`: {e:#}",
jan.display(),
args.join(" ")
);
false
}
}
}
fn handle_client(mut stream: UnixStream, state: Arc<Mutex<DaemonState>>, started: Instant) {
let reader = BufReader::new(
stream
.try_clone()
.unwrap_or_else(|_| stream.try_clone().expect("clone daemon control socket")),
);
let line = match reader.lines().next() {
Some(Ok(l)) => l,
_ => return,
};
let line = line.trim();
let (cmd, rest) = match line.split_once(char::is_whitespace) {
Some((c, r)) => (c, r.trim()),
None => (line, ""),
};
let reply = match cmd.to_ascii_lowercase().as_str() {
"ping" => "pong".to_string(),
"status" => state.lock().unwrap().status_line(started),
"status-detail" => match serde_json::to_string(&state.lock().unwrap().status_detail()) {
Ok(json) => format!("ok detail {json}"),
Err(e) => format!("error {e}"),
},
"wakeups" => match serde_json::to_string(&state.lock().unwrap().wakeup_leaves()) {
Ok(json) => format!("ok wakeups {json}"),
Err(e) => format!("error {e}"),
},
"reload" | "refresh" => match state.lock().unwrap().refresh() {
Ok(n) => format!("ok schedules={n}"),
Err(e) => format!("error {e:#}"),
},
"list" => match serde_json::to_string(&state.lock().unwrap().cached_entries()) {
Ok(json) => format!("ok list {json}"),
Err(e) => format!("error {e}"),
},
"disable" => {
if rest.is_empty() {
"error usage: disable <agent>".to_string()
} else {
match state.lock().unwrap().disable_agent(rest) {
Ok(name) => format!("ok disabled={name}"),
Err(e) => format!("error {e:#}"),
}
}
}
"enable" => {
if rest.is_empty() {
"error usage: enable <agent>".to_string()
} else {
match state.lock().unwrap().enable_agent(rest) {
Ok(name) => format!("ok enabled={name}"),
Err(e) => format!("error {e:#}"),
}
}
}
"disabled" => match serde_json::to_string(&state.lock().unwrap().disabled_names()) {
Ok(json) => format!("ok disabled {json}"),
Err(e) => format!("error {e}"),
},
"stop" => {
STOP_REQUESTED.store(true, Ordering::SeqCst);
"ok stopping".to_string()
}
other => format!("error unknown command `{other}`"),
};
let _ = writeln!(stream, "{reply}");
}
fn accept_control(state: Arc<Mutex<DaemonState>>, started: Instant) {
let listener = match UnixListener::bind(&socket_path()) {
Ok(l) => l,
Err(e) => {
eprintln!("jan cron daemon: bind {}: {e:#}", socket_path().display());
return;
}
};
if let Err(e) = listener.set_nonblocking(true) {
eprintln!("jan cron daemon: set_nonblocking: {e:#}");
return;
}
while !STOP_REQUESTED.load(Ordering::SeqCst) {
match listener.accept() {
Ok((stream, _)) => {
let state = Arc::clone(&state);
thread::spawn(move || handle_client(stream, state, started));
}
Err(e) if e.kind() == std::io::ErrorKind::WouldBlock => {
thread::sleep(Duration::from_millis(50));
}
Err(e) => {
eprintln!("jan cron daemon: accept: {e:#}");
thread::sleep(Duration::from_millis(100));
}
}
}
}
fn sleep_until_next_tick(start: Instant, tick_index: u64) {
let target = start + Duration::from_millis(tick_index * TICK_MS);
let now = Instant::now();
if target > now {
thread::sleep(target - now);
}
}
pub fn run_foreground(verbose: bool) -> Result<i32> {
STOP_REQUESTED.store(false, Ordering::SeqCst);
fs::create_dir_all(runtime_dir())
.with_context(|| format!("create {}", runtime_dir().display()))?;
remove_runtime_files();
write_pid_file()?;
let jan = jan_bin()?;
let state = Arc::new(Mutex::new(DaemonState::empty()));
{
let n = state
.lock()
.unwrap()
.refresh()
.context("load cron schedules")?;
if verbose {
eprintln!("jan cron daemon: cached {n} scheduled script(s) (refresh to reload)");
}
}
if let Err(e) = crate::runtime_daemon::start_daemon_background() {
if verbose {
eprintln!("jan cron daemon: runtime pool not started ({e:#})");
}
}
let started = Instant::now();
let state_bg = Arc::clone(&state);
let control = thread::spawn(move || accept_control(state_bg, started));
let (wakeups, events_connected, wakeup_metrics) = {
let st = state.lock().unwrap();
(
Arc::clone(&st.wakeups),
Arc::clone(&st.events_connected),
Arc::clone(&st.wakeup_metrics),
)
};
let events = thread::spawn(move || {
unifier_events::run_listener(wakeups, events_connected, wakeup_metrics, verbose);
});
let loop_start = Instant::now();
let mut tick_index = 0u64;
while !STOP_REQUESTED.load(Ordering::SeqCst) {
sleep_until_next_tick(loop_start, tick_index);
let now = TickTime::now_local();
state.lock().unwrap().tick(&now, &jan, verbose);
tick_index += 1;
}
remove_runtime_files();
let _ = control.join();
let _ = events.join();
Ok(0)
}
pub fn send_command(cmd: &str) -> Result<String> {
let path = socket_path();
if !path.exists() {
bail!(
"jan cron daemon is not running (no socket at {})",
path.display()
);
}
let mut stream =
UnixStream::connect(&path).with_context(|| format!("connect to {}", path.display()))?;
stream.set_read_timeout(Some(Duration::from_secs(5)))?;
stream.set_write_timeout(Some(Duration::from_secs(5)))?;
writeln!(stream, "{cmd}").context("write daemon command")?;
let mut reader = BufReader::new(stream);
let mut reply = String::new();
reader.read_line(&mut reply).context("read daemon reply")?;
Ok(reply.trim().to_string())
}
pub fn fetch_cached_entries() -> Result<Vec<CachedCronEntry>> {
let reply = send_command("list")
.with_context(|| "jan cron --list reads the daemon cache; run `jan cron start` first")?;
if let Some(rest) = reply.strip_prefix("error ") {
if rest.contains("unknown command") && rest.contains("list") {
bail!(
"jan cron daemon is outdated (no `list` cache command); \
run `jan cron stop` then `jan cron start` with this jan binary"
);
}
bail!("daemon list failed: {rest}");
}
let json = reply
.strip_prefix("ok list ")
.ok_or_else(|| anyhow::anyhow!("unexpected daemon list reply: {reply}"))?;
serde_json::from_str(json).with_context(|| format!("parse daemon list JSON: {json}"))
}
pub fn daemon_running() -> bool {
match send_command("ping") {
Ok(r) => r == "pong",
Err(_) => false,
}
}
pub fn start_daemon_background(verbose: bool) -> Result<()> {
if daemon_running() {
return Ok(());
}
let jan = jan_bin()?;
let mut cmd = Command::new(&jan);
cmd.args(["--no-log", "cron", "daemon"]);
if verbose {
cmd.arg("-v");
}
cmd.stdin(Stdio::null())
.stdout(Stdio::null())
.stderr(Stdio::null());
cmd.spawn()
.with_context(|| format!("spawn `{} cron daemon`", jan.display()))?;
for _ in 0..50 {
if daemon_running() {
return Ok(());
}
thread::sleep(Duration::from_millis(100));
}
bail!(
"jan cron daemon failed to start (no response on {})",
socket_path().display()
);
}
pub fn stop_daemon() -> Result<()> {
if !socket_path().exists() {
println!("(jan cron daemon is not running)");
return Ok(());
}
let reply = send_command("stop")?;
for _ in 0..30 {
if !socket_path().exists() {
println!("stopped jan cron daemon");
return Ok(());
}
thread::sleep(Duration::from_millis(100));
}
bail!("jan cron daemon did not stop: {reply}");
}
#[derive(Debug, Clone, Copy, Default)]
pub struct StatusOpts {
pub compact: bool,
pub json: bool,
}
pub fn fetch_status_detail() -> Result<StatusDetail> {
let reply = send_command("status-detail")?;
let Some(json) = reply.strip_prefix("ok detail ") else {
if let Some(err) = reply.strip_prefix("error ") {
bail!("{err}");
}
bail!("unexpected daemon reply: {reply}");
};
serde_json::from_str(json).with_context(|| format!("parse status-detail: {json}"))
}
fn fetch_wakeup_leaves() -> Result<Vec<WakeupLeaf>> {
let reply = send_command("wakeups")?;
let Some(json) = reply.strip_prefix("ok wakeups ") else {
if let Some(err) = reply.strip_prefix("error ") {
bail!("{err}");
}
bail!("unexpected daemon reply: {reply}");
};
serde_json::from_str(json).with_context(|| format!("parse wakeups: {json}"))
}
fn print_job_views(title: &str, jobs: &[JobView]) {
println!("{title}");
if jobs.is_empty() {
println!(" (none)");
return;
}
for j in jobs {
let id = j
.wakeup_id
.as_deref()
.map(|s| format!(" id={s}"))
.unwrap_or_default();
let sys = j
.system
.as_deref()
.map(|s| format!(" system={s}"))
.unwrap_or_default();
println!(
" {} trigger={}{sys} started={}{id}",
j.leaf,
j.trigger.as_str(),
j.started_at
);
}
}
fn print_recent(records: &[SpawnRecord]) {
println!("recent:");
if records.is_empty() {
println!(" (none)");
return;
}
for r in records.iter().take(16) {
let sys = r
.system
.as_deref()
.map(|s| format!(" system={s}"))
.unwrap_or_default();
if r.overlap_skip {
println!(
" {} trigger={}{sys} {} overlap_skip",
r.leaf,
r.trigger.as_str(),
r.started_at
);
continue;
}
let code = r
.exit_code
.map(|c| format!("exit={c}"))
.unwrap_or_else(|| "exit=?".into());
let finished = r.finished_at.as_deref().unwrap_or("-");
println!(
" {} trigger={}{sys} {} → {} {code}",
r.leaf,
r.trigger.as_str(),
r.started_at,
finished
);
}
}
fn print_wakeups(leaves: &[WakeupLeaf]) {
println!("wake targets:");
if leaves.is_empty() {
println!(" (none)");
return;
}
for w in leaves {
let cron = if w.has_cron { "cron" } else { "event-only" };
let status = if w.disabled { "disabled" } else { "enabled" };
let sys = w.system.as_deref().unwrap_or("-");
println!(
" {}|{}|{}|{}|{}",
w.leaf,
w.chain.join(" "),
sys,
cron,
status
);
}
}
pub fn daemon_status_opts(opts: StatusOpts) -> Result<i32> {
match send_command("status") {
Ok(reply) if reply.starts_with("ok ") => {
if opts.json {
let detail = fetch_status_detail().unwrap_or(StatusDetail {
running: Vec::new(),
deferred: Vec::new(),
recent: Vec::new(),
wakeups: Vec::new(),
});
let mut obj = serde_json::Map::new();
obj.insert(
"summary".into(),
serde_json::Value::String(reply.clone()),
);
obj.insert(
"detail".into(),
serde_json::to_value(&detail).unwrap_or(serde_json::Value::Null),
);
println!("{}", serde_json::to_string_pretty(&obj)?);
return Ok(0);
}
println!("jan cron daemon running ({reply})");
if !opts.compact {
match fetch_status_detail() {
Ok(detail) => {
print_job_views("running:", &detail.running);
print_job_views("deferred:", &detail.deferred);
print_recent(&detail.recent);
}
Err(e) => {
eprintln!("(status detail unavailable: {e:#})");
}
}
}
Ok(0)
}
Ok(reply) => {
println!("jan cron daemon: {reply}");
Ok(1)
}
Err(e) => {
println!("jan cron daemon is not running ({e:#})");
Ok(1)
}
}
}
pub fn list_wakeup_leaves() -> Result<i32> {
let leaves = fetch_wakeup_leaves().with_context(|| {
"jan cron wakeups talks to the daemon; run `jan cron start` first"
})?;
print_wakeups(&leaves);
eprintln!("({} wakeup target(s))", leaves.len());
Ok(0)
}
pub fn watch_status(interval_ms: u64) -> Result<i32> {
loop {
println!("---- {}", now_rfc3339());
let code = daemon_status_opts(StatusOpts {
compact: false,
json: false,
})?;
if code != 0 {
return Ok(code);
}
thread::sleep(Duration::from_millis(interval_ms.max(200)));
}
}
pub fn reload_daemon() -> Result<()> {
let reply = send_command("refresh")?;
if reply.starts_with("ok ") {
println!("refreshed jan cron schedule cache ({reply})");
Ok(())
} else {
bail!("refresh failed: {reply}");
}
}
pub fn disable_agent(target: &str) -> Result<String> {
let reply = send_command(&format!("disable {target}"))
.with_context(|| "jan cron disable talks to the daemon; run `jan cron start` first")?;
if let Some(name) = reply.strip_prefix("ok disabled=") {
println!("disabled agent `{name}` (cron + unifier event wakeups)");
return Ok(name.to_string());
}
if let Some(err) = reply.strip_prefix("error ") {
bail!("{err}");
}
bail!("unexpected daemon reply: {reply}");
}
pub fn enable_agent(target: &str) -> Result<String> {
let reply = send_command(&format!("enable {target}"))
.with_context(|| "jan cron enable talks to the daemon; run `jan cron start` first")?;
if let Some(name) = reply.strip_prefix("ok enabled=") {
println!("enabled agent `{name}`");
return Ok(name.to_string());
}
if let Some(err) = reply.strip_prefix("error ") {
bail!("{err}");
}
bail!("unexpected daemon reply: {reply}");
}
pub fn list_disabled_agents() -> Result<Vec<String>> {
match send_command("disabled") {
Ok(reply) => {
if let Some(json) = reply.strip_prefix("ok disabled ") {
let names: Vec<String> = serde_json::from_str(json)
.with_context(|| format!("parse disabled list: {json}"))?;
return Ok(names);
}
if let Some(err) = reply.strip_prefix("error ") {
if err.contains("unknown command") {
let mut names: Vec<_> = load_disabled().into_iter().collect();
names.sort();
return Ok(names);
}
bail!("{err}");
}
bail!("unexpected daemon reply: {reply}");
}
Err(_) => {
let mut names: Vec<_> = load_disabled().into_iter().collect();
names.sort();
Ok(names)
}
}
}
fn systemd_unit_body(jan: &Path) -> String {
format!(
r#"[Unit]
Description=Jan cron scheduler daemon (100ms ticks)
After=default.target
[Service]
Type=simple
ExecStart={} --no-log cron daemon --foreground
Restart=on-failure
RestartSec=5
[Install]
WantedBy=default.target
"#,
jan.display()
)
}
pub fn install_systemd(dry_run: bool) -> Result<i32> {
let jan = jan_bin()?;
let unit_path = systemd_unit_path();
let body = systemd_unit_body(&jan);
if dry_run {
println!("# dry-run: would write {}:", unit_path.display());
print!("{body}");
println!("# dry-run: would run systemctl --user daemon-reload");
println!("# dry-run: would run systemctl --user enable --now {SERVICE_NAME}.service");
return Ok(0);
}
if let Some(parent) = unit_path.parent() {
fs::create_dir_all(parent).with_context(|| format!("create {}", parent.display()))?;
}
fs::write(&unit_path, &body).with_context(|| format!("write {}", unit_path.display()))?;
run_systemctl(&["--user", "daemon-reload"])?;
run_systemctl(&[
"--user",
"enable",
"--now",
&format!("{SERVICE_NAME}.service"),
])?;
println!(
"installed and started {} (unit: {})",
SERVICE_NAME,
unit_path.display()
);
println!(
" ExecStart={} --no-log cron daemon --foreground",
jan.display()
);
Ok(0)
}
pub fn uninstall_systemd(dry_run: bool) -> Result<i32> {
let unit_path = systemd_unit_path();
if dry_run {
println!("# dry-run: would run systemctl --user disable --now {SERVICE_NAME}.service");
if unit_path.is_file() {
println!("# dry-run: would remove {}", unit_path.display());
}
return Ok(0);
}
let _ = run_systemctl(&["--user", "stop", &format!("{SERVICE_NAME}.service")]);
let _ = run_systemctl(&["--user", "disable", &format!("{SERVICE_NAME}.service")]);
if unit_path.is_file() {
fs::remove_file(&unit_path).with_context(|| format!("remove {}", unit_path.display()))?;
}
let _ = run_systemctl(&["--user", "daemon-reload"]);
stop_daemon().ok();
println!("removed {SERVICE_NAME} systemd user service");
Ok(0)
}
fn run_systemctl(args: &[&str]) -> Result<()> {
let out = Command::new("systemctl")
.args(args)
.output()
.with_context(|| format!("spawn systemctl {}", args.join(" ")))?;
if !out.status.success() {
bail!(
"systemctl {} failed ({}): {}",
args.join(" "),
out.status,
String::from_utf8_lossy(&out.stderr).trim()
);
}
Ok(())
}
#[cfg(test)]
mod tests {
use super::*;
use crate::cron::{CronExpr, TickTime};
#[test]
fn systemd_unit_contains_execstart() {
let body = systemd_unit_body(Path::new("/usr/local/bin/jan"));
assert!(body.contains("ExecStart=/usr/local/bin/jan --no-log cron daemon --foreground"));
assert!(body.contains("WantedBy=default.target"));
}
#[test]
fn cached_schedule_matches_without_reparsing() {
let sched = ParsedSchedule {
chain: vec!["scripts".into(), "misc".into(), "tick".into()],
system: None,
cron_raw: vec!["* * * * * *".into()],
exprs: vec![CronExpr::parse("* * * * * *").unwrap()],
};
let now = TickTime {
decisecond: 0,
second: 12,
minute: 30,
hour: 10,
day: 5,
month: 8,
dow: 2,
};
assert!(sched.matches_tick(&now));
let mid = TickTime {
decisecond: 5,
..now
};
assert!(!sched.matches_tick(&mid));
}
#[test]
fn cached_entries_serialize_roundtrip() {
let entries = vec![CachedCronEntry {
chain: vec!["scripts".into(), "misc".into(), "tick".into()],
cron: vec!["30 10 * * *".into(), "* * * * * *".into()],
disabled: true,
}];
let json = serde_json::to_string(&entries).unwrap();
let back: Vec<CachedCronEntry> = serde_json::from_str(&json).unwrap();
assert_eq!(back, entries);
assert_eq!(back[0].chain_str(), "scripts misc tick");
assert!(back[0].disabled);
}
#[test]
fn disabled_file_roundtrip() {
let dir = tempfile::tempdir().unwrap();
std::env::set_var("JAN_CONFIG_DIR", dir.path());
let mut set = HashSet::new();
set.insert("ping-agent".into());
set.insert("mute-tracker".into());
save_disabled(&set).unwrap();
let loaded = load_disabled();
assert_eq!(loaded.len(), 2);
assert!(loaded.contains("ping-agent"));
assert!(is_disabled_name("ping-agent"));
assert!(!is_disabled_name("report"));
std::env::remove_var("JAN_CONFIG_DIR");
}
#[test]
fn job_gate_defers_when_at_capacity() {
let mut gate = JobGate {
max_concurrent: 1,
max_deferred: 8,
allow_overlap: true,
running_total: 0,
running_by_leaf: HashMap::new(),
running_jobs: Vec::new(),
deferred: VecDeque::new(),
recent: VecDeque::new(),
deferred_drops: 0,
skipped_overlap: 0,
spawned: 0,
};
let started = gate.note_started(&PendingSpawn {
chain: vec!["a".into()],
system: None,
extra_args: vec![],
env: vec![],
trigger: TriggerKind::Cron,
wakeup_id: None,
});
assert_eq!(gate.running_total, 1);
gate.defer(
PendingSpawn {
chain: vec!["scripts".into(), "agents".into(), "b".into()],
system: Some("agents".into()),
extra_args: vec![],
env: vec![],
trigger: TriggerKind::Mailbox,
wakeup_id: Some("mid".into()),
},
false,
);
assert_eq!(gate.deferred.len(), 1);
gate.note_finished("a", &started, Some(0));
assert_eq!(gate.running_total, 0);
assert_eq!(gate.leaf_running("a"), 0);
assert_eq!(gate.recent.len(), 1);
assert_eq!(gate.recent[0].exit_code, Some(0));
}
#[test]
fn job_gate_overlap_skip_counter() {
let mut gate = JobGate {
max_concurrent: 8,
max_deferred: 8,
allow_overlap: false,
running_total: 0,
running_by_leaf: HashMap::new(),
running_jobs: Vec::new(),
deferred: VecDeque::new(),
recent: VecDeque::new(),
deferred_drops: 0,
skipped_overlap: 0,
spawned: 0,
};
let _ = gate.note_started(&PendingSpawn {
chain: vec!["pong-agent".into()],
system: Some("agents".into()),
extra_args: vec![],
env: vec![],
trigger: TriggerKind::Event,
wakeup_id: Some("e1".into()),
});
assert_eq!(gate.leaf_running("pong-agent"), 1);
gate.note_overlap_skip(&PendingSpawn {
chain: vec!["pong-agent".into()],
system: Some("agents".into()),
extra_args: vec![],
env: vec![],
trigger: TriggerKind::Event,
wakeup_id: Some("e2".into()),
});
assert_eq!(gate.skipped_overlap, 1);
assert!(gate.recent.back().unwrap().overlap_skip);
}
#[test]
fn job_gate_deferred_queue_drops_oldest() {
let mut gate = JobGate {
max_concurrent: 1,
max_deferred: 2,
allow_overlap: true,
running_total: 1,
running_by_leaf: HashMap::new(),
running_jobs: Vec::new(),
deferred: VecDeque::new(),
recent: VecDeque::new(),
deferred_drops: 0,
skipped_overlap: 0,
spawned: 0,
};
for name in ["a", "b", "c"] {
gate.defer(
PendingSpawn {
chain: vec![name.into()],
system: None,
extra_args: vec![],
env: vec![],
trigger: TriggerKind::Cron,
wakeup_id: None,
},
false,
);
}
assert_eq!(gate.deferred.len(), 2);
assert_eq!(gate.deferred_drops, 1);
assert_eq!(gate.deferred.front().unwrap().chain[0], "b");
}
}