use crate::{ScheduleError, TimerDirective, TimerRegistration};
use std::time::Duration;
#[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)]
pub enum TimerPolicy {
Once,
AfterCompletion {
cadence_ns: u64,
},
Watchdog {
cadence_ns: u64,
},
}
impl TimerPolicy {
#[must_use]
pub const fn label(self) -> &'static str {
match self {
Self::Once => "once",
Self::AfterCompletion { .. } => "after_completion",
Self::Watchdog { .. } => "watchdog",
}
}
#[must_use]
pub const fn cadence_ns(self) -> Option<u64> {
match self {
Self::Once => None,
Self::AfterCompletion { cadence_ns } | Self::Watchdog { cadence_ns } => {
Some(cadence_ns)
}
}
}
#[must_use]
pub const fn initial_mode(self) -> TimerSchedulingMode {
match self {
Self::Once => TimerSchedulingMode::Once,
Self::AfterCompletion { .. } => TimerSchedulingMode::AfterCompletion,
Self::Watchdog { .. } => TimerSchedulingMode::Watchdog,
}
}
}
#[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)]
pub enum TimerSchedulingMode {
Once,
AfterCompletion,
Deadline,
Retry,
Continuation,
Watchdog,
}
impl TimerSchedulingMode {
#[must_use]
pub const fn label(self) -> &'static str {
match self {
Self::Once => "once",
Self::AfterCompletion => "after_completion",
Self::Deadline => "deadline",
Self::Retry => "retry",
Self::Continuation => "continuation",
Self::Watchdog => "watchdog",
}
}
}
#[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)]
pub enum TimerDirectiveSnapshot {
Stop,
ContinueImmediately,
RetryAfter {
delay_ns: u64,
},
ScheduleAt {
deadline_ns: u64,
},
RecurAfter {
delay_ns: u64,
},
}
impl TimerDirectiveSnapshot {
#[must_use]
pub const fn scheduling_mode(self) -> Option<TimerSchedulingMode> {
match self {
Self::Stop => None,
Self::ContinueImmediately => Some(TimerSchedulingMode::Continuation),
Self::RetryAfter { .. } => Some(TimerSchedulingMode::Retry),
Self::ScheduleAt { .. } => Some(TimerSchedulingMode::Deadline),
Self::RecurAfter { .. } => Some(TimerSchedulingMode::AfterCompletion),
}
}
}
impl TryFrom<TimerDirective> for TimerDirectiveSnapshot {
type Error = ScheduleError;
fn try_from(value: TimerDirective) -> Result<Self, Self::Error> {
Ok(match value {
TimerDirective::Stop => Self::Stop,
TimerDirective::ContinueImmediately => Self::ContinueImmediately,
TimerDirective::RetryAfter(delay) => Self::RetryAfter {
delay_ns: duration_ns(delay)?,
},
TimerDirective::ScheduleAt(deadline_ns) => Self::ScheduleAt { deadline_ns },
TimerDirective::RecurAfter(delay) => Self::RecurAfter {
delay_ns: duration_ns(delay)?,
},
})
}
}
impl From<TimerDirectiveSnapshot> for TimerDirective {
fn from(value: TimerDirectiveSnapshot) -> Self {
match value {
TimerDirectiveSnapshot::Stop => Self::Stop,
TimerDirectiveSnapshot::ContinueImmediately => Self::ContinueImmediately,
TimerDirectiveSnapshot::RetryAfter { delay_ns } => {
Self::RetryAfter(Duration::from_nanos(delay_ns))
}
TimerDirectiveSnapshot::ScheduleAt { deadline_ns } => Self::ScheduleAt(deadline_ns),
TimerDirectiveSnapshot::RecurAfter { delay_ns } => {
Self::RecurAfter(Duration::from_nanos(delay_ns))
}
}
}
}
fn duration_ns(duration: Duration) -> Result<u64, ScheduleError> {
u64::try_from(duration.as_nanos()).map_err(|_| ScheduleError::DelayOutOfRange)
}
#[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)]
pub struct PreArmedSuccessor {
pub generation: u64,
pub deadline_ns: u64,
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub struct TimerSchedulingSnapshot {
pub configured_policy: TimerPolicy,
pub current_mode: TimerSchedulingMode,
pub latest_directive: Option<TimerDirectiveSnapshot>,
pub latest_requested_delay_ns: Option<u64>,
pub latest_armed_delay_ns: Option<u64>,
pub next_deadline_ns: Option<u64>,
pub pre_armed_successor: Option<PreArmedSuccessor>,
}
impl TimerSchedulingSnapshot {
#[must_use]
pub const fn new(configured_policy: TimerPolicy) -> Self {
Self {
configured_policy,
current_mode: configured_policy.initial_mode(),
latest_directive: None,
latest_requested_delay_ns: None,
latest_armed_delay_ns: None,
next_deadline_ns: None,
pre_armed_successor: None,
}
}
}
#[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)]
pub enum TimerRegistrationStatus {
Unregistered,
Scheduled,
Running,
}
impl TimerRegistrationStatus {
#[must_use]
pub const fn label(self) -> &'static str {
match self {
Self::Unregistered => "unregistered",
Self::Scheduled => "scheduled",
Self::Running => "running",
}
}
}
impl From<TimerRegistration> for TimerRegistrationStatus {
fn from(value: TimerRegistration) -> Self {
match value {
TimerRegistration::Unregistered => Self::Unregistered,
TimerRegistration::Scheduled { .. } => Self::Scheduled,
TimerRegistration::Running { .. } => Self::Running,
}
}
}
#[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)]
pub enum TimerProcessCondition {
Disabled,
Idle,
Active,
Retrying,
Failed,
MissingRegistration,
}
impl TimerProcessCondition {
#[must_use]
pub const fn label(self) -> &'static str {
match self {
Self::Disabled => "disabled",
Self::Idle => "idle",
Self::Active => "active",
Self::Retrying => "retrying",
Self::Failed => "failed",
Self::MissingRegistration => "missing_registration",
}
}
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub struct TimerStateSnapshot {
pub enabled: bool,
pub registration: TimerRegistrationStatus,
pub condition: TimerProcessCondition,
pub generation: u64,
pub in_flight: bool,
}
impl Default for TimerStateSnapshot {
fn default() -> Self {
Self {
enabled: true,
registration: TimerRegistrationStatus::Unregistered,
condition: TimerProcessCondition::Idle,
generation: 0,
in_flight: false,
}
}
}
#[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)]
pub enum TimerCompletionOutcome {
Success,
NoWork,
RetryableFailure,
InvariantFailure,
}
impl TimerCompletionOutcome {
#[must_use]
pub const fn label(self) -> &'static str {
match self {
Self::Success => "success",
Self::NoWork => "no_work",
Self::RetryableFailure => "retryable_failure",
Self::InvariantFailure => "invariant_failure",
}
}
}
#[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)]
pub enum TimerLastOutcome {
Completed(TimerCompletionOutcome),
Interrupted,
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub struct TimerCompletion {
pub outcome: TimerCompletionOutcome,
pub work_count: u64,
}
impl TimerCompletion {
#[must_use]
pub const fn success(work_count: u64) -> Self {
Self {
outcome: TimerCompletionOutcome::Success,
work_count,
}
}
#[must_use]
pub const fn no_work() -> Self {
Self {
outcome: TimerCompletionOutcome::NoWork,
work_count: 0,
}
}
#[must_use]
pub const fn retryable_failure(work_count: u64) -> Self {
Self {
outcome: TimerCompletionOutcome::RetryableFailure,
work_count,
}
}
#[must_use]
pub const fn invariant_failure(work_count: u64) -> Self {
Self {
outcome: TimerCompletionOutcome::InvariantFailure,
work_count,
}
}
}
#[derive(Clone, Copy, Debug, Default, Eq, PartialEq)]
pub struct TimerOutcomeSnapshot {
last_outcome: Option<TimerLastOutcome>,
last_work_count: Option<u64>,
last_success_at_ns: Option<u64>,
last_failure_at_ns: Option<u64>,
last_interrupted_at_ns: Option<u64>,
consecutive_expected_failures: u64,
}
impl TimerOutcomeSnapshot {
pub const fn record_completion(&mut self, completion: TimerCompletion, completed_at_ns: u64) {
self.last_outcome = Some(TimerLastOutcome::Completed(completion.outcome));
self.last_work_count = Some(completion.work_count);
match completion.outcome {
TimerCompletionOutcome::Success | TimerCompletionOutcome::NoWork => {
self.last_success_at_ns = Some(completed_at_ns);
self.consecutive_expected_failures = 0;
}
TimerCompletionOutcome::RetryableFailure => {
self.last_failure_at_ns = Some(completed_at_ns);
self.consecutive_expected_failures =
self.consecutive_expected_failures.saturating_add(1);
}
TimerCompletionOutcome::InvariantFailure => {
self.last_failure_at_ns = Some(completed_at_ns);
self.consecutive_expected_failures = 0;
}
}
}
pub const fn record_interruption(&mut self, observed_at_ns: u64) {
self.last_outcome = Some(TimerLastOutcome::Interrupted);
self.last_work_count = None;
self.last_interrupted_at_ns = Some(observed_at_ns);
}
#[must_use]
pub const fn last_outcome(self) -> Option<TimerLastOutcome> {
self.last_outcome
}
#[must_use]
pub const fn last_work_count(self) -> Option<u64> {
self.last_work_count
}
#[must_use]
pub const fn last_success_at_ns(self) -> Option<u64> {
self.last_success_at_ns
}
#[must_use]
pub const fn last_failure_at_ns(self) -> Option<u64> {
self.last_failure_at_ns
}
#[must_use]
pub const fn last_interrupted_at_ns(self) -> Option<u64> {
self.last_interrupted_at_ns
}
#[must_use]
pub const fn consecutive_expected_failures(self) -> u64 {
self.consecutive_expected_failures
}
}
#[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)]
pub struct TimerEpoch {
pub id: u64,
pub started_at_ns: u64,
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn expected_failure_streak_saturates() {
let mut outcomes = TimerOutcomeSnapshot {
consecutive_expected_failures: u64::MAX,
..TimerOutcomeSnapshot::default()
};
outcomes.record_completion(TimerCompletion::retryable_failure(0), 10);
assert_eq!(outcomes.consecutive_expected_failures(), u64::MAX);
}
}