use super::OptionValueExt;
use meerkat_machine_dsl::machine;
machine! {
machine ScheduleLifecycleMachine {
version: 1,
rust: "self" / "catalog::dsl::schedule_lifecycle",
state {
schedule_id: ScheduleId,
lifecycle_phase: ScheduleLifecycleState,
revision: u64,
trigger_key: TriggerKey,
target_binding_key: TargetBindingId,
misfire_policy: Enum<MisfirePolicy>,
overlap_policy: Enum<OverlapPolicy>,
missing_target_policy: Enum<MissingTargetPolicy>,
planning_horizon_days: u64,
planning_horizon_occurrences: u64,
planning_cursor_utc_ms: Option<u64>,
next_occurrence_ordinal: u64,
superseded_ack_ids: Set<OccurrenceId>,
}
init(Active) {
schedule_id = "schedule-0",
revision = 1,
trigger_key = "trigger-0",
target_binding_key = "target-0",
misfire_policy = MisfirePolicy::Skip,
overlap_policy = OverlapPolicy::SkipIfRunning,
missing_target_policy = MissingTargetPolicy::MarkMisfired,
planning_horizon_days = 30,
planning_horizon_occurrences = 64,
planning_cursor_utc_ms = None,
next_occurrence_ordinal = 0,
superseded_ack_ids = EmptySet,
}
terminal [Deleted]
phase ScheduleLifecycleState {
Active,
Paused,
Deleted,
}
input ScheduleLifecycleInput {
Create {
schedule_id: ScheduleId,
trigger_key: TriggerKey,
target_binding_key: TargetBindingId,
misfire_policy: Enum<MisfirePolicy>,
overlap_policy: Enum<OverlapPolicy>,
missing_target_policy: Enum<MissingTargetPolicy>,
planning_horizon_days: Option<u64>,
planning_horizon_occurrences: Option<u64>,
},
Revise {
trigger_key: TriggerKey,
target_binding_key: TargetBindingId,
misfire_policy: Enum<MisfirePolicy>,
overlap_policy: Enum<OverlapPolicy>,
missing_target_policy: Enum<MissingTargetPolicy>,
planning_horizon_days: u64,
planning_horizon_occurrences: u64,
at_utc_ms: u64,
},
UpdatePlanningConfig { planning_horizon_days: u64, planning_horizon_occurrences: u64 },
RecordPlanningWindow {
planning_cursor_utc_ms: u64,
next_occurrence_ordinal: u64,
},
SyncTargetSnapshot { target_binding_key: TargetBindingId },
Pause { at_utc_ms: u64 },
Resume { at_utc_ms: u64 },
Delete { at_utc_ms: u64 },
ConfirmOccurrencesSuperseded { occurrence_id: OccurrenceId, superseding_revision: u64 },
}
effect ScheduleLifecycleEffect {
EmitScheduleNotice { new_state: ScheduleLifecycleState, revision: u64 },
SupersedePendingOccurrences { superseding_revision: u64, at_utc_ms: u64 },
PlanningWindowRecorded { planning_cursor_utc_ms: u64, next_occurrence_ordinal: u64 },
}
invariant revision_is_positive {
self.revision > 0
}
invariant deleted_has_no_planning_cursor {
self.lifecycle_phase != Phase::Deleted || self.planning_cursor_utc_ms == None
}
invariant planning_cursor_requires_occurrence_progress {
self.planning_cursor_utc_ms == None || self.next_occurrence_ordinal > 0
}
disposition EmitScheduleNotice => external seam SurfaceResultAlignment,
disposition SupersedePendingOccurrences => routed [OccurrenceLifecycleMachine] seam NoOwnerRealization,
disposition PlanningWindowRecorded => local seam NoOwnerRealization,
transition CreateSchedule {
on input Create {
schedule_id,
trigger_key,
target_binding_key,
misfire_policy,
overlap_policy,
missing_target_policy,
planning_horizon_days,
planning_horizon_occurrences
}
guard { self.lifecycle_phase == Phase::Active }
update {
self.schedule_id = schedule_id;
self.trigger_key = trigger_key;
self.target_binding_key = target_binding_key;
self.misfire_policy = misfire_policy;
self.overlap_policy = overlap_policy;
self.missing_target_policy = missing_target_policy;
if planning_horizon_days != None {
self.planning_horizon_days = planning_horizon_days.get("value");
}
if planning_horizon_occurrences != None {
self.planning_horizon_occurrences = planning_horizon_occurrences.get("value");
}
}
to Active
emit EmitScheduleNotice { new_state: self.lifecycle_phase, revision: self.revision }
}
transition ReviseActive {
on input Revise {
trigger_key,
target_binding_key,
misfire_policy,
overlap_policy,
missing_target_policy,
planning_horizon_days,
planning_horizon_occurrences,
at_utc_ms
}
guard { self.lifecycle_phase == Phase::Active }
update {
self.trigger_key = trigger_key;
self.target_binding_key = target_binding_key;
self.misfire_policy = misfire_policy;
self.overlap_policy = overlap_policy;
self.missing_target_policy = missing_target_policy;
self.planning_horizon_days = planning_horizon_days;
self.planning_horizon_occurrences = planning_horizon_occurrences;
self.revision += 1;
self.planning_cursor_utc_ms = None;
}
to Active
emit EmitScheduleNotice { new_state: self.lifecycle_phase, revision: self.revision }
emit SupersedePendingOccurrences { superseding_revision: self.revision, at_utc_ms: at_utc_ms }
}
transition RevisePaused {
on input Revise {
trigger_key,
target_binding_key,
misfire_policy,
overlap_policy,
missing_target_policy,
planning_horizon_days,
planning_horizon_occurrences,
at_utc_ms
}
guard { self.lifecycle_phase == Phase::Paused }
update {
self.trigger_key = trigger_key;
self.target_binding_key = target_binding_key;
self.misfire_policy = misfire_policy;
self.overlap_policy = overlap_policy;
self.missing_target_policy = missing_target_policy;
self.planning_horizon_days = planning_horizon_days;
self.planning_horizon_occurrences = planning_horizon_occurrences;
self.revision += 1;
self.planning_cursor_utc_ms = None;
}
to Paused
emit EmitScheduleNotice { new_state: self.lifecycle_phase, revision: self.revision }
emit SupersedePendingOccurrences { superseding_revision: self.revision, at_utc_ms: at_utc_ms }
}
transition UpdatePlanningConfigActive {
on input UpdatePlanningConfig { planning_horizon_days, planning_horizon_occurrences }
guard { self.lifecycle_phase == Phase::Active }
update {
self.planning_horizon_days = planning_horizon_days;
self.planning_horizon_occurrences = planning_horizon_occurrences;
}
to Active
emit EmitScheduleNotice { new_state: self.lifecycle_phase, revision: self.revision }
}
transition UpdatePlanningConfigPaused {
on input UpdatePlanningConfig { planning_horizon_days, planning_horizon_occurrences }
guard { self.lifecycle_phase == Phase::Paused }
update {
self.planning_horizon_days = planning_horizon_days;
self.planning_horizon_occurrences = planning_horizon_occurrences;
}
to Paused
emit EmitScheduleNotice { new_state: self.lifecycle_phase, revision: self.revision }
}
transition RecordPlanningWindowActive {
on input RecordPlanningWindow { planning_cursor_utc_ms, next_occurrence_ordinal }
guard "planning_window_advances_ordinal" { self.lifecycle_phase == Phase::Active && next_occurrence_ordinal > 0 }
guard "planning_cursor_advances" {
self.planning_cursor_utc_ms == None
|| planning_cursor_utc_ms > self.planning_cursor_utc_ms.get("value")
}
update {
self.planning_cursor_utc_ms = Some(planning_cursor_utc_ms);
self.next_occurrence_ordinal = next_occurrence_ordinal;
}
to Active
emit EmitScheduleNotice { new_state: self.lifecycle_phase, revision: self.revision }
emit PlanningWindowRecorded { planning_cursor_utc_ms: planning_cursor_utc_ms, next_occurrence_ordinal: next_occurrence_ordinal }
}
transition SyncTargetSnapshotActive {
on input SyncTargetSnapshot { target_binding_key }
guard { self.lifecycle_phase == Phase::Active }
update {
self.target_binding_key = target_binding_key;
}
to Active
}
transition SyncTargetSnapshotPaused {
on input SyncTargetSnapshot { target_binding_key }
guard { self.lifecycle_phase == Phase::Paused }
update {
self.target_binding_key = target_binding_key;
}
to Paused
}
transition PauseActiveOrPaused {
on input Pause { at_utc_ms }
guard { self.lifecycle_phase == Phase::Active || self.lifecycle_phase == Phase::Paused }
update {}
to Paused
emit EmitScheduleNotice { new_state: self.lifecycle_phase, revision: self.revision }
}
transition ResumeActiveOrPaused {
on input Resume { at_utc_ms }
guard { self.lifecycle_phase == Phase::Active || self.lifecycle_phase == Phase::Paused }
update {}
to Active
emit EmitScheduleNotice { new_state: self.lifecycle_phase, revision: self.revision }
}
transition DeleteActive {
on input Delete { at_utc_ms }
guard { self.lifecycle_phase == Phase::Active }
update {
self.revision += 1;
self.planning_cursor_utc_ms = None;
}
to Deleted
emit EmitScheduleNotice { new_state: self.lifecycle_phase, revision: self.revision }
emit SupersedePendingOccurrences { superseding_revision: self.revision, at_utc_ms: at_utc_ms }
}
transition DeletePaused {
on input Delete { at_utc_ms }
guard { self.lifecycle_phase == Phase::Paused }
update {
self.revision += 1;
self.planning_cursor_utc_ms = None;
}
to Deleted
emit EmitScheduleNotice { new_state: self.lifecycle_phase, revision: self.revision }
emit SupersedePendingOccurrences { superseding_revision: self.revision, at_utc_ms: at_utc_ms }
}
transition DeleteDeleted {
on input Delete { at_utc_ms }
guard { self.lifecycle_phase == Phase::Deleted }
update {}
to Deleted
}
transition ConfirmOccurrencesSupersededActive {
on input ConfirmOccurrencesSuperseded { occurrence_id, superseding_revision }
guard { self.lifecycle_phase == Phase::Active }
update {
self.superseded_ack_ids.insert(occurrence_id);
}
to Active
}
transition ConfirmOccurrencesSupersededPaused {
on input ConfirmOccurrencesSuperseded { occurrence_id, superseding_revision }
guard { self.lifecycle_phase == Phase::Paused }
update {
self.superseded_ack_ids.insert(occurrence_id);
}
to Paused
}
transition ConfirmOccurrencesSupersededDeleted {
on input ConfirmOccurrencesSuperseded { occurrence_id, superseding_revision }
guard { self.lifecycle_phase == Phase::Deleted }
update {
self.superseded_ack_ids.insert(occurrence_id);
}
to Deleted
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum MisfirePolicy {
Skip,
CatchUpWithin,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum OverlapPolicy {
AllowConcurrent,
SkipIfRunning,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum MissingTargetPolicy {
MarkMisfired,
Skip,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ScheduleId(pub String);
impl<T: Into<String>> From<T> for ScheduleId {
fn from(s: T) -> Self {
Self(s.into())
}
}
#[derive(Debug, Clone, PartialEq, Eq, Hash, Ord, PartialOrd)]
pub struct TriggerKey(pub String);
impl<T: Into<String>> From<T> for TriggerKey {
fn from(s: T) -> Self {
Self(s.into())
}
}
#[derive(Debug, Clone, PartialEq, Eq, Hash, Ord, PartialOrd)]
pub struct TargetBindingId(pub String);
impl<T: Into<String>> From<T> for TargetBindingId {
fn from(s: T) -> Self {
Self(s.into())
}
}
#[derive(Debug, Clone, PartialEq, Eq, Hash, Ord, PartialOrd)]
pub struct ClaimOwner(pub String);
impl<T: Into<String>> From<T> for ClaimOwner {
fn from(s: T) -> Self {
Self(s.into())
}
}
#[derive(Debug, Clone, PartialEq, Eq, Hash, Ord, PartialOrd)]
pub struct CorrelationId(pub String);
impl<T: Into<String>> From<T> for CorrelationId {
fn from(s: T) -> Self {
Self(s.into())
}
}
#[derive(Debug, Clone, PartialEq, Eq, Hash, Ord, PartialOrd)]
pub struct OccurrenceId(pub String);
impl<T: Into<String>> From<T> for OccurrenceId {
fn from(s: T) -> Self {
Self(s.into())
}
}
#[cfg(test)]
#[allow(clippy::expect_used)]
mod tests {
use super::*;
#[test]
fn delete_from_deleted_is_noop() {
let mut auth = ScheduleLifecycleMachineAuthority::new();
let first = ScheduleLifecycleMachineMutator::apply(
&mut auth,
ScheduleLifecycleInput::Delete { at_utc_ms: 100 },
)
.expect("first Delete from Active must succeed");
assert_eq!(first.to_phase, ScheduleLifecycleState::Deleted);
let revision_after_delete = auth.state().revision;
let planning_after_delete = auth.state().planning_cursor_utc_ms;
let ordinal_after_delete = auth.state().next_occurrence_ordinal;
let second = ScheduleLifecycleMachineMutator::apply(
&mut auth,
ScheduleLifecycleInput::Delete { at_utc_ms: 200 },
)
.expect("Delete from Deleted must be idempotent, not a transition error");
assert_eq!(second.from_phase, ScheduleLifecycleState::Deleted);
assert_eq!(second.to_phase, ScheduleLifecycleState::Deleted);
assert!(
second.effects.is_empty(),
"Delete from Deleted must emit zero effects, got {:?}",
second.effects
);
assert_eq!(auth.state().revision, revision_after_delete);
assert_eq!(auth.state().planning_cursor_utc_ms, planning_after_delete);
assert_eq!(auth.state().next_occurrence_ordinal, ordinal_after_delete);
}
}