use pforge_runtime::Handler;
use std::path::{Path, PathBuf};
use crate::core::types::StateLock;
use crate::core::{parser, resolver, state, unattended};
use crate::tripwire::drift;
use super::types::*;
pub struct DriftHandler;
const UNATTENDED: drift::DriftOptions = drift::DriftOptions {
run_task_checks: false,
};
#[derive(Default)]
struct Scan {
findings: Vec<DriftFindingOutput>,
unchecked: Vec<String>,
census: Vec<serde_json::Value>,
declined: Vec<String>,
inspected: usize,
skipped: usize,
}
impl Scan {
fn absorb(&mut self, machine: &str, report: drift::DriftReport) {
let drift::DriftReport { findings, census } = report;
for f in &findings {
self.findings.push(DriftFindingOutput {
resource: f.resource_id.clone(),
expected_hash: f.expected_hash.clone(),
actual_hash: f.actual_hash.clone(),
detail: f.detail.clone(),
});
}
for id in census.skipped_ids(drift::SkipReason::TaskChecksDisabled) {
self.declined.push(format!(
"{id}: completion_check not executed on {machine}; its assertion \
was not evaluated"
));
}
self.inspected += census.inspected_total();
self.skipped += census.skipped_total();
let mut value = census.to_json();
if let Some(obj) = value.as_object_mut() {
obj.insert("machine".to_string(), serde_json::json!(machine));
}
self.census.push(value);
}
fn finish(self, mut unattended_skipped: Vec<String>) -> DriftOutput {
unattended_skipped.extend(self.declined);
DriftOutput {
drifted: !self.findings.is_empty(),
findings: self.findings,
unchecked: self.unchecked,
census: self.census,
resources_inspected: self.inspected,
resources_skipped: self.skipped,
unattended_skipped,
}
}
}
fn machine_lock(
state_dir: &Path,
machine_name: &str,
unchecked: &mut Vec<String>,
) -> Result<Option<StateLock>, pforge_runtime::Error> {
match state::load_lock(state_dir, machine_name) {
Ok(Some(lock)) => Ok(Some(lock)),
Ok(None) => {
unchecked.push(format!("{machine_name}: no state recorded (never applied)"));
Ok(None)
}
Err(e) => Err(pforge_runtime::Error::Handler(format!(
"cannot read state for machine '{machine_name}' in {}: {e}",
state_dir.display()
))),
}
}
#[async_trait::async_trait]
impl Handler for DriftHandler {
type Input = DriftInput;
type Output = DriftOutput;
type Error = pforge_runtime::Error;
async fn handle(&self, input: Self::Input) -> pforge_runtime::Result<Self::Output> {
let path = PathBuf::from(&input.path);
let state_dir = super::paths::resolve_state_dir(&path, input.state_dir.as_deref());
let parsed = parser::parse_and_validate(&path).map_err(pforge_runtime::Error::Handler)?;
let (config, unattended_skipped) = unattended::sanitize_config(&parsed);
let resolved = resolver::resolve_all(
&config.resources,
&config.params,
&config.machines,
&config.secrets,
);
if let Some(wanted) = input.machine.as_deref() {
if !config.machines.contains_key(wanted) {
let known: Vec<&str> = config.machines.keys().map(String::as_str).collect();
return Err(pforge_runtime::Error::Handler(format!(
"unknown machine `{wanted}`; the config declares: {}",
known.join(", ")
)));
}
}
let mut scan = Scan::default();
for (machine_name, machine) in &config.machines {
if input.machine.as_ref().is_some_and(|m| m != machine_name) {
continue;
}
let Some(lock) = machine_lock(&state_dir, machine_name, &mut scan.unchecked)? else {
continue;
};
let report = drift::detect_drift_full_reported(&lock, machine, &resolved, UNATTENDED);
scan.absorb(machine_name, report);
}
Ok(scan.finish(unattended_skipped))
}
}