use crate::contract::{ActionContract, WorkerContract};
use super::error::AwlScaffoldError;
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct Connection {
pub node: Option<String>,
pub actions: Vec<String>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ConnectionPlan {
pub task_queue: String,
pub connections: Vec<Connection>,
pub servable: Vec<ServableAction>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ServableAction {
pub name: String,
pub node: Option<String>,
pub input_schema: serde_json::Value,
pub output_schema: serde_json::Value,
}
pub fn plan(contract: &WorkerContract) -> Result<ConnectionPlan, AwlScaffoldError> {
let servable = servable_actions(contract)?;
let plan = ConnectionPlan {
task_queue: contract.task_queue.clone(),
connections: connections(&servable),
servable,
};
Ok(plan)
}
fn servable_actions(contract: &WorkerContract) -> Result<Vec<ServableAction>, AwlScaffoldError> {
let servable = contract
.actions
.iter()
.filter(|action| action.worker_owed())
.map(|action| servable_action(&contract.task_queue, action))
.collect::<Result<Vec<_>, _>>()?;
if servable.is_empty() {
return Err(AwlScaffoldError::NoServableAction {
task_queue: contract.task_queue.clone(),
});
}
Ok(servable)
}
fn servable_action(
task_queue: &str,
action: &ActionContract,
) -> Result<ServableAction, AwlScaffoldError> {
validate_action_name(task_queue, &action.name)?;
Ok(ServableAction {
name: action.name.clone(),
node: action.node.clone(),
input_schema: action.input_schema.clone(),
output_schema: action.output_schema.clone(),
})
}
fn validate_action_name(task_queue: &str, name: &str) -> Result<(), AwlScaffoldError> {
if matches!(name, "self" | "crate" | "super" | "Self") {
return Err(AwlScaffoldError::ActionNameUnnameable {
task_queue: task_queue.to_owned(),
action: name.to_owned(),
});
}
let mut characters = name.chars();
let starts = characters
.next()
.is_some_and(|first| first.is_ascii_alphabetic() || first == '_');
let continues = characters.all(|rest| rest.is_ascii_alphanumeric() || rest == '_');
if starts && continues {
return Ok(());
}
Err(AwlScaffoldError::ActionNameNotAnIdentifier {
task_queue: task_queue.to_owned(),
action: name.to_owned(),
})
}
fn connections(servable: &[ServableAction]) -> Vec<Connection> {
let mut nodes: Vec<&str> = Vec::new();
for action in servable {
if let Some(node) = action.node.as_deref()
&& !nodes.contains(&node)
{
nodes.push(node);
}
}
if nodes.is_empty() {
return vec![Connection {
node: None,
actions: servable.iter().map(|action| action.name.clone()).collect(),
}];
}
nodes
.into_iter()
.map(|node| Connection {
node: Some(node.to_owned()),
actions: servable
.iter()
.filter(|action| dispatch_can_reach(action.node.as_deref(), Some(node)))
.map(|action| action.name.clone())
.collect(),
})
.collect()
}
fn dispatch_can_reach(action_node: Option<&str>, worker_node: Option<&str>) -> bool {
match action_node {
None => true,
Some(pin) => worker_node == Some(pin),
}
}