use std::future::Future;
use std::sync::Arc;
use tracing::Instrument as _;
use zeph_config::DurableConfig;
use zeph_durable::{
DurableBackendEnum, DurableContext, EffectIntentSubClass, ExecutionId, ExecutionKind,
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 = local_backend
.open_execution(exec_id, ExecutionKind::ScheduledJob)
.await
.map_err(|e| 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}")))?;
ctx.step(desc, |_handle| async move {
fire().await.map_err(|e| StepError::new(e.to_string()))
})
.await
.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!(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 = 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(fast_config());
let (writer, handle) = JournalWriter::new(local, &cfg);
let task = tokio::spawn(writer.run()); (SchedulerDurableAdapter::new(backend, handle, cfg), task)
}
#[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"
);
}
}
}