use std::path::Path;
use std::time::Duration;
use serde::{Deserialize, Serialize};
use super::backpressure::BackpressureStrategy;
use crate::api::{Error, Result};
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum SchedulerKind {
#[default]
Default,
SingleThread,
WorkerPool,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct RuntimeConfig {
#[serde(default = "default_workers")]
pub workers: usize,
#[serde(default = "default_batch_size")]
pub batch_size: usize,
#[serde(default)]
pub backpressure: BackpressureStrategy,
#[serde(default)]
pub scheduler: SchedulerKind,
#[serde(default = "default_true")]
pub metrics: bool,
#[serde(default)]
pub backpressure_retries: Option<u32>,
#[serde(default = "default_yield_ms")]
pub backpressure_yield_ms: u64,
#[serde(default)]
pub idle_sleep_ms: Option<u64>,
#[serde(default)]
pub overflow_sink: Option<u32>,
}
fn default_workers() -> usize {
1
}
fn default_batch_size() -> usize {
512
}
fn default_true() -> bool {
true
}
fn default_yield_ms() -> u64 {
0
}
impl Default for RuntimeConfig {
fn default() -> Self {
Self {
workers: default_workers(),
batch_size: default_batch_size(),
backpressure: BackpressureStrategy::default(),
scheduler: SchedulerKind::default(),
metrics: true,
backpressure_retries: None,
backpressure_yield_ms: 0,
idle_sleep_ms: None,
overflow_sink: None,
}
}
}
#[derive(Debug, Clone, Default, Deserialize)]
struct RuntimeFile {
#[serde(default)]
runtime: RuntimeConfig,
}
impl RuntimeConfig {
pub fn from_toml_str(toml_src: &str) -> Result<Self> {
let file: RuntimeFile =
toml::from_str(toml_src).map_err(|e| Error::config(format!("runtime TOML: {e}")))?;
file.runtime.validate()?;
Ok(file.runtime)
}
pub fn from_toml_path(path: impl AsRef<Path>) -> Result<Self> {
let path = path.as_ref();
let text = std::fs::read_to_string(path).map_err(|e| {
Error::config(format!(
"failed to read runtime config '{}': {e}",
path.display()
))
})?;
Self::from_toml_str(&text)
}
pub fn validate(&self) -> Result<()> {
if self.workers == 0 {
return Err(Error::config("runtime.workers must be ≥ 1"));
}
if self.batch_size == 0 {
return Err(Error::config("runtime.batch_size must be ≥ 1"));
}
if matches!(self.scheduler, SchedulerKind::WorkerPool) && self.workers < 1 {
return Err(Error::config("worker pool requires workers ≥ 1"));
}
Ok(())
}
pub fn idle_sleep(&self) -> Option<Duration> {
self.idle_sleep_ms.map(Duration::from_millis)
}
pub fn backpressure_yield(&self) -> Duration {
Duration::from_millis(self.backpressure_yield_ms)
}
pub fn with_workers(mut self, n: usize) -> Self {
self.workers = n.max(1);
self
}
pub fn with_batch_size(mut self, n: usize) -> Self {
self.batch_size = n.max(1);
self
}
pub fn with_backpressure(mut self, bp: BackpressureStrategy) -> Self {
self.backpressure = bp;
self
}
pub fn with_scheduler(mut self, kind: SchedulerKind) -> Self {
self.scheduler = kind;
self
}
pub fn with_metrics(mut self, on: bool) -> Self {
self.metrics = on;
self
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn parses_toml_example() {
let cfg = RuntimeConfig::from_toml_str(
r#"
[runtime]
workers = 4
batch_size = 512
backpressure = "block"
scheduler = "default"
metrics = true
"#,
)
.unwrap();
assert_eq!(cfg.workers, 4);
assert_eq!(cfg.batch_size, 512);
assert_eq!(cfg.backpressure, BackpressureStrategy::Block);
assert!(cfg.metrics);
}
#[test]
fn rejects_zero_batch() {
let cfg = RuntimeConfig {
batch_size: 0,
..RuntimeConfig::default()
};
assert!(cfg.validate().is_err());
}
}