From 6a04d7f7930a1d15f6ff9050492bab0f02b87f59 Mon Sep 17 00:00:00 2001 From: Sunil Thakare Date: Wed, 22 Jul 2026 14:09:05 +0530 Subject: [PATCH] WIP: recover runtime implementation after accidental git clean --- .../systemd/nx9-auth.service | 0 deploy.sh => scripts/deploy.sh | 0 src/db/models/refresh_token.rs | 13 ++ src/runtime/cancellation.rs | 42 ++++ src/runtime/hooks.rs | 82 +++++++ src/runtime/lifecycle.rs | 16 ++ src/runtime/mod.rs | 27 +++ src/runtime/signals.rs | 67 ++++++ src/runtime/state.rs | 202 ++++++++++++++++++ src/runtime/workers.rs | 114 ++++++++++ 10 files changed, 563 insertions(+) rename nx9-auth.service => deploy/systemd/nx9-auth.service (100%) rename deploy.sh => scripts/deploy.sh (100%) create mode 100644 src/db/models/refresh_token.rs create mode 100644 src/runtime/cancellation.rs create mode 100644 src/runtime/hooks.rs create mode 100644 src/runtime/lifecycle.rs create mode 100644 src/runtime/mod.rs create mode 100644 src/runtime/signals.rs create mode 100644 src/runtime/state.rs create mode 100644 src/runtime/workers.rs diff --git a/nx9-auth.service b/deploy/systemd/nx9-auth.service similarity index 100% rename from nx9-auth.service rename to deploy/systemd/nx9-auth.service diff --git a/deploy.sh b/scripts/deploy.sh similarity index 100% rename from deploy.sh rename to scripts/deploy.sh diff --git a/src/db/models/refresh_token.rs b/src/db/models/refresh_token.rs new file mode 100644 index 0000000..67dd3d8 --- /dev/null +++ b/src/db/models/refresh_token.rs @@ -0,0 +1,13 @@ +use serde::{Deserialize, Serialize}; +use sqlx::FromRow; + +#[derive(Debug, Clone, Serialize, Deserialize, FromRow)] +pub struct RefreshToken { + pub id: String, + pub user_id: String, + #[serde(skip_serializing)] + pub token_hash: String, + pub expires_at: String, + pub created_at: String, + pub revoked: bool, +} diff --git a/src/runtime/cancellation.rs b/src/runtime/cancellation.rs new file mode 100644 index 0000000..5f0b0c2 --- /dev/null +++ b/src/runtime/cancellation.rs @@ -0,0 +1,42 @@ +//! Runtime-wide shutdown coordination using `CancellationToken`. + +use tokio_util::sync::CancellationToken; + +#[derive(Clone)] +pub struct ShutdownCoordinator { + root: CancellationToken, +} + +impl ShutdownCoordinator { + pub fn new() -> Self { + Self { + root: CancellationToken::new(), + } + } + + pub fn token(&self) -> &CancellationToken { + &self.root + } + + pub fn child_token(&self) -> CancellationToken { + self.root.child_token() + } + + pub fn cancel(&self) { + self.root.cancel(); + } + + pub fn is_cancelled(&self) -> bool { + self.root.is_cancelled() + } + + pub async fn cancelled(&self) { + self.root.cancelled().await; + } +} + +impl Default for ShutdownCoordinator { + fn default() -> Self { + Self::new() + } +} diff --git a/src/runtime/hooks.rs b/src/runtime/hooks.rs new file mode 100644 index 0000000..8ccc6de --- /dev/null +++ b/src/runtime/hooks.rs @@ -0,0 +1,82 @@ +//! Prioritized, idempotent shutdown hooks. + +use anyhow::Result; + +#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord)] +pub enum ShutdownPriority { + First = 0, + Normal = 1, + Last = 2, +} + +#[async_trait::async_trait] +pub trait ShutdownHook: Send + Sync { + fn name(&self) -> &'static str; + + fn priority(&self) -> ShutdownPriority { + ShutdownPriority::Normal + } + + async fn shutdown(&self) -> Result<()>; +} + +pub struct HookRegistry { + hooks: Vec>, +} + +impl HookRegistry { + pub fn new() -> Self { + Self { hooks: Vec::new() } + } + + pub fn register(&mut self, hook: Box) { + tracing::debug!(hook = hook.name(), priority = ?hook.priority(), "shutdown hook registered"); + self.hooks.push(hook); + } + + pub async fn execute_all(&self) { + if self.hooks.is_empty() { + return; + } + + let mut indices: Vec = (0..self.hooks.len()).collect(); + indices.sort_by_key(|&i| self.hooks[i].priority()); + + for i in indices { + let hook = &self.hooks[i]; + let start = std::time::Instant::now(); + tracing::info!(hook = hook.name(), priority = ?hook.priority(), "executing shutdown hook"); + match hook.shutdown().await { + Ok(()) => { + tracing::info!( + hook = hook.name(), + duration_ms = start.elapsed().as_millis(), + "shutdown hook completed successfully" + ); + } + Err(e) => { + tracing::error!( + hook = hook.name(), + duration_ms = start.elapsed().as_millis(), + error = %e, + "shutdown hook failed" + ); + } + } + } + } + + pub fn len(&self) -> usize { + self.hooks.len() + } + + pub fn is_empty(&self) -> bool { + self.hooks.is_empty() + } +} + +impl Default for HookRegistry { + fn default() -> Self { + Self::new() + } +} diff --git a/src/runtime/lifecycle.rs b/src/runtime/lifecycle.rs new file mode 100644 index 0000000..0f5d6a4 --- /dev/null +++ b/src/runtime/lifecycle.rs @@ -0,0 +1,16 @@ +//! Common lifecycle contract for runtime components. + +use anyhow::Result; + +/// Lifecycle contract for the application runtime. +#[async_trait::async_trait] +pub trait Lifecycle { + /// Initialize subsystems. + async fn initialize(&mut self) -> Result<()>; + + /// Start serving. + async fn start(&mut self) -> Result<()>; + + /// Perform graceful shutdown. + async fn shutdown(&mut self) -> Result<()>; +} diff --git a/src/runtime/mod.rs b/src/runtime/mod.rs new file mode 100644 index 0000000..00282d3 --- /dev/null +++ b/src/runtime/mod.rs @@ -0,0 +1,27 @@ +//! Unified Enterprise Runtime Lifecycle for `nx9-auth`. +//! +//! Provides application assembly (`ApplicationBuilder`), lifecycle management +//! (`Application`, `Lifecycle`), atomic state machine (`RuntimeState`), +//! signal coordination (`SignalManager`), prioritized shutdown hooks +//! (`ShutdownHook`), worker management (`WorkerManager`), and operational +//! metrics (`RuntimeMetrics`). + +pub mod application; +pub mod builder; +pub mod cancellation; +pub mod hooks; +pub mod lifecycle; +pub mod metrics; +pub mod signals; +pub mod state; +pub mod workers; + +pub use application::Application; +pub use builder::ApplicationBuilder; +pub use cancellation::ShutdownCoordinator; +pub use hooks::{HookRegistry, ShutdownHook, ShutdownPriority}; +pub use lifecycle::Lifecycle; +pub use metrics::RuntimeMetrics; +pub use signals::SignalManager; +pub use state::{AtomicRuntimeState, RuntimeState}; +pub use workers::{TaskGroup, WorkerManager}; diff --git a/src/runtime/signals.rs b/src/runtime/signals.rs new file mode 100644 index 0000000..56ea635 --- /dev/null +++ b/src/runtime/signals.rs @@ -0,0 +1,67 @@ +//! Unix signal handling for graceful and forced shutdown. + +use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering}; +use std::sync::Arc; + +#[derive(Clone)] +pub struct SignalManager { + signal_count: Arc, + force_shutdown: Arc, +} + +impl SignalManager { + pub fn new() -> Self { + Self { + signal_count: Arc::new(AtomicUsize::new(0)), + force_shutdown: Arc::new(AtomicBool::new(false)), + } + } + + pub fn is_force_shutdown(&self) -> bool { + self.force_shutdown.load(Ordering::Acquire) + } + + pub fn signal_count(&self) -> usize { + self.signal_count.load(Ordering::Acquire) + } + + pub fn record_signal(&self) -> usize { + let count = self.signal_count.fetch_add(1, Ordering::AcqRel) + 1; + if count >= 2 { + self.force_shutdown.store(true, Ordering::Release); + } + count + } +} + +impl Default for SignalManager { + fn default() -> Self { + Self::new() + } +} + +pub async fn wait_for_shutdown_signal() -> &'static str { + let ctrl_c = async { + tokio::signal::ctrl_c() + .await + .expect("failed to install SIGINT handler"); + "SIGINT" + }; + + #[cfg(unix)] + let sigterm = async { + tokio::signal::unix::signal(tokio::signal::unix::SignalKind::terminate()) + .expect("failed to install SIGTERM handler") + .recv() + .await; + "SIGTERM" + }; + + #[cfg(not(unix))] + let sigterm = std::future::pending::<&str>(); + + tokio::select! { + name = ctrl_c => name, + name = sigterm => name, + } +} diff --git a/src/runtime/state.rs b/src/runtime/state.rs new file mode 100644 index 0000000..40acea0 --- /dev/null +++ b/src/runtime/state.rs @@ -0,0 +1,202 @@ +//! Lock-free atomic runtime state machine. +//! +//! Tracks the application through granular lifecycle phases using `AtomicU8` +//! with `compare_exchange` transitions. This avoids mutex contention and +//! provides deterministic, race-free state management across async tasks. + +use std::fmt; +use std::sync::atomic::{AtomicU8, Ordering}; + +/// Granular runtime lifecycle states for enterprise production observability. +/// +/// Each transition is deterministic and logged. Only forward transitions +/// are permitted during normal operation; the state machine never moves +/// backward. +#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord)] +#[repr(u8)] +pub enum RuntimeState { + /// Runtime is being configured and dependencies are being assembled. + Initializing = 0, + /// Subsystems are starting (database, HTTP listener, workers). + Starting = 1, + /// Application is fully operational and serving requests. + Running = 2, + /// Shutdown signal received; HTTP listener stopped, draining active requests. + Draining = 3, + /// Active requests drained; cancelling and awaiting background workers. + StoppingWorkers = 4, + /// Workers stopped; executing registered shutdown hooks. + ExecutingHooks = 5, + /// Hooks executed; closing database connection pools and flushing logs. + ClosingResources = 6, + /// All resources released; process is ready to exit. + Stopped = 7, +} + +impl RuntimeState { + /// Convert a raw `u8` value back into a `RuntimeState`. + fn from_u8(val: u8) -> Option { + match val { + 0 => Some(Self::Initializing), + 1 => Some(Self::Starting), + 2 => Some(Self::Running), + 3 => Some(Self::Draining), + 4 => Some(Self::StoppingWorkers), + 5 => Some(Self::ExecutingHooks), + 6 => Some(Self::ClosingResources), + 7 => Some(Self::Stopped), + _ => None, + } + } + + /// Returns `true` if the runtime is in any shutdown phase. + pub fn is_shutting_down(self) -> bool { + (self as u8) >= (Self::Draining as u8) + } +} + +impl fmt::Display for RuntimeState { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + match self { + Self::Initializing => write!(f, "Initializing"), + Self::Starting => write!(f, "Starting"), + Self::Running => write!(f, "Running"), + Self::Draining => write!(f, "Draining"), + Self::StoppingWorkers => write!(f, "StoppingWorkers"), + Self::ExecutingHooks => write!(f, "ExecutingHooks"), + Self::ClosingResources => write!(f, "ClosingResources"), + Self::Stopped => write!(f, "Stopped"), + } + } +} + +/// Lock-free atomic runtime state container. +/// +/// Uses `AtomicU8` with `compare_exchange` to ensure deterministic, +/// race-free state transitions without mutex contention. +pub struct AtomicRuntimeState { + state: AtomicU8, +} + +impl AtomicRuntimeState { + /// Create a new state machine in the `Initializing` state. + pub fn new() -> Self { + Self { + state: AtomicU8::new(RuntimeState::Initializing as u8), + } + } + + /// Read the current state (acquire ordering for visibility). + pub fn load(&self) -> RuntimeState { + RuntimeState::from_u8(self.state.load(Ordering::Acquire)) + .unwrap_or(RuntimeState::Stopped) + } + + /// Attempt an atomic state transition from `expected` to `new`. + /// + /// Returns `Ok(new)` if the transition succeeded, or `Err(actual)` if the + /// current state did not match `expected`. + pub fn transition( + &self, + expected: RuntimeState, + new: RuntimeState, + ) -> Result { + match self.state.compare_exchange( + expected as u8, + new as u8, + Ordering::AcqRel, + Ordering::Acquire, + ) { + Ok(_) => { + tracing::info!(from = %expected, to = %new, "runtime state transition"); + Ok(new) + } + Err(actual) => { + let actual_state = + RuntimeState::from_u8(actual).unwrap_or(RuntimeState::Stopped); + tracing::debug!( + expected = %expected, + actual = %actual_state, + target = %new, + "state transition skipped (unexpected current state)" + ); + Err(actual_state) + } + } + } + + /// Unconditionally advance the state. Used during forced shutdown when + /// intermediate states may have been skipped. + pub fn force_set(&self, new: RuntimeState) { + let prev = self.state.swap(new as u8, Ordering::AcqRel); + let prev_state = RuntimeState::from_u8(prev).unwrap_or(RuntimeState::Stopped); + if prev_state != new { + tracing::info!(from = %prev_state, to = %new, "runtime state forced"); + } + } + + /// Attempt to transition from `Running` to `Draining`. + /// + /// Returns `true` if this call initiated shutdown (first caller wins). + /// Returns `false` if shutdown was already in progress or the runtime + /// was not yet `Running`. + pub fn initiate_shutdown(&self) -> bool { + self.transition(RuntimeState::Running, RuntimeState::Draining) + .is_ok() + } +} + +impl Default for AtomicRuntimeState { + fn default() -> Self { + Self::new() + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn test_state_transitions() { + let state = AtomicRuntimeState::new(); + assert_eq!(state.load(), RuntimeState::Initializing); + + assert!( + state + .transition(RuntimeState::Initializing, RuntimeState::Starting) + .is_ok() + ); + assert_eq!(state.load(), RuntimeState::Starting); + + assert!( + state + .transition(RuntimeState::Starting, RuntimeState::Running) + .is_ok() + ); + assert_eq!(state.load(), RuntimeState::Running); + + // Duplicate shutdown protection + assert!(state.initiate_shutdown()); + assert!(!state.initiate_shutdown()); // second call fails + assert_eq!(state.load(), RuntimeState::Draining); + } + + #[test] + fn test_force_set() { + let state = AtomicRuntimeState::new(); + state.force_set(RuntimeState::ClosingResources); + assert_eq!(state.load(), RuntimeState::ClosingResources); + } + + #[test] + fn test_is_shutting_down() { + assert!(!RuntimeState::Initializing.is_shutting_down()); + assert!(!RuntimeState::Starting.is_shutting_down()); + assert!(!RuntimeState::Running.is_shutting_down()); + assert!(RuntimeState::Draining.is_shutting_down()); + assert!(RuntimeState::StoppingWorkers.is_shutting_down()); + assert!(RuntimeState::ExecutingHooks.is_shutting_down()); + assert!(RuntimeState::ClosingResources.is_shutting_down()); + assert!(RuntimeState::Stopped.is_shutting_down()); + } +} diff --git a/src/runtime/workers.rs b/src/runtime/workers.rs new file mode 100644 index 0000000..5d1d6f1 --- /dev/null +++ b/src/runtime/workers.rs @@ -0,0 +1,114 @@ +//! Background worker management using `tokio::task::JoinSet`. + +use std::collections::HashMap; +use std::future::Future; +use std::time::Duration; + +use tokio::task::JoinSet; + +pub struct TaskGroup { + name: String, + tasks: JoinSet<()>, +} + +impl TaskGroup { + pub fn new(name: impl Into) -> Self { + Self { + name: name.into(), + tasks: JoinSet::new(), + } + } + + pub fn name(&self) -> &str { + &self.name + } + + pub fn spawn(&mut self, future: F) + where + F: Future + Send + 'static, + { + self.tasks.spawn(future); + } + + pub fn len(&self) -> usize { + self.tasks.len() + } + + pub fn is_empty(&self) -> bool { + self.tasks.is_empty() + } + + pub fn abort_all(&mut self) { + self.tasks.abort_all(); + } + + pub async fn shutdown(&mut self, timeout: Duration) { + if self.tasks.is_empty() { + return; + } + + let task_count = self.tasks.len(); + tracing::info!( + group = %self.name, + tasks = task_count, + timeout_secs = timeout.as_secs(), + "waiting for task group to complete" + ); + + let result = tokio::time::timeout(timeout, async { + while self.tasks.join_next().await.is_some() {} + }) + .await; + + if result.is_err() { + let remaining = self.tasks.len(); + tracing::warn!( + group = %self.name, + remaining, + "task group timed out, aborting remaining tasks" + ); + self.tasks.abort_all(); + while self.tasks.join_next().await.is_some() {} + } + } +} + +pub struct WorkerManager { + groups: HashMap, +} + +impl WorkerManager { + pub fn new() -> Self { + Self { + groups: HashMap::new(), + } + } + + pub fn group(&mut self, name: &str) -> &mut TaskGroup { + self.groups + .entry(name.to_string()) + .or_insert_with(|| TaskGroup::new(name)) + } + + pub fn active_tasks(&self) -> usize { + self.groups.values().map(|g| g.len()).sum() + } + + pub fn abort_all(&mut self) { + for group in self.groups.values_mut() { + group.abort_all(); + } + } + + pub async fn shutdown_all(&mut self, timeout: Duration) { + for group in self.groups.values_mut() { + group.shutdown(timeout).await; + } + } +} + +impl Default for WorkerManager { + fn default() -> Self { + Self::new() + } +}