use trusty_common::memory_core::palace::{Palace, PalaceId, RoomType};
use trusty_common::memory_core::retrieval::recall_with_default_embedder;
use uuid::Uuid;
pub(crate) const PROBE_SENTINEL_CONTENT: &str =
"__trusty_memory_health_sentinel__ issue-#1142 self-heal probe";
pub(crate) const PROBE_SENTINEL_PREFIX: &str = "__trusty_memory_health_sentinel__";
use crate::AppState;
use super::{to_value, HEALTH_PROBE_PALACE};
use crate::transport::api_error::ApiError;
#[derive(Debug, Clone, Copy, Default, serde::Deserialize)]
pub struct HealthQuery {
#[serde(default)]
pub probe: bool,
#[serde(default)]
pub deep: bool,
}
impl HealthQuery {
fn wants_deep_probe(&self) -> bool {
self.probe || self.deep
}
}
#[derive(serde::Serialize)]
pub struct HealthResponse {
pub status: String,
#[serde(skip_serializing_if = "Option::is_none")]
pub detail: Option<String>,
pub version: &'static str,
pub rss_mb: u64,
pub disk_bytes: u64,
pub cpu_pct: f32,
pub uptime_secs: u64,
#[serde(skip_serializing_if = "Option::is_none")]
pub socket: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub open_fds: Option<u64>,
#[serde(skip_serializing_if = "Option::is_none")]
pub fd_soft_limit: Option<u64>,
#[serde(skip_serializing_if = "Option::is_none")]
pub update_available: Option<String>,
pub daemon_state: String,
pub worker: WorkerHealth,
#[serde(skip_serializing_if = "Vec::is_empty")]
pub unopenable_palaces: Vec<UnopenablePalace>,
#[serde(skip_serializing_if = "Vec::is_empty")]
pub drawer_degraded_palaces: Vec<String>,
}
#[derive(serde::Serialize)]
pub struct UnopenablePalace {
pub id: String,
pub reason: String,
}
#[derive(serde::Serialize)]
pub struct WorkerHealth {
pub in_flight: usize,
#[serde(skip_serializing_if = "Option::is_none")]
pub oldest_age_secs: Option<u64>,
pub wedged: bool,
#[serde(skip_serializing_if = "Option::is_none")]
pub wedged_reason: Option<WedgeReason>,
#[serde(skip_serializing_if = "Option::is_none")]
pub stalled_lock: Option<StalledLockHealth>,
pub stall_tracking_ok: bool,
}
#[derive(serde::Serialize, Debug, Clone, Copy, PartialEq, Eq)]
#[serde(rename_all = "snake_case")]
pub enum WedgeReason {
Lock,
Pool,
}
#[derive(serde::Serialize)]
pub struct StalledLockHealth {
pub palace: String,
pub lock: crate::lock_stall::PalaceLock,
pub age_secs: u64,
}
pub async fn health(state: &AppState, query: HealthQuery) -> Result<serde_json::Value, ApiError> {
let (rss_mb, cpu_pct) = {
let mut metrics = state.sys_metrics.lock().await;
metrics.sample()
};
let disk_bytes = state.disk_bytes.load(std::sync::atomic::Ordering::Relaxed);
let uptime_secs = state.started_at.elapsed().as_secs();
let socket = crate::transport::uds::socket_path()
.ok()
.map(|p| p.display().to_string());
let open_fds = crate::fd_metrics::count_open_fds();
let fd_soft_limit = crate::fd_metrics::fd_soft_limit();
let wedge_threshold = state.wedge_threshold;
let oldest_age = state.worker_liveness.oldest_age();
state.lock_stalls.sweep_if_due(
&state.registry,
crate::lock_stall::probe_interval(wedge_threshold),
);
let now = std::time::Instant::now();
let stalled = state.lock_stalls.oldest_stall_at(now);
let lock_wedged = stalled.as_ref().is_some_and(|s| s.age > wedge_threshold);
let pool_wedged = oldest_age.is_some_and(|age| age > wedge_threshold);
let stall_tracking_degraded = state.lock_stalls.degraded_at(now);
let worker = WorkerHealth {
in_flight: state.worker_liveness.in_flight(),
oldest_age_secs: oldest_age.map(|d| d.as_secs()),
wedged: lock_wedged || pool_wedged,
wedged_reason: match (lock_wedged, pool_wedged) {
(true, _) => Some(WedgeReason::Lock),
(false, true) => Some(WedgeReason::Pool),
(false, false) => None,
},
stalled_lock: stalled.as_ref().map(|s| StalledLockHealth {
palace: s.palace.clone(),
lock: s.lock,
age_secs: s.age.as_secs(),
}),
stall_tracking_ok: stall_tracking_degraded.is_none(),
};
let (status, detail) = if query.wants_deep_probe() {
match run_health_round_trip(state).await {
Ok(()) => ("ok".to_string(), None),
Err(err) => {
tracing::warn!("/health round-trip degraded: {err}");
("degraded".to_string(), Some(err.to_string()))
}
}
} else {
("ok".to_string(), None)
};
let (status, detail) = if let Some(s) = stalled.as_ref().filter(|_| lock_wedged) {
tracing::warn!(
palace = %s.palace,
lock = ?s.lock,
age_secs = s.age.as_secs(),
"/health: palace handle lock held past the wedge threshold"
);
(
"wedged".to_string(),
Some(format!(
"palace '{}' {:?} lock has been unavailable for {}s (threshold {}s) — its \
holder is not releasing it, and writers queued behind it time out",
s.palace,
s.lock,
s.age.as_secs(),
wedge_threshold.as_secs()
)),
)
} else if worker.wedged {
let secs = worker.oldest_age_secs.unwrap_or_default();
tracing::warn!(
oldest_age_secs = secs,
in_flight = worker.in_flight,
"/health: worker pool appears wedged"
);
(
"wedged".to_string(),
Some(format!(
"oldest in-flight palace operation has been running {secs}s \
(threshold {}s, {} in flight) — workers are not making progress",
wedge_threshold.as_secs(),
worker.in_flight
)),
)
} else if let Some(reason) = stall_tracking_degraded.filter(|_| status == "ok") {
tracing::warn!("/health: {reason}");
("degraded".to_string(), Some(reason))
} else {
(status, detail)
};
let update_available = state.update_available.lock().ok().and_then(|g| g.clone());
let daemon_state = match state.readiness() {
crate::DaemonReadiness::Ready => "ready",
crate::DaemonReadiness::Warming
if trusty_common::memory_core::retrieval::shared_embedder_initialized() =>
{
"ready"
}
crate::DaemonReadiness::Warming => "warming",
}
.to_string();
let mut unopenable_palaces: Vec<UnopenablePalace> = state
.registry
.unopenable()
.into_iter()
.map(|(id, reason)| UnopenablePalace {
id: id.as_str().to_string(),
reason,
})
.collect();
unopenable_palaces.sort_by(|a, b| a.id.cmp(&b.id));
let mut drawer_degraded_palaces: Vec<String> = state
.registry
.list()
.into_iter()
.filter(|id| {
state
.registry
.peek(id)
.is_some_and(|handle| handle.drawer_load_degraded)
})
.map(|id| id.as_str().to_string())
.collect();
drawer_degraded_palaces.sort();
if !drawer_degraded_palaces.is_empty() {
tracing::warn!(
palaces = ?drawer_degraded_palaces,
"/health: open palaces are serving a partial drawer corpus"
);
}
to_value(HealthResponse {
status,
detail,
version: env!("CARGO_PKG_VERSION"),
rss_mb,
disk_bytes,
cpu_pct,
uptime_secs,
socket,
open_fds,
fd_soft_limit,
update_available,
daemon_state,
worker,
unopenable_palaces,
drawer_degraded_palaces,
})
}
#[derive(Debug, thiserror::Error)]
pub(crate) enum HealthProbeError {
#[error("open palace failed: {0}")]
OpenPalace(String),
#[error("provision health probe palace failed: {0}")]
EnsureProbePalace(String),
#[error("store failed: {0}")]
Store(String),
#[error("recall failed: {0}")]
Recall(String),
#[error("recall did not return the probe drawer (id={0})")]
ProbeMissing(Uuid),
#[error("delete probe drawer failed: {0}")]
Delete(String),
}
pub(crate) fn ensure_health_probe_palace(state: &AppState) -> Result<(), HealthProbeError> {
let id = PalaceId::new(HEALTH_PROBE_PALACE);
if state.registry.get(&id).is_some() {
return Ok(());
}
if state.registry.open_palace(&state.data_root, &id).is_ok() {
return Ok(());
}
let palace = Palace {
id: id.clone(),
name: HEALTH_PROBE_PALACE.to_string(),
description: Some(
"Internal health-probe palace (issue #185). Hidden from listings; \
holds short-lived round-trip drawers cleaned up on every probe."
.to_string(),
),
created_at: chrono::Utc::now(),
data_dir: state.data_root.join(HEALTH_PROBE_PALACE),
};
state
.registry
.create_palace(&state.data_root, palace)
.map_err(|e| HealthProbeError::EnsureProbePalace(format!("{e:#}")))?;
Ok(())
}
pub(crate) async fn seed_probe_sentinel_if_absent(
handle: &std::sync::Arc<trusty_common::memory_core::PalaceHandle>,
) -> Result<bool, HealthProbeError> {
let sentinel_present = handle
.drawers
.read()
.iter()
.any(|d| d.content().starts_with(PROBE_SENTINEL_PREFIX));
if sentinel_present {
return Ok(false);
}
use trusty_common::memory_core::retrieval::RememberOptions;
handle
.remember_with_options(
PROBE_SENTINEL_CONTENT.to_string(),
RoomType::General,
vec!["healthcheck".to_string(), "sentinel".to_string()],
0.0,
RememberOptions {
force: true,
..Default::default()
},
)
.await
.map_err(|e| HealthProbeError::EnsureProbePalace(format!("seed sentinel: {e:#}")))?;
tracing::info!(
palace = HEALTH_PROBE_PALACE,
self_heal = true,
"health probe: seeded sentinel drawer (issue #1142 self-heal)"
);
Ok(true)
}
pub(crate) async fn run_health_round_trip(state: &AppState) -> Result<(), HealthProbeError> {
ensure_health_probe_palace(state)?;
let probe_id = PalaceId::new(HEALTH_PROBE_PALACE);
let handle = state
.registry
.open_palace(&state.data_root, &probe_id)
.map_err(|e| HealthProbeError::OpenPalace(format!("{e:#}")))?;
if seed_probe_sentinel_if_absent(&handle).await? {
return Ok(());
}
run_health_round_trip_inner(handle, |handle, query| async move {
recall_with_default_embedder(&handle, &query, 5)
.await
.map_err(|e| HealthProbeError::Recall(format!("{e:#}")))
})
.await
}
pub(crate) async fn run_health_round_trip_inner<F, Fut>(
handle: std::sync::Arc<trusty_common::memory_core::PalaceHandle>,
recall: F,
) -> Result<(), HealthProbeError>
where
F: FnOnce(std::sync::Arc<trusty_common::memory_core::PalaceHandle>, String) -> Fut,
Fut: std::future::Future<
Output = Result<Vec<trusty_common::memory_core::retrieval::RecallResult>, HealthProbeError>,
>,
{
let probe_token = Uuid::new_v4();
let probe_content = format!("__trusty_memory_healthcheck__ probe {probe_token}");
let drawer_id = handle
.remember(
probe_content.clone(),
RoomType::General,
vec!["healthcheck".to_string()],
0.0,
)
.await
.map_err(|e| HealthProbeError::Store(format!("{e:#}")))?;
let recall_result = recall(handle.clone(), probe_content).await;
let delete_result = handle.forget(drawer_id).await;
match recall_result {
Ok(hits) => {
if !hits.iter().any(|hit| hit.drawer.id == drawer_id) {
return Err(HealthProbeError::ProbeMissing(drawer_id));
}
}
Err(e) => return Err(e),
}
delete_result.map_err(|e| HealthProbeError::Delete(format!("{e:#}")))?;
Ok(())
}