lib

Core libraries for Radroots
git clone https://radroots.dev/git/lib.git
Log | Files | Refs | README

task.rs (9400B)


      1 //! Validated metadata for every supervisor-owned task.
      2 
      3 use core::fmt;
      4 
      5 use serde::Serialize;
      6 
      7 pub const TASK_NAME_MAX_BYTES: usize = 64;
      8 
      9 /// A stable, non-secret identifier for one supervised task.
     10 #[derive(Clone, Debug, Hash, PartialEq, Eq, PartialOrd, Ord, Serialize)]
     11 #[serde(transparent)]
     12 pub struct TaskName(String);
     13 
     14 impl TaskName {
     15     pub fn new(value: impl AsRef<str>) -> Result<Self, TaskMetadataError> {
     16         let value = value.as_ref();
     17         if !valid_task_name(value) {
     18             return Err(TaskMetadataError::InvalidTaskName);
     19         }
     20         Ok(Self(value.to_owned()))
     21     }
     22 
     23     #[must_use]
     24     pub fn as_str(&self) -> &str {
     25         &self.0
     26     }
     27 }
     28 
     29 impl fmt::Display for TaskName {
     30     fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
     31         formatter.write_str(self.as_str())
     32     }
     33 }
     34 
     35 /// Service impact and lifetime class for a supervised task.
     36 #[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize)]
     37 #[serde(rename_all = "snake_case")]
     38 pub enum TaskClassification {
     39     /// A long-lived authoritative task whose failure or early success fails the service.
     40     Critical,
     41     /// A bounded optional capability whose failure is observable but not service-fatal.
     42     Optional,
     43     /// A finite authoritative operation whose successful completion is expected.
     44     OneShot,
     45 }
     46 
     47 impl TaskClassification {
     48     #[must_use]
     49     pub const fn completion_expectation(self) -> TaskCompletionExpectation {
     50         match self {
     51             Self::Critical => TaskCompletionExpectation::AfterCancellation,
     52             Self::Optional => TaskCompletionExpectation::MayComplete,
     53             Self::OneShot => TaskCompletionExpectation::CompletesOnce,
     54         }
     55     }
     56 
     57     /// Returns whether a returned error or panic fails the service.
     58     #[must_use]
     59     pub const fn failure_is_fatal(self) -> bool {
     60         matches!(self, Self::Critical | Self::OneShot)
     61     }
     62 
     63     #[must_use]
     64     pub const fn requires_shutdown_phase(self) -> bool {
     65         !matches!(self, Self::OneShot)
     66     }
     67 }
     68 
     69 /// When successful completion is valid for a task class.
     70 #[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize)]
     71 #[serde(rename_all = "snake_case")]
     72 pub enum TaskCompletionExpectation {
     73     AfterCancellation,
     74     MayComplete,
     75     CompletesOnce,
     76 }
     77 
     78 /// Ordered shutdown phase assigned to each long-lived task.
     79 #[derive(Clone, Copy, Debug, Hash, PartialEq, Eq, PartialOrd, Ord, Serialize)]
     80 #[serde(rename_all = "snake_case")]
     81 pub enum ShutdownPhase {
     82     RejectNewMutations,
     83     CancelIngress,
     84     DrainOperations,
     85     PersistRecoverableWork,
     86     CloseNetwork,
     87     CloseSqlite,
     88     CloseSockets,
     89 }
     90 
     91 /// Complete static identity and lifecycle ownership for one task.
     92 #[derive(Clone, Debug, PartialEq, Eq, Serialize)]
     93 pub struct TaskMetadata {
     94     name: TaskName,
     95     classification: TaskClassification,
     96     shutdown_phase: Option<ShutdownPhase>,
     97 }
     98 
     99 impl TaskMetadata {
    100     pub fn new(
    101         name: TaskName,
    102         classification: TaskClassification,
    103         shutdown_phase: Option<ShutdownPhase>,
    104     ) -> Result<Self, TaskMetadataError> {
    105         if classification.requires_shutdown_phase() != shutdown_phase.is_some() {
    106             return Err(TaskMetadataError::InvalidShutdownPhaseAssignment);
    107         }
    108         Ok(Self {
    109             name,
    110             classification,
    111             shutdown_phase,
    112         })
    113     }
    114 
    115     #[must_use]
    116     pub const fn name(&self) -> &TaskName {
    117         &self.name
    118     }
    119 
    120     #[must_use]
    121     pub const fn classification(&self) -> TaskClassification {
    122         self.classification
    123     }
    124 
    125     #[must_use]
    126     pub const fn shutdown_phase(&self) -> Option<ShutdownPhase> {
    127         self.shutdown_phase
    128     }
    129 }
    130 
    131 /// Validation failure for static supervisor metadata.
    132 #[derive(Clone, Copy, Debug, PartialEq, Eq)]
    133 pub enum TaskMetadataError {
    134     InvalidTaskName,
    135     InvalidShutdownPhaseAssignment,
    136 }
    137 
    138 impl fmt::Display for TaskMetadataError {
    139     fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
    140         formatter.write_str(match self {
    141             Self::InvalidTaskName => "supervised task name is invalid",
    142             Self::InvalidShutdownPhaseAssignment => {
    143                 "supervised task shutdown phase does not match its classification"
    144             }
    145         })
    146     }
    147 }
    148 
    149 impl std::error::Error for TaskMetadataError {}
    150 
    151 fn valid_task_name(value: &str) -> bool {
    152     let mut bytes = value.bytes();
    153     let Some(first) = bytes.next() else {
    154         return false;
    155     };
    156     value.len() <= TASK_NAME_MAX_BYTES
    157         && first.is_ascii_lowercase()
    158         && bytes.all(|byte| byte.is_ascii_lowercase() || byte.is_ascii_digit() || byte == b'_')
    159 }
    160 
    161 #[cfg(test)]
    162 mod tests {
    163     use super::*;
    164 
    165     #[test]
    166     fn task_names_are_bounded_stable_and_secret_safe() {
    167         for valid in ["a", "admin_listener", "outbox_worker_01"] {
    168             let name = TaskName::new(valid).unwrap();
    169             assert_eq!(name.as_str(), valid);
    170             assert_eq!(name.to_string(), valid);
    171         }
    172         assert!(TaskName::new("a".repeat(TASK_NAME_MAX_BYTES)).is_ok());
    173         for invalid in [
    174             "",
    175             "AdminListener",
    176             "1_worker",
    177             "admin-listener",
    178             "admin.listener",
    179             "admin listener",
    180             "worker/secret",
    181             "café",
    182         ] {
    183             assert_eq!(
    184                 TaskName::new(invalid),
    185                 Err(TaskMetadataError::InvalidTaskName)
    186             );
    187         }
    188         assert!(TaskName::new("a".repeat(TASK_NAME_MAX_BYTES + 1)).is_err());
    189         assert_eq!(
    190             TaskName::new("a".repeat(4 * 1024 * 1024)),
    191             Err(TaskMetadataError::InvalidTaskName)
    192         );
    193         assert_eq!(
    194             TaskMetadataError::InvalidTaskName.to_string(),
    195             "supervised task name is invalid"
    196         );
    197         assert_eq!(
    198             TaskMetadataError::InvalidShutdownPhaseAssignment.to_string(),
    199             "supervised task shutdown phase does not match its classification"
    200         );
    201     }
    202 
    203     #[test]
    204     fn classification_semantics_are_exhaustive() {
    205         let expected = [
    206             (
    207                 TaskClassification::Critical,
    208                 TaskCompletionExpectation::AfterCancellation,
    209                 true,
    210                 true,
    211             ),
    212             (
    213                 TaskClassification::Optional,
    214                 TaskCompletionExpectation::MayComplete,
    215                 false,
    216                 true,
    217             ),
    218             (
    219                 TaskClassification::OneShot,
    220                 TaskCompletionExpectation::CompletesOnce,
    221                 true,
    222                 false,
    223             ),
    224         ];
    225         for (classification, completion, fatal_failure, requires_shutdown) in expected {
    226             assert_eq!(classification.completion_expectation(), completion);
    227             assert_eq!(classification.failure_is_fatal(), fatal_failure);
    228             assert_eq!(classification.requires_shutdown_phase(), requires_shutdown);
    229         }
    230     }
    231 
    232     #[test]
    233     fn shutdown_phase_assignment_matches_task_lifetime() {
    234         let name = || TaskName::new("relay_subscription").unwrap();
    235         assert!(
    236             TaskMetadata::new(
    237                 name(),
    238                 TaskClassification::Critical,
    239                 Some(ShutdownPhase::CancelIngress),
    240             )
    241             .is_ok()
    242         );
    243         let optional = TaskMetadata::new(
    244             name(),
    245             TaskClassification::Optional,
    246             Some(ShutdownPhase::CloseNetwork),
    247         )
    248         .unwrap();
    249         assert_eq!(optional.shutdown_phase(), Some(ShutdownPhase::CloseNetwork));
    250         assert_eq!(
    251             TaskMetadata::new(name(), TaskClassification::Optional, None),
    252             Err(TaskMetadataError::InvalidShutdownPhaseAssignment)
    253         );
    254         assert_eq!(
    255             TaskMetadata::new(
    256                 name(),
    257                 TaskClassification::OneShot,
    258                 Some(ShutdownPhase::CancelIngress),
    259             ),
    260             Err(TaskMetadataError::InvalidShutdownPhaseAssignment)
    261         );
    262         assert!(TaskMetadata::new(name(), TaskClassification::OneShot, None).is_ok());
    263     }
    264 
    265     #[test]
    266     fn serde_debug_and_shutdown_order_are_stable() {
    267         let phases = [
    268             ShutdownPhase::RejectNewMutations,
    269             ShutdownPhase::CancelIngress,
    270             ShutdownPhase::DrainOperations,
    271             ShutdownPhase::PersistRecoverableWork,
    272             ShutdownPhase::CloseNetwork,
    273             ShutdownPhase::CloseSqlite,
    274             ShutdownPhase::CloseSockets,
    275         ];
    276         assert!(phases.windows(2).all(|pair| pair[0] < pair[1]));
    277         assert_eq!(
    278             serde_json::to_string(&phases).unwrap(),
    279             r#"["reject_new_mutations","cancel_ingress","drain_operations","persist_recoverable_work","close_network","close_sqlite","close_sockets"]"#
    280         );
    281 
    282         let metadata = TaskMetadata::new(
    283             TaskName::new("admin_listener").unwrap(),
    284             TaskClassification::Critical,
    285             Some(ShutdownPhase::CloseSockets),
    286         )
    287         .unwrap();
    288         assert_eq!(
    289             serde_json::to_string(&metadata).unwrap(),
    290             r#"{"name":"admin_listener","classification":"critical","shutdown_phase":"close_sockets"}"#
    291         );
    292         assert_eq!(
    293             format!("{metadata:?}"),
    294             "TaskMetadata { name: TaskName(\"admin_listener\"), classification: Critical, shutdown_phase: Some(CloseSockets) }"
    295         );
    296     }
    297 }