use std::sync::Arc;
use crate::server::state::AppState;
#[derive(Clone, Copy, Default)]
pub struct ReloadOpts {
pub kafka_restart_jitter: bool,
}
pub async fn reload_engine(state: &AppState) -> Result<(), crate::errors::OrionError> {
reload_engine_with_opts(state, ReloadOpts::default()).await
}
pub async fn resync_from_db(state: &AppState) -> Result<(), crate::errors::OrionError> {
state
.connector_registry
.reload(state.repos.connectors.as_ref())
.await?;
state.caches.sql_pool_cache.evict_all().await;
state.caches.mongo_pool_cache.evict_all().await;
state.caches.cache_pool.evict_all_pools().await;
reload_engine_with_opts(
state,
ReloadOpts {
kafka_restart_jitter: true,
},
)
.await
}
#[tracing::instrument(skip(state, opts))]
pub async fn reload_engine_with_opts(
state: &AppState,
opts: ReloadOpts,
) -> Result<(), crate::errors::OrionError> {
let start = std::time::Instant::now();
let result = async {
let channels = state.repos.channels.list_active().await?;
let channels = crate::engine::filter_channels(channels, &state.config.channel_filter);
let active_workflows = state.repos.workflows.list_active().await?;
let (workflows, engine_issues) =
crate::engine::build_engine_workflows(&channels, &active_workflows);
let current_engine = state.engine.load();
let new_engine = Arc::new(
current_engine
.with_new_workflows(workflows)
.map_err(crate::errors::OrionError::Engine)?,
);
state
.channel_registry
.reload(
&channels,
&state.connector_registry,
&state.caches.cache_pool,
&state.datalogic,
&state.config.trace_storage,
engine_issues,
)
.await;
state.engine.store(new_engine);
crate::metrics::set_active_workflows(active_workflows.len() as f64);
if state.config.kafka.enabled {
restart_kafka_consumer_if_needed(state, &channels, opts).await;
}
tracing::info!(
workflow_count = active_workflows.len(),
channel_count = channels.len(),
"Engine reloaded"
);
Ok(())
}
.await;
let duration = start.elapsed().as_secs_f64();
crate::metrics::record_engine_reload_duration(duration);
match &result {
Ok(()) => crate::metrics::record_engine_reload("success"),
Err(_) => crate::metrics::record_engine_reload("failure"),
}
result
}
async fn restart_kafka_consumer_if_needed(
state: &AppState,
channels: &[crate::storage::models::Channel],
opts: ReloadOpts,
) {
use std::collections::HashSet;
let all_topics = crate::kafka::merge_kafka_topics(&state.config.kafka, channels);
let new_topic_set: HashSet<String> = all_topics.iter().map(|t| t.topic.clone()).collect();
let mut handle_guard = state.kafka.consumer_handle.lock().await;
if let Some(ref existing_handle) = *handle_guard
&& *existing_handle.topics() == new_topic_set
{
tracing::info!("Kafka topics unchanged, pausing consumer during engine swap");
if let Err(e) = existing_handle.pause() {
tracing::warn!(error = %e, "Failed to pause Kafka consumer, falling back to full restart");
} else {
tokio::time::sleep(std::time::Duration::from_millis(100)).await;
if let Err(e) = existing_handle.resume() {
tracing::error!(error = %e, "Failed to resume Kafka consumer after engine reload");
} else {
tracing::info!("Kafka consumer resumed after engine reload");
return;
}
}
}
if opts.kafka_restart_jitter {
let jitter_ms = rand::random_range(0..=5000u64);
tracing::info!(
jitter_ms,
"Jittering Kafka consumer restart (epoch-driven reload)"
);
tokio::time::sleep(std::time::Duration::from_millis(jitter_ms)).await;
}
if let Some(ref existing_handle) = *handle_guard {
let _ = existing_handle.pause(); tokio::time::sleep(std::time::Duration::from_millis(50)).await;
}
if let Some(old_handle) = handle_guard.take() {
tracing::info!("Shutting down Kafka consumer for topic refresh...");
old_handle.shutdown().await;
}
match try_start_ingest(state, channels) {
Ok(Some(new_handle)) => {
tracing::info!(
topics = ?new_topic_set,
"Kafka consumer restarted with updated topics"
);
*handle_guard = Some(new_handle);
state.kafka.ingest_status.set_degraded(false);
}
Ok(None) => {
tracing::info!("No Kafka topics configured or from DB, consumer not started");
state.kafka.ingest_status.set_degraded(false);
}
Err(e) => {
state.kafka.ingest_status.set_degraded(true);
crate::metrics::record_error("kafka_restart");
tracing::error!(
error = %e,
"Failed to restart Kafka consumer; ingestion is down — retrying with backoff"
);
drop(handle_guard);
spawn_kafka_restart_supervisor(state);
}
}
}
fn try_start_ingest(
state: &AppState,
channels: &[crate::storage::models::Channel],
) -> Result<Option<crate::kafka::consumer::ConsumerHandle>, String> {
crate::bootstrap::start_kafka_ingest(
&state.config.kafka,
channels,
state.engine.clone(),
state.channel_registry.clone(),
state.datalogic.clone(),
state.kafka.producer.clone(),
state
.cluster
.enabled
.then(|| state.cluster.instance_id.as_str()),
)
.map_err(|e| e.to_string())
}
pub fn spawn_kafka_restart_supervisor(state: &AppState) {
if !state.kafka.ingest_status.claim_supervisor() {
return;
}
let state = state.clone();
tokio::spawn(async move {
let mut backoff_ms = crate::kafka::consumer::INITIAL_RETRY_BACKOFF_MS;
loop {
tokio::time::sleep(std::time::Duration::from_millis(backoff_ms)).await;
let mut handle_guard = state.kafka.consumer_handle.lock().await;
if !state.ready.load(std::sync::atomic::Ordering::Acquire) {
state.kafka.ingest_status.release_supervisor();
break;
}
if handle_guard.is_some() {
state.kafka.ingest_status.release_supervisor();
break;
}
let channels = match state.repos.channels.list_active().await {
Ok(channels) => {
crate::engine::filter_channels(channels, &state.config.channel_filter)
}
Err(e) => {
tracing::warn!(
error = %e,
backoff_ms,
"Kafka restart supervisor could not list channels; will retry"
);
drop(handle_guard);
backoff_ms = crate::kafka::consumer::next_backoff_ms(backoff_ms);
continue;
}
};
let started = try_start_ingest(&state, &channels);
match started {
Ok(Some(handle)) => {
*handle_guard = Some(handle);
state.kafka.ingest_status.release_supervisor();
state.kafka.ingest_status.set_degraded(false);
tracing::info!("Kafka consumer restored by restart supervisor");
break;
}
Ok(None) => {
state.kafka.ingest_status.release_supervisor();
state.kafka.ingest_status.set_degraded(false);
tracing::info!("No Kafka topics remain; restart supervisor standing down");
break;
}
Err(e) => {
crate::metrics::record_error("kafka_restart");
tracing::error!(
error = %e,
backoff_ms,
"Kafka consumer restart failed; ingestion still down"
);
drop(handle_guard);
backoff_ms = crate::kafka::consumer::next_backoff_ms(backoff_ms);
}
}
}
});
}