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,
scope: crate::cluster::EpochScope,
) -> Result<(), crate::errors::OrionError> {
if scope.touches_connectors() {
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 _reload_guard = state.reload_lock.lock().await;
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 current_engine = state.engine.load();
let (workflows, engine_issues) =
crate::engine::build_engine_workflows(&channels, &active_workflows, &*current_engine);
let new_engine = Arc::new(
current_engine
.with_new_workflows(workflows)
.map_err(crate::errors::OrionError::Engine)?,
);
state
.channel_registry
.reload(
&channels,
crate::channel::ReloadDeps {
connector_registry: &state.connector_registry,
cache_pool: &state.caches.cache_pool,
datalogic: &state.datalogic,
jwks: &state.jwks,
global_trace_storage: &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");
state
.reload_degraded
.store(false, std::sync::atomic::Ordering::Release);
}
Err(e) => {
crate::metrics::record_engine_reload("failure");
state
.reload_degraded
.store(true, std::sync::atomic::Ordering::Release);
tracing::error!(
error = %e,
"Engine reload failed: this node is still serving the previous \
generation, so its running config no longer matches the database"
);
}
}
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,
crate::bootstrap::IngestDeps {
engine: state.engine.clone(),
channel_registry: state.channel_registry.clone(),
datalogic: state.datalogic.clone(),
vars: state.vars.clone(),
kafka_producer: state.kafka.producer.clone(),
instance_id: state
.cluster
.enabled
.then(|| state.cluster.instance_id.clone()),
trace_repo: state.repos.traces.clone(),
persistence_queue: state.trace_persistence_queue.clone(),
max_result_size_bytes: state.config.trace_queue.max_result_size_bytes,
},
)
.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);
}
}
}
});
}