#![allow(dead_code)]
use serde::{Deserialize, Serialize};
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub(crate) enum OperationKind {
FullIndex,
IncrementalIndex,
WatchCycle,
SemanticBackfill,
HistoricalSalvage,
LexicalRebuild,
Repair,
}
impl OperationKind {
pub(crate) fn is_repair(self) -> bool {
matches!(self, Self::Repair | Self::HistoricalSalvage)
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub(crate) enum OperationState {
Missing,
Building,
Repairing,
Stalled,
Stale,
WaitingOnLock,
Ready,
}
impl OperationState {
pub(crate) fn is_healthy(self) -> bool {
matches!(self, Self::Building | Self::Repairing | Self::Ready)
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub(crate) enum ProgressNextStep {
None,
WaitBounded,
AttachOrWait,
WaitForLockOwner,
RestartOperation,
StartOperation,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub(crate) struct ActiveLock {
pub owner: String,
pub acquired_at_ms: i64,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub(crate) struct ProgressReport {
pub operation: OperationKind,
pub phase: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub subphase: Option<String>,
pub current: u64,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub total: Option<u64>,
pub units: String,
pub started_at_ms: i64,
pub last_forward_progress_at_ms: i64,
pub heartbeat_at_ms: i64,
pub stall_threshold_ms: i64,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub active_lock: Option<ActiveLock>,
}
impl ProgressReport {
fn is_complete(&self) -> bool {
matches!(self.total, Some(t) if t > 0 && self.current >= t)
}
pub(crate) fn resolve(&self, now_ms: i64) -> ProgressSnapshot {
let elapsed_ms = (now_ms - self.started_at_ms).max(0);
let since_forward_ms = (now_ms - self.last_forward_progress_at_ms).max(0);
let since_heartbeat_ms = (now_ms - self.heartbeat_at_ms).max(0);
let threshold = self.stall_threshold_ms.max(0);
let complete = self.is_complete();
let heartbeat_dead = since_heartbeat_ms > threshold;
let no_forward = since_forward_ms >= threshold;
let (state, stall_reason) = if complete {
(OperationState::Ready, None)
} else if heartbeat_dead {
(
OperationState::Stale,
Some(format!(
"no heartbeat for {since_heartbeat_ms}ms (> {threshold}ms); operation likely died mid-run"
)),
)
} else if no_forward {
match &self.active_lock {
Some(lock) => (
OperationState::WaitingOnLock,
Some(format!(
"no forward progress for {since_forward_ms}ms; blocked on lock held by {}",
lock.owner
)),
),
None => (
OperationState::Stalled,
Some(format!(
"heartbeat alive but no forward progress for {since_forward_ms}ms (> {threshold}ms)"
)),
),
}
} else if self.operation.is_repair() {
(OperationState::Repairing, None)
} else {
(OperationState::Building, None)
};
let next_step = match state {
OperationState::Ready => ProgressNextStep::None,
OperationState::Building | OperationState::Repairing => ProgressNextStep::WaitBounded,
OperationState::WaitingOnLock => ProgressNextStep::WaitForLockOwner,
OperationState::Stalled => ProgressNextStep::AttachOrWait,
OperationState::Stale => ProgressNextStep::RestartOperation,
OperationState::Missing => ProgressNextStep::StartOperation,
};
ProgressSnapshot {
operation: self.operation,
phase: self.phase.clone(),
subphase: self.subphase.clone(),
current: self.current,
total: self.total,
units: self.units.clone(),
elapsed_ms,
since_forward_progress_ms: since_forward_ms,
since_heartbeat_ms,
stall_threshold_ms: threshold,
stalled: matches!(state, OperationState::Stalled),
stall_reason,
active_lock: self.active_lock.clone(),
state,
next_step,
}
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub(crate) struct ProgressSnapshot {
pub operation: OperationKind,
pub phase: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub subphase: Option<String>,
pub current: u64,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub total: Option<u64>,
pub units: String,
pub elapsed_ms: i64,
pub since_forward_progress_ms: i64,
pub since_heartbeat_ms: i64,
pub stall_threshold_ms: i64,
pub stalled: bool,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub stall_reason: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub active_lock: Option<ActiveLock>,
pub state: OperationState,
pub next_step: ProgressNextStep,
}
impl ProgressSnapshot {
pub(crate) fn missing(operation: OperationKind, units: impl Into<String>) -> Self {
Self {
operation,
phase: "absent".to_string(),
subphase: None,
current: 0,
total: None,
units: units.into(),
elapsed_ms: 0,
since_forward_progress_ms: 0,
since_heartbeat_ms: 0,
stall_threshold_ms: 0,
stalled: false,
stall_reason: None,
active_lock: None,
state: OperationState::Missing,
next_step: ProgressNextStep::StartOperation,
}
}
}
#[cfg(test)]
mod tests {
use super::*;
const MIN: i64 = 60_000;
fn report(now: i64, kind: OperationKind) -> ProgressReport {
ProgressReport {
operation: kind,
phase: "scanning".to_string(),
subphase: Some("connectors".to_string()),
current: 100,
total: Some(1000),
units: "sessions".to_string(),
started_at_ms: now - 10 * MIN,
last_forward_progress_at_ms: now,
heartbeat_at_ms: now,
stall_threshold_ms: 5 * MIN,
active_lock: None,
}
}
#[test]
fn enums_serialize_as_snake_case() {
assert_eq!(
serde_json::to_string(&OperationKind::SemanticBackfill).unwrap(),
"\"semantic_backfill\""
);
assert_eq!(
serde_json::to_string(&OperationState::WaitingOnLock).unwrap(),
"\"waiting_on_lock\""
);
assert_eq!(
serde_json::to_string(&ProgressNextStep::AttachOrWait).unwrap(),
"\"attach_or_wait\""
);
}
#[test]
fn true_forward_progress_reads_as_building() {
let now = 1_000 * MIN;
let snap = report(now, OperationKind::FullIndex).resolve(now);
assert_eq!(snap.state, OperationState::Building);
assert!(!snap.stalled);
assert!(snap.stall_reason.is_none());
assert_eq!(snap.next_step, ProgressNextStep::WaitBounded);
assert_eq!(snap.elapsed_ms, 10 * MIN);
}
#[test]
fn repair_progress_reads_as_repairing() {
let now = 1_000 * MIN;
for kind in [OperationKind::Repair, OperationKind::HistoricalSalvage] {
let snap = report(now, kind).resolve(now);
assert_eq!(snap.state, OperationState::Repairing, "{kind:?}");
assert!(snap.state.is_healthy());
}
}
#[test]
fn heartbeat_without_forward_progress_is_stalled_not_building() {
let now = 1_000 * MIN;
let mut r = report(now, OperationKind::LexicalRebuild);
r.heartbeat_at_ms = now;
r.last_forward_progress_at_ms = now - 6 * MIN;
let snap = r.resolve(now);
assert_eq!(snap.state, OperationState::Stalled);
assert!(snap.stalled);
assert!(
snap.stall_reason
.as_deref()
.unwrap()
.contains("no forward progress")
);
assert_eq!(snap.next_step, ProgressNextStep::AttachOrWait);
}
#[test]
fn dead_heartbeat_past_threshold_is_stale_with_restart() {
let now = 1_000 * MIN;
let mut r = report(now, OperationKind::SemanticBackfill);
r.heartbeat_at_ms = now - 6 * MIN;
r.last_forward_progress_at_ms = now - 6 * MIN;
let snap = r.resolve(now);
assert_eq!(snap.state, OperationState::Stale);
assert!(!snap.stalled, "stale is distinct from a live stall");
assert_eq!(snap.next_step, ProgressNextStep::RestartOperation);
}
#[test]
fn no_progress_with_lock_is_waiting_on_lock_not_stalled() {
let now = 1_000 * MIN;
let mut r = report(now, OperationKind::FullIndex);
r.last_forward_progress_at_ms = now - 6 * MIN;
r.heartbeat_at_ms = now;
r.active_lock = Some(ActiveLock {
owner: "host-b/pid-4242".to_string(),
acquired_at_ms: now - 7 * MIN,
});
let snap = r.resolve(now);
assert_eq!(snap.state, OperationState::WaitingOnLock);
assert!(!snap.stalled);
assert_eq!(snap.next_step, ProgressNextStep::WaitForLockOwner);
assert!(
snap.stall_reason
.as_deref()
.unwrap()
.contains("host-b/pid-4242")
);
}
#[test]
fn completed_publish_is_ready() {
let now = 1_000 * MIN;
let mut r = report(now, OperationKind::FullIndex);
r.current = 1000;
r.total = Some(1000);
r.last_forward_progress_at_ms = now - 30 * MIN;
r.heartbeat_at_ms = now - 30 * MIN;
let snap = r.resolve(now);
assert_eq!(snap.state, OperationState::Ready);
assert!(!snap.stalled);
assert_eq!(snap.next_step, ProgressNextStep::None);
}
#[test]
fn missing_operation_snapshot_recommends_start() {
let snap = ProgressSnapshot::missing(OperationKind::FullIndex, "sessions");
assert_eq!(snap.state, OperationState::Missing);
assert_eq!(snap.next_step, ProgressNextStep::StartOperation);
assert!(!snap.state.is_healthy());
}
#[test]
fn indeterminate_total_never_reports_ready() {
let now = 1_000 * MIN;
let mut r = report(now, OperationKind::WatchCycle);
r.total = None; r.current = 9999;
let snap = r.resolve(now);
assert_ne!(snap.state, OperationState::Ready);
assert_eq!(snap.state, OperationState::Building);
}
#[test]
fn snapshot_round_trips_through_json() {
let now = 1_000 * MIN;
let snap = report(now, OperationKind::IncrementalIndex).resolve(now);
let json = serde_json::to_string(&snap).unwrap();
assert!(json.contains("\"operation\":\"incremental_index\""));
assert!(json.contains("\"state\":\"building\""));
assert!(json.contains("\"next_step\":\"wait_bounded\""));
let parsed: ProgressSnapshot = serde_json::from_str(&json).unwrap();
assert_eq!(parsed, snap);
}
}