runtime_supervision.rs (10548B)
1 //! Sealed, bounded Myc critical-task supervision over the shared host runner. 2 3 use core::{fmt, future::Future, pin::Pin}; 4 use std::error::Error; 5 6 use radroots_service_host::{ 7 CancellationToken, HostError, HostErrorKind, ShutdownPhase, SupervisionFailureKind, 8 TaskClassification, TaskMetadata, TaskName, TaskSupervisor, 9 }; 10 11 use crate::{MycLogRecord, MycProcessResult}; 12 13 /// Exact Myc runtime-supervision contract version. 14 pub const MYC_RUNTIME_SUPERVISION_CONTRACT_VERSION: u32 = 1; 15 /// Maximum number of critical tasks admitted into one Myc runtime graph. 16 pub const MYC_CRITICAL_TASK_MAX_COUNT: usize = 32; 17 18 type CriticalTaskFuture = 19 Pin<Box<dyn Future<Output = Result<(), MycCriticalTaskError>> + Send + 'static>>; 20 type CriticalTaskFactory = 21 Box<dyn FnOnce(MycTaskCancellation) -> CriticalTaskFuture + Send + 'static>; 22 23 // TaskMetadata requires a phase for critical work, but the Step 158 supervisor 24 // does not execute ordered shutdown. Step 159 replaces this inert placeholder 25 // while composing the exact durability-aware phase inventory. 26 const STEP_158_PLACEHOLDER_SHUTDOWN_PHASE: ShutdownPhase = ShutdownPhase::DrainOperations; 27 28 /// Read-only cooperative cancellation evidence supplied to one critical task. 29 #[derive(Clone)] 30 pub struct MycTaskCancellation { 31 inner: CancellationToken, 32 } 33 34 impl MycTaskCancellation { 35 pub(crate) const fn from_host(inner: CancellationToken) -> Self { 36 Self { inner } 37 } 38 39 #[cfg(any(target_os = "linux", target_os = "macos"))] 40 pub(crate) fn uncancelled() -> Self { 41 Self { 42 inner: CancellationToken::new(), 43 } 44 } 45 46 /// Returns whether coordinated cancellation has already been requested. 47 #[must_use] 48 pub fn is_cancelled(&self) -> bool { 49 self.inner.is_cancelled() 50 } 51 52 /// Waits until the owning Myc supervisor requests coordinated cancellation. 53 pub async fn cancelled(&self) { 54 self.inner.cancelled().await; 55 } 56 57 #[cfg(all(test, any(target_os = "linux", target_os = "macos")))] 58 pub(crate) fn test_pair() -> (Self, CancellationToken) { 59 let token = CancellationToken::new(); 60 ( 61 Self { 62 inner: token.clone(), 63 }, 64 token, 65 ) 66 } 67 } 68 69 impl fmt::Debug for MycTaskCancellation { 70 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 71 formatter.write_str("MycTaskCancellation([sealed])") 72 } 73 } 74 75 /// Source-free failure returned by one caller-supplied critical task. 76 #[derive(Clone, Copy, Debug, Default, PartialEq, Eq)] 77 pub struct MycCriticalTaskError; 78 79 impl MycCriticalTaskError { 80 /// Constructs the sole safe task-failure classification. 81 #[must_use] 82 pub const fn failed() -> Self { 83 Self 84 } 85 } 86 87 impl fmt::Display for MycCriticalTaskError { 88 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 89 formatter.write_str("Myc critical task failed") 90 } 91 } 92 93 impl Error for MycCriticalTaskError {} 94 95 /// One sealed critical task without a caller-controlled name or detachable handle. 96 pub struct MycCriticalTask { 97 factory: CriticalTaskFactory, 98 } 99 100 impl MycCriticalTask { 101 /// Wraps one authoritative task for owned supervision. 102 pub fn new<F, Fut>(task: F) -> Self 103 where 104 F: FnOnce(MycTaskCancellation) -> Fut + Send + 'static, 105 Fut: Future<Output = Result<(), MycCriticalTaskError>> + Send + 'static, 106 { 107 Self { 108 factory: Box::new(move |cancellation| Box::pin(task(cancellation))), 109 } 110 } 111 } 112 113 impl fmt::Debug for MycCriticalTask { 114 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 115 formatter.write_str("MycCriticalTask([sealed])") 116 } 117 } 118 119 /// Stable source-free failure classes for the Myc critical-task graph. 120 #[derive(Clone, Copy, Debug, PartialEq, Eq)] 121 pub enum MycRuntimeSupervisionErrorKind { 122 EmptyTaskSet, 123 TooManyTasks, 124 TaskRegistration, 125 TaskReturnedError, 126 TaskPanicked, 127 UnexpectedCompletion, 128 UnexpectedCancellation, 129 JoinFailed, 130 } 131 132 impl MycRuntimeSupervisionErrorKind { 133 /// Returns the stable machine-readable failure code. 134 #[must_use] 135 pub const fn code(self) -> &'static str { 136 match self { 137 Self::EmptyTaskSet => "runtime_task_set_empty", 138 Self::TooManyTasks => "runtime_task_set_too_large", 139 Self::TaskRegistration => "runtime_task_registration_failed", 140 Self::TaskReturnedError => "runtime_task_returned_error", 141 Self::TaskPanicked => "runtime_task_panicked", 142 Self::UnexpectedCompletion => "runtime_task_completed_early", 143 Self::UnexpectedCancellation => "runtime_task_cancelled_unexpectedly", 144 Self::JoinFailed => "runtime_task_join_failed", 145 } 146 } 147 } 148 149 /// Redacted failure returned only after the shared supervisor joins owned work. 150 #[derive(Clone, Copy, PartialEq, Eq)] 151 pub struct MycRuntimeSupervisionError { 152 kind: MycRuntimeSupervisionErrorKind, 153 } 154 155 impl MycRuntimeSupervisionError { 156 const fn new(kind: MycRuntimeSupervisionErrorKind) -> Self { 157 Self { kind } 158 } 159 160 /// Returns the stable failure classification. 161 #[must_use] 162 pub const fn kind(self) -> MycRuntimeSupervisionErrorKind { 163 self.kind 164 } 165 166 /// Returns the stable machine-readable failure code. 167 #[must_use] 168 pub const fn code(self) -> &'static str { 169 self.kind.code() 170 } 171 172 /// Returns the fixed nonzero process result for every fatal task-graph outcome. 173 #[must_use] 174 pub const fn process_result(self) -> MycProcessResult { 175 MycProcessResult::UnexpectedInternal 176 } 177 178 /// Returns the fixed safe diagnostic for every fatal task-graph outcome. 179 #[must_use] 180 pub const fn diagnostic(self) -> MycLogRecord { 181 MycLogRecord::critical_task_failed() 182 } 183 } 184 185 impl fmt::Display for MycRuntimeSupervisionError { 186 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 187 formatter.write_str("Myc critical-task supervision failed") 188 } 189 } 190 191 impl fmt::Debug for MycRuntimeSupervisionError { 192 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 193 formatter 194 .debug_struct("MycRuntimeSupervisionError") 195 .field("kind", &self.kind) 196 .finish() 197 } 198 } 199 200 impl Error for MycRuntimeSupervisionError {} 201 202 /// One nonforgeable, bounded set of critical tasks awaiting owned execution. 203 #[must_use = "the supervised runtime must be run so authoritative tasks are joined"] 204 pub struct MycSupervisedRuntime { 205 tasks: Box<[MycCriticalTask]>, 206 } 207 208 impl MycSupervisedRuntime { 209 /// Validates and retains between one and 32 critical tasks. 210 /// 211 /// Iterator ingestion stops after the maximum plus one item. 212 pub fn new( 213 tasks: impl IntoIterator<Item = MycCriticalTask>, 214 ) -> Result<Self, MycRuntimeSupervisionError> { 215 let mut bounded = Vec::with_capacity(MYC_CRITICAL_TASK_MAX_COUNT); 216 for task in tasks.into_iter().take(MYC_CRITICAL_TASK_MAX_COUNT + 1) { 217 if bounded.len() == MYC_CRITICAL_TASK_MAX_COUNT { 218 return Err(MycRuntimeSupervisionError::new( 219 MycRuntimeSupervisionErrorKind::TooManyTasks, 220 )); 221 } 222 bounded.push(task); 223 } 224 if bounded.is_empty() { 225 return Err(MycRuntimeSupervisionError::new( 226 MycRuntimeSupervisionErrorKind::EmptyTaskSet, 227 )); 228 } 229 Ok(Self { 230 tasks: bounded.into_boxed_slice(), 231 }) 232 } 233 234 /// Returns the exact number of retained authoritative tasks. 235 #[must_use] 236 pub fn task_count(&self) -> usize { 237 self.tasks.len() 238 } 239 240 /// Runs the single owned graph until a fatal outcome coordinates peer 241 /// cancellation and every task join has been observed. 242 pub async fn run(self) -> Result<(), MycRuntimeSupervisionError> { 243 let metadata = (0..self.tasks.len()) 244 .map(task_metadata) 245 .collect::<Result<Vec<_>, _>>()?; 246 let mut supervisor = TaskSupervisor::new(); 247 for (metadata, task) in metadata.into_iter().zip(self.tasks.into_vec()) { 248 let factory = task.factory; 249 if supervisor 250 .spawn(metadata, move |cancellation| async move { 251 factory(MycTaskCancellation { 252 inner: cancellation, 253 }) 254 .await 255 .map_err(|_| HostError::new(HostErrorKind::TaskFailure)) 256 }) 257 .is_err() 258 { 259 supervisor.request_cancellation(); 260 let _ = supervisor.supervise().await; 261 return Err(MycRuntimeSupervisionError::new( 262 MycRuntimeSupervisionErrorKind::TaskRegistration, 263 )); 264 } 265 } 266 supervisor 267 .supervise() 268 .await 269 .map(|_| ()) 270 .map_err(|failure| MycRuntimeSupervisionError::new(map_failure_kind(failure.kind()))) 271 } 272 } 273 274 impl fmt::Debug for MycSupervisedRuntime { 275 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 276 formatter 277 .debug_struct("MycSupervisedRuntime") 278 .field("task_count", &self.tasks.len()) 279 .field("tasks", &"[sealed]") 280 .finish() 281 } 282 } 283 284 fn task_metadata(index: usize) -> Result<TaskMetadata, MycRuntimeSupervisionError> { 285 let name = TaskName::new(format!("critical_task_{index:02}")).map_err(|_| { 286 MycRuntimeSupervisionError::new(MycRuntimeSupervisionErrorKind::TaskRegistration) 287 })?; 288 TaskMetadata::new( 289 name, 290 TaskClassification::Critical, 291 Some(STEP_158_PLACEHOLDER_SHUTDOWN_PHASE), 292 ) 293 .map_err(|_| MycRuntimeSupervisionError::new(MycRuntimeSupervisionErrorKind::TaskRegistration)) 294 } 295 296 const fn map_failure_kind(kind: SupervisionFailureKind) -> MycRuntimeSupervisionErrorKind { 297 match kind { 298 SupervisionFailureKind::TaskReturnedError => { 299 MycRuntimeSupervisionErrorKind::TaskReturnedError 300 } 301 SupervisionFailureKind::TaskPanicked => MycRuntimeSupervisionErrorKind::TaskPanicked, 302 SupervisionFailureKind::UnexpectedCompletion => { 303 MycRuntimeSupervisionErrorKind::UnexpectedCompletion 304 } 305 SupervisionFailureKind::UnexpectedCancellation => { 306 MycRuntimeSupervisionErrorKind::UnexpectedCancellation 307 } 308 SupervisionFailureKind::JoinFailed => MycRuntimeSupervisionErrorKind::JoinFailed, 309 } 310 }