use std::future::Future;
use std::sync::Arc;
use tracing::Instrument as _;
use zeph_config::DurableConfig;
use zeph_durable::{
DurableBackendEnum, DurableContext, DurableError, EffectIntentSubClass, ExecutionId,
ExecutionKind, ExecutionStatus, JournalWriterHandle, LocalBackend, StepDescriptor, StepError,
};
use crate::error::SchedulerError;
#[must_use]
pub fn derive_execution_id(job_name: &str, slot_ms: i64) -> ExecutionId {
let name_bytes = job_name.as_bytes();
let mut payload = Vec::with_capacity(8 + name_bytes.len() + 1 + 8);
payload.extend_from_slice(&(name_bytes.len() as u64).to_le_bytes());
payload.extend_from_slice(name_bytes);
payload.push(0u8);
payload.extend_from_slice(&slot_ms.to_le_bytes());
ExecutionId::derive(b"zeph.scheduler.fire.v1", &payload)
}
#[derive(Clone)]
pub struct SchedulerDurableAdapter {
backend: Arc<DurableBackendEnum>,
writer: JournalWriterHandle,
config: Arc<DurableConfig>,
}
impl SchedulerDurableAdapter {
#[must_use]
pub fn new(
backend: Arc<DurableBackendEnum>,
writer: JournalWriterHandle,
config: Arc<DurableConfig>,
) -> Self {
Self {
backend,
writer,
config,
}
}
}
pub async fn fire_with_durable<F, Fut>(
adapter: &SchedulerDurableAdapter,
job_name: &str,
slot_ms: i64,
fire: F,
) -> Result<(), SchedulerError>
where
F: FnOnce() -> Fut + Send + 'static,
Fut: Future<Output = Result<(), SchedulerError>> + Send + 'static,
{
let span = tracing::info_span!("sched.durable.fire", job = job_name, slot_ms = slot_ms,);
async move {
let exec_id = derive_execution_id(job_name, slot_ms);
let local_backend: Arc<LocalBackend> = match &*adapter.backend {
DurableBackendEnum::Local(lb) => lb.clone(),
_ => {
return Err(SchedulerError::TaskFailed(
"durable scheduler adapter requires a LocalBackend".into(),
));
}
};
let (is_resume, _lock) = match local_backend
.open_execution_exclusive(exec_id, ExecutionKind::ScheduledJob)
.await
{
Ok(result) => result,
Err(DurableError::ExecutionLocked {
execution_id,
holder_pid,
}) => {
tracing::info!(
execution_id = %execution_id,
holder_pid,
job = job_name,
slot_ms,
"sched.durable.fire: execution already open in another process; skipping \
this fire"
);
return Ok(());
}
Err(e) => {
return Err(SchedulerError::TaskFailed(format!(
"durable open failed: {e}"
)));
}
};
let ctx = DurableContext::new(
exec_id,
ExecutionKind::ScheduledJob,
is_resume,
adapter.backend.clone(),
adapter.writer.clone(),
&adapter.config,
);
let desc = StepDescriptor::exactly_once_guarded(
"scheduler_fire",
EffectIntentSubClass::CostBearingOrBoundaryIdempotent,
None,
fingerprint(job_name, slot_ms),
)
.map_err(|e| SchedulerError::TaskFailed(format!("durable descriptor error: {e}")))?;
let step_result = ctx
.step(desc, |_handle| async move {
fire().await.map_err(|e| StepError::new(e.to_string()))
})
.await;
let terminal_status = if step_result.is_ok() {
ExecutionStatus::Completed
} else {
ExecutionStatus::Failed
};
if let Err(finalize_err) = ctx.finalize(terminal_status).await {
tracing::warn!(
error = %finalize_err,
job = job_name,
slot_ms,
"sched.durable.fire: failed to finalize scheduled-job execution"
);
}
step_result.map_err(|e| SchedulerError::TaskFailed(format!("durable step failed: {e}")))
}
.instrument(span)
.await
}
fn fingerprint(job_name: &str, slot_ms: i64) -> Vec<u8> {
let name_bytes = job_name.as_bytes();
let mut buf = Vec::with_capacity(8 + name_bytes.len() + 8);
buf.extend_from_slice(&(name_bytes.len() as u64).to_le_bytes());
buf.extend_from_slice(name_bytes);
buf.extend_from_slice(&slot_ms.to_le_bytes());
buf
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn derive_execution_id_is_stable() {
let a = derive_execution_id("weekly-digest", 1_700_000_000_000);
let b = derive_execution_id("weekly-digest", 1_700_000_000_000);
assert_eq!(a, b, "same inputs must produce the same ExecutionId");
}
#[test]
fn derive_execution_id_differs_by_slot() {
let a = derive_execution_id("weekly-digest", 1_700_000_000_000);
let c = derive_execution_id("weekly-digest", 1_700_000_001_000);
assert_ne!(
a, c,
"different slot_ms must produce a different ExecutionId"
);
}
#[test]
fn derive_execution_id_differs_by_name() {
let a = derive_execution_id("job-a", 1_000);
let b = derive_execution_id("job-b", 1_000);
assert_ne!(a, b, "different job names must produce different ids");
}
#[test]
fn derive_execution_id_no_prefix_collision() {
let a = derive_execution_id("ab", 1_000);
let b = derive_execution_id("abc", 1_000);
assert_ne!(
a, b,
"length-delimited framing must prevent prefix collisions"
);
}
#[test]
fn no_adapter_means_no_durable_overhead() {
assert_ne!(
derive_execution_id("job", 0),
derive_execution_id("other", 0)
);
}
#[cfg(feature = "sqlite")]
mod with_backend {
use std::sync::atomic::{AtomicU32, Ordering};
use zeph_durable::config::DurableConfig;
use zeph_durable::{DurableBackendEnum, JournalWriter, LocalBackend};
use super::*;
fn fast_config() -> DurableConfig {
DurableConfig {
journal_flush_interval_ms: 5,
journal_ack_timeout_ms: 2000,
..DurableConfig::default()
}
}
async fn make_adapter() -> (SchedulerDurableAdapter, tokio::task::JoinHandle<()>) {
let (_local, adapter, task) = make_adapter_with_local(fast_config()).await;
(adapter, task)
}
async fn make_adapter_with_local(
config: DurableConfig,
) -> (
Arc<LocalBackend>,
SchedulerDurableAdapter,
tokio::task::JoinHandle<()>,
) {
let local = Arc::new(LocalBackend::open(":memory:", 1_048_576).await.unwrap());
local.init().await.unwrap();
let backend = Arc::new(DurableBackendEnum::Local(local.clone()));
let cfg = Arc::new(config);
let (writer, handle) = JournalWriter::new(local.clone(), &cfg);
let task = tokio::spawn(writer.run()); (
local,
SchedulerDurableAdapter::new(backend, handle, cfg),
task,
)
}
async fn execution_status(local: &LocalBackend, exec: ExecutionId) -> String {
let (status,): (String,) = zeph_db::query_as(zeph_db::sql!(
"SELECT status FROM durable_executions WHERE execution_id = ?"
))
.bind(exec.as_uuid().to_string())
.fetch_one(local.pool())
.await
.unwrap();
status
}
#[tokio::test]
async fn fire_with_durable_replays_without_reinvocation() {
let (adapter, _task) = make_adapter().await;
let count = Arc::new(AtomicU32::new(0));
let c = count.clone();
fire_with_durable(&adapter, "test-job", 1_000, move || async move {
c.fetch_add(1, Ordering::Relaxed);
Ok(())
})
.await
.unwrap();
assert_eq!(count.load(Ordering::Relaxed), 1, "first fire must execute");
let c = count.clone();
fire_with_durable(&adapter, "test-job", 1_000, move || async move {
c.fetch_add(1, Ordering::Relaxed);
Ok(())
})
.await
.unwrap();
assert_eq!(
count.load(Ordering::Relaxed),
1,
"replay must not re-invoke the closure"
);
}
#[tokio::test]
async fn fire_with_durable_different_slot_runs_again() {
let (adapter, _task) = make_adapter().await;
let count = Arc::new(AtomicU32::new(0));
let c = count.clone();
fire_with_durable(&adapter, "test-job", 1_000, move || async move {
c.fetch_add(1, Ordering::Relaxed);
Ok(())
})
.await
.unwrap();
let c = count.clone();
fire_with_durable(&adapter, "test-job", 2_000, move || async move {
c.fetch_add(1, Ordering::Relaxed);
Ok(())
})
.await
.unwrap();
assert_eq!(
count.load(Ordering::Relaxed),
2,
"different slot_ms must produce a separate execution and run"
);
}
#[tokio::test]
async fn fire_with_durable_finalizes_the_execution_as_completed() {
let (local, adapter, _task) = make_adapter_with_local(fast_config()).await;
fire_with_durable(&adapter, "test-job", 1_000, || async { Ok(()) })
.await
.unwrap();
let exec = derive_execution_id("test-job", 1_000);
assert_eq!(execution_status(&local, exec).await, "completed");
}
#[tokio::test]
async fn fire_with_durable_finalizes_the_execution_as_failed_on_fire_error() {
let (local, adapter, _task) = make_adapter_with_local(fast_config()).await;
let err = fire_with_durable(&adapter, "test-job", 1_000, || async {
Err(SchedulerError::TaskFailed("boom".to_string()))
})
.await
.expect_err("the fire body's error must propagate");
assert!(
matches!(&err, SchedulerError::TaskFailed(_)),
"unexpected error: {err:?}"
);
let exec = derive_execution_id("test-job", 1_000);
assert_eq!(execution_status(&local, exec).await, "failed");
}
#[tokio::test]
async fn fire_with_durable_skips_gracefully_when_execution_is_locked() {
let dir = tempfile::tempdir().unwrap();
let url = dir.path().join("durable.db").to_string_lossy().into_owned();
let owner = LocalBackend::open(&url, 1_048_576).await.unwrap();
owner.init().await.unwrap();
let contender = Arc::new(LocalBackend::open(&url, 1_048_576).await.unwrap());
let backend = Arc::new(DurableBackendEnum::Local(contender.clone()));
let cfg = Arc::new(fast_config());
let (writer, handle) = JournalWriter::new(contender.clone(), &cfg);
let _task = tokio::spawn(writer.run()); let adapter = SchedulerDurableAdapter::new(backend, handle, cfg);
let exec_id = derive_execution_id("test-job", 1_000);
let (_, _lock) = owner
.open_execution_exclusive(exec_id, ExecutionKind::ScheduledJob)
.await
.unwrap();
let count = Arc::new(AtomicU32::new(0));
let c = count.clone();
let result = fire_with_durable(&adapter, "test-job", 1_000, move || async move {
c.fetch_add(1, Ordering::Relaxed);
Ok(())
})
.await;
assert!(
result.is_ok(),
"ExecutionLocked must degrade to a graceful Ok(()) skip, got {result:?}"
);
assert_eq!(
count.load(Ordering::Relaxed),
0,
"a locked execution must not re-invoke the fire body"
);
assert_eq!(execution_status(&owner, exec_id).await, "running");
}
}
}