use std::collections::{BTreeMap, BTreeSet};
use serde::{Deserialize, Serialize};
use crate::claude_runtime_state::{
ClaudeCronJob, ClaudeRuntimeExecutionState, ClaudeRuntimeManifest, ClaudeWakeup,
};
use crate::{Error, Result};
fn default_timezone() -> String {
"UTC".to_string()
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct ClaudeRuntimeSchedulerState {
#[serde(default)]
pub activated_at_unix: Option<i64>,
#[serde(default = "default_timezone")]
pub timezone: String,
#[serde(default)]
pub cron_jobs: Vec<ClaudeCronScheduleState>,
#[serde(default)]
pub wakeups: Vec<ClaudeWakeupScheduleState>,
#[serde(default)]
pub deliveries: Vec<ClaudeRuntimeDeliveryState>,
}
impl Default for ClaudeRuntimeSchedulerState {
fn default() -> Self {
Self {
activated_at_unix: None,
timezone: default_timezone(),
cron_jobs: Vec::new(),
wakeups: Vec::new(),
deliveries: Vec::new(),
}
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct ClaudeCronScheduleState {
pub id: String,
pub next_due_unix: i64,
#[serde(default)]
pub expires_at_unix: Option<i64>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct ClaudeWakeupScheduleState {
pub tool_use_id: String,
pub due_unix: i64,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum ClaudeRuntimeTriggerKind {
Queue,
Cron,
Wakeup,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct ClaudeRuntimeTrigger {
pub due_unix: i64,
pub kind: ClaudeRuntimeTriggerKind,
pub id: String,
pub prompt: Option<String>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct ClaudeRuntimeDeliveryState {
pub trigger: ClaudeRuntimeTrigger,
pub retry_at_unix: i64,
}
impl ClaudeRuntimeTrigger {
fn order_key(&self) -> (i64, ClaudeRuntimeTriggerKind, &str) {
(self.due_unix, self.kind, self.id.as_str())
}
}
impl ClaudeRuntimeManifest {
pub fn activate_scheduler(&mut self, now_unix: i64) -> Result<()> {
if self.execution_state == ClaudeRuntimeExecutionState::Active {
return Err(Error::Other(
"Claude runtime scheduler is already active".into(),
));
}
let (scheduler, retained_crons) = build_scheduler(self, now_unix, false)?;
self.active_crons = retained_crons;
self.scheduler = scheduler;
self.execution_state = ClaudeRuntimeExecutionState::Active;
Ok(())
}
pub fn reconcile_scheduler(&mut self, now_unix: i64) -> Result<()> {
self.require_active()?;
let (scheduler, retained_crons) = build_scheduler(self, now_unix, true)?;
self.active_crons = retained_crons;
self.scheduler = scheduler;
Ok(())
}
pub fn next_due(&self, now_unix: i64) -> Result<Option<ClaudeRuntimeTrigger>> {
self.require_active()?;
let mut due = scheduler_triggers(self, now_unix, false)?;
due.extend(self.scheduler.deliveries.iter().map(|delivery| {
let mut trigger = delivery.trigger.clone();
trigger.due_unix = delivery.retry_at_unix;
trigger
}));
due.sort_by(|a, b| a.order_key().cmp(&b.order_key()));
Ok(due.into_iter().next())
}
pub fn claim_due(&mut self, now_unix: i64) -> Result<Vec<ClaudeRuntimeTrigger>> {
self.require_active()?;
let mut claimed: Vec<_> = self
.scheduler
.deliveries
.iter()
.filter(|delivery| delivery.retry_at_unix <= now_unix)
.map(|delivery| delivery.trigger.clone())
.collect();
let mut newly_claimed = scheduler_triggers(self, now_unix, true)?;
newly_claimed.sort_by(|a, b| a.order_key().cmp(&b.order_key()));
let claimed_crons: BTreeSet<String> = newly_claimed
.iter()
.filter(|item| item.kind == ClaudeRuntimeTriggerKind::Cron)
.map(|item| item.id.clone())
.collect();
let claimed_queue: BTreeSet<String> = newly_claimed
.iter()
.filter(|item| item.kind == ClaudeRuntimeTriggerKind::Queue)
.map(|item| item.id.clone())
.collect();
let claimed_wakeups: BTreeSet<String> = newly_claimed
.iter()
.filter(|item| item.kind == ClaudeRuntimeTriggerKind::Wakeup)
.map(|item| item.id.clone())
.collect();
let mut next_crons = Vec::new();
let mut retained_jobs = Vec::new();
for job in std::mem::take(&mut self.active_crons) {
let Some(cursor) = self.scheduler.cron_jobs.iter().find(|c| c.id == job.id) else {
return Err(Error::Other(format!(
"active Claude cron `{}` has no scheduler cursor",
job.id
)));
};
if cursor
.expires_at_unix
.is_some_and(|expiry| now_unix >= expiry)
{
continue;
}
if !claimed_crons.contains(&job.id) {
next_crons.push(cursor.clone());
retained_jobs.push(job);
continue;
}
if !job.recurring {
continue;
}
let parsed = CronSchedule::parse(&job.schedule)?;
let next = parsed.next_after(now_unix)?;
if cursor.expires_at_unix.is_some_and(|expiry| next >= expiry) {
continue;
}
next_crons.push(ClaudeCronScheduleState {
id: job.id.clone(),
next_due_unix: next,
expires_at_unix: cursor.expires_at_unix,
});
retained_jobs.push(job);
}
self.active_crons = retained_jobs;
self.scheduler.cron_jobs = next_crons;
self.pending_wakeups
.retain(|wakeup| !claimed_wakeups.contains(&wakeup.tool_use_id));
self.scheduler
.wakeups
.retain(|wakeup| !claimed_wakeups.contains(&wakeup.tool_use_id));
if !claimed_queue.is_empty() {
let expected: BTreeSet<String> = (0..self.queue.pending.len())
.map(|index| queue_trigger_id(&self.queue, index))
.collect::<Result<_>>()?;
if claimed_queue != expected {
return Err(Error::Other(
"Claude queue claim did not cover the complete pending FIFO".into(),
));
}
self.queue.pending.clear();
}
for trigger in &newly_claimed {
self.scheduler.deliveries.push(ClaudeRuntimeDeliveryState {
trigger: trigger.clone(),
retry_at_unix: now_unix,
});
}
self.scheduler
.deliveries
.sort_by(|a, b| a.trigger.order_key().cmp(&b.trigger.order_key()));
claimed.extend(newly_claimed);
claimed.sort_by(|a, b| a.order_key().cmp(&b.order_key()));
Ok(claimed)
}
pub fn complete_delivery(&mut self, kind: ClaudeRuntimeTriggerKind, id: &str) -> Result<()> {
self.require_active()?;
let before = self.scheduler.deliveries.len();
self.scheduler
.deliveries
.retain(|delivery| delivery.trigger.kind != kind || delivery.trigger.id != id);
if self.scheduler.deliveries.len() == before {
return Err(Error::Other(format!(
"unknown Claude runtime delivery `{kind:?}` `{id}`"
)));
}
Ok(())
}
pub fn defer_delivery(
&mut self,
kind: ClaudeRuntimeTriggerKind,
id: &str,
retry_at_unix: i64,
) -> Result<()> {
self.require_active()?;
let Some(delivery) = self
.scheduler
.deliveries
.iter_mut()
.find(|delivery| delivery.trigger.kind == kind && delivery.trigger.id == id)
else {
return Err(Error::Other(format!(
"unknown Claude runtime delivery `{kind:?}` `{id}`"
)));
};
delivery.retry_at_unix = retry_at_unix;
Ok(())
}
fn require_active(&self) -> Result<()> {
if self.execution_state != ClaudeRuntimeExecutionState::Active {
return Err(Error::Other(
"Claude runtime scheduler is paused; activate it explicitly first".into(),
));
}
if self.scheduler.activated_at_unix.is_none() {
return Err(Error::Other(
"active Claude runtime manifest has no scheduler activation metadata".into(),
));
}
Ok(())
}
}
fn build_scheduler(
manifest: &ClaudeRuntimeManifest,
now_unix: i64,
preserve_existing: bool,
) -> Result<(ClaudeRuntimeSchedulerState, Vec<ClaudeCronJob>)> {
for index in 0..manifest.queue.pending.len() {
queue_trigger_id(&manifest.queue, index)?;
}
let old_crons: BTreeMap<&str, &ClaudeCronScheduleState> = manifest
.scheduler
.cron_jobs
.iter()
.map(|cursor| (cursor.id.as_str(), cursor))
.collect();
let old_wakeups: BTreeMap<&str, &ClaudeWakeupScheduleState> = manifest
.scheduler
.wakeups
.iter()
.map(|cursor| (cursor.tool_use_id.as_str(), cursor))
.collect();
let mut seen = BTreeSet::new();
let mut cron_jobs = Vec::new();
let mut retained_crons = Vec::new();
for job in &manifest.active_crons {
if !seen.insert(job.id.as_str()) {
return Err(Error::Other(format!(
"duplicate active Claude cron id `{}`",
job.id
)));
}
let parsed = CronSchedule::parse(&job.schedule)?;
let expires_at_unix = expiry_for(job)?;
if expires_at_unix.is_some_and(|expiry| now_unix >= expiry) {
continue;
}
let next_due_unix = if preserve_existing {
old_crons
.get(job.id.as_str())
.map(|cursor| cursor.next_due_unix)
.unwrap_or(parsed.next_after(now_unix)?)
} else {
parsed.next_after(now_unix)?
};
if expires_at_unix.is_some_and(|expiry| next_due_unix >= expiry) {
continue;
}
cron_jobs.push(ClaudeCronScheduleState {
id: job.id.clone(),
next_due_unix,
expires_at_unix,
});
retained_crons.push(job.clone());
}
let mut wakeups = Vec::new();
let mut seen_wakeups = BTreeSet::new();
for wakeup in &manifest.pending_wakeups {
if !seen_wakeups.insert(wakeup.tool_use_id.as_str()) {
return Err(Error::Other(format!(
"duplicate pending Claude wakeup id `{}`",
wakeup.tool_use_id
)));
}
let due_unix = if preserve_existing {
old_wakeups
.get(wakeup.tool_use_id.as_str())
.map(|cursor| cursor.due_unix)
.unwrap_or(wakeup_due(wakeup, now_unix)?)
} else {
wakeup_due(wakeup, now_unix)?
};
wakeups.push(ClaudeWakeupScheduleState {
tool_use_id: wakeup.tool_use_id.clone(),
due_unix,
});
}
cron_jobs.sort_by(|a, b| a.id.cmp(&b.id));
wakeups.sort_by(|a, b| a.tool_use_id.cmp(&b.tool_use_id));
Ok((
ClaudeRuntimeSchedulerState {
activated_at_unix: Some(
manifest
.scheduler
.activated_at_unix
.filter(|_| preserve_existing)
.unwrap_or(now_unix),
),
timezone: default_timezone(),
cron_jobs,
wakeups,
deliveries: if preserve_existing {
manifest.scheduler.deliveries.clone()
} else {
Vec::new()
},
},
retained_crons,
))
}
fn expiry_for(job: &ClaudeCronJob) -> Result<Option<i64>> {
let Some(seconds) = job.expires_after_seconds else {
return Ok(None);
};
let created = parse_created(job.created_at.as_deref(), "cron", &job.id)?;
let seconds = i64::try_from(seconds).map_err(|_| {
Error::Other(format!(
"Claude cron `{}` expiry exceeds supported Unix time",
job.id
))
})?;
created.checked_add(seconds).map(Some).ok_or_else(|| {
Error::Other(format!(
"Claude cron `{}` expiry overflows Unix time",
job.id
))
})
}
fn wakeup_due(wakeup: &ClaudeWakeup, now_unix: i64) -> Result<i64> {
let created = parse_created(wakeup.created_at.as_deref(), "wakeup", &wakeup.tool_use_id)?;
let delay = i64::try_from(wakeup.delay_seconds).map_err(|_| {
Error::Other(format!(
"Claude wakeup `{}` delay exceeds supported Unix time",
wakeup.tool_use_id
))
})?;
let due = created.checked_add(delay).ok_or_else(|| {
Error::Other(format!(
"Claude wakeup `{}` due time overflows Unix time",
wakeup.tool_use_id
))
})?;
Ok(due.max(now_unix))
}
fn parse_created(value: Option<&str>, kind: &str, id: &str) -> Result<i64> {
let value = value.ok_or_else(|| {
Error::Other(format!(
"Claude {kind} `{id}` has no creation timestamp for deterministic activation"
))
})?;
crate::sidecar::rfc3339_to_ms(value)
.map(|ms| ms.div_euclid(1000))
.ok_or_else(|| {
Error::Other(format!(
"Claude {kind} `{id}` has invalid RFC3339 timestamp `{value}`"
))
})
}
fn scheduler_triggers(
manifest: &ClaudeRuntimeManifest,
now_unix: i64,
only_due: bool,
) -> Result<Vec<ClaudeRuntimeTrigger>> {
let jobs: BTreeMap<&str, &ClaudeCronJob> = manifest
.active_crons
.iter()
.map(|job| (job.id.as_str(), job))
.collect();
let wakeups: BTreeMap<&str, &ClaudeWakeup> = manifest
.pending_wakeups
.iter()
.map(|wakeup| (wakeup.tool_use_id.as_str(), wakeup))
.collect();
let mut out = Vec::new();
let queue_due = manifest.scheduler.activated_at_unix.ok_or_else(|| {
Error::Other("active Claude runtime manifest has no scheduler activation metadata".into())
})?;
if !only_due || queue_due <= now_unix {
for (index, prompt) in manifest.queue.pending.iter().enumerate() {
out.push(ClaudeRuntimeTrigger {
due_unix: queue_due,
kind: ClaudeRuntimeTriggerKind::Queue,
id: queue_trigger_id(&manifest.queue, index)?,
prompt: Some(prompt.clone()),
});
}
}
for cursor in &manifest.scheduler.cron_jobs {
let job = jobs.get(cursor.id.as_str()).ok_or_else(|| {
Error::Other(format!(
"scheduler cursor references missing Claude cron `{}`",
cursor.id
))
})?;
if cursor
.expires_at_unix
.is_some_and(|expiry| now_unix >= expiry)
{
continue;
}
if !only_due || cursor.next_due_unix <= now_unix {
out.push(ClaudeRuntimeTrigger {
due_unix: cursor.next_due_unix,
kind: ClaudeRuntimeTriggerKind::Cron,
id: cursor.id.clone(),
prompt: Some(job.prompt.clone()),
});
}
}
for cursor in &manifest.scheduler.wakeups {
let wakeup = wakeups.get(cursor.tool_use_id.as_str()).ok_or_else(|| {
Error::Other(format!(
"scheduler cursor references missing Claude wakeup `{}`",
cursor.tool_use_id
))
})?;
if !only_due || cursor.due_unix <= now_unix {
out.push(ClaudeRuntimeTrigger {
due_unix: cursor.due_unix,
kind: ClaudeRuntimeTriggerKind::Wakeup,
id: cursor.tool_use_id.clone(),
prompt: wakeup.prompt.clone().or_else(|| wakeup.reason.clone()),
});
}
}
Ok(out)
}
fn queue_trigger_id(
queue: &crate::claude_runtime_state::ClaudeQueueState,
index: usize,
) -> Result<String> {
let pending = u64::try_from(queue.pending.len())
.map_err(|_| Error::Other("Claude pending queue length exceeds u64".into()))?;
let index = u64::try_from(index)
.map_err(|_| Error::Other("Claude pending queue index exceeds u64".into()))?;
let first_ordinal = queue.enqueued.checked_sub(pending).ok_or_else(|| {
Error::Other(format!(
"Claude queue has {} pending prompts but only {} enqueue records",
queue.pending.len(),
queue.enqueued
))
})?;
let ordinal = first_ordinal
.checked_add(index)
.ok_or_else(|| Error::Other("Claude queue ordinal overflows u64".into()))?;
Ok(format!("queue-{ordinal:020}"))
}
#[derive(Debug, Clone)]
struct CronField {
allowed: Vec<bool>,
}
impl CronField {
fn parse(text: &str, min: u32, max: u32, dow: bool) -> Result<Self> {
if text.is_empty() {
return Err(Error::Other("empty cron field".into()));
}
let mut allowed = vec![false; (max - min + 1) as usize];
for item in text.split(',') {
if item.is_empty() {
return Err(Error::Other(format!(
"invalid empty item in cron field `{text}`"
)));
}
let mut parts = item.split('/');
let base = parts.next().unwrap_or_default();
let step = parts
.next()
.map(|value| value.parse::<u32>())
.transpose()
.map_err(|_| Error::Other(format!("invalid cron step in `{item}`")))?
.unwrap_or(1);
if parts.next().is_some() || step == 0 {
return Err(Error::Other(format!("invalid cron step in `{item}`")));
}
let (start, end) = if base == "*" {
(min, max)
} else if let Some((start, end)) = base.split_once('-') {
(
parse_cron_num(start, min, max, dow)?,
parse_cron_num(end, min, max, dow)?,
)
} else {
let start = parse_cron_num(base, min, max, dow)?;
(start, if item.contains('/') { max } else { start })
};
if start > end {
return Err(Error::Other(format!(
"descending cron range `{base}` is unsupported"
)));
}
let mut value = start;
while value <= end {
let normalized = if dow && value == 7 { 0 } else { value };
allowed[(normalized - min) as usize] = true;
let Some(next) = value.checked_add(step) else {
break;
};
value = next;
}
}
if !allowed.iter().any(|allowed| *allowed) {
return Err(Error::Other(format!(
"cron field `{text}` matches no values"
)));
}
Ok(Self { allowed })
}
fn contains(&self, value: u32, min: u32) -> bool {
self.allowed
.get((value - min) as usize)
.copied()
.unwrap_or(false)
}
fn unrestricted(&self) -> bool {
self.allowed.iter().all(|allowed| *allowed)
}
}
fn parse_cron_num(text: &str, min: u32, max: u32, dow: bool) -> Result<u32> {
let value = text
.parse::<u32>()
.map_err(|_| Error::Other(format!("invalid cron number `{text}`")))?;
let upper = if dow { 7 } else { max };
if value < min || value > upper {
return Err(Error::Other(format!(
"cron number `{value}` is outside {min}..={upper}"
)));
}
Ok(value)
}
#[derive(Debug, Clone)]
struct CronSchedule {
minute: CronField,
hour: CronField,
day_of_month: CronField,
month: CronField,
day_of_week: CronField,
}
impl CronSchedule {
fn parse(schedule: &str) -> Result<Self> {
let fields: Vec<&str> = schedule.split_whitespace().collect();
if fields.len() != 5 {
return Err(Error::Other(format!(
"invalid Claude cron `{schedule}`: expected exactly 5 fields"
)));
}
Ok(Self {
minute: CronField::parse(fields[0], 0, 59, false)?,
hour: CronField::parse(fields[1], 0, 23, false)?,
day_of_month: CronField::parse(fields[2], 1, 31, false)?,
month: CronField::parse(fields[3], 1, 12, false)?,
day_of_week: CronField::parse(fields[4], 0, 6, true)?,
})
}
fn next_after(&self, after_unix: i64) -> Result<i64> {
let start_minute = after_unix
.div_euclid(60)
.checked_add(1)
.ok_or_else(|| Error::Other("cron search overflows Unix time".into()))?;
const SEARCH_MINUTES: i64 = 8 * 366 * 24 * 60;
for delta in 0..SEARCH_MINUTES {
let unix = start_minute
.checked_add(delta)
.and_then(|minute| minute.checked_mul(60))
.ok_or_else(|| Error::Other("cron search overflows Unix time".into()))?;
if self.matches(unix) {
return Ok(unix);
}
}
Err(Error::Other(
"cron expression has no matching UTC minute within eight years".into(),
))
}
fn matches(&self, unix: i64) -> bool {
let days = unix.div_euclid(86_400);
let seconds = unix.rem_euclid(86_400);
let (year, month, day) = civil_from_days(days);
let _ = year;
let hour = (seconds / 3600) as u32;
let minute = ((seconds % 3600) / 60) as u32;
let dow = (days + 4).rem_euclid(7) as u32;
let dom_match = self.day_of_month.contains(day, 1);
let dow_match = self.day_of_week.contains(dow, 0);
let day_match = match (
self.day_of_month.unrestricted(),
self.day_of_week.unrestricted(),
) {
(true, true) => true,
(true, false) => dow_match,
(false, true) => dom_match,
(false, false) => dom_match || dow_match,
};
self.minute.contains(minute, 0)
&& self.hour.contains(hour, 0)
&& self.month.contains(month, 1)
&& day_match
}
}
fn civil_from_days(days: i64) -> (i64, u32, u32) {
let z = days + 719_468;
let era = if z >= 0 { z } else { z - 146_096 }.div_euclid(146_097);
let doe = z - era * 146_097;
let yoe = (doe - doe / 1460 + doe / 36_524 - doe / 146_096) / 365;
let mut year = yoe + era * 400;
let doy = doe - (365 * yoe + yoe / 4 - yoe / 100);
let mp = (5 * doy + 2) / 153;
let day = doy - (153 * mp + 2) / 5 + 1;
let month = if mp < 10 { mp + 3 } else { mp - 9 };
year += (month <= 2) as i64;
(year, month as u32, day as u32)
}