use std::panic::{AssertUnwindSafe, catch_unwind};
use std::time::Instant;
use tracing::{error, warn};
use crate::bridge::dispatch::BridgeResponse;
use crate::bridge::envelope::{ErrorCode, Payload, Response, Status};
use crate::data::core_health::{CoreHealthWatchdog, MAX_CONSECUTIVE_PANICS, PANIC_WINDOW_SECS};
use crate::data::eventfd::EventFd;
use crate::data::executor::core_loop::CoreLoop;
const IDLE_POLL_TIMEOUT_MS: i32 = 100;
const MAX_TASKS_PER_ITERATION: usize = 256;
pub(super) fn run_event_loop(
core: &mut CoreLoop,
core_id: usize,
efd: &EventFd,
checkpoint_interval: std::time::Duration,
) {
let mut watchdog = CoreHealthWatchdog::new();
let mut last_checkpoint = Instant::now();
let mut last_event_emit = Instant::now();
let mut heartbeat_interval = heartbeat_interval_with_jitter();
loop {
efd.poll_wait(IDLE_POLL_TIMEOUT_MS);
while efd.drain() > 0 {}
if watchdog.is_degraded() {
drain_and_reject(core, core_id);
continue;
}
let mut tasks_processed = 0usize;
loop {
let result = catch_unwind(AssertUnwindSafe(|| core.tick()));
match result {
Ok(0) => break, Ok(_) => {
watchdog.record_success();
tasks_processed += 1;
if tasks_processed >= MAX_TASKS_PER_ITERATION {
break; }
}
Err(panic_payload) => {
let msg = panic_message(&panic_payload);
error!(
core_id,
panic_count = watchdog.consecutive_panics + 1,
message = %msg,
"data plane core caught panic during tick"
);
let is_degraded = watchdog.record_panic();
if is_degraded {
error!(
core_id,
threshold = MAX_CONSECUTIVE_PANICS,
window_secs = PANIC_WINDOW_SECS,
"core entered DEGRADED mode — rejecting all requests"
);
drain_and_reject(core, core_id);
}
break; }
}
}
if last_checkpoint.elapsed() >= checkpoint_interval {
if let Err(e) = core.checkpoint_vector_indexes() {
warn!(
core = core_id,
error = %e,
"periodic vector checkpoint failed; the coordinated checkpoint \
will clamp its reported LSN if it fails there too"
);
}
last_checkpoint = Instant::now();
}
core.maybe_run_maintenance();
if tasks_processed > 0 {
last_event_emit = Instant::now();
} else if last_event_emit.elapsed() >= heartbeat_interval {
core.emit_heartbeat();
last_event_emit = Instant::now();
heartbeat_interval = heartbeat_interval_with_jitter();
}
}
}
fn drain_and_reject(core: &mut CoreLoop, core_id: usize) {
core.drain_requests();
while let Some(task) = core.task_queue.pop_front() {
let response = Response {
request_id: task.request_id(),
status: Status::Error,
attempt: 1,
partial: false,
payload: Payload::empty(),
watermark_lsn: core.watermark,
error_code: Some(Box::new(ErrorCode::Internal {
detail: format!("core-{core_id} is degraded after repeated panics"),
})),
read_set_valid: None,
read_version_lsn: crate::types::Lsn::ZERO,
write_set: Vec::new(),
};
if let Err(e) = core
.response_tx
.try_push(BridgeResponse { inner: response })
{
warn!(core_id, error = %e, "failed to send degraded-rejection response");
}
}
}
fn heartbeat_interval_with_jitter() -> std::time::Duration {
let seed = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap_or_default()
.as_nanos() as u64;
let mut x = seed;
x ^= x >> 30;
x = x.wrapping_mul(0xbf58476d1ce4e5b9);
x ^= x >> 27;
let jitter_ms = (x % 201) as i64 - 100;
std::time::Duration::from_millis((1000 + jitter_ms) as u64)
}
fn panic_message(payload: &Box<dyn std::any::Any + Send>) -> String {
if let Some(s) = payload.downcast_ref::<&str>() {
(*s).to_string()
} else if let Some(s) = payload.downcast_ref::<String>() {
s.clone()
} else {
"non-string panic payload".to_string()
}
}