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