1use super::module_boundary::{
4 CORE_MODULE_MCP_TIMEOUT, SCHEDULING_DISPATCH_MCP_TOOL, call_module_mcp_tool_json,
5 mcp_required_error, module_uses_mcp,
6};
7use super::*;
8
9pub fn evaluate_schedules_at_tick(
10 schedules: &[ScheduleDefinition],
11 tick_ms: u64,
12) -> Result<ScheduleEvaluation, ScheduleValidationError> {
13 validate_schedule_tick_ms_supported(tick_ms)?;
14 validate_schedules(schedules)?;
15 let mut due_triggers = Vec::new();
16 let mut lookback_budget = CRON_LOOKBACK_BUDGET_PER_REQUEST;
19 for schedule in schedules.iter().filter(|s| s.enabled) {
20 let canonical_schedule_id = canonical_schedule_id(&schedule.schedule_id);
21 let interval = parse_schedule_interval(&schedule.interval).ok_or_else(|| {
22 ScheduleValidationError::InvalidInterval {
23 schedule_id: canonical_schedule_id.clone(),
24 interval: schedule.interval.clone(),
25 }
26 })?;
27 let timezone = parse_schedule_timezone(&schedule.timezone).ok_or_else(|| {
28 ScheduleValidationError::InvalidTimezone {
29 schedule_id: canonical_schedule_id.clone(),
30 timezone: schedule.timezone.clone(),
31 }
32 })?;
33 let Some(due_tick_ms) = latest_due_tick_at_or_before(
34 &canonical_schedule_id,
35 &interval,
36 &timezone,
37 schedule.jitter_ms,
38 tick_ms,
39 &mut lookback_budget,
40 )
41 .map_err(|_| ScheduleValidationError::LookbackBudgetExceeded {
42 schedule_id: canonical_schedule_id.clone(),
43 })?
44 else {
45 continue;
46 };
47 if due_tick_ms != tick_ms {
48 continue;
49 }
50 due_triggers.push(ScheduleTrigger {
51 schedule_id: canonical_schedule_id,
52 interval: schedule.interval.clone(),
53 timezone: schedule.timezone.clone(),
54 due_tick_ms,
55 });
56 }
57
58 due_triggers.sort_by(|left, right| {
59 left.due_tick_ms
60 .cmp(&right.due_tick_ms)
61 .then_with(|| left.schedule_id.cmp(&right.schedule_id))
62 .then_with(|| left.interval.cmp(&right.interval))
63 .then_with(|| left.timezone.cmp(&right.timezone))
64 });
65
66 Ok(ScheduleEvaluation {
67 tick_ms,
68 due_triggers,
69 })
70}
71
72pub(crate) fn validate_schedules(
73 schedules: &[ScheduleDefinition],
74) -> Result<(), ScheduleValidationError> {
75 let mut seen = BTreeSet::new();
76 for schedule in schedules {
77 let canonical_schedule_id = canonical_schedule_id(&schedule.schedule_id);
78 if canonical_schedule_id.is_empty() {
79 return Err(ScheduleValidationError::EmptyScheduleId);
80 }
81 if !seen.insert(canonical_schedule_id.clone()) {
82 return Err(ScheduleValidationError::DuplicateScheduleId(
83 canonical_schedule_id,
84 ));
85 }
86 if parse_schedule_interval(&schedule.interval).is_none() {
87 return Err(ScheduleValidationError::InvalidInterval {
88 schedule_id: canonical_schedule_id,
89 interval: schedule.interval.clone(),
90 });
91 }
92 if parse_schedule_timezone(&schedule.timezone).is_none() {
93 return Err(ScheduleValidationError::InvalidTimezone {
94 schedule_id: canonical_schedule_id,
95 timezone: schedule.timezone.clone(),
96 });
97 }
98 }
99 Ok(())
100}
101
102impl MobkitRuntimeHandle {
103 fn parse_scheduling_runtime_injection_response(
104 response: Value,
105 ) -> Result<Option<(String, String)>, RuntimeBoundaryError> {
106 let Some(injection) = response
107 .as_object()
108 .and_then(|payload| payload.get("runtime_injection"))
109 .and_then(Value::as_object)
110 .cloned()
111 else {
112 return Ok(None);
113 };
114 let member_id = injection
115 .get("member_id")
116 .and_then(Value::as_str)
117 .map(str::trim)
118 .filter(|value| !value.is_empty())
119 .ok_or_else(|| {
120 RuntimeBoundaryError::Mcp(McpBoundaryError::InvalidToolPayload {
121 module_id: "scheduling".to_string(),
122 tool: SCHEDULING_DISPATCH_MCP_TOOL.to_string(),
123 reason: "runtime_injection.member_id must be a non-empty string".to_string(),
124 })
125 })?;
126 let message = injection
127 .get("message")
128 .and_then(Value::as_str)
129 .map(str::trim)
130 .filter(|value| !value.is_empty())
131 .ok_or_else(|| {
132 RuntimeBoundaryError::Mcp(McpBoundaryError::InvalidToolPayload {
133 module_id: "scheduling".to_string(),
134 tool: SCHEDULING_DISPATCH_MCP_TOOL.to_string(),
135 reason: "runtime_injection.message must be a non-empty string".to_string(),
136 })
137 })?;
138 Ok(Some((member_id.to_string(), message.to_string())))
139 }
140
141 fn scheduling_runtime_injection_for_dispatch(
142 &self,
143 schedule_id: &str,
144 interval: &str,
145 timezone: &str,
146 due_tick_ms: u64,
147 tick_ms: u64,
148 claim_key: &str,
149 ) -> Result<Option<(String, String)>, RuntimeBoundaryError> {
150 let Some((scheduling_module, pre_spawn)) = self.module_and_prespawn("scheduling") else {
151 return Ok(None);
152 };
153 if !self.is_module_loaded("scheduling") {
154 return Ok(None);
155 }
156 if !module_uses_mcp(scheduling_module, pre_spawn) {
157 return Err(mcp_required_error(
158 "scheduling",
159 SCHEDULING_DISPATCH_MCP_TOOL,
160 ));
161 }
162 let response = call_module_mcp_tool_json(
163 scheduling_module,
164 pre_spawn,
165 SCHEDULING_DISPATCH_MCP_TOOL,
166 &serde_json::json!({
167 "schedule_id": schedule_id,
168 "interval": interval,
169 "timezone": timezone,
170 "due_tick_ms": due_tick_ms,
171 "tick_ms": tick_ms,
172 "claim_key": claim_key,
173 }),
174 CORE_MODULE_MCP_TIMEOUT,
175 )?;
176 Self::parse_scheduling_runtime_injection_response(response)
177 }
178
179 fn next_scheduling_dispatch_sequence(&mut self) -> u64 {
180 Self::next_sequence(&mut self.scheduling_dispatch_sequence)
181 }
182 pub fn evaluate_schedule_tick(
183 &self,
184 schedules: &[ScheduleDefinition],
185 tick_ms: u64,
186 ) -> Result<ScheduleEvaluation, ScheduleValidationError> {
187 evaluate_schedules_at_tick(schedules, tick_ms)
188 }
189
190 pub fn dispatch_schedule_tick(
191 &mut self,
192 schedules: &[ScheduleDefinition],
193 tick_ms: u64,
194 ) -> Result<ScheduleDispatchReport, ScheduleValidationError> {
195 validate_schedule_tick_ms_supported(tick_ms)?;
196 validate_schedules(schedules)?;
197 self.prune_schedule_claims(tick_ms);
198 self.prune_scheduling_last_due_ticks(tick_ms);
199 let mut due_triggers = Vec::new();
200 let mut lookback_budget = CRON_LOOKBACK_BUDGET_PER_REQUEST;
203 for schedule in schedules.iter().filter(|s| s.enabled) {
204 let canonical_schedule_id = canonical_schedule_id(&schedule.schedule_id);
205 let interval = parse_schedule_interval(&schedule.interval).ok_or_else(|| {
206 ScheduleValidationError::InvalidInterval {
207 schedule_id: canonical_schedule_id.clone(),
208 interval: schedule.interval.clone(),
209 }
210 })?;
211 let timezone = parse_schedule_timezone(&schedule.timezone).ok_or_else(|| {
212 ScheduleValidationError::InvalidTimezone {
213 schedule_id: canonical_schedule_id.clone(),
214 timezone: schedule.timezone.clone(),
215 }
216 })?;
217 let Some(due_tick_ms) = latest_due_tick_at_or_before(
218 &canonical_schedule_id,
219 &interval,
220 &timezone,
221 schedule.jitter_ms,
222 tick_ms,
223 &mut lookback_budget,
224 )
225 .map_err(|_| ScheduleValidationError::LookbackBudgetExceeded {
226 schedule_id: canonical_schedule_id.clone(),
227 })?
228 else {
229 continue;
230 };
231 let last_due_tick = self
232 .scheduling_last_due_ticks
233 .get(&canonical_schedule_id)
234 .copied();
235 if schedule.catch_up {
236 if last_due_tick.is_some_and(|last| last >= due_tick_ms) {
237 continue;
238 }
239 } else if last_due_tick
240 .is_some_and(|last| last >= due_tick_ms && due_tick_ms != tick_ms)
241 {
242 continue;
243 }
244 due_triggers.push((schedule, canonical_schedule_id, due_tick_ms));
245 }
246 due_triggers.sort_by(
247 |(left_schedule, left_schedule_id, left_due_tick),
248 (right_schedule, right_schedule_id, right_due_tick)| {
249 left_due_tick
250 .cmp(right_due_tick)
251 .then_with(|| left_schedule_id.cmp(right_schedule_id))
252 .then_with(|| left_schedule.interval.cmp(&right_schedule.interval))
253 .then_with(|| left_schedule.timezone.cmp(&right_schedule.timezone))
254 },
255 );
256 let mut dispatched = Vec::new();
257 let mut skipped_claims = Vec::new();
258 let scheduling_signal = self.scheduling_supervisor_signal();
259 let mut supervisor_restart_emitted = false;
260
261 for (trigger, canonical_schedule_id, due_tick_ms) in &due_triggers {
262 let claim_key = format!("{canonical_schedule_id}:{due_tick_ms}");
263 if !self.record_schedule_claim(claim_key.clone(), tick_ms) {
264 skipped_claims.push(claim_key);
265 continue;
266 }
267 self.scheduling_last_due_ticks
268 .insert(canonical_schedule_id.clone(), *due_tick_ms);
269 self.prune_scheduling_last_due_ticks(tick_ms);
270
271 let event_sequence = self.next_scheduling_dispatch_sequence();
272 let event_id =
273 format!("evt-schedule-{canonical_schedule_id}-{due_tick_ms}-{event_sequence}");
274 insert_event_sorted(
275 &mut self.merged_events,
276 EventEnvelope {
277 event_id: event_id.clone(),
278 source: "module".to_string(),
279 timestamp_ms: tick_ms,
280 event: UnifiedEvent::Module(ModuleEvent {
281 module: "scheduling".to_string(),
282 event_type: "dispatch".to_string(),
283 payload: serde_json::json!({
284 "schedule_id": canonical_schedule_id,
285 "interval": trigger.interval,
286 "timezone": trigger.timezone,
287 "tick_ms": tick_ms,
288 "due_tick_ms": due_tick_ms,
289 "claim_key": claim_key,
290 "supervisor_signal": scheduling_signal,
291 }),
292 }),
293 },
294 );
295
296 if let Some(signal) = &scheduling_signal
297 && signal.restart_observed
298 && !supervisor_restart_emitted
299 {
300 insert_event_sorted(
301 &mut self.merged_events,
302 EventEnvelope {
303 event_id: format!("evt-scheduling-supervisor-{tick_ms}-{event_sequence}"),
304 source: "module".to_string(),
305 timestamp_ms: tick_ms,
306 event: UnifiedEvent::Module(ModuleEvent {
307 module: "scheduling".to_string(),
308 event_type: "supervisor.restart".to_string(),
309 payload: serde_json::json!({
310 "module_id": signal.module_id,
311 "latest_state": signal.latest_state,
312 "latest_attempt": signal.latest_attempt,
313 "restart_observed": signal.restart_observed,
314 }),
315 }),
316 },
317 );
318 supervisor_restart_emitted = true;
319 }
320
321 let mut runtime_injection = None;
322 let mut runtime_injection_error = None;
323 match self.scheduling_runtime_injection_for_dispatch(
324 canonical_schedule_id,
325 &trigger.interval,
326 &trigger.timezone,
327 *due_tick_ms,
328 tick_ms,
329 &claim_key,
330 ) {
331 Ok(Some((member_id, message))) => {
332 let injection_event_id =
333 format!("evt-runtime-injection-{tick_ms}-{event_sequence}");
334 insert_event_sorted(
335 &mut self.merged_events,
336 EventEnvelope {
337 event_id: injection_event_id.clone(),
338 source: "module".to_string(),
339 timestamp_ms: tick_ms,
340 event: UnifiedEvent::Module(ModuleEvent {
341 module: "runtime".to_string(),
342 event_type: "injection.dispatch".to_string(),
343 payload: serde_json::json!({
344 "schedule_id": canonical_schedule_id,
345 "claim_key": claim_key,
346 "member_id": member_id,
347 "message": message,
348 }),
349 }),
350 },
351 );
352 runtime_injection = Some(ScheduleRuntimeInjection {
353 member_id,
354 message,
355 injection_event_id,
356 });
357 }
358 Ok(None) => {}
359 Err(error) => {
360 runtime_injection_error = Some(format!("{error:?}"));
361 insert_event_sorted(
362 &mut self.merged_events,
363 EventEnvelope {
364 event_id: format!(
365 "evt-runtime-injection-failed-{tick_ms}-{event_sequence}"
366 ),
367 source: "module".to_string(),
368 timestamp_ms: tick_ms,
369 event: UnifiedEvent::Module(ModuleEvent {
370 module: "runtime".to_string(),
371 event_type: "runtime.injection.failed".to_string(),
372 payload: serde_json::json!({
373 "schedule_id": canonical_schedule_id,
374 "claim_key": claim_key,
375 "error": format!("{error:?}"),
376 }),
377 }),
378 },
379 );
380 }
381 }
382
383 dispatched.push(ScheduleDispatch {
384 claim_key,
385 schedule_id: canonical_schedule_id.clone(),
386 interval: trigger.interval.clone(),
387 timezone: trigger.timezone.clone(),
388 due_tick_ms: *due_tick_ms,
389 tick_ms,
390 event_id,
391 supervisor_signal: scheduling_signal.clone(),
392 runtime_injection,
393 runtime_injection_error,
394 });
395 }
396
397 Ok(ScheduleDispatchReport {
398 tick_ms,
399 due_count: due_triggers.len(),
400 dispatched,
401 skipped_claims,
402 })
403 }
404 fn record_schedule_claim(&mut self, claim_key: String, tick_ms: u64) -> bool {
405 if !self.scheduling_claims.insert(claim_key.clone()) {
406 return false;
407 }
408 self.scheduling_claim_ticks
409 .entry(tick_ms)
410 .or_default()
411 .push(claim_key);
412 true
413 }
414
415 fn prune_schedule_claims(&mut self, current_tick_ms: u64) {
416 let cutoff_tick = current_tick_ms.saturating_sub(SCHEDULING_CLAIM_RETENTION_WINDOW_MS);
417 let expired_ticks = self
418 .scheduling_claim_ticks
419 .keys()
420 .copied()
421 .take_while(|tick| *tick < cutoff_tick)
422 .collect::<Vec<_>>();
423 for tick in expired_ticks {
424 if let Some(keys) = self.scheduling_claim_ticks.remove(&tick) {
425 for key in keys {
426 self.scheduling_claims.remove(&key);
427 }
428 }
429 }
430
431 while self.scheduling_claims.len() > SCHEDULING_CLAIMS_MAX_RETAINED {
432 let Some(oldest_tick) = self.scheduling_claim_ticks.keys().next().copied() else {
433 break;
434 };
435 if let Some(keys) = self.scheduling_claim_ticks.remove(&oldest_tick) {
436 for key in keys {
437 self.scheduling_claims.remove(&key);
438 }
439 } else {
440 break;
441 }
442 }
443 }
444
445 fn prune_scheduling_last_due_ticks(&mut self, current_tick_ms: u64) {
446 let cutoff_tick = current_tick_ms.saturating_sub(SCHEDULING_CLAIM_RETENTION_WINDOW_MS);
447 self.scheduling_last_due_ticks
448 .retain(|_, due_tick| *due_tick >= cutoff_tick);
449
450 while self.scheduling_last_due_ticks.len() > SCHEDULING_LAST_DUE_MAX_RETAINED {
451 let Some(oldest_schedule_id) = self
452 .scheduling_last_due_ticks
453 .iter()
454 .min_by(|(left_id, left_due), (right_id, right_due)| {
455 left_due.cmp(right_due).then_with(|| left_id.cmp(right_id))
456 })
457 .map(|(schedule_id, _)| schedule_id.clone())
458 else {
459 break;
460 };
461 self.scheduling_last_due_ticks.remove(&oldest_schedule_id);
462 }
463 }
464 fn scheduling_supervisor_signal(&self) -> Option<SchedulingSupervisorSignal> {
465 let module_transitions = self
466 .supervisor_report
467 .transitions
468 .iter()
469 .filter(|transition| transition.module_id == "scheduling")
470 .collect::<Vec<_>>();
471 let latest = module_transitions.last()?;
472 let restart_observed = module_transitions
473 .iter()
474 .any(|transition| transition.to == ModuleHealthState::Restarting);
475 Some(SchedulingSupervisorSignal {
476 module_id: latest.module_id.clone(),
477 latest_state: latest.to.clone(),
478 latest_attempt: latest.attempt,
479 restart_observed,
480 })
481 }
482}
483
484#[derive(Debug, Clone, PartialEq, Eq)]
485enum ParsedInterval {
486 Marker { interval_ms: u64 },
487 Cron(CronExpression),
488}
489
490impl ParsedInterval {
491 fn jitter_base_interval_ms(&self) -> u64 {
492 match self {
493 Self::Marker { interval_ms } => *interval_ms,
494 Self::Cron(_) => 60_000,
496 }
497 }
498}
499
500#[derive(Debug, Clone, PartialEq, Eq)]
501enum ParsedTimezone {
502 FixedOffsetMs(i64),
503 Iana(chrono_tz::Tz),
504}
505
506#[derive(Debug, Clone, PartialEq, Eq)]
507struct CronExpression {
508 minute: CronFieldSet,
509 hour: CronFieldSet,
510 day_of_month: CronFieldSet,
511 month: CronFieldSet,
512 day_of_week: CronFieldSet,
513}
514
515#[derive(Debug, Clone, PartialEq, Eq)]
516struct CronFieldSet {
517 any: bool,
518 min: u32,
519 allowed: Vec<bool>,
520}
521
522impl CronExpression {
523 fn parse(expression: &str) -> Option<Self> {
524 let fields = expression.split_whitespace().collect::<Vec<_>>();
525 if fields.len() != 5 {
526 return None;
527 }
528 let parsed = Self {
529 minute: parse_cron_field(fields[0], 0, 59, false)?,
530 hour: parse_cron_field(fields[1], 0, 23, false)?,
531 day_of_month: parse_cron_field(fields[2], 1, 31, false)?,
532 month: parse_cron_field(fields[3], 1, 12, false)?,
533 day_of_week: parse_cron_field(fields[4], 0, 7, true)?,
534 };
535
536 if parsed.day_of_week.any
539 && !parsed.day_of_month.any
540 && !parsed.has_possible_day_of_month_for_selected_months()
541 {
542 return None;
543 }
544
545 Some(parsed)
546 }
547
548 fn matches(&self, local: &LocalDateTimeFields) -> bool {
549 if !self.minute.matches(local.minute)
550 || !self.hour.matches(local.hour)
551 || !self.month.matches(local.month)
552 {
553 return false;
554 }
555
556 let dom_match = self.day_of_month.matches(local.day_of_month);
557 let dow_match = self.day_of_week.matches(local.day_of_week);
558
559 if self.day_of_month.any && self.day_of_week.any {
560 true
561 } else if self.day_of_month.any {
562 dow_match
563 } else if self.day_of_week.any {
564 dom_match
565 } else {
566 dom_match || dow_match
567 }
568 }
569
570 fn has_possible_day_of_month_for_selected_months(&self) -> bool {
571 for month in 1..=12 {
572 if !self.month.matches(month) {
573 continue;
574 }
575 let max_day = max_day_for_month_with_feb_29(month);
576 for day in 1..=max_day {
577 if self.day_of_month.matches(day) {
578 return true;
579 }
580 }
581 }
582 false
583 }
584}
585
586impl CronFieldSet {
587 fn matches(&self, value: u32) -> bool {
588 if value < self.min {
589 return false;
590 }
591 let idx = (value - self.min) as usize;
592 self.allowed.get(idx).copied().unwrap_or(false)
593 }
594}
595
596#[derive(Debug, Clone, PartialEq, Eq)]
597struct LocalDateTimeFields {
598 minute: u32,
599 hour: u32,
600 day_of_month: u32,
601 month: u32,
602 day_of_week: u32,
603 second: u32,
604 subsec_nanos: u32,
605}
606
607fn parse_cron_field(
608 field: &str,
609 min: u32,
610 max: u32,
611 map_sunday_seven_to_zero: bool,
612) -> Option<CronFieldSet> {
613 let mut allowed = vec![false; (max - min + 1) as usize];
614
615 for raw_token in field.split(',') {
616 let token = raw_token.trim();
617 if token.is_empty() {
618 return None;
619 }
620 let (base, step) = match token.split_once('/') {
621 Some((base, step)) => {
622 let step = step.parse::<u32>().ok()?;
623 if step == 0 {
624 return None;
625 }
626 (base.trim(), step)
627 }
628 None => (token, 1),
629 };
630
631 if base == "*" {
632 let mut value = min;
633 while value <= max {
634 let mapped = normalize_cron_value(value, map_sunday_seven_to_zero);
635 let idx = (mapped - min) as usize;
636 allowed[idx] = true;
637 match value.checked_add(step) {
638 Some(next) => value = next,
639 None => break,
640 }
641 }
642 continue;
643 }
644
645 if let Some((start, end)) = base.split_once('-') {
646 let start = parse_cron_raw_value(start.trim(), min, max)?;
647 let end = parse_cron_raw_value(end.trim(), min, max)?;
648 if start > end {
649 return None;
650 }
651 let mut value = start;
652 while value <= end {
653 let mapped = normalize_cron_value(value, map_sunday_seven_to_zero);
654 let idx = (mapped - min) as usize;
655 allowed[idx] = true;
656 match value.checked_add(step) {
657 Some(next) => value = next,
658 None => break,
659 }
660 }
661 continue;
662 }
663
664 let value = parse_cron_value(base, min, max, map_sunday_seven_to_zero)?;
665 let idx = (value - min) as usize;
666 allowed[idx] = true;
667 }
668
669 if allowed.iter().all(|allowed| !allowed) {
670 return None;
671 }
672
673 let any = cron_field_is_semantic_wildcard(min, max, map_sunday_seven_to_zero, &allowed);
674 Some(CronFieldSet { any, min, allowed })
675}
676
677fn cron_field_is_semantic_wildcard(
678 min: u32,
679 max: u32,
680 map_sunday_seven_to_zero: bool,
681 allowed: &[bool],
682) -> bool {
683 let mut covered = vec![false; allowed.len()];
684 for raw in min..=max {
685 let mapped = normalize_cron_value(raw, map_sunday_seven_to_zero);
686 if mapped < min || mapped > max {
687 return false;
688 }
689 let mapped_idx = (mapped - min) as usize;
690 covered[mapped_idx] = true;
691 }
692
693 covered
694 .iter()
695 .enumerate()
696 .filter(|(_, is_semantic_value)| **is_semantic_value)
697 .all(|(idx, _)| allowed.get(idx).copied().unwrap_or(false))
698}
699
700fn parse_cron_value(raw: &str, min: u32, max: u32, map_sunday_seven_to_zero: bool) -> Option<u32> {
701 let value = normalize_cron_value(
702 parse_cron_raw_value(raw, min, max)?,
703 map_sunday_seven_to_zero,
704 );
705 if value < min || value > max {
706 return None;
707 }
708 Some(value)
709}
710
711fn parse_cron_raw_value(raw: &str, min: u32, max: u32) -> Option<u32> {
712 let value = raw.parse::<u32>().ok()?;
713 if value < min || value > max {
714 return None;
715 }
716 Some(value)
717}
718
719fn normalize_cron_value(value: u32, map_sunday_seven_to_zero: bool) -> u32 {
720 if map_sunday_seven_to_zero && value == 7 {
721 0
722 } else {
723 value
724 }
725}
726
727fn max_day_for_month_with_feb_29(month: u32) -> u32 {
728 match month {
729 2 => 29,
730 4 | 6 | 9 | 11 => 30,
731 _ => 31,
732 }
733}
734
735fn parse_interval_marker_ms(interval: &str) -> Option<u64> {
736 let marker = interval.trim().to_ascii_lowercase();
737 let marker = marker.strip_prefix("*/")?;
738 if marker.len() < 2 {
739 return None;
740 }
741 let (count_part, unit_part) = marker.split_at(marker.len() - 1);
742 let count = count_part.parse::<u64>().ok()?;
743 if count == 0 {
744 return None;
745 }
746 let unit_ms = match unit_part {
747 "s" => 1_000,
748 "m" => 60_000,
749 "h" => 3_600_000,
750 "d" => 86_400_000,
751 _ => return None,
752 };
753 count.checked_mul(unit_ms)
754}
755
756fn parse_schedule_interval(interval: &str) -> Option<ParsedInterval> {
757 parse_interval_marker_ms(interval)
758 .map(|interval_ms| ParsedInterval::Marker { interval_ms })
759 .or_else(|| CronExpression::parse(interval.trim()).map(ParsedInterval::Cron))
760}
761
762fn deterministic_jitter_offset_ms(schedule_id: &str, jitter_ms: u64, interval_ms: u64) -> u64 {
763 if jitter_ms == 0 || interval_ms <= 1 {
764 return 0;
765 }
766 let mut hash = 1_469_598_103_934_665_603_u64;
767 for byte in schedule_id.bytes() {
768 hash ^= byte as u64;
769 hash = hash.wrapping_mul(1_099_511_628_211);
770 }
771 let max_jitter = jitter_ms.min(interval_ms.saturating_sub(1));
772 hash % (max_jitter + 1)
773}
774
775fn parse_schedule_timezone(timezone: &str) -> Option<ParsedTimezone> {
776 let timezone = timezone.trim();
777 if timezone.is_empty() {
778 return None;
779 }
780 parse_timezone_offset_ms(timezone)
781 .map(ParsedTimezone::FixedOffsetMs)
782 .or_else(|| {
783 timezone
784 .parse::<chrono_tz::Tz>()
785 .ok()
786 .map(ParsedTimezone::Iana)
787 })
788}
789
790fn parse_timezone_offset_ms(timezone: &str) -> Option<i64> {
791 let tz = timezone.trim();
792 if tz.is_empty() {
793 return None;
794 }
795 if tz.eq_ignore_ascii_case("utc") || tz == "Z" {
796 return Some(0);
797 }
798 let offset = tz
799 .strip_prefix("UTC")
800 .or_else(|| tz.strip_prefix("utc"))
801 .or_else(|| tz.strip_prefix("GMT"))
802 .or_else(|| tz.strip_prefix("gmt"))
803 .unwrap_or(tz);
804 parse_hhmm_offset(offset)
805}
806
807fn parse_hhmm_offset(offset: &str) -> Option<i64> {
808 if offset.is_empty() {
809 return Some(0);
810 }
811 let sign = if offset.starts_with('+') {
812 1_i64
813 } else if offset.starts_with('-') {
814 -1_i64
815 } else {
816 return None;
817 };
818 let body = &offset[1..];
819 let (hours, minutes) = if let Some((h, m)) = body.split_once(':') {
820 (h, m)
821 } else if body.len() == 4 {
822 body.split_at(2)
823 } else {
824 return None;
825 };
826 let hours = hours.parse::<i64>().ok()?;
827 let minutes = minutes.parse::<i64>().ok()?;
828 if hours > 23 || minutes > 59 {
829 return None;
830 }
831 let total_minutes = hours.saturating_mul(60).saturating_add(minutes);
832 Some(sign.saturating_mul(total_minutes).saturating_mul(60_000))
833}
834
835fn utc_datetime_from_tick_ms(tick_ms: u64) -> Option<chrono::DateTime<Utc>> {
836 let tick_ms = i64::try_from(tick_ms).ok()?;
837 chrono::DateTime::<Utc>::from_timestamp_millis(tick_ms)
838}
839
840fn local_fields_at_tick(timezone: &ParsedTimezone, tick_ms: u64) -> Option<LocalDateTimeFields> {
841 let utc = utc_datetime_from_tick_ms(tick_ms)?;
842 let (minute, hour, day_of_month, month, day_of_week, second, subsec_nanos) = match timezone {
843 ParsedTimezone::FixedOffsetMs(offset_ms) => {
844 let offset_seconds = i32::try_from(offset_ms / 1_000).ok()?;
845 let offset = chrono::FixedOffset::east_opt(offset_seconds)?;
846 let local = utc.with_timezone(&offset);
847 (
848 local.minute(),
849 local.hour(),
850 local.day(),
851 local.month(),
852 local.weekday().num_days_from_sunday(),
853 local.second(),
854 local.nanosecond(),
855 )
856 }
857 ParsedTimezone::Iana(timezone) => {
858 let local = utc.with_timezone(timezone);
859 (
860 local.minute(),
861 local.hour(),
862 local.day(),
863 local.month(),
864 local.weekday().num_days_from_sunday(),
865 local.second(),
866 local.nanosecond(),
867 )
868 }
869 };
870 Some(LocalDateTimeFields {
871 minute,
872 hour,
873 day_of_month,
874 month,
875 day_of_week,
876 second,
877 subsec_nanos,
878 })
879}
880
881fn timezone_offset_ms_at_tick(timezone: &ParsedTimezone, tick_ms: u64) -> Option<i64> {
882 match timezone {
883 ParsedTimezone::FixedOffsetMs(offset) => Some(*offset),
884 ParsedTimezone::Iana(tz) => {
885 let utc = utc_datetime_from_tick_ms(tick_ms)?;
886 let local = utc.with_timezone(tz);
887 Some(i64::from(local.offset().fix().local_minus_utc()).saturating_mul(1_000))
888 }
889 }
890}
891
892fn latest_due_marker_tick_at_or_before(
893 interval_ms: u64,
894 timezone: &ParsedTimezone,
895 tick_ms: u64,
896) -> Option<u64> {
897 match timezone {
898 ParsedTimezone::FixedOffsetMs(timezone_offset_ms) => {
899 latest_due_marker_tick_at_or_before_with_offset(
900 interval_ms,
901 *timezone_offset_ms,
902 tick_ms,
903 )
904 }
905 ParsedTimezone::Iana(_) => {
906 let mut timezone_offset_ms = timezone_offset_ms_at_tick(timezone, tick_ms)?;
907 for _ in 0..4 {
908 let due_tick = latest_due_marker_tick_at_or_before_with_offset(
909 interval_ms,
910 timezone_offset_ms,
911 tick_ms,
912 )?;
913 let due_offset_ms = timezone_offset_ms_at_tick(timezone, due_tick)?;
914 if due_offset_ms == timezone_offset_ms {
915 return Some(due_tick);
916 }
917 timezone_offset_ms = due_offset_ms;
918 }
919 latest_due_marker_tick_at_or_before_with_offset(
920 interval_ms,
921 timezone_offset_ms,
922 tick_ms,
923 )
924 }
925 }
926}
927
928fn latest_due_marker_tick_at_or_before_with_offset(
929 interval_ms: u64,
930 timezone_offset_ms: i64,
931 tick_ms: u64,
932) -> Option<u64> {
933 let local_tick = i128::from(tick_ms) + i128::from(timezone_offset_ms);
934 if local_tick < 0 {
935 return None;
936 }
937 let local_tick = local_tick as u64;
938 let latest_due_local_tick = local_tick - (local_tick % interval_ms);
939 let due_tick = i128::from(latest_due_local_tick) - i128::from(timezone_offset_ms);
940 if due_tick < 0 {
941 return None;
942 }
943 Some(due_tick as u64)
944}
945
946fn canonical_schedule_id(schedule_id: &str) -> String {
947 schedule_id.trim().to_string()
948}
949
950fn validate_schedule_tick_ms_supported(tick_ms: u64) -> Result<(), ScheduleValidationError> {
951 if tick_ms > i64::MAX as u64 {
952 return Err(ScheduleValidationError::InvalidTickMs(tick_ms));
953 }
954 Ok(())
955}
956
957#[derive(Debug)]
963struct LookbackBudgetExhausted;
964
965fn latest_due_cron_tick_at_or_before(
966 cron: &CronExpression,
967 timezone: &ParsedTimezone,
968 tick_ms: u64,
969 budget: &mut u64,
970) -> Result<Option<u64>, LookbackBudgetExhausted> {
971 let mut candidate = tick_ms - (tick_ms % 60_000);
972 for _ in 0..=CRON_LOOKBACK_MINUTES {
973 if *budget == 0 {
974 return Err(LookbackBudgetExhausted);
975 }
976 *budget -= 1;
977 let Some(fields) = local_fields_at_tick(timezone, candidate) else {
978 return Ok(None);
979 };
980 if fields.second == 0 && fields.subsec_nanos == 0 && cron.matches(&fields) {
981 return Ok(Some(candidate));
982 }
983 let Some(next) = candidate.checked_sub(60_000) else {
984 return Ok(None);
985 };
986 candidate = next;
987 }
988 Ok(None)
989}
990
991fn latest_due_tick_at_or_before(
992 schedule_id: &str,
993 interval: &ParsedInterval,
994 timezone: &ParsedTimezone,
995 jitter_ms: u64,
996 tick_ms: u64,
997 budget: &mut u64,
998) -> Result<Option<u64>, LookbackBudgetExhausted> {
999 let jitter_offset_ms =
1000 deterministic_jitter_offset_ms(schedule_id, jitter_ms, interval.jitter_base_interval_ms());
1001 let Some(tick_without_jitter) = tick_ms.checked_sub(jitter_offset_ms) else {
1002 return Ok(None);
1003 };
1004 let due_without_jitter = match interval {
1005 ParsedInterval::Marker { interval_ms } => {
1006 latest_due_marker_tick_at_or_before(*interval_ms, timezone, tick_without_jitter)
1007 }
1008 ParsedInterval::Cron(cron) => {
1009 latest_due_cron_tick_at_or_before(cron, timezone, tick_without_jitter, budget)?
1010 }
1011 };
1012 Ok(due_without_jitter.and_then(|tick| tick.checked_add(jitter_offset_ms)))
1013}
1014
1015#[cfg(test)]
1016#[allow(clippy::expect_used, clippy::unwrap_used, clippy::panic)]
1017mod budget_tests {
1018 use super::*;
1019
1020 #[test]
1025 fn sparse_leap_day_cron_still_resolves_with_ample_budget() {
1026 let cron = match parse_schedule_interval("0 0 29 2 *").expect("valid leap-day cron") {
1027 ParsedInterval::Cron(cron) => cron,
1028 ParsedInterval::Marker { .. } => panic!("expected a cron interval"),
1029 };
1030 let tz = parse_schedule_timezone("UTC").expect("utc");
1031 let tick_ms = 1_750_032_000_000;
1034 let mut budget = CRON_LOOKBACK_BUDGET_PER_REQUEST;
1035 let due = latest_due_cron_tick_at_or_before(&cron, &tz, tick_ms, &mut budget)
1036 .expect("ample budget must not be exhausted");
1037 let due = due.expect("a Feb 29 must exist within the lookback window");
1038 let fields = local_fields_at_tick(&tz, due).expect("local fields");
1040 assert_eq!((fields.month, fields.day_of_month), (2, 29));
1041 assert!(
1042 budget < CRON_LOOKBACK_BUDGET_PER_REQUEST,
1043 "budget was consumed"
1044 );
1045 }
1046
1047 #[test]
1052 fn exhausted_budget_returns_marker_for_sparse_cron() {
1053 let cron = match parse_schedule_interval("0 0 29 2 *").expect("valid leap-day cron") {
1054 ParsedInterval::Cron(cron) => cron,
1055 ParsedInterval::Marker { .. } => panic!("expected a cron interval"),
1056 };
1057 let tz = parse_schedule_timezone("UTC").expect("utc");
1058 let tick_ms = 1_750_032_000_000;
1059 let mut budget = 10u64;
1061 let result = latest_due_cron_tick_at_or_before(&cron, &tz, tick_ms, &mut budget);
1062 assert!(
1063 result.is_err(),
1064 "a tiny budget must surface LookbackBudgetExhausted, not stall or silently miss"
1065 );
1066 assert_eq!(budget, 0, "budget must be fully consumed before failing");
1067 }
1068
1069 #[test]
1073 fn budget_is_shared_across_schedules_in_a_request() {
1074 let cron = match parse_schedule_interval("0 0 29 2 *").expect("valid leap-day cron") {
1075 ParsedInterval::Cron(cron) => cron,
1076 ParsedInterval::Marker { .. } => panic!("expected a cron interval"),
1077 };
1078 let tz = parse_schedule_timezone("UTC").expect("utc");
1079 let tick_ms = 1_750_032_000_000;
1080 let mut budget = CRON_LOOKBACK_BUDGET_PER_REQUEST;
1082 latest_due_cron_tick_at_or_before(&cron, &tz, tick_ms, &mut budget)
1083 .expect("first sparse cron resolves within budget");
1084 let remaining = budget;
1085 assert!(remaining < CRON_LOOKBACK_BUDGET_PER_REQUEST);
1086 budget = 5;
1089 let second = latest_due_cron_tick_at_or_before(&cron, &tz, tick_ms, &mut budget);
1090 assert!(
1091 second.is_err(),
1092 "shared budget must fail the second sparse cron closed"
1093 );
1094 }
1095}