myc

Self-custodial remote signer for Radroots apps
git clone https://radroots.dev/git/myc.git
Log | Files | Refs | README | LICENSE

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 }