Skip to main content

saddle_runtime/
application.rs

1use std::{
2    future::Future,
3    sync::Arc,
4    sync::atomic::{AtomicBool, Ordering},
5    time::{Duration, Instant},
6};
7
8use saddle_core::{ComponentLifecycle, ErrorKind, Result, SaddleError};
9use saddle_observability::root_diagnostic::RecordedComponentLifecycle;
10
11use crate::{RecordedComponentBoundary, RequestLifecycle};
12
13enum ApplicationComponent {
14    Legacy(Arc<dyn ComponentLifecycle>),
15    Recorded(Arc<dyn RecordedComponentLifecycle>),
16}
17type PrevalidatedLegacyComponent = Arc<dyn ComponentLifecycle>;
18impl ApplicationComponent {
19    fn name(&self) -> &'static str {
20        match self {
21            Self::Legacy(component) => component.name(),
22            Self::Recorded(component) => component.name(),
23        }
24    }
25}
26
27static RUNTIME_STARTED: AtomicBool = AtomicBool::new(false);
28const WORKER_THREADS: usize = 2;
29const MAX_IO_EVENTS_PER_TICK: usize = 5;
30const DEFAULT_START_TIMEOUT: Duration = Duration::from_secs(30);
31const DEFAULT_SHUTDOWN_TIMEOUT: Duration = Duration::from_secs(30);
32
33/// Fixed wall-clock limits for the managed component lifecycle.
34#[derive(Clone, Copy, Debug, Eq, PartialEq)]
35pub struct LifecycleTimeouts {
36    start: Duration,
37    shutdown: Duration,
38}
39
40impl LifecycleTimeouts {
41    pub fn from_millis(start_ms: u64, shutdown_ms: u64) -> Option<Self> {
42        if start_ms == 0 || shutdown_ms == 0 {
43            return None;
44        }
45        Some(Self {
46            start: Duration::from_millis(start_ms),
47            shutdown: Duration::from_millis(shutdown_ms),
48        })
49    }
50}
51
52impl Default for LifecycleTimeouts {
53    fn default() -> Self {
54        Self {
55            start: DEFAULT_START_TIMEOUT,
56            shutdown: DEFAULT_SHUTDOWN_TIMEOUT,
57        }
58    }
59}
60
61/// A complete Saddle application hosted by the process-wide async runtime.
62///
63/// This is an assembly API, not a general-purpose async executor: it exposes no
64/// Tokio handle, task spawning, runtime configuration, or arbitrary `block_on`.
65pub struct Application {
66    components: Vec<ApplicationComponent>,
67    recorded_boundary: Option<RecordedComponentBoundary>,
68    requests: RequestLifecycle,
69    deployment_resource_budget: Option<saddle_admission::DeploymentResourceBudget>,
70    ingress_bridge_issued: AtomicBool,
71    lifecycle_timeouts: LifecycleTimeouts,
72    shutdown_deadline: Arc<std::sync::Mutex<Option<Instant>>>,
73    lifecycle_observer: Option<(saddle_observability::Observer, String)>,
74    #[cfg(all(target_arch = "x86_64", target_os = "linux"))]
75    pending_driver_finalizer: crate::post_driver::PendingDriverFinalizerSlot,
76}
77
78impl Application {
79    /// Creates an empty application assembly.
80    pub fn new() -> Self {
81        Self {
82            components: Vec::new(),
83            recorded_boundary: None,
84            requests: RequestLifecycle::new(),
85            deployment_resource_budget: None,
86            ingress_bridge_issued: AtomicBool::new(false),
87            lifecycle_timeouts: LifecycleTimeouts::default(),
88            shutdown_deadline: Arc::new(std::sync::Mutex::new(None)),
89            lifecycle_observer: None,
90            #[cfg(all(target_arch = "x86_64", target_os = "linux"))]
91            pending_driver_finalizer: crate::post_driver::PendingDriverFinalizerSlot::new(),
92        }
93    }
94
95    /// Installs the frozen process lifecycle policy before component startup.
96    #[doc(hidden)]
97    pub fn set_lifecycle_timeouts(&mut self, timeouts: LifecycleTimeouts) {
98        self.lifecycle_timeouts = timeouts;
99    }
100
101    #[doc(hidden)]
102    pub fn install_lifecycle_observer(
103        &mut self,
104        observer: saddle_observability::Observer,
105        application: &str,
106    ) {
107        self.lifecycle_observer = Some((observer, application.to_owned()));
108    }
109
110    /// Installs the same application identity with the selected source output
111    /// before any internal request admission owner can run.
112    #[doc(hidden)]
113    pub fn install_lifecycle_observer_with_output(
114        &mut self,
115        observer: saddle_observability::Observer,
116        application: &str,
117        output: Option<saddle_observability::EmergencyDiagnosticHandle>,
118    ) {
119        assert!(self.requests.install_admission_source(application.into(), output),
120            "one Application installs one admission source binding");
121        self.lifecycle_observer = Some((observer, application.to_owned()));
122    }
123
124    /// Formal process assembly accepts only the output selected and opened
125    /// during checked startup. The legacy optional installation stays outside
126    /// this capability boundary.
127    #[doc(hidden)]
128    pub fn install_lifecycle_observer_with_source_output(
129        &mut self,
130        observer: saddle_observability::Observer,
131        application: &str,
132        output: saddle_observability::SourceOutput,
133    ) {
134        assert!(self.components.is_empty(),
135            "formal source output must precede component registration");
136        let label = saddle_core::ContextLabel::checked(application)
137            .expect("formal Application requires checked identity");
138        assert!(self.requests.install_admission_source(
139            application.into(), Some(output.handle().clone())),
140            "one Application installs one admission source binding");
141        self.recorded_boundary = Some(RecordedComponentBoundary::new(label, output));
142        self.lifecycle_observer = Some((observer, application.to_owned()));
143    }
144
145    #[cfg(all(target_arch = "x86_64", target_os = "linux"))]
146    pub(crate) fn install_prevalidated_components(
147        &mut self,
148        components: Vec<PrevalidatedLegacyComponent>,
149    ) {
150        debug_assert!(self.components.is_empty());
151        assert!(self.recorded_boundary.is_none(),
152            "prevalidated legacy components cannot replace formal recorded components");
153        self.components = components.into_iter().map(ApplicationComponent::Legacy).collect();
154    }
155
156    #[cfg(all(target_arch = "x86_64", target_os = "linux"))]
157    #[doc(hidden)]
158    pub fn pending_driver_finalizer(&self) -> crate::post_driver::PendingDriverFinalizerSlot {
159        self.pending_driver_finalizer.clone()
160    }
161
162    #[cfg(all(test, target_arch = "x86_64", target_os = "linux"))]
163    pub(crate) fn post_driver_is_unarmed_for_test(&self) -> bool {
164        self.pending_driver_finalizer.is_unarmed_for_test()
165    }
166
167    #[cfg(all(target_arch = "x86_64", target_os = "linux"))]
168    #[doc(hidden)]
169    #[allow(clippy::result_large_err)]
170    pub fn commit_post_driver_install(
171        &self,
172        binding: saddle_admission::VerifiedPostDriverInstallBinding,
173    ) {
174        self.pending_driver_finalizer
175            .commit_verified_install(binding)
176    }
177
178    pub(crate) fn reserved_post_driver_submit(
179        &self,
180    ) -> crate::post_driver::MustSubmitDriverFinalizer {
181        self.pending_driver_finalizer.reserved_submit_handle()
182    }
183
184    /// Returns the request lifecycle shared with Saddle's Service adapter.
185    pub fn request_lifecycle(&self) -> RequestLifecycle {
186        self.requests.clone()
187    }
188
189    /// Returns a read-only observer of the framework's unique lifecycle
190    /// state. The observer cannot admit requests or mutate readiness.
191    pub fn health(&self) -> crate::ApplicationHealth {
192        self.requests.health()
193    }
194
195    /// Reserves the fixed alpha.1 Ingress execution bridge attached to this
196    /// application's existing 0.2 request lifecycle.
197    #[doc(hidden)]
198    pub fn managed_ingress_bridge(
199        &self,
200        capacity: usize,
201    ) -> Option<crate::alpha1_ingress::ManagedIngressBridge> {
202        if capacity == 0 {
203            return None;
204        }
205        if self
206            .ingress_bridge_issued
207            .compare_exchange(false, true, Ordering::AcqRel, Ordering::Acquire)
208            .is_err()
209        {
210            return None;
211        }
212        crate::alpha1_ingress::ManagedIngressBridge::new(self.requests.clone(), capacity)
213    }
214
215    /// Registers a framework component for managed startup and shutdown.
216    ///
217    /// Components start in registration order and stop in reverse order.
218    pub fn register<C>(&mut self, component: C) -> Result<()>
219    where
220        C: ComponentLifecycle + 'static,
221    {
222        self.register_shared(Arc::new(component))
223    }
224
225    /// Registers an already shared framework component.
226    pub fn register_shared(&mut self, component: Arc<dyn ComponentLifecycle>) -> Result<()> {
227        if self.recorded_boundary.is_some() {
228            return Err(SaddleError::new(ErrorKind::Infrastructure,
229                "runtime.recorded_component_required",
230                "formal Application requires a recorded component"));
231        }
232        if self
233            .components
234            .iter()
235            .any(|registered| registered.name() == component.name())
236        {
237            return Err(SaddleError::new(
238                ErrorKind::Conflict,
239                "runtime.duplicate_component",
240                format!("component '{}' is already registered", component.name()),
241            ));
242        }
243        self.components.push(ApplicationComponent::Legacy(component));
244        Ok(())
245    }
246
247    /// Formal registration requires a selected source output before startup.
248    pub fn register_recorded<C>(&mut self, component: C) -> Result<()>
249    where C: RecordedComponentLifecycle + 'static {
250        self.register_recorded_shared(Arc::new(component))
251    }
252
253    /// Shares a recorded component with a lifecycle fault owner without
254    /// reopening the legacy registration path.
255    pub fn register_recorded_shared(
256        &mut self, component: Arc<dyn RecordedComponentLifecycle>,
257    ) -> Result<()> {
258        if self.recorded_boundary.is_none() {
259            return Err(SaddleError::new(ErrorKind::Infrastructure,
260                "runtime.source_output_required",
261                "formal component output has not been selected").with_unconfirmed_source());
262        }
263        if self.components.iter().any(|registered| registered.name() == component.name()) {
264            return Err(SaddleError::new(ErrorKind::Conflict,
265                "runtime.duplicate_component",
266                format!("component '{}' is already registered", component.name())));
267        }
268        self.components.push(ApplicationComponent::Recorded(component));
269        Ok(())
270    }
271
272    /// Runs the application on Saddle's single process-wide async runtime.
273    ///
274    /// The call blocks the process entry thread until SIGINT or, on Unix,
275    /// SIGTERM. Shutdown first closes request admission, then waits for every
276    /// admitted request, and finally stops components in reverse order.
277    pub fn run(self) -> Result<()> {
278        Self::run_with(|| async move { Ok(self) })
279    }
280
281    /// Creates the application inside Saddle's process-wide async runtime and
282    /// then runs it until shutdown.
283    ///
284    /// This is the framework assembly path for components whose initialization
285    /// performs async I/O. Business code is not given a runtime handle or an
286    /// executor through this API.
287    pub fn run_with<F, Fut>(bootstrap: F) -> Result<()>
288    where
289        F: FnOnce() -> Fut + Send + 'static,
290        Fut: Future<Output = Result<Self>> + Send + 'static,
291    {
292        if RUNTIME_STARTED
293            .compare_exchange(false, true, Ordering::AcqRel, Ordering::Acquire)
294            .is_err()
295        {
296            return Err(SaddleError::new(
297                ErrorKind::Conflict,
298                "runtime.already_started",
299                "the Saddle runtime has already started in this process",
300            ));
301        }
302
303        let runtime = build_runtime()?;
304
305        #[cfg(all(target_arch = "x86_64", target_os = "linux"))]
306        {
307            Self::run_with_owned_runtime(runtime, bootstrap)
308        }
309
310        #[cfg(not(all(target_arch = "x86_64", target_os = "linux")))]
311        runtime.block_on(async {
312            let signal = ShutdownSignal::register()?;
313            bootstrap_and_run(bootstrap, signal.wait()).await
314        })
315    }
316
317    /// Runs the formal process while retaining its one frozen deployment
318    /// budget inside Runtime assembly. No read or replacement surface escapes.
319    #[doc(hidden)]
320    pub fn run_with_deployment_resource_budget<F, Fut>(
321        budget: saddle_admission::DeploymentResourceBudget,
322        bootstrap: F,
323    ) -> Result<()>
324    where
325        F: FnOnce() -> Fut + Send + 'static,
326        Fut: Future<Output = Result<Self>> + Send + 'static,
327    {
328        Self::run_with(move || async move {
329            let mut application = bootstrap().await?;
330            if application.deployment_resource_budget.is_some() {
331                return Err(SaddleError::new(
332                    ErrorKind::Conflict,
333                    "runtime.deployment_resource_budget_already_installed",
334                    "the deployment resource budget was already installed",
335                ));
336            }
337            application.deployment_resource_budget = Some(budget);
338            Ok(application)
339        })
340    }
341
342    #[cfg(all(target_arch = "x86_64", target_os = "linux"))]
343    pub(crate) fn claim_process_runtime() -> Result<()> {
344        if RUNTIME_STARTED
345            .compare_exchange(false, true, Ordering::AcqRel, Ordering::Acquire)
346            .is_err()
347        {
348            return Err(SaddleError::new(
349                ErrorKind::Conflict,
350                "runtime.already_started",
351                "the Saddle runtime has already started in this process",
352            ));
353        }
354        Ok(())
355    }
356
357    #[cfg(all(target_arch = "x86_64", target_os = "linux"))]
358    pub(crate) fn run_with_owned_runtime<F, Fut>(
359        runtime: tokio::runtime::Runtime,
360        bootstrap: F,
361    ) -> Result<()>
362    where
363        F: FnOnce() -> Fut,
364        Fut: Future<Output = Result<Self>>,
365    {
366        let outcome = runtime.block_on(async {
367            let signal = ShutdownSignal::register()?;
368            let application = bootstrap().await?;
369            let finalizer = application.pending_driver_finalizer();
370            let shutdown_deadline = Arc::clone(&application.shutdown_deadline);
371            let lifecycle_observer = application.lifecycle_observer_handle();
372            let result = application.run_until_shutdown(signal.wait()).await;
373            Ok::<_, SaddleError>((finalizer, shutdown_deadline, lifecycle_observer, result))
374        });
375        match outcome {
376            Ok((finalizer, shutdown_deadline, lifecycle_observer, result)) => {
377                let deadline = *shutdown_deadline
378                    .lock()
379                    .unwrap_or_else(|poisoned| poisoned.into_inner());
380                finalizer.finish(runtime, result, deadline, lifecycle_observer)
381            }
382            Err(error) => {
383                drop(runtime);
384                Err(error)
385            }
386        }
387    }
388
389    pub(crate) async fn run_until_shutdown<F>(self, shutdown: F) -> Result<()>
390    where
391        F: Future<Output = Result<()>>,
392    {
393        tokio::pin!(shutdown);
394        let signal_before_start = tokio::select! {
395            biased;
396            signal_result = &mut shutdown => Some(signal_result),
397            _ = std::future::ready(()) => None,
398        };
399        if let Some(signal_result) = signal_before_start {
400            self.requests.begin_draining();
401            self.requests.wait_until_drained().await;
402            self.requests.mark_stopped();
403            return signal_result.map_err(|error| crate::diagnostics::shutdown_signal_source(
404                error, self.lifecycle_observer.as_ref().map(|(_, application)| application.as_str()),
405            ));
406        }
407
408        let mut started = 0;
409
410        for component in &self.components {
411            let start_started = Instant::now();
412            let start = crate::diagnostics::task(
413                async {
414                    match component {
415                        ApplicationComponent::Legacy(legacy) => legacy.start().await,
416                        ApplicationComponent::Recorded(component) => self.recorded_boundary
417                            .as_ref().expect("recorded component requires selected output")
418                            .start(component.as_ref()).await.map_err(|failure| failure.into_safe()),
419                    }
420                },
421                saddle_core::DiagnosticStage::StartupListener,
422                "runtime.component_start",
423            );
424            tokio::pin!(start);
425            let mut shutdown_during_start = None;
426            let start_result = tokio::select! {
427                biased;
428                signal_result = &mut shutdown => {
429                    let deadline = Instant::now() + self.lifecycle_timeouts.shutdown;
430                    self.set_shutdown_deadline(deadline);
431                    shutdown_during_start = Some((signal_result, deadline));
432                    // A component may have partially initialized before its
433                    // start future yielded. Bound that in-progress start by
434                    // both lifecycle budgets, then include it in rollback.
435                    let start_deadline = std::cmp::min(
436                        start_started + self.lifecycle_timeouts.start,
437                        deadline,
438                    );
439                    match tokio::time::timeout_at(start_deadline.into(), start).await {
440                        Ok(result) => result,
441                        Err(elapsed) => Err(match component {
442                            ApplicationComponent::Legacy(_) => lifecycle_timeout_error("component_start"),
443                            ApplicationComponent::Recorded(_) => self.recorded_boundary.as_ref()
444                                .expect("selected output").timeout_failure(elapsed,
445                                    saddle_observability::root_diagnostic::ComponentSourceKind::Start,
446                                    None).into_safe(),
447                        }),
448                    }
449                }
450                start_result = tokio::time::timeout(self.lifecycle_timeouts.start, &mut start) => {
451                    start_result.unwrap_or_else(|elapsed| Err(match component {
452                        ApplicationComponent::Legacy(_) => lifecycle_timeout_error("component_start"),
453                        ApplicationComponent::Recorded(_) => self.recorded_boundary.as_ref()
454                            .expect("selected output").timeout_failure(elapsed,
455                                saddle_observability::root_diagnostic::ComponentSourceKind::Start,
456                                None).into_safe(),
457                    }))
458                },
459            };
460
461            if let Err(error) = start_result {
462                // Fix cleanup's budget BEFORE any diagnostic capture/submit.
463                let deadline = shutdown_during_start
464                    .as_ref()
465                    .map(|(_, deadline)| *deadline)
466                    .unwrap_or_else(|| Instant::now() + self.lifecycle_timeouts.shutdown);
467                self.set_shutdown_deadline(deadline);
468                let error = if matches!(component, ApplicationComponent::Recorded(_)) {
469                    error
470                } else {
471                    crate::diagnostics::component_start_source(error,
472                        self.lifecycle_observer.as_ref().map(|(_, application)| application.as_str()))
473                };
474                self.record_timeout(&error, start_started.elapsed());
475                self.requests.begin_draining();
476                let cleanup_count = if error.code() == "runtime.lifecycle_timeout.component_start" {
477                    started + 1
478                } else {
479                    started
480                };
481                let drain_result = timeout_at(
482                    deadline,
483                    self.requests.wait_until_drained(),
484                    "request_drain",
485                )
486                .await;
487                let cleanup_source_unavailable = match drain_result {
488                    Ok(()) => self
489                        .shutdown_components(cleanup_count, deadline, Some(&error))
490                        .await
491                        .err()
492                        .is_some_and(|cleanup| cleanup.cleanup_source_unavailable()),
493                    Err(cleanup) => {
494                        let cleanup = crate::diagnostics::request_drain_source(
495                            cleanup,
496                            Some(&error),
497                            self.lifecycle_observer.as_ref().map(|(_, application)| application.as_str()),
498                        );
499                        cleanup.cleanup_source_unavailable()
500                    }
501                };
502                self.requests.mark_stopped();
503                return Err(if cleanup_source_unavailable {
504                    error.with_unconfirmed_cleanup_source()
505                } else { error });
506            }
507            started += 1;
508
509            if let Some((signal_result, deadline)) = shutdown_during_start {
510                self.requests.begin_draining();
511                let drain_result = timeout_at(
512                    deadline,
513                    self.requests.wait_until_drained(),
514                    "request_drain",
515                )
516                .await;
517                let signal_result = signal_result.map_err(|e| {
518                    crate::diagnostics::shutdown_signal_source(
519                        e, self.lifecycle_observer.as_ref().map(|(_, application)| application.as_str()),
520                    )
521                });
522                let shutdown_result = if drain_result.is_ok() {
523                    self.shutdown_components(started, deadline, signal_result.as_ref().err())
524                        .await
525                } else {
526                    drain_result.map_err(|e| {
527                        crate::diagnostics::request_drain_source(
528                            e,
529                            signal_result.as_ref().err(),
530                            self.lifecycle_observer.as_ref().map(|(_, application)| application.as_str()),
531                        )
532                    })
533                };
534                self.requests.mark_stopped();
535                return combine_lifecycle_results(signal_result, shutdown_result);
536            }
537        }
538
539        self.requests.mark_ready();
540        let signal_result = shutdown.await;
541        let deadline = Instant::now() + self.lifecycle_timeouts.shutdown;
542        self.set_shutdown_deadline(deadline);
543        self.requests.begin_draining();
544        let signal_result = signal_result.map_err(|e| {
545            crate::diagnostics::shutdown_signal_source(
546                e, self.lifecycle_observer.as_ref().map(|(_, application)| application.as_str()),
547            )
548        });
549        let drain_result = timeout_at(
550            deadline,
551            self.requests.wait_until_drained(),
552            "request_drain",
553        )
554        .await;
555        if let Err(error) = &drain_result {
556            self.record_timeout(error, self.lifecycle_timeouts.shutdown);
557        }
558        let shutdown_result = if drain_result.is_ok() {
559            self.shutdown_components(started, deadline, signal_result.as_ref().err())
560                .await
561        } else {
562            drain_result.map_err(|e| {
563                crate::diagnostics::request_drain_source(
564                    e,
565                    signal_result.as_ref().err(),
566                    self.lifecycle_observer.as_ref().map(|(_, application)| application.as_str()),
567                )
568            })
569        };
570        self.requests.mark_stopped();
571
572        combine_lifecycle_results(signal_result, shutdown_result)
573    }
574
575    fn set_shutdown_deadline(&self, deadline: Instant) {
576        *self
577            .shutdown_deadline
578            .lock()
579            .unwrap_or_else(|poisoned| poisoned.into_inner()) = Some(deadline);
580    }
581
582    pub(crate) fn shutdown_deadline_handle(&self) -> Arc<std::sync::Mutex<Option<Instant>>> {
583        Arc::clone(&self.shutdown_deadline)
584    }
585
586    pub(crate) fn lifecycle_observer_handle(
587        &self,
588    ) -> Option<(saddle_observability::Observer, String)> {
589        self.lifecycle_observer.clone()
590    }
591
592    fn record_timeout(&self, error: &SaddleError, elapsed: Duration) {
593        let stage = match error.code() {
594            "runtime.lifecycle_timeout.component_start" => {
595                saddle_observability::LifecycleTimeoutStage::ComponentStart
596            }
597            "runtime.lifecycle_timeout.request_drain" => {
598                saddle_observability::LifecycleTimeoutStage::RequestDrain
599            }
600            "runtime.lifecycle_timeout.component_shutdown" => {
601                saddle_observability::LifecycleTimeoutStage::ComponentShutdown
602            }
603            _ => return,
604        };
605        if let Some((observer, application)) = &self.lifecycle_observer {
606            observer.record_lifecycle_timeout(
607                application.as_str(),
608                stage,
609                u64::try_from(elapsed.as_millis()).unwrap_or(u64::MAX),
610            );
611        }
612    }
613
614    async fn shutdown_components(
615        &self,
616        started: usize,
617        deadline: Instant,
618        primary: Option<&SaddleError>,
619    ) -> Result<()> {
620        let mut first_error = None;
621        let mut cleanup_source_unavailable = false;
622        for component in self.components[..started].iter().rev() {
623            let result = match tokio::time::timeout_at(
624                deadline.into(),
625                crate::diagnostics::task(
626                    async {
627                        let parent = primary.or(first_error.as_ref())
628                            .and_then(|error| error.diagnostic().map(|d| d.occurrence()));
629                        match component {
630                            ApplicationComponent::Legacy(legacy) =>
631                                legacy.shutdown_with_primary(parent).await,
632                            ApplicationComponent::Recorded(component) => self.recorded_boundary
633                                .as_ref().expect("recorded component requires selected output")
634                                .shutdown(component.as_ref(), parent).await
635                                .map_err(|failure| {
636                                    let safe = failure.into_safe();
637                                    if safe.source_unavailable() {
638                                        safe.with_unconfirmed_cleanup_source()
639                                    } else { safe }
640                                }),
641                        }
642                    },
643                    saddle_core::DiagnosticStage::ShutdownComponent,
644                    "runtime.component_shutdown",
645                ),
646            )
647            .await
648            {
649                Ok(result) => result,
650                Err(elapsed) => Err(match component {
651                    ApplicationComponent::Legacy(_) => lifecycle_timeout_error("component_shutdown"),
652                    ApplicationComponent::Recorded(_) => {
653                        let safe = self.recorded_boundary.as_ref()
654                            .expect("selected output").timeout_failure(elapsed,
655                                saddle_observability::root_diagnostic::ComponentSourceKind::Cleanup,
656                                primary.or(first_error.as_ref())
657                                    .and_then(|error| error.diagnostic().map(|d| d.occurrence())))
658                            .into_safe();
659                        if safe.source_unavailable() { safe.with_unconfirmed_cleanup_source() }
660                        else { safe }
661                    }
662                }),
663            };
664            if let Err(error) = result {
665                let error = if matches!(component, ApplicationComponent::Recorded(_)) {
666                    error
667                } else {
668                    crate::diagnostics::component_cleanup_source(error,
669                        primary.or(first_error.as_ref()),
670                        self.lifecycle_observer.as_ref().map(|(_, application)| application.as_str()))
671                };
672                cleanup_source_unavailable |= error.cleanup_source_unavailable();
673                self.record_timeout(&error, self.lifecycle_timeouts.shutdown);
674                if first_error.is_none() {
675                    first_error = Some(error);
676                }
677                if Instant::now() >= deadline {
678                    break;
679                }
680            }
681        }
682        first_error.map_or(Ok(()), |error| Err(if cleanup_source_unavailable {
683            error.with_unconfirmed_cleanup_source()
684        } else { error }))
685    }
686}
687
688pub(crate) fn combine_lifecycle_results(primary: Result<()>, cleanup: Result<()>) -> Result<()> {
689    match (primary, cleanup) {
690        (Err(primary), Err(cleanup)) if cleanup.cleanup_source_unavailable() => {
691            Err(primary.with_unconfirmed_cleanup_source())
692        }
693        (Err(primary), _) => Err(primary),
694        (Ok(()), result) => result,
695    }
696}
697
698async fn timeout_at<T>(
699    deadline: Instant,
700    future: impl Future<Output = T>,
701    stage: &'static str,
702) -> Result<T> {
703    let remaining = deadline.saturating_duration_since(Instant::now());
704    tokio::time::timeout(remaining, future)
705        .await
706        .map_err(|_| lifecycle_timeout_error(stage))
707}
708
709fn lifecycle_timeout_error(stage: &'static str) -> SaddleError {
710    SaddleError::new(
711        ErrorKind::Infrastructure,
712        match stage {
713            "component_start" => "runtime.lifecycle_timeout.component_start",
714            "request_drain" => "runtime.lifecycle_timeout.request_drain",
715            "component_shutdown" => "runtime.lifecycle_timeout.component_shutdown",
716            _ => "runtime.lifecycle_timeout.post_driver",
717        },
718        format!("managed lifecycle stage '{stage}' exceeded its wall-clock deadline"),
719    )
720}
721
722#[cfg(any(test, not(all(target_arch = "x86_64", target_os = "linux"))))]
723async fn bootstrap_and_run<F, Fut, S>(bootstrap: F, shutdown: S) -> Result<()>
724where
725    F: FnOnce() -> Fut,
726    Fut: Future<Output = Result<Application>>,
727    S: Future<Output = Result<()>>,
728{
729    let application = bootstrap().await?;
730    application.run_until_shutdown(shutdown).await
731}
732
733impl Default for Application {
734    fn default() -> Self {
735        Self::new()
736    }
737}
738
739fn build_runtime() -> Result<tokio::runtime::Runtime> {
740    tokio::runtime::Builder::new_multi_thread()
741        .worker_threads(WORKER_THREADS)
742        .max_io_events_per_tick(MAX_IO_EVENTS_PER_TICK)
743        .enable_all()
744        .build()
745        .map_err(|_| {
746            SaddleError::new(
747                ErrorKind::Infrastructure,
748                "runtime.initialization_failed",
749                "failed to initialize the Saddle async runtime",
750            )
751        })
752}
753
754#[cfg(all(target_arch = "x86_64", target_os = "linux"))]
755pub(crate) fn claim_owned_runtime() -> Result<tokio::runtime::Runtime> {
756    Application::claim_process_runtime()?;
757    build_runtime()
758}
759
760#[cfg(unix)]
761pub(crate) struct ShutdownSignal {
762    interrupt: tokio::signal::unix::Signal,
763    terminate: tokio::signal::unix::Signal,
764}
765
766#[cfg(unix)]
767impl ShutdownSignal {
768    /// Registers both listeners synchronously before any component starts.
769    pub(crate) fn register() -> Result<Self> {
770        Ok(Self {
771            interrupt: tokio::signal::unix::signal(tokio::signal::unix::SignalKind::interrupt())
772                .map_err(|_| signal_error())?,
773            terminate: tokio::signal::unix::signal(tokio::signal::unix::SignalKind::terminate())
774                .map_err(|_| signal_error())?,
775        })
776    }
777
778    pub(crate) async fn wait(mut self) -> Result<()> {
779        tokio::select! {
780            _ = self.interrupt.recv() => Ok(()),
781            _ = self.terminate.recv() => Ok(()),
782        }
783    }
784}
785
786#[cfg(windows)]
787struct ShutdownSignal {
788    ctrl_c: tokio::signal::windows::CtrlC,
789    ctrl_break: tokio::signal::windows::CtrlBreak,
790}
791
792#[cfg(windows)]
793impl ShutdownSignal {
794    /// Registers both listeners synchronously before any component starts.
795    fn register() -> Result<Self> {
796        Ok(Self {
797            ctrl_c: tokio::signal::windows::ctrl_c().map_err(|_| signal_error())?,
798            ctrl_break: tokio::signal::windows::ctrl_break().map_err(|_| signal_error())?,
799        })
800    }
801
802    async fn wait(mut self) -> Result<()> {
803        tokio::select! {
804            _ = self.ctrl_c.recv() => Ok(()),
805            _ = self.ctrl_break.recv() => Ok(()),
806        }
807    }
808}
809
810fn signal_error() -> SaddleError {
811    SaddleError::new(
812        ErrorKind::Infrastructure,
813        "runtime.signal_registration_failed",
814        "failed to register the application shutdown signal",
815    )
816}
817
818#[cfg(test)]
819mod tests {
820    use std::sync::Mutex;
821
822    use saddle_core::LifecycleFuture;
823
824    use super::*;
825    use crate::ApplicationPhase;
826
827    struct RecordedStartFailure;
828    struct RecordedPendingStart;
829    struct RecordedPendingCleanup;
830    impl RecordedComponentLifecycle for RecordedPendingCleanup {
831        fn name(&self) -> &'static str { "recorded-pending-cleanup" }
832        fn start<'a>(&'a self, _: &'a saddle_core::ContextLabel,
833            _: &'a saddle_observability::SourceOutput)
834            -> saddle_observability::root_diagnostic::RecordedLifecycleFuture<'a> {
835            Box::pin(async { Ok(()) })
836        }
837        fn shutdown<'a>(&'a self, _: &'a saddle_core::ContextLabel,
838            _: &'a saddle_observability::SourceOutput)
839            -> saddle_observability::root_diagnostic::RecordedLifecycleFuture<'a> {
840            Box::pin(std::future::pending())
841        }
842        fn shutdown_with_primary<'a>(&'a self, _: &'a saddle_core::ContextLabel,
843            _: &'a saddle_observability::SourceOutput,
844            _: Option<saddle_core::DiagnosticOccurrence>)
845            -> saddle_observability::root_diagnostic::RecordedLifecycleFuture<'a> {
846            Box::pin(std::future::pending())
847        }
848    }
849    impl RecordedComponentLifecycle for RecordedPendingStart {
850        fn name(&self) -> &'static str { "recorded-pending-start" }
851        fn start<'a>(&'a self, _: &'a saddle_core::ContextLabel,
852            _: &'a saddle_observability::SourceOutput)
853            -> saddle_observability::root_diagnostic::RecordedLifecycleFuture<'a> {
854            Box::pin(std::future::pending())
855        }
856        fn shutdown<'a>(&'a self, _: &'a saddle_core::ContextLabel,
857            _: &'a saddle_observability::SourceOutput)
858            -> saddle_observability::root_diagnostic::RecordedLifecycleFuture<'a> {
859            Box::pin(async { Ok(()) })
860        }
861        fn shutdown_with_primary<'a>(&'a self, _: &'a saddle_core::ContextLabel,
862            _: &'a saddle_observability::SourceOutput,
863            _: Option<saddle_core::DiagnosticOccurrence>)
864            -> saddle_observability::root_diagnostic::RecordedLifecycleFuture<'a> {
865            Box::pin(async { Ok(()) })
866        }
867    }
868    impl RecordedComponentLifecycle for RecordedStartFailure {
869        fn name(&self) -> &'static str { "recorded-start-failure" }
870        fn start<'a>(&'a self, application: &'a saddle_core::ContextLabel,
871            output: &'a saddle_observability::SourceOutput)
872            -> saddle_observability::root_diagnostic::RecordedLifecycleFuture<'a> {
873            Box::pin(async move {
874                saddle_observability::root_diagnostic::RecordedSaddleError::component_result::<(), _>(
875                    Err(std::io::Error::other("M_ACTUAL_START_ORIGINAL")), application, output,
876                    saddle_observability::root_diagnostic::ComponentSourceKind::Start, None,
877                    ErrorKind::Infrastructure, "runtime.recorded_start_test", "safe start failure")
878            })
879        }
880        fn shutdown<'a>(&'a self, _: &'a saddle_core::ContextLabel,
881            _: &'a saddle_observability::SourceOutput)
882            -> saddle_observability::root_diagnostic::RecordedLifecycleFuture<'a> {
883            Box::pin(async { Ok(()) })
884        }
885        fn shutdown_with_primary<'a>(&'a self, _: &'a saddle_core::ContextLabel,
886            _: &'a saddle_observability::SourceOutput,
887            _: Option<saddle_core::DiagnosticOccurrence>)
888            -> saddle_observability::root_diagnostic::RecordedLifecycleFuture<'a> {
889            Box::pin(async { Ok(()) })
890        }
891    }
892
893    #[test]
894    fn formal_application_requires_output_and_reads_original_before_return() {
895        let mut absent = Application::new();
896        assert_eq!(absent.register_recorded(RecordedStartFailure).unwrap_err().code(),
897            "runtime.source_output_required");
898        let root = std::env::temp_dir().join(format!("saddle-m-recorded-{}-{}",
899            std::process::id(), std::time::SystemTime::now()
900                .duration_since(std::time::UNIX_EPOCH).unwrap().as_nanos()));
901        std::fs::create_dir(&root).unwrap();
902        let mut output = saddle_observability::EmergencyDiagnostics::start_checked(
903            &saddle_observability::FileLoggingConfig::new(root.clone(),
904                saddle_observability::Rotation::Daily)).unwrap();
905        let target = output.target().to_owned();
906        let observer = saddle_observability::Observer::with_writer(
907            Default::default(), std::io::sink()).unwrap();
908        let mut app = Application::new();
909        app.install_lifecycle_observer_with_source_output(observer, "m-formal-app",
910            output.source_output().unwrap());
911        assert_eq!(app.register(RecordingComponent {
912            name: "legacy", events: Arc::new(Mutex::new(Vec::new())),
913            start_error: false, shutdown_error: false,
914        }).unwrap_err().code(), "runtime.recorded_component_required");
915        app.register_recorded(RecordedStartFailure).unwrap();
916        let failed = test_runtime().block_on(app.run_until_shutdown(
917            std::future::pending::<Result<()>>())).unwrap_err();
918        assert_eq!(failed.code(), "runtime.recorded_start_test");
919        let id = failed.diagnostic().unwrap().id().to_string();
920        let raw = std::fs::read_to_string(&target).unwrap();
921        assert!(raw.contains("M_ACTUAL_START_ORIGINAL"));
922        assert!(raw.contains(&id));
923        assert!(!failed.source_unavailable());
924        while output.shutdown() == saddle_observability::DiagnosticShutdown::Pending {
925            std::thread::yield_now();
926        }
927        drop(output);
928        std::fs::remove_file(target).unwrap();
929        std::fs::remove_dir(root).unwrap();
930    }
931
932    #[test]
933    fn formal_application_timeout_records_actual_elapsed_before_return() {
934        let root = std::env::temp_dir().join(format!("saddle-m-timeout-{}-{}",
935            std::process::id(), std::time::SystemTime::now()
936                .duration_since(std::time::UNIX_EPOCH).unwrap().as_nanos()));
937        std::fs::create_dir(&root).unwrap();
938        let mut output = saddle_observability::EmergencyDiagnostics::start_checked(
939            &saddle_observability::FileLoggingConfig::new(root.clone(),
940                saddle_observability::Rotation::Daily)).unwrap();
941        let target = output.target().to_owned();
942        let observer = saddle_observability::Observer::with_writer(
943            Default::default(), std::io::sink()).unwrap();
944        let mut app = Application::new();
945        app.set_lifecycle_timeouts(LifecycleTimeouts::from_millis(20, 100).unwrap());
946        app.install_lifecycle_observer_with_source_output(observer, "m-timeout-app",
947            output.source_output().unwrap());
948        app.register_recorded(RecordedPendingStart).unwrap();
949        let failed = test_runtime().block_on(app.run_until_shutdown(
950            std::future::pending::<Result<()>>())).unwrap_err();
951        assert_eq!(failed.code(), "runtime.lifecycle_timeout.component_start");
952        assert!(!failed.source_unavailable());
953        let raw = std::fs::read_to_string(&target).unwrap();
954        assert!(raw.contains("deadline has elapsed"));
955        assert!(raw.contains(&failed.diagnostic().unwrap().id().to_string()));
956        while output.shutdown() == saddle_observability::DiagnosticShutdown::Pending {
957            std::thread::yield_now();
958        }
959        drop(output);
960        std::fs::remove_file(target).unwrap();
961        std::fs::remove_dir(root).unwrap();
962    }
963
964    #[test]
965    fn formal_application_cleanup_timeout_records_actual_elapsed() {
966        let root = std::env::temp_dir().join(format!("saddle-m-cleanup-{}-{}",
967            std::process::id(), std::time::SystemTime::now()
968                .duration_since(std::time::UNIX_EPOCH).unwrap().as_nanos()));
969        std::fs::create_dir(&root).unwrap();
970        let mut output = saddle_observability::EmergencyDiagnostics::start_checked(
971            &saddle_observability::FileLoggingConfig::new(root.clone(),
972                saddle_observability::Rotation::Daily)).unwrap();
973        let target = output.target().to_owned();
974        let observer = saddle_observability::Observer::with_writer(
975            Default::default(), std::io::sink()).unwrap();
976        let mut app = Application::new();
977        app.set_lifecycle_timeouts(LifecycleTimeouts::from_millis(100, 20).unwrap());
978        app.install_lifecycle_observer_with_source_output(observer, "m-cleanup-app",
979            output.source_output().unwrap());
980        app.register_recorded(RecordedPendingCleanup).unwrap();
981        let shutdown = shutdown_when_ready(app.request_lifecycle());
982        let failed = test_runtime().block_on(app.run_until_shutdown(shutdown)).unwrap_err();
983        assert_eq!(failed.code(), "runtime.lifecycle_timeout.component_shutdown");
984        assert!(!failed.cleanup_source_unavailable());
985        let raw = std::fs::read_to_string(&target).unwrap();
986        assert!(raw.contains("deadline has elapsed"));
987        assert!(raw.contains(&failed.diagnostic().unwrap().id().to_string()));
988        while output.shutdown() == saddle_observability::DiagnosticShutdown::Pending {
989            std::thread::yield_now();
990        }
991        drop(output);
992        std::fs::remove_file(target).unwrap();
993        std::fs::remove_dir(root).unwrap();
994    }
995
996    struct RecordingComponent {
997        name: &'static str,
998        events: Arc<Mutex<Vec<String>>>,
999        start_error: bool,
1000        shutdown_error: bool,
1001    }
1002
1003    struct BlockingStartComponent {
1004        events: Arc<Mutex<Vec<String>>>,
1005        started: Mutex<Option<tokio::sync::oneshot::Sender<()>>>,
1006        release: Mutex<Option<tokio::sync::oneshot::Receiver<()>>>,
1007    }
1008
1009    struct BlockingShutdownComponent {
1010        events: Arc<Mutex<Vec<String>>>,
1011    }
1012
1013    struct HealthAwareListener {
1014        health: crate::ApplicationHealth,
1015        events: Arc<Mutex<Vec<String>>>,
1016    }
1017
1018    impl ComponentLifecycle for RecordingComponent {
1019        fn name(&self) -> &'static str {
1020            self.name
1021        }
1022
1023        fn start(&self) -> LifecycleFuture<'_> {
1024            Box::pin(async move {
1025                self.events
1026                    .lock()
1027                    .unwrap()
1028                    .push(format!("start:{}", self.name));
1029                if self.start_error {
1030                    Err(test_error("start failed"))
1031                } else {
1032                    Ok(())
1033                }
1034            })
1035        }
1036
1037        fn shutdown(&self) -> LifecycleFuture<'_> {
1038            Box::pin(async move {
1039                self.events
1040                    .lock()
1041                    .unwrap()
1042                    .push(format!("shutdown:{}", self.name));
1043                if self.shutdown_error {
1044                    Err(test_error("shutdown failed"))
1045                } else {
1046                    Ok(())
1047                }
1048            })
1049        }
1050    }
1051
1052    impl ComponentLifecycle for BlockingStartComponent {
1053        fn name(&self) -> &'static str {
1054            "blocking"
1055        }
1056
1057        fn start(&self) -> LifecycleFuture<'_> {
1058            Box::pin(async move {
1059                self.events
1060                    .lock()
1061                    .unwrap()
1062                    .push("start:blocking".to_owned());
1063                let started = self.started.lock().unwrap().take().unwrap();
1064                let release = self.release.lock().unwrap().take().unwrap();
1065                started.send(()).unwrap();
1066                release.await.unwrap();
1067                Ok(())
1068            })
1069        }
1070
1071        fn shutdown(&self) -> LifecycleFuture<'_> {
1072            Box::pin(async move {
1073                self.events
1074                    .lock()
1075                    .unwrap()
1076                    .push("shutdown:blocking".to_owned());
1077                Ok(())
1078            })
1079        }
1080    }
1081
1082    impl ComponentLifecycle for BlockingShutdownComponent {
1083        fn name(&self) -> &'static str {
1084            "blocking-shutdown"
1085        }
1086
1087        fn start(&self) -> LifecycleFuture<'_> {
1088            Box::pin(async move {
1089                self.events
1090                    .lock()
1091                    .unwrap()
1092                    .push("start:blocking-shutdown".into());
1093                Ok(())
1094            })
1095        }
1096
1097        fn shutdown(&self) -> LifecycleFuture<'_> {
1098            Box::pin(async move {
1099                self.events
1100                    .lock()
1101                    .unwrap()
1102                    .push("shutdown:blocking-shutdown".into());
1103                std::future::pending().await
1104            })
1105        }
1106    }
1107
1108    impl ComponentLifecycle for HealthAwareListener {
1109        fn name(&self) -> &'static str {
1110            "health-aware-listener"
1111        }
1112
1113        fn start(&self) -> LifecycleFuture<'_> {
1114            Box::pin(async move {
1115                let snapshot = self.health.snapshot();
1116                assert!(snapshot.is_live());
1117                assert!(!snapshot.is_ready());
1118                assert_eq!(snapshot.phase(), ApplicationPhase::Starting);
1119                self.events
1120                    .lock()
1121                    .unwrap()
1122                    .push("listener:accepting".into());
1123                Ok(())
1124            })
1125        }
1126
1127        fn shutdown(&self) -> LifecycleFuture<'_> {
1128            Box::pin(async move {
1129                let snapshot = self.health.snapshot();
1130                assert!(snapshot.is_live());
1131                assert!(!snapshot.is_ready());
1132                assert_eq!(snapshot.phase(), ApplicationPhase::Draining);
1133                self.events.lock().unwrap().push("listener:stopped".into());
1134                Ok(())
1135            })
1136        }
1137    }
1138
1139    fn component(name: &'static str, events: &Arc<Mutex<Vec<String>>>) -> RecordingComponent {
1140        RecordingComponent {
1141            name,
1142            events: Arc::clone(events),
1143            start_error: false,
1144            shutdown_error: false,
1145        }
1146    }
1147
1148    fn test_error(message: &'static str) -> SaddleError {
1149        SaddleError::new(ErrorKind::Infrastructure, "test.failure", message)
1150    }
1151
1152    fn test_runtime() -> tokio::runtime::Runtime {
1153        tokio::runtime::Builder::new_current_thread()
1154            .enable_time()
1155            .build()
1156            .expect("test runtime must build")
1157    }
1158
1159    async fn shutdown_when_ready(requests: RequestLifecycle) -> Result<()> {
1160        while requests.phase() != crate::ApplicationPhase::Ready {
1161            tokio::task::yield_now().await;
1162        }
1163        Ok(())
1164    }
1165
1166    #[test]
1167    fn components_start_in_order_and_shutdown_in_reverse() {
1168        let events = Arc::new(Mutex::new(Vec::new()));
1169        let mut application = Application::new();
1170        application.register(component("db", &events)).unwrap();
1171        application.register(component("service", &events)).unwrap();
1172        let shutdown = shutdown_when_ready(application.request_lifecycle());
1173
1174        test_runtime()
1175            .block_on(application.run_until_shutdown(shutdown))
1176            .unwrap();
1177
1178        assert_eq!(
1179            *events.lock().unwrap(),
1180            [
1181                "start:db",
1182                "start:service",
1183                "shutdown:service",
1184                "shutdown:db"
1185            ]
1186        );
1187    }
1188
1189    #[test]
1190    fn health_uses_the_unique_lifecycle_and_clears_ready_before_listener_shutdown() {
1191        test_runtime().block_on(async {
1192            let events = Arc::new(Mutex::new(Vec::new()));
1193            let mut application = Application::new();
1194            let health = application.health();
1195            let initial = health.snapshot();
1196            assert!(initial.is_live());
1197            assert!(!initial.is_ready());
1198            assert_eq!(initial.phase(), ApplicationPhase::Starting);
1199
1200            application
1201                .register(HealthAwareListener {
1202                    health: health.clone(),
1203                    events: Arc::clone(&events),
1204                })
1205                .unwrap();
1206            let shutdown_health = health.clone();
1207            application
1208                .run_until_shutdown(async move {
1209                    loop {
1210                        let snapshot = shutdown_health.snapshot();
1211                        if snapshot.is_ready() {
1212                            assert!(snapshot.is_live());
1213                            assert_eq!(snapshot.phase(), ApplicationPhase::Ready);
1214                            return Ok(());
1215                        }
1216                        tokio::task::yield_now().await;
1217                    }
1218                })
1219                .await
1220                .unwrap();
1221
1222            let stopped = health.snapshot();
1223            assert!(!stopped.is_live());
1224            assert!(!stopped.is_ready());
1225            assert_eq!(stopped.phase(), ApplicationPhase::Stopped);
1226            assert_eq!(
1227                *events.lock().unwrap(),
1228                ["listener:accepting", "listener:stopped"]
1229            );
1230        });
1231    }
1232
1233    #[test]
1234    fn async_bootstrap_runs_before_early_shutdown_prevents_component_start() {
1235        let events = Arc::new(Mutex::new(Vec::new()));
1236        let bootstrap_events = Arc::clone(&events);
1237
1238        test_runtime()
1239            .block_on(bootstrap_and_run(
1240                move || async move {
1241                    bootstrap_events
1242                        .lock()
1243                        .unwrap()
1244                        .push("bootstrap".to_owned());
1245                    let mut application = Application::new();
1246                    application.register(component("component", &bootstrap_events))?;
1247                    Ok(application)
1248                },
1249                std::future::ready(Ok(())),
1250            ))
1251            .unwrap();
1252
1253        assert_eq!(*events.lock().unwrap(), ["bootstrap"]);
1254    }
1255
1256    #[test]
1257    fn failed_async_bootstrap_does_not_start_components() {
1258        let error = test_runtime()
1259            .block_on(bootstrap_and_run(
1260                || async { Err(test_error("bootstrap failed")) },
1261                std::future::pending(),
1262            ))
1263            .unwrap_err();
1264        assert_eq!(error.message(), "bootstrap failed");
1265    }
1266
1267    #[test]
1268    fn startup_failure_rolls_back_only_started_components() {
1269        let events = Arc::new(Mutex::new(Vec::new()));
1270        let mut application = Application::new();
1271        application.register(component("first", &events)).unwrap();
1272        let mut failing = component("failing", &events);
1273        failing.start_error = true;
1274        application.register(failing).unwrap();
1275        application.register(component("never", &events)).unwrap();
1276        let shutdown = shutdown_when_ready(application.request_lifecycle());
1277
1278        let error = test_runtime()
1279            .block_on(application.run_until_shutdown(shutdown))
1280            .unwrap_err();
1281
1282        assert_eq!(error.message(), "start failed");
1283        assert_eq!(
1284            *events.lock().unwrap(),
1285            ["start:first", "start:failing", "shutdown:first"]
1286        );
1287    }
1288
1289    #[test]
1290    fn startup_primary_retains_unconfirmed_cleanup_fact() {
1291        let events = Arc::new(Mutex::new(Vec::new()));
1292        let mut application = Application::new();
1293        let mut cleanup = component("cleanup", &events);
1294        cleanup.shutdown_error = true;
1295        application.register(cleanup).unwrap();
1296        let mut startup = component("startup", &events);
1297        startup.start_error = true;
1298        application.register(startup).unwrap();
1299        let primary = test_runtime().block_on(application.run_until_shutdown(
1300            std::future::pending::<Result<()>>(),
1301        )).unwrap_err();
1302        assert_eq!(primary.code(), "test.failure");
1303        assert_eq!(primary.message(), "start failed");
1304        assert!(primary.source_unavailable());
1305        assert!(primary.cleanup_source_unavailable());
1306    }
1307
1308    #[test]
1309    fn shutdown_continues_after_a_component_error() {
1310        let events = Arc::new(Mutex::new(Vec::new()));
1311        let mut application = Application::new();
1312        application.register(component("first", &events)).unwrap();
1313        let mut failing = component("second", &events);
1314        failing.shutdown_error = true;
1315        application.register(failing).unwrap();
1316        let shutdown = shutdown_when_ready(application.request_lifecycle());
1317
1318        let error = test_runtime()
1319            .block_on(application.run_until_shutdown(shutdown))
1320            .unwrap_err();
1321
1322        assert_eq!(error.message(), "shutdown failed");
1323        assert_eq!(
1324            *events.lock().unwrap(),
1325            [
1326                "start:first",
1327                "start:second",
1328                "shutdown:second",
1329                "shutdown:first"
1330            ]
1331        );
1332    }
1333
1334    #[test]
1335    fn component_cleanup_original_precedes_output_close() {
1336        const CHILD: &str = "SADDLE_COMPONENT_CLEANUP_ORIGINAL_CHILD";
1337        if let Some(path) = std::env::var_os(CHILD) {
1338            use saddle_observability::{EmergencyDiagnostics, FileLoggingConfig, Rotation};
1339            let output = EmergencyDiagnostics::start(&FileLoggingConfig::new(
1340                std::path::PathBuf::from(&path), Rotation::Daily,
1341            )).unwrap();
1342            assert!(crate::diagnostics::install_output(output.handle()).is_ok());
1343            let events = Arc::new(Mutex::new(Vec::new()));
1344            let mut application = Application::new();
1345            application.install_lifecycle_observer(
1346                saddle_observability::Observer::with_writer(Default::default(), std::io::sink()).unwrap(),
1347                "test-app",
1348            );
1349            let mut cleanup = component("cleanup", &events);
1350            cleanup.shutdown_error = true;
1351            application.register(cleanup).unwrap();
1352            let mut startup = component("startup", &events);
1353            startup.start_error = true;
1354            application.register(startup).unwrap();
1355            let primary = test_runtime().block_on(application.run_until_shutdown(
1356                std::future::pending::<Result<()>>(),
1357            )).unwrap_err();
1358            let primary_id = primary.diagnostic().unwrap().id();
1359            assert!(!primary.source_unavailable());
1360            assert!(!primary.cleanup_source_unavailable());
1361            let mut signal_application = Application::new();
1362            signal_application.install_lifecycle_observer(
1363                saddle_observability::Observer::with_writer(Default::default(), std::io::sink()).unwrap(),
1364                "test-app",
1365            );
1366            let signal = test_runtime().block_on(signal_application.run_until_shutdown(
1367                async { Err(signal_error()) },
1368            )).unwrap_err();
1369            assert!(!signal.source_unavailable());
1370            let signal_id = signal.diagnostic().unwrap().id();
1371            let mut drain_application = Application::new();
1372            drain_application.set_lifecycle_timeouts(LifecycleTimeouts::from_millis(100, 10).unwrap());
1373            drain_application.install_lifecycle_observer(
1374                saddle_observability::Observer::with_writer(Default::default(), std::io::sink()).unwrap(),
1375                "test-app",
1376            );
1377            let requests = drain_application.request_lifecycle();
1378            let drain = test_runtime().block_on(drain_application.run_until_shutdown(async move {
1379                while requests.phase() != crate::ApplicationPhase::Ready {
1380                    tokio::task::yield_now().await;
1381                }
1382                let request = requests.try_accept().unwrap();
1383                tokio::spawn(async move {
1384                    tokio::time::sleep(Duration::from_millis(50)).await;
1385                    drop(request);
1386                });
1387                Ok(())
1388            })).unwrap_err();
1389            assert_eq!(drain.code(), "runtime.lifecycle_timeout.request_drain");
1390            assert!(!drain.source_unavailable());
1391            let drain_id = drain.diagnostic().unwrap().id();
1392            let exit = crate::diagnostics::close_output(
1393                output, Some(Instant::now() + Duration::from_secs(2)),
1394            );
1395            assert_eq!(exit.shutdown, saddle_observability::DiagnosticShutdown::Finished);
1396            let text = std::fs::read_to_string(
1397                std::path::PathBuf::from(path).join("saddle.emergency.log"),
1398            ).unwrap();
1399            let rows: Vec<serde_json::Value> = text.lines()
1400                .map(|line| serde_json::from_str(line).unwrap()).collect();
1401            let field = |rows: &[serde_json::Value], occurrence: &serde_json::Value, channel: &str| -> Option<String> {
1402                let segments: Vec<_> = rows.iter().filter(|row| row["event"] == "request_error_original"
1403                    && row["occurrence"] == *occurrence && row["channel"] == channel
1404                    && row["cause_depth"] == 0).collect();
1405                if segments.last()?["state"] != "field_end" { return None; }
1406                Some(segments.iter().map(|row| row["payload"].as_str().unwrap()).collect::<String>())
1407            };
1408            let signal_header = rows.iter().find(|row| row["event"] == "request_error_original"
1409                && row["channel"] == "context"
1410                && row["occurrence"]["diagnostic_id"] == signal_id)
1411                .expect("signal original written before output close");
1412            let signal_context: serde_json::Value = serde_json::from_str(
1413                signal_header["payload"].as_str().unwrap(),
1414            ).unwrap();
1415            assert_eq!(signal_context["stage"], "shutdown_signal");
1416            assert_eq!(signal_context["context"]["application"]["value"], "test-app");
1417            assert_eq!(field(&rows, &signal_header["occurrence"], "description").unwrap(),
1418                "runtime.signal_registration_failed: failed to register the application shutdown signal");
1419            let drain_header = rows.iter().find(|row| row["event"] == "request_error_original"
1420                && row["channel"] == "context"
1421                && row["occurrence"]["diagnostic_id"] == drain_id)
1422                .expect("drain original written before output close");
1423            let drain_context: serde_json::Value = serde_json::from_str(
1424                drain_header["payload"].as_str().unwrap(),
1425            ).unwrap();
1426            assert_eq!(drain_context["stage"], "request_drain");
1427            assert_eq!(drain_context["context"]["application"]["value"], "test-app");
1428            assert_eq!(field(&rows, &drain_header["occurrence"], "description").unwrap(),
1429                "runtime.lifecycle_timeout.request_drain: managed lifecycle stage 'request_drain' exceeded its wall-clock deadline");
1430            let startup_header = rows.iter().find(|row| row["event"] == "request_error_original"
1431                && row["channel"] == "context"
1432                && row["occurrence"]["diagnostic_id"] == primary_id)
1433                .expect("startup original written before cleanup");
1434            let startup_context: serde_json::Value = serde_json::from_str(
1435                startup_header["payload"].as_str().unwrap(),
1436            ).unwrap();
1437            assert_eq!(startup_context["stage"], "component_start");
1438            assert_eq!(startup_context["context"]["application"]["value"], "test-app");
1439            assert_eq!(startup_context["context"]["request"]["state"], "not_established");
1440            let startup_occurrence = &startup_header["occurrence"];
1441            assert_eq!(field(&rows, startup_occurrence, "description").unwrap(), "test.failure: start failed");
1442            assert!(field(&rows, startup_occurrence, "debug").unwrap().contains("start failed"));
1443            assert!(rows.iter().any(|row| row["occurrence"] == *startup_occurrence
1444                && row["channel"] == "terminal" && row["state"] == "exposed_chain_complete"));
1445            let mut missing_startup = rows.clone();
1446            missing_startup.retain(|row| !(row["occurrence"] == *startup_occurrence
1447                && row["channel"] == "description" && row["state"] == "field_end"));
1448            assert!(field(&missing_startup, startup_occurrence, "description").is_none());
1449            let header = rows.iter().find(|row| row["event"] == "request_error_original"
1450                && row["channel"] == "context"
1451                && row["occurrence"]["primary_diagnostic_id"] == primary_id)
1452                .expect("cleanup original linked to startup primary");
1453            let context: serde_json::Value = serde_json::from_str(header["payload"].as_str().unwrap()).unwrap();
1454            assert_eq!(context["context"]["application"]["state"], "present");
1455            assert_eq!(context["context"]["application"]["value"], "test-app");
1456            assert_eq!(context["stage"], "component_cleanup");
1457            let occurrence = &header["occurrence"];
1458            assert_eq!(field(&rows, occurrence, "description").unwrap(), "test.failure: shutdown failed");
1459            let debug = field(&rows, occurrence, "debug").unwrap();
1460            assert!(debug.contains("test.failure"));
1461            assert!(debug.contains("shutdown failed"));
1462            let mut missing_debug = rows.clone();
1463            missing_debug.retain(|row| !(row["occurrence"] == *occurrence
1464                && row["channel"] == "debug" && row["state"] == "field_end"));
1465            assert!(field(&missing_debug, occurrence, "debug").is_none());
1466            assert!(rows.iter().any(|row| row["occurrence"] == *occurrence
1467                && row["channel"] == "terminal" && row["state"] == "exposed_chain_complete"));
1468            return;
1469        }
1470        let path = std::env::temp_dir().join(format!(
1471            "saddle-component-cleanup-{}", std::process::id(),
1472        ));
1473        std::fs::create_dir(&path).unwrap();
1474        let mut child = std::process::Command::new(std::env::current_exe().unwrap())
1475            .args(["--exact", "application::tests::component_cleanup_original_precedes_output_close"])
1476            .env(CHILD, &path).spawn().unwrap();
1477        let started = Instant::now();
1478        let status = loop {
1479            if let Some(status) = child.try_wait().unwrap() { break status; }
1480            if started.elapsed() > Duration::from_secs(10) {
1481                child.kill().unwrap();
1482                child.wait().unwrap();
1483                panic!("component cleanup child exceeded watchdog");
1484            }
1485            std::thread::sleep(Duration::from_millis(10));
1486        };
1487        assert!(status.success());
1488        std::fs::remove_file(path.join("saddle.emergency.log")).unwrap();
1489        std::fs::remove_dir(path).unwrap();
1490    }
1491
1492    #[test]
1493    fn duplicate_component_names_are_rejected() {
1494        let events = Arc::new(Mutex::new(Vec::new()));
1495        let mut application = Application::new();
1496        application.register(component("db", &events)).unwrap();
1497
1498        let error = application.register(component("db", &events)).unwrap_err();
1499        assert_eq!(error.code(), "runtime.duplicate_component");
1500    }
1501
1502    #[test]
1503    fn application_shutdown_waits_for_an_admitted_request() {
1504        test_runtime().block_on(async {
1505            let application = Application::new();
1506            let requests = application.request_lifecycle();
1507            let (release, released) = tokio::sync::oneshot::channel();
1508
1509            let shutdown = async move {
1510                shutdown_when_ready(requests.clone()).await?;
1511                let request = requests
1512                    .try_accept()
1513                    .expect("application is ready before waiting for shutdown");
1514                tokio::spawn(async move {
1515                    released.await.unwrap();
1516                    drop(request);
1517                });
1518                Ok(())
1519            };
1520            let running = tokio::spawn(application.run_until_shutdown(shutdown));
1521
1522            tokio::task::yield_now().await;
1523            assert!(!running.is_finished());
1524            release.send(()).unwrap();
1525            running.await.unwrap().unwrap();
1526        });
1527    }
1528
1529    #[test]
1530    fn signal_failure_before_start_prevents_component_startup() {
1531        let events = Arc::new(Mutex::new(Vec::new()));
1532        let mut application = Application::new();
1533        application.register(component("service", &events)).unwrap();
1534
1535        let error = test_runtime()
1536            .block_on(application.run_until_shutdown(async { Err(signal_error()) }))
1537            .unwrap_err();
1538
1539        assert_eq!(error.code(), "runtime.signal_registration_failed");
1540        assert!(error.source_unavailable());
1541        assert!(events.lock().unwrap().is_empty());
1542    }
1543
1544    #[test]
1545    fn signal_primary_keeps_independent_cleanup_failure_fact() {
1546        let signal = test_error("signal failed").with_unconfirmed_source();
1547        let cleanup = test_error("cleanup failed").with_unconfirmed_cleanup_source();
1548        let result = super::combine_lifecycle_results(Err(signal), Err(cleanup)).unwrap_err();
1549        assert_eq!(result.message(), "signal failed");
1550        assert!(result.source_unavailable());
1551        assert!(result.cleanup_source_unavailable());
1552    }
1553
1554    #[test]
1555    fn shutdown_during_startup_stops_starting_and_rolls_back() {
1556        test_runtime().block_on(async {
1557            let events = Arc::new(Mutex::new(Vec::new()));
1558            let (started_tx, started_rx) = tokio::sync::oneshot::channel();
1559            let (release_tx, release_rx) = tokio::sync::oneshot::channel();
1560            let (shutdown_tx, shutdown_rx) = tokio::sync::oneshot::channel();
1561            let mut application = Application::new();
1562            application
1563                .register(BlockingStartComponent {
1564                    events: Arc::clone(&events),
1565                    started: Mutex::new(Some(started_tx)),
1566                    release: Mutex::new(Some(release_rx)),
1567                })
1568                .unwrap();
1569            application.register(component("never", &events)).unwrap();
1570
1571            let running = tokio::spawn(application.run_until_shutdown(async move {
1572                shutdown_rx.await.unwrap();
1573                Ok(())
1574            }));
1575            started_rx.await.unwrap();
1576            shutdown_tx.send(()).unwrap();
1577            tokio::task::yield_now().await;
1578            release_tx.send(()).unwrap();
1579
1580            running.await.unwrap().unwrap();
1581            assert_eq!(
1582                *events.lock().unwrap(),
1583                ["start:blocking", "shutdown:blocking"]
1584            );
1585        });
1586    }
1587
1588    #[test]
1589    fn managed_runtime_provides_an_async_io_driver() {
1590        build_runtime()
1591            .unwrap()
1592            .block_on(async {
1593                tokio::net::TcpListener::bind((std::net::Ipv4Addr::LOCALHOST, 0)).await
1594            })
1595            .expect("service listeners require the managed async I/O driver");
1596    }
1597
1598    #[test]
1599    fn blocked_component_start_times_out_and_rolls_back_started_components() {
1600        test_runtime().block_on(async {
1601            let events = Arc::new(Mutex::new(Vec::new()));
1602            let (started_tx, started_rx) = tokio::sync::oneshot::channel();
1603            let (_release_tx, release_rx) = tokio::sync::oneshot::channel();
1604            let mut application = Application::new();
1605            application.set_lifecycle_timeouts(LifecycleTimeouts::from_millis(10, 100).unwrap());
1606            application.register(component("first", &events)).unwrap();
1607            application
1608                .register(BlockingStartComponent {
1609                    events: Arc::clone(&events),
1610                    started: Mutex::new(Some(started_tx)),
1611                    release: Mutex::new(Some(release_rx)),
1612                })
1613                .unwrap();
1614
1615            let running = tokio::spawn(application.run_until_shutdown(std::future::pending()));
1616            started_rx.await.unwrap();
1617            let error = running.await.unwrap().unwrap_err();
1618            assert_eq!(error.code(), "runtime.lifecycle_timeout.component_start");
1619            assert_eq!(
1620                *events.lock().unwrap(),
1621                [
1622                    "start:first",
1623                    "start:blocking",
1624                    "shutdown:blocking",
1625                    "shutdown:first"
1626                ]
1627            );
1628        });
1629    }
1630
1631    #[test]
1632    fn blocked_component_shutdown_uses_one_total_deadline_and_is_not_clean() {
1633        test_runtime().block_on(async {
1634            let events = Arc::new(Mutex::new(Vec::new()));
1635            let mut application = Application::new();
1636            application.set_lifecycle_timeouts(LifecycleTimeouts::from_millis(100, 10).unwrap());
1637            application
1638                .register(BlockingShutdownComponent {
1639                    events: Arc::clone(&events),
1640                })
1641                .unwrap();
1642            let shutdown = shutdown_when_ready(application.request_lifecycle());
1643            let error = application.run_until_shutdown(shutdown).await.unwrap_err();
1644            assert_eq!(error.code(), "runtime.lifecycle_timeout.component_shutdown");
1645            assert_eq!(
1646                *events.lock().unwrap(),
1647                ["start:blocking-shutdown", "shutdown:blocking-shutdown"]
1648            );
1649        });
1650    }
1651}