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 }