use std::path::Path;
use aion_awl::{CompiledHarness, CompiledHarnessKind, CwdSource};
use aion_core::{ClusterEvent, PutOutcome};
use aion_package::AwlSource;
use aion_store::{DesiredState, WorkerDeployment, WorkerDeploymentListing};
use crate::ServerState;
use crate::worker::supervisor::{Convergence, converge_and_report};
use super::documents;
use super::listener::worker_dial_address;
use super::narrate;
use super::outcome::{AutoWorkerDecision, AutoWorkerOutcome};
use super::queues::{self, HarnessQueue};
use super::record;
use super::retire;
pub(super) const OPERATION: &str = "worker.auto_provision";
pub async fn provision(
state: &ServerState,
root: &Path,
awl: &AwlSource,
workflow_type: &str,
) -> Vec<AutoWorkerOutcome> {
let queues = match queues::harness_queues(awl.document()) {
Ok(queues) => queues,
Err(error) => return finish(state, vec![document_failure(&error, workflow_type)]),
};
let listing = match state
.worker_deployment_store()
.list_worker_deployments()
.await
{
Ok(listing) => listing,
Err(error) => {
return finish(
state,
vec![AutoWorkerOutcome::document_level(
workflow_type,
format!(
"the durable worker deployments could not be listed, so this server \
cannot tell what already serves the queues `{workflow_type}` declares, \
and provisioned nothing: {error}"
),
)],
);
}
};
let mut outcomes = Vec::with_capacity(queues.len().saturating_add(1));
for queue in &queues {
outcomes.push(provision_queue(state, root, awl, workflow_type, queue, &listing).await);
}
outcomes.extend(retire::withdraw_undeclared(state, workflow_type, &queues, &listing).await);
retire::prune_snapshots(state, root).await;
finish(state, outcomes)
}
fn finish(state: &ServerState, outcomes: Vec<AutoWorkerOutcome>) -> Vec<AutoWorkerOutcome> {
for outcome in &outcomes {
log_outcome(outcome);
}
state.worker_supervisor().record_auto_provision(&outcomes);
outcomes
}
fn document_failure(error: &queues::HarnessQueueError, workflow_type: &str) -> AutoWorkerOutcome {
match error {
queues::HarnessQueueError::Parse { .. } => {
AutoWorkerOutcome::document_level(workflow_type, error.to_string())
}
queues::HarnessQueueError::Harness { task_queue, .. }
| queues::HarnessQueueError::NoServiceableAction { task_queue, .. } => {
AutoWorkerOutcome::new(
task_queue.clone(),
workflow_type,
AutoWorkerDecision::Failed,
None,
error.to_string(),
)
}
}
}
async fn provision_queue(
state: &ServerState,
root: &Path,
awl: &AwlSource,
workflow_type: &str,
queue: &HarnessQueue,
listing: &WorkerDeploymentListing,
) -> AutoWorkerOutcome {
let task_queue = queue.task_queue.as_str();
let name = record::auto_name(task_queue);
let namespace = state.runtime_config().default_namespace.clone();
let node = state.cluster_self_node();
if let Some(skip) = operator_record_skip(listing, task_queue, &namespace, node, workflow_type) {
return skip;
}
if let Some(why) = harness_unrunnable_here(&queue.harness) {
return narrate::unrunnable_harness(task_queue, workflow_type, &why);
}
let dial = match worker_dial_address(&state.runtime_config().outbox) {
Ok(dial) => dial,
Err(dark) => return narrate::dark_outbox(task_queue, workflow_type, &dark),
};
let staged = match documents::stage(root, awl, workflow_type) {
Ok(staged) => staged,
Err(error) => {
return narrate::failure(
task_queue,
workflow_type,
format!(
"the deployed document could not be staged for a built-in agent worker on \
`{task_queue}`, so nothing was minted: {error}"
),
);
}
};
let binary = match crate::worker::capture_binary_identity() {
Ok(binary) => binary,
Err(error) => {
return narrate::failure(
task_queue,
workflow_type,
format!(
"the running server executable could not be identified, so no built-in agent \
worker was minted for `{task_queue}`: {error}"
),
);
}
};
let current = listing
.deployments
.iter()
.find(|record| record.name == name);
let desired = current.map_or(DesiredState::Running, |record| record.desired);
let mut requested = record::new_deployment(task_queue, &staged.path, &dial, &namespace, binary);
requested.desired = desired;
if let Some(current) = current
&& current.artifact == requested.artifact
&& current.task_queue == requested.task_queue
&& current.namespaces == requested.namespaces
{
let converged = converge(state, &name, Convergence::Idempotent).await;
return narrate::unchanged(
&narrate::Subject {
task_queue,
workflow_type,
name: &name,
namespace: &namespace,
document: &staged.path,
},
desired,
converged,
);
}
mint(
state,
MintRequest {
subject: narrate::Subject {
task_queue,
workflow_type,
name: &name,
namespace: &namespace,
document: &staged.path,
},
dial: &dial,
},
requested,
)
.await
}
struct MintRequest<'a> {
subject: narrate::Subject<'a>,
dial: &'a str,
}
async fn mint(
state: &ServerState,
request: MintRequest<'_>,
requested: aion_store::NewWorkerDeployment,
) -> AutoWorkerOutcome {
let stopped = requested.desired == DesiredState::Stopped;
let record = match WorkerDeployment::new(requested, chrono::Utc::now()) {
Ok(record) => record,
Err(error) => {
return narrate::failure(
request.subject.task_queue,
request.subject.workflow_type,
format!(
"the auto worker-deployment record for `{}` is invalid: {error}",
request.subject.task_queue
),
);
}
};
let result = match state
.worker_deployment_store()
.put_worker_deployment(record)
.await
{
Ok(result) => result,
Err(error) => {
return narrate::failure(
request.subject.task_queue,
request.subject.workflow_type,
format!(
"the auto worker-deployment record for `{}` could not be persisted: {error}",
request.subject.task_queue
),
);
}
};
publish_put(state, &result);
let mode = match result.outcome {
PutOutcome::Created => Convergence::Idempotent,
PutOutcome::Replaced => Convergence::Replacing,
};
let converged = converge(state, request.subject.name, mode).await;
let decision = match (stopped, result.outcome, converged) {
(true, _, _) => AutoWorkerDecision::OperatorStopped,
(false, _, false) => AutoWorkerDecision::RecordedNotRunning,
(false, PutOutcome::Created, true) => AutoWorkerDecision::Minted,
(false, PutOutcome::Replaced, true) => AutoWorkerDecision::Reminted,
};
narrate::written(&request.subject, decision, request.dial)
}
fn operator_record_skip(
listing: &WorkerDeploymentListing,
task_queue: &str,
namespace: &str,
node: Option<&str>,
workflow_type: &str,
) -> Option<AutoWorkerOutcome> {
let operator = listing.deployments.iter().find(|record| {
record.task_queue == task_queue
&& !record::is_auto_name(&record.name)
&& record.namespaces.contains(namespace)
&& match record.node.as_deref() {
None => true,
Some(pinned) => node == Some(pinned),
}
})?;
Some(AutoWorkerOutcome::new(
task_queue,
workflow_type,
AutoWorkerDecision::OperatorRecord,
Some(operator.name.clone()),
format!(
"worker deployment `{}` already serves task queue `{task_queue}` in namespace \
`{namespace}` and was authored by an operator, so nothing was minted or replaced: an \
explicit record wins over the document's declaration. Delete or rename it if you want \
this server to provision the queue itself",
operator.name
),
))
}
async fn converge(state: &ServerState, name: &str, mode: Convergence) -> bool {
converge_and_report(state.worker_supervisor(), name, mode, OPERATION)
.await
.is_ok()
}
fn publish_put(state: &ServerState, result: &aion_store::WorkerDeploymentPutResult) {
let name = result.deployment.name.clone();
let outcome = result.outcome;
let desired_state = result.deployment.desired;
let binary_version = result.deployment.binary.version.clone();
let binary_content_hash = result.deployment.binary.content_hash.clone();
drop(
state
.cluster_publisher()
.emit(|meta| ClusterEvent::WorkerDeploymentPut {
meta,
name,
outcome,
desired_state,
binary_version,
binary_content_hash,
}),
);
}
pub(super) fn log_outcome(outcome: &AutoWorkerOutcome) {
let task_queue = outcome.task_queue.as_str();
let workflow_type = outcome.workflow_type.as_str();
let decision = outcome.decision.token();
let detail = outcome.detail.as_str();
if outcome.decision.is_refusal() {
tracing::error!(
operation = OPERATION,
workflow_type,
task_queue,
decision,
"{detail}"
);
} else {
tracing::info!(
operation = OPERATION,
workflow_type,
task_queue,
decision,
"{detail}"
);
}
}
#[cfg(test)]
#[path = "provision_tests.rs"]
mod tests;
fn harness_unrunnable_here(harness: &CompiledHarness) -> Option<String> {
let (command, cwd) = match &harness.kind {
CompiledHarnessKind::Acp { command, cwd, .. } => (Some(command.clone()), cwd),
CompiledHarnessKind::Norn { binary, cwd } => (binary.clone(), cwd),
};
if let Some(command) = command {
if let Some(why) = executable_missing(&command) {
return Some(why);
}
} else if let Some(why) = resolves_nowhere_on_path("norn") {
return Some(why);
}
if let CwdSource::Literal(path) = cwd
&& path.is_absolute()
&& !path.is_dir()
{
return Some(format!(
"the harness working directory `{}` is not a directory on this host",
path.display()
));
}
None
}
fn executable_missing(command: &Path) -> Option<String> {
let Ok(metadata) = std::fs::metadata(command) else {
return Some(format!(
"the harness agent executable `{}` does not exist on this host",
command.display()
));
};
if !metadata.is_file() {
return Some(format!(
"the harness agent path `{}` is not a file on this host",
command.display()
));
}
#[cfg(unix)]
{
use std::os::unix::fs::PermissionsExt as _;
if metadata.permissions().mode() & 0o111 == 0 {
return Some(format!(
"the harness agent `{}` exists on this host but is not executable",
command.display()
));
}
}
None
}
fn resolves_nowhere_on_path(name: &str) -> Option<String> {
let path = std::env::var_os("PATH")?;
let resolves =
std::env::split_paths(&path).any(|dir| executable_missing(&dir.join(name)).is_none());
if resolves {
None
} else {
Some(format!(
"the harness names no binary and `{name}` resolves nowhere on this server's PATH"
))
}
}