use super::census::{DriftCensus, SkipReason};
use super::ignore::should_ignore_drift;
use super::{DriftFinding, DRIFT_QUERY_TIMEOUT_SECS};
use crate::core::types::{Machine, Resource, ResourceStatus, ResourceType, TaskMode};
#[derive(Debug, Clone, Copy)]
pub struct DriftOptions {
pub run_task_checks: bool,
}
impl Default for DriftOptions {
fn default() -> Self {
Self {
run_task_checks: true,
}
}
}
pub(super) fn owns(resource: &Resource) -> bool {
resource.resource_type == ResourceType::Task
&& resource.completion_check.is_some()
&& resource.task_mode.as_ref() != Some(&TaskMode::Service)
}
pub(super) fn detect_task_drift(
lock: &crate::core::types::StateLock,
machine: &Machine,
resources: &indexmap::IndexMap<String, Resource>,
opts: DriftOptions,
census: &mut DriftCensus,
) -> Vec<DriftFinding> {
let mut findings = Vec::new();
for (id, rl) in &lock.resources {
if rl.resource_type != ResourceType::Task {
continue;
}
let Some(resource) = resources.get(id).filter(|r| owns(r)) else {
continue;
};
if let Some(reason) = skip_reason(rl, id, resources, opts) {
census.skipped(id, &rl.resource_type, reason);
continue;
}
census.inspected(id, &rl.resource_type);
if let Some(f) = check_task_drift(id, resource, machine) {
findings.push(f);
}
}
findings
}
fn skip_reason(
rl: &crate::core::types::ResourceLock,
id: &str,
resources: &indexmap::IndexMap<String, Resource>,
opts: DriftOptions,
) -> Option<SkipReason> {
if rl.status != ResourceStatus::Converged && rl.status != ResourceStatus::Drifted {
return Some(SkipReason::NotConverged);
}
if should_ignore_drift(id, resources) {
return Some(SkipReason::IgnoreDrift);
}
if !opts.run_task_checks {
return Some(SkipReason::TaskChecksDisabled);
}
None
}
pub(super) fn check_task_drift(
resource_id: &str,
resource: &Resource,
machine: &Machine,
) -> Option<DriftFinding> {
let script = match crate::core::codegen::check_script(resource) {
Ok(s) => s,
Err(e) => {
return Some(finding(
resource_id,
"ERROR",
format!("codegen failed: {e}"),
))
}
};
match crate::transport::exec_script_timeout(machine, &script, Some(DRIFT_QUERY_TIMEOUT_SECS)) {
Ok(out) if out.success() => None,
Ok(out) => Some(finding(
resource_id,
"completion_check: FAIL",
format!(
"completion_check fails on {}: {}",
machine.hostname,
marker(&out)
),
)),
Err(e) => Some(finding(
resource_id,
"ERROR",
format!("transport error: {e}"),
)),
}
}
fn marker(out: &crate::transport::ExecOutput) -> String {
let last = out
.stdout
.lines()
.rev()
.find(|l| !l.trim().is_empty())
.or_else(|| out.stderr.lines().find(|l| !l.trim().is_empty()))
.unwrap_or("no output")
.trim();
last.chars().take(200).collect()
}
fn finding(resource_id: &str, actual: &str, detail: String) -> DriftFinding {
DriftFinding {
resource_id: resource_id.to_string(),
resource_type: ResourceType::Task,
expected_hash: "completion_check: pass".to_string(),
actual_hash: actual.to_string(),
detail,
}
}