#![cfg(all(unix, feature = "scheduler"))]
use anyhow::Context as _;
use crate::bootstrap::{load_config_or_default, resolve_config_path};
pub(crate) async fn handle_serve(
config_path: Option<&std::path::Path>,
foreground: bool,
catch_up: bool,
vault_override: Option<&str>,
vault_key_override: Option<&std::path::Path>,
vault_path_override: Option<&std::path::Path>,
) -> anyhow::Result<()> {
let config_file = resolve_config_path(config_path);
let config = load_config_or_default(&config_file)?;
let daemon_cfg = build_daemon_config(&config);
if foreground {
let (vault_key_path, vault_secrets_path) = crate::bootstrap::resolve_vault_paths(
&config,
vault_override,
vault_key_override,
vault_path_override,
);
run_foreground(daemon_cfg, &config, &vault_key_path, &vault_secrets_path).await
} else {
let config_str = config_file.to_string_lossy();
let extra = build_detach_args(
&config_str,
catch_up,
vault_override,
vault_key_override,
vault_path_override,
);
let extra_refs: Vec<&str> = extra.iter().map(String::as_str).collect();
zeph_scheduler::detach_and_run(&daemon_cfg, &extra_refs)
.context("failed to detach scheduler daemon")
}
}
fn build_detach_args(
config_str: &str,
catch_up: bool,
vault_override: Option<&str>,
vault_key_override: Option<&std::path::Path>,
vault_path_override: Option<&std::path::Path>,
) -> Vec<String> {
let mut extra: Vec<String> = vec![
"--config".to_owned(),
config_str.to_owned(),
"serve".to_owned(),
"--foreground".to_owned(),
];
if !catch_up {
extra.push("--no-catch-up".to_owned());
}
if let Some(backend) = vault_override {
extra.push("--vault".to_owned());
extra.push(backend.to_owned());
}
if let Some(key_path) = vault_key_override {
extra.push("--vault-key".to_owned());
extra.push(key_path.to_string_lossy().into_owned());
}
if let Some(vault_path) = vault_path_override {
extra.push("--vault-path".to_owned());
extra.push(vault_path.to_string_lossy().into_owned());
}
extra
}
pub(crate) fn handle_stop(
config_path: Option<&std::path::Path>,
timeout_secs: u64,
) -> anyhow::Result<()> {
let config_file = resolve_config_path(config_path);
let config = load_config_or_default(&config_file)?;
let daemon_cfg = build_daemon_config(&config);
zeph_scheduler::stop_daemon(&daemon_cfg, timeout_secs)
.context("failed to stop scheduler daemon")
}
pub(crate) async fn handle_status(
config_path: Option<&std::path::Path>,
json: bool,
n: usize,
) -> anyhow::Result<()> {
let config_file = resolve_config_path(config_path);
let config = load_config_or_default(&config_file)?;
let daemon_cfg = build_daemon_config(&config);
let db_url = crate::db_url::resolve_db_url(&config);
let status = zeph_scheduler::daemon_status(&daemon_cfg, db_url, n)
.await
.context("failed to read daemon status")?;
if json {
println!(
"{}",
serde_json::to_string_pretty(&status).context("failed to serialize daemon status")?
);
} else {
print_status_human(&status);
}
Ok(())
}
fn build_daemon_config(config: &zeph_core::config::Config) -> zeph_scheduler::DaemonConfig {
let sched = &config.scheduler.daemon;
zeph_scheduler::DaemonConfig {
pid_file: std::path::PathBuf::from(&sched.pid_file),
log_file: std::path::PathBuf::from(&sched.log_file),
catch_up: sched.catch_up,
tick_secs: sched.tick_secs,
shutdown_grace_secs: sched.shutdown_grace_secs,
handler_timeout_secs: sched.handler_timeout_secs,
}
}
async fn run_foreground(
daemon_cfg: zeph_scheduler::DaemonConfig,
config: &zeph_core::config::Config,
vault_key_path: &std::path::Path,
vault_secrets_path: &std::path::Path,
) -> anyhow::Result<()> {
let db_url = crate::db_url::resolve_db_url(config);
let store = zeph_scheduler::JobStore::open(db_url)
.await
.context("failed to open scheduler store")?;
store
.init()
.await
.context("failed to init scheduler store")?;
let (shutdown_tx, shutdown_rx) = tokio::sync::watch::channel(false);
let sched_cancel = tokio_util::sync::CancellationToken::new();
let sched_supervisor = zeph_common::TaskSupervisor::new(sched_cancel.clone());
{
let signal_fut = async move {
use tokio::signal::unix::{SignalKind, signal};
let mut sigterm =
signal(SignalKind::terminate()).expect("failed to register SIGTERM handler");
let mut sigint =
signal(SignalKind::interrupt()).expect("failed to register SIGINT handler");
tokio::select! {
_ = sigterm.recv() => tracing::info!("received SIGTERM"),
_ = sigint.recv() => tracing::info!("received SIGINT"),
}
let _ = shutdown_tx.send(true);
};
let cell = std::sync::Arc::new(parking_lot::Mutex::new(Some(signal_fut)));
sched_supervisor.spawn(zeph_common::TaskDescriptor {
name: "sched_daemon_signal",
restart: zeph_common::RestartPolicy::RunOnce,
factory: move || {
let f = cell.lock().take();
async move {
if let Some(f) = f {
f.await;
}
}
},
});
}
let (mut scheduler, ctrl_tx) = zeph_scheduler::Scheduler::new(store, shutdown_rx);
scheduler = scheduler.with_reentry_defense(
config.scheduler.security.enabled,
config.scheduler.security.injection_pattern_check,
config.scheduler.security.attenuate_after_external_read,
);
if let Some(adapter) = build_durable_adapter(
config,
&sched_supervisor,
vault_key_path,
vault_secrets_path,
)
.await?
{
scheduler = scheduler.with_durable(adapter);
}
if config.agent.auto_update_check {
let (update_tx, _update_rx) = tokio::sync::mpsc::channel(4);
let handler = zeph_scheduler::UpdateCheckHandler::new(env!("CARGO_PKG_VERSION"), update_tx);
scheduler.register_handler(&zeph_scheduler::TaskKind::UpdateCheck, Box::new(handler));
}
crate::scheduler::load_config_tasks(&config.scheduler.tasks, &ctrl_tx);
let result = zeph_scheduler::run_foreground(scheduler, &daemon_cfg)
.await
.context("scheduler daemon exited with error");
sched_cancel.cancel();
sched_supervisor
.shutdown_all(std::time::Duration::from_secs(5))
.await;
result
}
async fn build_durable_adapter(
config: &zeph_core::config::Config,
sched_supervisor: &zeph_common::TaskSupervisor,
vault_key_path: &std::path::Path,
vault_secrets_path: &std::path::Path,
) -> anyhow::Result<Option<zeph_scheduler::durable::SchedulerDurableAdapter>> {
if !(config.durable.enabled && config.durable.scheduler) {
return Ok(None);
}
let cipher =
crate::commands::durable::load_write_cipher(config, vault_key_path, vault_secrets_path)?;
let hmac_keys =
crate::commands::durable::load_write_hmac_key(config, vault_key_path, vault_secrets_path)?;
let hwm_keys =
crate::commands::durable::load_write_hwm_key(config, vault_key_path, vault_secrets_path)?;
let durable_url = crate::commands::durable::resolve_durable_db_url(config);
let local = zeph_durable::LocalBackend::open(&durable_url, config.durable.max_payload_bytes)
.await
.context("failed to open scheduler durable backend")?;
local
.init()
.await
.context("failed to init scheduler durable schema")?;
let local = if let Some(cipher) = cipher {
local.with_cipher(cipher)
} else {
local
};
let local = if let Some(key) = hmac_keys.current {
local.with_hmac_key(key)
} else {
local
};
let local = if let Some(key) = hmac_keys.previous {
local.with_previous_hmac_key(key)
} else {
local
};
let local = if let Some(slot) = hwm_keys.current {
local.with_hwm_key(slot.epoch, slot.key)
} else {
local
};
let local = if let Some(slot) = hwm_keys.previous {
local.with_previous_hwm_key(slot.epoch, slot.key)
} else {
local
};
let local = std::sync::Arc::new(local);
let backend = std::sync::Arc::new(zeph_durable::DurableBackendEnum::Local(local.clone()));
let durable_cfg = std::sync::Arc::new(config.durable.clone());
let (writer_actor, writer_handle) = zeph_durable::JournalWriter::new(local, &durable_cfg);
{
let writer_fut = writer_actor.run();
let cell = std::sync::Arc::new(parking_lot::Mutex::new(Some(writer_fut)));
sched_supervisor.spawn(zeph_common::TaskDescriptor {
name: "journal_writer",
restart: zeph_common::RestartPolicy::RunOnce,
factory: move || {
let f = cell.lock().take();
async move {
if let Some(f) = f {
f.await;
}
}
},
});
}
{
let retention_backend = std::sync::Arc::clone(&backend);
let retention_policy = durable_cfg.retention.clone();
sched_supervisor.spawn(zeph_common::TaskDescriptor {
name: "durable.retention_sweep",
restart: zeph_common::RestartPolicy::Restart {
max: 5,
base_delay: std::time::Duration::from_secs(5),
},
factory: move || {
zeph_durable::DurableRetentionService::new(
std::sync::Arc::clone(&retention_backend),
retention_policy.clone(),
)
.run()
},
});
}
Ok(Some(zeph_scheduler::durable::SchedulerDurableAdapter::new(
backend,
writer_handle,
durable_cfg,
)))
}
fn print_status_human(status: &zeph_scheduler::DaemonStatus) {
let running = if status.running {
"running"
} else {
"not running"
};
let pid_str = status
.pid
.map(|p| format!(" (pid {p})"))
.unwrap_or_default();
println!("daemon: {running}{pid_str}");
println!("pid_file: {}", status.pid_file.display());
println!("log_file: {}", status.log_file.display());
println!("tasks: {}", status.task_count);
if !status.recent_runs.is_empty() {
println!("last runs:");
for run in &status.recent_runs {
let last_run = if run.last_run.is_empty() {
"never"
} else {
run.last_run.as_str()
};
println!(
" {:<24} last: {:<25} next: {}",
run.name, last_run, run.next_run
);
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use zeph_durable::Journal as _;
#[test]
fn build_detach_args_forwards_all_three_vault_overrides_when_present() {
let args = build_detach_args(
"/some/config.toml",
true,
Some("age"),
Some(std::path::Path::new("/override/vault-key.txt")),
Some(std::path::Path::new("/override/secrets.age")),
);
assert!(args.iter().any(|a| a == "--vault"));
assert!(args.iter().any(|a| a == "age"));
assert!(args.iter().any(|a| a == "--vault-key"));
assert!(args.iter().any(|a| a == "/override/vault-key.txt"));
assert!(args.iter().any(|a| a == "--vault-path"));
assert!(args.iter().any(|a| a == "/override/secrets.age"));
let key_idx = args.iter().position(|a| a == "--vault-key").unwrap();
assert_eq!(args[key_idx + 1], "/override/vault-key.txt");
let path_idx = args.iter().position(|a| a == "--vault-path").unwrap();
assert_eq!(args[path_idx + 1], "/override/secrets.age");
}
#[test]
fn build_detach_args_omits_vault_flags_when_no_override_given() {
let args = build_detach_args("/some/config.toml", true, None, None, None);
assert!(
!args.iter().any(|a| a == "--vault"),
"no --vault override must not appear in the child argv: {args:?}"
);
assert!(
!args.iter().any(|a| a == "--vault-key"),
"no --vault-key override must not appear in the child argv: {args:?}"
);
assert!(
!args.iter().any(|a| a == "--vault-path"),
"no --vault-path override must not appear in the child argv: {args:?}"
);
}
#[test]
fn build_detach_args_forwards_single_override_independently() {
let args = build_detach_args(
"/some/config.toml",
true,
None,
Some(std::path::Path::new("/override/vault-key.txt")),
None,
);
assert!(args.iter().any(|a| a == "--vault-key"));
assert!(!args.iter().any(|a| a == "--vault-path"));
assert!(!args.iter().any(|a| a == "--vault"));
}
#[test]
fn build_detach_args_always_includes_config_and_foreground() {
let args = build_detach_args("/x/config.toml", true, None, None, None);
assert_eq!(args[0], "--config");
assert_eq!(args[1], "/x/config.toml");
assert!(args.iter().any(|a| a == "serve"));
assert!(args.iter().any(|a| a == "--foreground"));
}
#[test]
fn build_detach_args_no_catch_up_flag_added_when_catch_up_false() {
let with_catch_up = build_detach_args("/x/config.toml", true, None, None, None);
assert!(!with_catch_up.iter().any(|a| a == "--no-catch-up"));
let without_catch_up = build_detach_args("/x/config.toml", false, None, None, None);
assert!(without_catch_up.iter().any(|a| a == "--no-catch-up"));
}
#[tokio::test]
async fn spawns_retention_sweep_alongside_journal_writer() {
let dir = tempfile::tempdir().unwrap();
let mut config = zeph_core::config::Config::default();
config.memory.sqlite_path = dir.path().join("zeph.db").to_string_lossy().into_owned();
config.durable.enabled = true;
config.durable.scheduler = true;
config.durable.encrypt_payload = false;
let cancel = tokio_util::sync::CancellationToken::new();
let sched_supervisor = zeph_common::TaskSupervisor::new(cancel);
let vault_dir = zeph_core::vault::default_vault_dir();
let adapter = build_durable_adapter(
&config,
&sched_supervisor,
&vault_dir.join("vault-key.txt"),
&vault_dir.join("secrets.age"),
)
.await
.expect("build_durable_adapter must succeed with a local, unencrypted backend");
assert!(
adapter.is_some(),
"durable.enabled && durable.scheduler must produce an adapter"
);
let names: Vec<String> = sched_supervisor
.snapshot()
.iter()
.map(|s| s.name.to_string())
.collect();
assert!(
names.contains(&"journal_writer".to_owned()),
"expected journal_writer among supervised tasks, got {names:?}"
);
assert!(
names.contains(&"durable.retention_sweep".to_owned()),
"expected durable.retention_sweep among supervised tasks, got {names:?}"
);
}
#[tokio::test]
async fn does_not_spawn_when_scheduler_adapter_disabled() {
let dir = tempfile::tempdir().unwrap();
let mut config = zeph_core::config::Config::default();
config.memory.sqlite_path = dir.path().join("zeph.db").to_string_lossy().into_owned();
config.durable.enabled = true;
config.durable.scheduler = false;
let cancel = tokio_util::sync::CancellationToken::new();
let sched_supervisor = zeph_common::TaskSupervisor::new(cancel);
let vault_dir = zeph_core::vault::default_vault_dir();
let adapter = build_durable_adapter(
&config,
&sched_supervisor,
&vault_dir.join("vault-key.txt"),
&vault_dir.join("secrets.age"),
)
.await
.expect("build_durable_adapter must succeed (returns None, not an error)");
assert!(adapter.is_none());
assert!(sched_supervisor.snapshot().is_empty());
}
#[allow(unsafe_code)]
#[tokio::test]
#[serial_test::serial]
async fn scheduler_daemon_reads_previous_key_control_entry_through_rotation_window() {
let vault_dir = tempfile::tempdir().unwrap();
let prev_xdg = std::env::var("XDG_CONFIG_HOME").ok();
unsafe {
std::env::set_var("XDG_CONFIG_HOME", vault_dir.path());
}
let vault_root = zeph_core::vault::default_vault_dir();
zeph_core::vault::AgeVaultProvider::init_vault(&vault_root).unwrap();
let mut provider = zeph_core::vault::AgeVaultProvider::load(
&vault_root.join("vault-key.txt"),
&vault_root.join("secrets.age"),
)
.unwrap();
let old_key_b64 = zeph_core::durable::generate_durable_key_b64();
let new_key_b64 = zeph_core::durable::generate_durable_key_b64();
provider
.set_secret_mut(
"ZEPH_DURABLE_KEY_PREVIOUS".to_owned(),
old_key_b64.clone(),
false,
)
.unwrap();
provider
.set_secret_mut("ZEPH_DURABLE_KEY".to_owned(), new_key_b64, false)
.unwrap();
provider.save().unwrap();
let dir = tempfile::tempdir().unwrap();
let mut config = zeph_core::config::Config::default();
config.memory.sqlite_path = dir.path().join("zeph.db").to_string_lossy().into_owned();
config.durable.enabled = true;
config.durable.scheduler = true;
config.durable.shared_db = true;
config.durable.key_id = 1;
config.durable.previous_key_id = Some(0);
let old_hmac_key = zeph_core::durable::derive_control_hmac_key_b64(&old_key_b64).unwrap();
let durable_url = crate::commands::durable::resolve_durable_db_url(&config);
let exec = zeph_durable::ExecutionId::new();
{
let pre_rotation_writer =
zeph_durable::LocalBackend::open(&durable_url, config.durable.max_payload_bytes)
.await
.unwrap()
.with_hmac_key(old_hmac_key);
pre_rotation_writer.init().await.unwrap();
pre_rotation_writer
.open_execution(exec, zeph_durable::ExecutionKind::AgentTurn)
.await
.unwrap();
let step_id = zeph_durable::StepId::new(0);
pre_rotation_writer
.append(zeph_durable::JournalEntry {
seq: None,
execution_id: exec,
kind: zeph_durable::ExecutionKind::AgentTurn,
step_id,
entry: zeph_durable::EntryKind::EffectIntent {
idempotency_key: zeph_durable::IdempotencyKey::derive(
exec,
step_id,
b"transfer",
),
effect: zeph_durable::EffectClass::ExactlyOnceGuarded,
hmac: None,
},
created_at_ms: 0,
})
.await
.unwrap();
}
let cancel = tokio_util::sync::CancellationToken::new();
let sched_supervisor = zeph_common::TaskSupervisor::new(cancel);
let adapter = build_durable_adapter(
&config,
&sched_supervisor,
&vault_root.join("vault-key.txt"),
&vault_root.join("secrets.age"),
)
.await;
assert!(
adapter.is_ok() && adapter.unwrap().is_some(),
"build_durable_adapter must succeed with a consistent rotation window declared"
);
let hmac_keys = crate::commands::durable::load_write_hmac_key(
&config,
&vault_root.join("vault-key.txt"),
&vault_root.join("secrets.age"),
)
.unwrap();
let reader =
zeph_durable::LocalBackend::open(&durable_url, config.durable.max_payload_bytes)
.await
.unwrap()
.with_hmac_key(hmac_keys.current.unwrap())
.with_previous_hmac_key(hmac_keys.previous.unwrap());
let read_result = reader.read_execution(exec).await;
unsafe {
match &prev_xdg {
Some(v) => std::env::set_var("XDG_CONFIG_HOME", v),
None => std::env::remove_var("XDG_CONFIG_HOME"),
}
}
assert!(
read_result.is_ok(),
"the scheduler-daemon read channel must verify a pre-rotation EffectIntent control \
entry through the rotation window"
);
}
#[allow(unsafe_code)]
#[tokio::test]
#[serial_test::serial]
async fn scheduler_daemon_resumes_previous_epoch_high_water_mark_through_rotation_window() {
let vault_dir = tempfile::tempdir().unwrap();
let prev_xdg = std::env::var("XDG_CONFIG_HOME").ok();
unsafe {
std::env::set_var("XDG_CONFIG_HOME", vault_dir.path());
}
let vault_root = zeph_core::vault::default_vault_dir();
zeph_core::vault::AgeVaultProvider::init_vault(&vault_root).unwrap();
let mut provider = zeph_core::vault::AgeVaultProvider::load(
&vault_root.join("vault-key.txt"),
&vault_root.join("secrets.age"),
)
.unwrap();
let old_key_b64 = zeph_core::durable::generate_durable_key_b64();
let new_key_b64 = zeph_core::durable::generate_durable_key_b64();
provider
.set_secret_mut(
"ZEPH_DURABLE_KEY_PREVIOUS".to_owned(),
old_key_b64.clone(),
false,
)
.unwrap();
provider
.set_secret_mut("ZEPH_DURABLE_KEY".to_owned(), new_key_b64, false)
.unwrap();
provider.save().unwrap();
let dir = tempfile::tempdir().unwrap();
let mut config = zeph_core::config::Config::default();
config.memory.sqlite_path = dir.path().join("zeph.db").to_string_lossy().into_owned();
config.durable.enabled = true;
config.durable.scheduler = true;
config.durable.key_id = 1;
config.durable.previous_key_id = Some(0);
let old_hwm_key = zeph_core::durable::derive_hwm_key_b64(&old_key_b64).unwrap();
let durable_url = crate::commands::durable::resolve_durable_db_url(&config);
let exec = zeph_durable::ExecutionId::new();
{
let pre_rotation_writer =
zeph_durable::LocalBackend::open(&durable_url, config.durable.max_payload_bytes)
.await
.unwrap()
.with_hwm_key(0, old_hwm_key);
pre_rotation_writer.init().await.unwrap();
pre_rotation_writer
.open_execution(exec, zeph_durable::ExecutionKind::AgentTurn)
.await
.unwrap();
let step_id = zeph_durable::StepId::new(0);
pre_rotation_writer
.append(zeph_durable::JournalEntry {
seq: None,
execution_id: exec,
kind: zeph_durable::ExecutionKind::AgentTurn,
step_id,
entry: zeph_durable::EntryKind::StepResult {
idempotency_key: zeph_durable::IdempotencyKey::derive(
exec,
step_id,
b"tool:test",
),
payload: bytes::Bytes::copy_from_slice(b"pre-rotation result"),
effect: zeph_durable::EffectClass::Idempotent,
payload_version: 1,
},
created_at_ms: 0,
})
.await
.unwrap();
}
let cancel = tokio_util::sync::CancellationToken::new();
let sched_supervisor = zeph_common::TaskSupervisor::new(cancel);
let adapter = build_durable_adapter(
&config,
&sched_supervisor,
&vault_root.join("vault-key.txt"),
&vault_root.join("secrets.age"),
)
.await;
assert!(
adapter.is_ok() && adapter.unwrap().is_some(),
"build_durable_adapter must succeed with a consistent rotation window declared"
);
let hwm_keys = crate::commands::durable::load_write_hwm_key(
&config,
&vault_root.join("vault-key.txt"),
&vault_root.join("secrets.age"),
)
.unwrap();
let current = hwm_keys.current.expect("current HWM slot must resolve");
let previous = hwm_keys
.previous
.expect("previous HWM slot must resolve while the window is open");
let reader =
zeph_durable::LocalBackend::open(&durable_url, config.durable.max_payload_bytes)
.await
.unwrap()
.with_hwm_key(current.epoch, current.key)
.with_previous_hwm_key(previous.epoch, previous.key);
let resume_result = reader
.open_execution(exec, zeph_durable::ExecutionKind::AgentTurn)
.await;
unsafe {
match &prev_xdg {
Some(v) => std::env::set_var("XDG_CONFIG_HOME", v),
None => std::env::remove_var("XDG_CONFIG_HOME"),
}
}
assert!(
resume_result.is_ok(),
"the scheduler-daemon read channel must resume a pre-rotation execution's \
high-water-mark through the rotation window: {resume_result:?}"
);
}
}