|
| 1 | +//! Demonstrates a supervised service role with initialization, running, |
| 2 | +//! cooperative stop, and observation output. |
| 3 | +
|
| 4 | +mod observation; |
| 5 | +mod service_task; |
| 6 | + |
| 7 | +// Import command result variants returned by the runtime handle. |
| 8 | +use rust_supervisor::control::command::CommandResult; |
| 9 | +// Import supervisor error values. |
| 10 | +use rust_supervisor::error::types::SupervisorError; |
| 11 | +// Import the supervisor runtime entry point. |
| 12 | +use rust_supervisor::runtime::supervisor::Supervisor; |
| 13 | +// Import supervisor specification values. |
| 14 | +use rust_supervisor::spec::supervisor::SupervisorSpec; |
| 15 | +// Import shutdown timing policy. |
| 16 | +use rust_supervisor::shutdown::stage::ShutdownPolicy; |
| 17 | +// Import duration values for the example timing budget. |
| 18 | +use std::time::Duration; |
| 19 | +// Import asynchronous channel helpers. |
| 20 | +use tokio::sync::mpsc; |
| 21 | + |
| 22 | +// Define the shared example result type. |
| 23 | +type ExampleResult = Result<(), rust_supervisor::error::types::SupervisorError>; |
| 24 | + |
| 25 | +// Use the Tokio runtime for the asynchronous example. |
| 26 | +#[tokio::main] |
| 27 | +// Return typed supervisor errors from the example. |
| 28 | +/// Runs the service role example. |
| 29 | +async fn main() -> ExampleResult { |
| 30 | + // Build a channel that receives service lifecycle facts. |
| 31 | + let (service_event_sender, mut service_events) = mpsc::unbounded_channel(); |
| 32 | + // Build one child declared as a service role. |
| 33 | + let service_child = service_task::service_child(service_event_sender); |
| 34 | + // Build a root supervisor with the service child. |
| 35 | + let mut spec = SupervisorSpec::root(vec![service_child]); |
| 36 | + // Keep enough event buffer for the full shutdown observation sequence. |
| 37 | + spec.event_channel_capacity = 32; |
| 38 | + // Use short shutdown windows so the example finishes quickly. |
| 39 | + let shutdown_policy = |
| 40 | + ShutdownPolicy::new(Duration::from_millis(250), Duration::from_millis(50), true); |
| 41 | + // Start the runtime with the service child. |
| 42 | + let handle = Supervisor::start_with_policy(spec, shutdown_policy).await?; |
| 43 | + // Subscribe to lifecycle event text before commands are sent. |
| 44 | + let mut runtime_events = handle.subscribe_events(); |
| 45 | + // Wait until the service reports initialization. |
| 46 | + observation::wait_for_initialization(&mut service_events).await; |
| 47 | + // Print the state after the service has initialized. |
| 48 | + observation::print_current_state("after-initialization", handle.current_state().await?); |
| 49 | + // Print the long-running operation hint. |
| 50 | + println!("service example: running until Ctrl+C"); |
| 51 | + // Build periodic observation ticks for the long-running service. |
| 52 | + let mut observation_interval = tokio::time::interval(Duration::from_secs(1)); |
| 53 | + // Keep the service running until the operator requests shutdown. |
| 54 | + loop { |
| 55 | + // Wait for either an operator signal or the next observation tick. |
| 56 | + tokio::select! { |
| 57 | + // Stop the example only when the operator sends Ctrl+C. |
| 58 | + signal = tokio::signal::ctrl_c() => { |
| 59 | + // Convert signal errors into the example error type. |
| 60 | + signal.map_err(|error| { |
| 61 | + SupervisorError::fatal_config(format!( |
| 62 | + "failed to receive Ctrl+C signal: {error}" |
| 63 | + )) |
| 64 | + })?; |
| 65 | + // Print the operator stop signal. |
| 66 | + println!("operator signal=ctrl_c"); |
| 67 | + // Leave the persistent running loop. |
| 68 | + break; |
| 69 | + } |
| 70 | + // Print periodic service observation. |
| 71 | + _ = observation_interval.tick() => { |
| 72 | + // Print service facts emitted while the service is running. |
| 73 | + observation::drain_service_events("while-running", &mut service_events); |
| 74 | + // Print the current runtime state. |
| 75 | + observation::print_current_state("while-running", handle.current_state().await?); |
| 76 | + // Print relevant runtime events that arrived during the tick. |
| 77 | + observation::drain_runtime_events(&mut runtime_events); |
| 78 | + } |
| 79 | + } |
| 80 | + } |
| 81 | + // Request cooperative shutdown for the whole supervisor tree. |
| 82 | + let shutdown = handle |
| 83 | + .shutdown_tree("operator", "service role example stopped by operator") |
| 84 | + .await?; |
| 85 | + // Print service facts emitted during cooperative stop. |
| 86 | + observation::drain_service_events("during-stop", &mut service_events); |
| 87 | + // Print runtime events that show command and shutdown observation. |
| 88 | + observation::drain_runtime_events(&mut runtime_events); |
| 89 | + // Print the shutdown result and per-child outcome. |
| 90 | + if let CommandResult::Shutdown { result } = shutdown { |
| 91 | + // Print the completed shutdown phase. |
| 92 | + observation::print_shutdown_result(result); |
| 93 | + } |
| 94 | + // Print the state after shutdown has completed. |
| 95 | + observation::print_current_state("after-shutdown", handle.current_state().await?); |
| 96 | + // Finish the example successfully. |
| 97 | + Ok(()) |
| 98 | + // End the service role example. |
| 99 | +} |
0 commit comments