mod budget;
mod git;
mod idle;
mod intervention;
mod metadata;
pub mod output_patterns;
use std::sync::Arc;
use std::time::Duration;
use crate::harness::{HarnessRegistry, HarnessSignals};
use idle::check_idle_sessions;
#[cfg(test)]
use idle::{check_session_idle, handle_active_session, handle_idle_session};
pub(crate) use metadata::refresh_exact_usage;
use metadata::{build_session_event, detect_and_store_output_metadata};
pub use output_patterns::detect_waiting_for_input;
use pulpo_common::event::PulpoEvent;
use tokio::sync::broadcast;
use tracing::info;
use pulpo_common::session::Session;
use crate::backend::Backend;
use crate::store::Store;
use git::update_git_info;
fn resolve_backend_id(session: &Session, backend: &dyn Backend) -> String {
session
.backend_session_id
.clone()
.unwrap_or_else(|| backend.session_id(&session.name))
}
pub(super) const fn harness_owns_state(session: &Session) -> bool {
session.harness_last_event_at.is_some()
}
static HARNESS_REGISTRY: std::sync::LazyLock<HarnessRegistry> =
std::sync::LazyLock::new(HarnessRegistry::default);
pub(super) fn owned_signals(session: &Session) -> HarnessSignals {
HARNESS_REGISTRY.owned_signals_for(session.harness.as_deref(), harness_owns_state(session))
}
#[cfg_attr(coverage, allow(unused_variables))]
async fn list_sessions_or_warn(store: &Store, context: &str) -> Vec<Session> {
match store.list_sessions().await {
Ok(sessions) => sessions,
#[allow(unused_variables)]
Err(error) => {
coverage_warn!("{context}: failed to list sessions: {error}");
Vec::new()
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum IdleAction {
Alert,
Kill,
}
#[derive(Debug, Clone)]
pub struct IdleConfig {
pub enabled: bool,
pub timeout_secs: u64,
pub action: IdleAction,
pub threshold_secs: u64,
}
impl Default for IdleConfig {
fn default() -> Self {
Self {
enabled: true,
timeout_secs: 600,
action: IdleAction::Alert,
threshold_secs: 60,
}
}
}
#[derive(Debug, Clone)]
pub struct WatchdogRuntimeConfig {
pub interval: Duration,
pub idle: IdleConfig,
pub extra_waiting_patterns: Vec<String>,
}
#[cfg_attr(coverage, allow(dead_code))]
pub struct ReadyContext {
pub event_tx: Option<broadcast::Sender<PulpoEvent>>,
pub node_name: String,
}
async fn run_watchdog_tick(
backend: &Arc<dyn Backend>,
store: &Store,
cfg: &WatchdogRuntimeConfig,
ready_ctx: &ReadyContext,
) {
if cfg.idle.enabled {
check_idle_sessions(
backend,
store,
&cfg.idle,
ready_ctx,
&cfg.extra_waiting_patterns,
)
.await;
}
budget::enforce_budgets(backend, store, ready_ctx).await;
update_git_info(store).await;
}
pub async fn run_watchdog_loop(
backend: Arc<dyn Backend>,
store: Store,
cfg: WatchdogRuntimeConfig,
mut shutdown_rx: tokio::sync::watch::Receiver<bool>,
ready_ctx: ReadyContext,
) {
let mut tick = tokio::time::interval(cfg.interval);
tick.tick().await;
loop {
tokio::select! {
_ = tick.tick() => {
run_watchdog_tick(&backend, &store, &cfg, &ready_ctx).await;
}
_ = shutdown_rx.changed() => {
info!("Watchdog shutting down");
break;
}
}
}
}
#[cfg(test)]
mod tests;