use std::sync::Arc;
use std::sync::atomic::AtomicI64;
use crate::config::ClusterConfig;
use crate::connector::cache_backend::CacheBackend;
use crate::errors::OrionError;
use crate::storage::DbPool;
use crate::storage::repositories::cluster::{ClusterRepository, SqlClusterRepository};
pub mod epoch_watcher;
pub mod job_lease;
pub use epoch_watcher::start_cluster_tasks;
pub use job_lease::JobLeaseGate;
pub struct ClusterRuntime {
pub enabled: bool,
pub instance_id: String,
pub redis: Option<redis::aio::ConnectionManager>,
pub default_cache: Option<Arc<dyn CacheBackend>>,
pub repo: Arc<dyn ClusterRepository>,
pub last_seen_epoch: AtomicI64,
pub last_seen_breaker_epoch: AtomicI64,
}
impl ClusterRuntime {
pub async fn bump_config_epoch(&self) -> Result<(), crate::errors::OrionError> {
match self.repo.bump_epoch().await {
Ok(epoch) => {
self.last_seen_epoch
.fetch_max(epoch, std::sync::atomic::Ordering::AcqRel);
Ok(())
}
Err(e) if self.enabled => Err(e),
Err(e) => {
tracing::warn!(error = %e, "Failed to bump config epoch (cluster disabled — ignored)");
Ok(())
}
}
}
}
impl From<&ClusterRuntime> for crate::channel::registry::ClusterBackends {
fn from(runtime: &ClusterRuntime) -> Self {
Self {
default_cache: runtime.default_cache.clone(),
redis: runtime.redis.clone(),
}
}
}
pub async fn init_cluster_runtime(
config: &ClusterConfig,
pool: &DbPool,
) -> Result<Arc<ClusterRuntime>, OrionError> {
let instance_id = config.effective_instance_id();
let (redis, default_cache) = if config.enabled {
let client =
redis::Client::open(config.redis_url.as_str()).map_err(|e| OrionError::Config {
message: format!("cluster.redis_url is invalid: {e}"),
})?;
let conn = client
.get_connection_manager()
.await
.map_err(|e| OrionError::Internal {
context: "Failed to connect to cluster Redis (cluster.redis_url)".to_string(),
source: Some(Box::new(e)),
})?;
let cache: Arc<dyn CacheBackend> = Arc::new(
crate::connector::cache_backend::RedisCacheBackend::new(conn.clone()),
);
(Some(conn), Some(cache))
} else {
(None, None)
};
let repo: Arc<dyn ClusterRepository> = Arc::new(SqlClusterRepository::new(pool.clone()));
let (epoch, breaker_epoch) = if config.enabled {
let row = repo.get_epoch().await?;
(row.epoch, row.breaker_epoch)
} else {
(0, 0)
};
Ok(Arc::new(ClusterRuntime {
enabled: config.enabled,
instance_id,
redis,
default_cache,
repo,
last_seen_epoch: AtomicI64::new(epoch),
last_seen_breaker_epoch: AtomicI64::new(breaker_epoch),
}))
}
#[cfg(test)]
mod tests {
use super::*;
async fn sqlite_pool() -> DbPool {
crate::storage::test_sqlite_pool().await
}
#[tokio::test]
async fn test_disabled_runtime_is_inert() {
let runtime = init_cluster_runtime(&ClusterConfig::default(), &sqlite_pool().await)
.await
.expect("disabled runtime never fails");
assert!(!runtime.enabled);
assert!(runtime.redis.is_none());
assert!(runtime.default_cache.is_none());
assert_eq!(runtime.instance_id.len(), 36); }
#[tokio::test]
async fn test_configured_instance_id_wins() {
let config = ClusterConfig {
instance_id: "node-7".to_string(),
..Default::default()
};
let runtime = init_cluster_runtime(&config, &sqlite_pool().await)
.await
.expect("runtime");
assert_eq!(runtime.instance_id, "node-7");
}
}