use std::sync::Arc;
use std::sync::atomic::AtomicBool;
use std::time::Instant;
use crate::api::{MetricsCollector, NullCollector, Pipeline, Result, StepOutcome};
use super::config::{RuntimeConfig, SchedulerKind};
use super::scheduler::{
SingleThreadScheduler, WorkerPoolScheduler, run_single_thread, shutdown_flag,
};
use super::telemetry::{RuntimeEvent, RuntimeMetricKey, emit_lifecycle};
#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord)]
pub enum RuntimePhase {
Built,
Validated,
Initialized,
Running,
Draining,
Shutdown,
CleanedUp,
}
impl RuntimePhase {
pub fn as_u8(self) -> u8 {
match self {
Self::Built => 0,
Self::Validated => 1,
Self::Initialized => 2,
Self::Running => 3,
Self::Draining => 4,
Self::Shutdown => 5,
Self::CleanedUp => 6,
}
}
}
#[derive(Debug, Clone, Default)]
pub struct RuntimeStats {
pub steps: u64,
pub elapsed: std::time::Duration,
pub phase: Option<RuntimePhase>,
}
pub struct Runtime<P: Pipeline> {
pipeline: P,
config: RuntimeConfig,
phase: RuntimePhase,
metrics: Arc<dyn MetricsCollector>,
shutdown: Arc<AtomicBool>,
stats: RuntimeStats,
}
impl<P: Pipeline> Runtime<P> {
pub fn build(pipeline: P, config: RuntimeConfig) -> Result<Self> {
config.validate()?;
let metrics: Arc<dyn MetricsCollector> = Arc::new(NullCollector);
let rt = Self {
pipeline,
config,
phase: RuntimePhase::Built,
metrics,
shutdown: shutdown_flag(),
stats: RuntimeStats::default(),
};
emit_lifecycle(rt.metrics.as_ref(), rt.config.metrics, &RuntimeEvent::Built);
Ok(rt)
}
pub fn with_metrics(mut self, metrics: impl MetricsCollector + 'static) -> Self {
self.metrics = Arc::new(metrics);
self
}
pub fn config(&self) -> &RuntimeConfig {
&self.config
}
pub fn phase(&self) -> RuntimePhase {
self.phase
}
pub fn stats(&self) -> &RuntimeStats {
&self.stats
}
pub fn shutdown_handle(&self) -> Arc<AtomicBool> {
Arc::clone(&self.shutdown)
}
pub fn request_shutdown(&self) {
self.shutdown
.store(true, std::sync::atomic::Ordering::SeqCst);
}
pub fn pipeline_mut(&mut self) -> &mut P {
&mut self.pipeline
}
pub fn pipeline(&self) -> &P {
&self.pipeline
}
fn set_phase(&mut self, phase: RuntimePhase) {
self.phase = phase;
if self.config.metrics {
self.metrics
.record_gauge(&RuntimeMetricKey::Phase, phase.as_u8() as f64);
}
}
pub fn validate(&mut self) -> Result<()> {
self.config.validate()?;
self.set_phase(RuntimePhase::Validated);
emit_lifecycle(
self.metrics.as_ref(),
self.config.metrics,
&RuntimeEvent::Validated,
);
Ok(())
}
pub fn initialize(&mut self) -> Result<()> {
if self.phase < RuntimePhase::Validated {
self.validate()?;
}
self.pipeline.init()?;
self.set_phase(RuntimePhase::Initialized);
emit_lifecycle(
self.metrics.as_ref(),
self.config.metrics,
&RuntimeEvent::Initialized,
);
Ok(())
}
pub fn run(&mut self) -> Result<RuntimeStats> {
if self.phase < RuntimePhase::Initialized {
self.initialize()?;
}
self.set_phase(RuntimePhase::Running);
emit_lifecycle(
self.metrics.as_ref(),
self.config.metrics,
&RuntimeEvent::Started,
);
if self.config.metrics {
self.metrics
.record_gauge(&RuntimeMetricKey::Workers, self.config.workers as f64);
}
let start = Instant::now();
match self.config.scheduler {
SchedulerKind::Default | SchedulerKind::SingleThread | SchedulerKind::WorkerPool => {
let mut sched = SingleThreadScheduler::with_shutdown(
self.config.clone(),
Arc::clone(&self.shutdown),
);
loop {
if self.shutdown.load(std::sync::atomic::Ordering::Relaxed) {
break;
}
let step_start = Instant::now();
let outcome = self.pipeline.step_outcome()?;
self.stats.steps += 1;
if self.config.metrics {
self.metrics.record_counter(&RuntimeMetricKey::Steps, 1);
self.metrics.record_histogram(
&RuntimeMetricKey::StepDurationNs,
step_start.elapsed().as_nanos() as f64,
);
}
match outcome {
StepOutcome::Exhausted => {
self.set_phase(RuntimePhase::Draining);
emit_lifecycle(
self.metrics.as_ref(),
self.config.metrics,
&RuntimeEvent::Draining,
);
break;
}
StepOutcome::Idle => {
if let Some(d) = self.config.idle_sleep() {
std::thread::sleep(d);
}
}
StepOutcome::BackPressured => {
if self.config.metrics {
self.metrics
.record_counter(&RuntimeMetricKey::BackpressureEvents, 1);
}
let y = self.config.backpressure_yield();
if y.is_zero() {
std::thread::yield_now();
} else {
std::thread::sleep(y);
}
}
StepOutcome::Progress => {}
}
let _ = &mut sched;
}
}
}
self.stats.elapsed = start.elapsed();
self.shutdown_pipeline()?;
self.cleanup()?;
self.stats.phase = Some(self.phase);
Ok(self.stats.clone())
}
pub fn shutdown_pipeline(&mut self) -> Result<()> {
self.set_phase(RuntimePhase::Shutdown);
emit_lifecycle(
self.metrics.as_ref(),
self.config.metrics,
&RuntimeEvent::Shutdown,
);
self.pipeline.shutdown()?;
Ok(())
}
pub fn cleanup(&mut self) -> Result<()> {
self.set_phase(RuntimePhase::CleanedUp);
emit_lifecycle(
self.metrics.as_ref(),
self.config.metrics,
&RuntimeEvent::Cleanup,
);
Ok(())
}
}
pub fn run_with_worker_factory<P, F>(config: RuntimeConfig, factory: F) -> Result<()>
where
P: Pipeline + 'static,
F: Fn(usize) -> Result<P> + Send + Sync + 'static,
{
let pool = WorkerPoolScheduler::new(config)?;
pool.run_factory(factory)
}
pub fn run_pipeline<P: Pipeline>(pipeline: &mut P, config: &RuntimeConfig) -> Result<()> {
pipeline.init()?;
let flag = shutdown_flag();
run_single_thread(pipeline, config, &flag)?;
pipeline.shutdown()?;
Ok(())
}