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 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 #[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 #[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 #[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 pub fn request_lifecycle(&self) -> RequestLifecycle {
186 self.requests.clone()
187 }
188
189 pub fn health(&self) -> crate::ApplicationHealth {
192 self.requests.health()
193 }
194
195 #[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 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 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 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 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 pub fn run(self) -> Result<()> {
278 Self::run_with(|| async move { Ok(self) })
279 }
280
281 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 #[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 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 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 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 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}