use std::collections::BTreeSet;
use crate::ServerState;
use super::documents::sanitise_workflow_type;
use super::outcome::AutoWorkerOutcome;
use super::provision::{OPERATION, log_outcome};
use super::record;
use super::retire::{WithdrawCause, withdraw};
pub async fn deployed_workflow_types(state: &ServerState) -> Result<BTreeSet<String>, String> {
let engine = state
.engine()
.map_err(|error| format!("no engine to read the package catalogue through: {error}"))?;
let packages = engine
.store()
.list_packages()
.await
.map_err(|error| format!("the package catalogue could not be listed: {error}"))?;
Ok(packages
.iter()
.map(|package| sanitise_workflow_type(&package.workflow_type))
.collect())
}
pub async fn withdraw_orphaned(
state: &ServerState,
deployed: &BTreeSet<String>,
) -> Vec<AutoWorkerOutcome> {
let listing = match state
.worker_deployment_store()
.list_worker_deployments()
.await
{
Ok(listing) => listing,
Err(error) => {
tracing::error!(
operation = OPERATION,
%error,
"the worker deployments could not be listed, so this boot cannot tell whether an \
auto-provisioned record outlived its workflow"
);
return Vec::new();
}
};
for poisoned in &listing.undecodable {
tracing::warn!(
operation = OPERATION,
record = %poisoned.name,
error = %poisoned.error,
"a worker deployment record could not be decoded, so whether it outlived its \
workflow cannot be judged; it was left as it is"
);
}
let mut outcomes = Vec::new();
for deployment in &listing.deployments {
if !record::is_auto_name(&deployment.name) {
continue;
}
let Some(workflow_type) = record::staged_workflow_type(&deployment.artifact) else {
continue;
};
if deployed.contains(&workflow_type) {
continue;
}
let outcome = withdraw(state, &workflow_type, deployment, WithdrawCause::NoPackage).await;
log_outcome(&outcome);
outcomes.push(outcome);
}
if !outcomes.is_empty() {
state.worker_supervisor().record_auto_provision(&outcomes);
}
outcomes
}
#[cfg(test)]
#[path = "orphans_tests.rs"]
mod tests;