use crate::criterion;
use crate::expression::{self, Exchange, ExpressionError, Scope, StepState, WorkflowState};
use crate::http::{HttpRequest, HttpResponse};
use crate::operation::{self, Source};
use crate::report::{
CriterionOutcome, ExecutionError, ExecutionReport, Outcome, Performed, StepRecord,
};
use crate::select;
use crate::select::SelectError;
use roas_arazzo::v1_1::{
Criterion, CriterionKind, CriterionType, Description, FailureActionType, Parameter,
ParameterLocation, ReusableOr, Step, SuccessActionType, ValueOrSelector, Workflow,
};
use serde_json::{Map, Value};
use std::collections::{BTreeMap, BTreeSet};
use std::time::{Duration, Instant};
#[derive(Clone, Copy, Debug)]
struct Limits {
steps: usize,
depth: usize,
retries: u32,
}
impl Default for Limits {
fn default() -> Self {
Self {
steps: 1_000,
depth: 8,
retries: 10,
}
}
}
#[derive(Clone, Debug, Default)]
pub struct Options {
workflow: Option<String>,
inputs: Map<String, Value>,
sources: BTreeMap<String, Source>,
base_urls: BTreeMap<String, String>,
headers: Vec<(String, String)>,
limits: Limits,
}
impl Options {
#[must_use]
pub fn new() -> Self {
Self::default()
}
#[must_use]
pub fn workflow(mut self, workflow_id: impl Into<String>) -> Self {
self.workflow = Some(workflow_id.into());
self
}
#[must_use]
pub fn input(mut self, name: impl Into<String>, value: impl Into<Value>) -> Self {
self.inputs.insert(name.into(), value.into());
self
}
#[must_use]
pub fn inputs(mut self, inputs: Value) -> Self {
if let Value::Object(inputs) = inputs {
self.inputs = inputs;
}
self
}
#[must_use]
pub fn source(
mut self,
name: impl Into<String>,
url: impl Into<String>,
document: Value,
) -> Self {
self.sources.insert(
name.into(),
Source {
url: url.into(),
document,
},
);
self
}
#[must_use]
pub fn base_url(mut self, source_name: impl Into<String>, url: impl Into<String>) -> Self {
self.base_urls.insert(source_name.into(), url.into());
self
}
#[must_use]
pub fn header(mut self, name: impl Into<String>, value: impl Into<String>) -> Self {
self.headers.push((name.into(), value.into()));
self
}
#[must_use]
pub fn max_steps(mut self, steps: usize) -> Self {
self.limits.steps = steps;
self
}
#[must_use]
pub fn max_depth(mut self, depth: usize) -> Self {
self.limits.depth = depth;
self
}
#[must_use]
pub fn max_retries(mut self, retries: u32) -> Self {
self.limits.retries = retries;
self
}
}
#[derive(Debug)]
pub enum Progress {
Send(HttpRequest),
Wait(Duration),
Done(Box<ExecutionReport>),
}
fn scope<'s>(
frame: &'s Frame<'_>,
steps: &'s BTreeMap<String, StepState>,
here: Option<&'s Exchange>,
finished: &'s BTreeMap<String, WorkflowState>,
ambient: &'s Ambient,
) -> Scope<'s> {
Scope {
inputs: &frame.inputs,
outputs: &frame.outputs,
steps,
workflows: finished,
here,
sources: &ambient.sources,
components: &ambient.components,
self_: ambient.self_.as_deref(),
declared_steps: &frame.declared,
declared_workflows: &ambient.workflows,
}
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
enum Then {
Advance,
Retry,
EndCaller,
}
struct Frame<'d> {
workflow: &'d Workflow,
inputs: Value,
order: Vec<usize>,
at: usize,
steps: BTreeMap<String, StepState>,
outputs: BTreeMap<String, Value>,
attempts: BTreeMap<String, u32>,
declared: BTreeSet<String>,
retries: BTreeMap<(String, usize), u32>,
caller: Option<(String, Then)>,
started: Instant,
detour: Option<usize>,
outcome: Outcome,
}
struct Completion {
exchange: Option<Exchange>,
given: BTreeMap<String, Value>,
default_pass: bool,
attempt: u32,
performed: Performed,
elapsed: Duration,
}
struct Pending {
step: usize,
attempt: u32,
exchange: Exchange,
started: Instant,
}
struct Ambient {
sources: Value,
components: Value,
self_: Option<String>,
workflows: BTreeSet<String>,
}
pub struct Run<'d> {
description: &'d Description,
options: &'d Options,
ambient: Ambient,
frames: Vec<Frame<'d>>,
queue: Vec<&'d Workflow>,
finished: BTreeMap<String, WorkflowState>,
pending: Option<Pending>,
wait: Option<Duration>,
records: Vec<StepRecord>,
taken: usize,
report: Option<Box<ExecutionReport>>,
}
impl<'d> Run<'d> {
pub fn start(
description: &'d Description,
options: &'d Options,
) -> Result<Self, ExecutionError> {
let wanted = match &options.workflow {
Some(id) => description
.workflows
.iter()
.find(|workflow| &workflow.workflow_id == id)
.ok_or_else(|| ExecutionError::UnknownWorkflow(id.clone()))?,
None => description
.workflows
.first()
.ok_or_else(|| ExecutionError::UnknownWorkflow(String::new()))?,
};
let sources = serde_json::to_value(
description
.source_descriptions
.iter()
.map(|source| (source.name.clone(), source))
.collect::<BTreeMap<_, _>>(),
)
.unwrap_or(Value::Null);
let components = description
.components
.as_ref()
.and_then(|components| serde_json::to_value(components).ok())
.unwrap_or(Value::Null);
let mut queue = ordered_workflows(description, wanted)?;
let first = queue.remove(0);
let mut run = Self {
description,
options,
ambient: Ambient {
sources,
components,
self_: description.self_.clone(),
workflows: description
.workflows
.iter()
.map(|workflow| workflow.workflow_id.clone())
.collect(),
},
frames: Vec::new(),
queue,
finished: BTreeMap::new(),
pending: None,
wait: None,
records: Vec::new(),
taken: 0,
report: None,
};
run.enter(first, Value::Object(options.inputs.clone()), None)?;
Ok(run)
}
pub fn advance(&mut self) -> Result<Progress, ExecutionError> {
if let Some(wait) = self.wait.take() {
return Ok(Progress::Wait(wait));
}
if let Some(pending) = &self.pending {
return Err(ExecutionError::Awaiting {
method: pending.exchange.request.method.clone(),
url: pending.exchange.request.url.clone(),
});
}
loop {
if let Some(report) = self.report.take() {
return Ok(Progress::Done(report));
}
let Some(frame) = self.frames.last() else {
return Ok(Progress::Done(Box::new(self.finish())));
};
if frame.at >= frame.order.len() {
self.leave()?;
continue;
}
self.taken += 1;
if self.taken > self.options.limits.steps {
return Err(ExecutionError::Limit {
limit: "step",
at: self.options.limits.steps,
});
}
let index = frame.order[frame.at];
let step = &frame.workflow.steps[index];
if step.workflow_id.is_some() {
self.call(index)?;
continue;
}
let (request, exchange) = self.build(index)?;
let attempt = self
.frames
.last()
.and_then(|frame| frame.attempts.get(&step.step_id).copied())
.unwrap_or(0)
+ 1;
self.pending = Some(Pending {
step: index,
attempt,
exchange,
started: Instant::now(),
});
return Ok(Progress::Send(request));
}
}
pub fn supply(&mut self, response: HttpResponse) -> Result<(), ExecutionError> {
let Some(mut pending) = self.pending.take() else {
return Err(ExecutionError::NotWaiting);
};
let elapsed = pending.started.elapsed();
let status = response.status;
pending.exchange.response_body = response.body_as_json();
pending.exchange.response = Some(response);
if self.frames.last().is_none() {
return Err(ExecutionError::NotWaiting);
}
let performed = Performed::Request {
method: pending.exchange.request.method.clone(),
url: pending.exchange.request.url.clone(),
status,
};
self.complete(
pending.step,
Completion {
exchange: Some(pending.exchange),
given: BTreeMap::new(),
default_pass: (200..400).contains(&status),
attempt: pending.attempt,
performed,
elapsed,
},
)
}
fn complete(&mut self, index: usize, done: Completion) -> Result<(), ExecutionError> {
let frame = self.frames.last().expect("a frame to complete in");
let step = &frame.workflow.steps[index];
let step_id = step.step_id.clone();
let workflow_id = frame.workflow.workflow_id.clone();
let mut state = StepState {
exchange: done.exchange.clone(),
outputs: done.given.clone(),
passed: true,
};
let (passed, criteria, outputs) = {
let mut ahead = frame.steps.clone();
ahead.insert(step_id.clone(), state.clone());
let scope = scope(
frame,
&ahead,
done.exchange.as_ref(),
&self.finished,
&self.ambient,
);
let mut criteria = Vec::with_capacity(step.success_criteria.len());
let mut passed = if step.success_criteria.is_empty() {
done.default_pass
} else {
true
};
for criterion in &step.success_criteria {
let holds = criterion::passes(criterion, &scope)?;
criteria.push(CriterionOutcome {
condition: criterion.condition.clone(),
passed: holds,
});
passed = passed && holds;
}
let outputs = if passed {
let mut outputs = done.given.clone();
outputs.extend(evaluate_outputs(&step.outputs, &scope)?);
outputs
} else {
BTreeMap::new()
};
(passed, criteria, outputs)
};
state.outputs = outputs.clone();
state.passed = passed;
let frame = self.frames.last_mut().expect("the frame is still there");
frame.steps.insert(step_id.clone(), state);
let action = self.decide(index, passed, done.exchange.as_ref())?;
let described = describe(&action);
self.records.push(StepRecord {
workflow_id,
step_id,
attempt: done.attempt,
performed: done.performed,
criteria,
passed,
outputs,
action: described,
elapsed: done.elapsed,
});
self.apply(action)
}
fn enter(
&mut self,
workflow: &'d Workflow,
inputs: Value,
caller: Option<(String, Then)>,
) -> Result<(), ExecutionError> {
if self.frames.len() >= self.options.limits.depth {
return Err(ExecutionError::Limit {
limit: "workflow depth",
at: self.options.limits.depth,
});
}
self.frames.push(Frame {
workflow,
inputs,
order: ordered_steps(workflow, self.description)?,
at: 0,
steps: BTreeMap::new(),
outputs: BTreeMap::new(),
declared: workflow
.steps
.iter()
.map(|step| step.step_id.clone())
.collect(),
attempts: BTreeMap::new(),
retries: BTreeMap::new(),
caller,
started: Instant::now(),
detour: None,
outcome: Outcome::Succeeded,
});
Ok(())
}
fn leave(&mut self) -> Result<(), ExecutionError> {
let frame = self.frames.pop().expect("a frame to leave");
let outputs = {
let scope = scope(&frame, &frame.steps, None, &self.finished, &self.ambient);
if frame.outcome == Outcome::Succeeded {
evaluate_outputs(&frame.workflow.outputs, &scope)?
} else {
evaluate_what_ran(&frame.workflow.outputs, &scope)?
}
};
self.finished.insert(
frame.workflow.workflow_id.clone(),
WorkflowState {
inputs: frame.inputs.clone(),
outputs: outputs.clone(),
},
);
let Some((step_id, then)) = frame.caller else {
if self.queue.is_empty() {
self.report = Some(Box::new(ExecutionReport {
workflow_id: frame.workflow.workflow_id.clone(),
outcome: frame.outcome,
outputs,
steps: std::mem::take(&mut self.records),
}));
} else {
let next = self.queue.remove(0);
let inputs = Value::Object(self.options.inputs.clone());
self.enter(next, inputs, None)?;
}
return Ok(());
};
let Some(parent) = self.frames.last() else {
return Ok(());
};
let index = parent.order[parent.at];
debug_assert_eq!(parent.workflow.steps[index].step_id, step_id);
match then {
Then::EndCaller => {
let parent = self.frames.last_mut().expect("the parent is still there");
parent.steps.insert(
step_id,
StepState {
exchange: None,
outputs,
passed: frame.outcome != Outcome::Failed,
},
);
parent.at = parent.order.len();
parent.outcome = frame.outcome;
Ok(())
}
Then::Retry => Ok(()),
Then::Advance => {
let elapsed = frame.started.elapsed();
let timed_out = parent.workflow.steps[index]
.timeout
.and_then(|timeout| u64::try_from(timeout).ok())
.is_some_and(|timeout| elapsed > Duration::from_millis(timeout));
let attempt = self
.frames
.last()
.and_then(|parent| parent.attempts.get(&step_id).copied())
.unwrap_or(0)
+ 1;
self.complete(
index,
Completion {
exchange: None,
given: outputs,
default_pass: frame.outcome != Outcome::Failed && !timed_out,
attempt,
performed: Performed::Workflow {
workflow_id: frame.workflow.workflow_id.clone(),
outcome: frame.outcome,
},
elapsed,
},
)
}
}
}
fn call(&mut self, index: usize) -> Result<(), ExecutionError> {
let frame = self.frames.last().expect("a frame to call from");
let step = &frame.workflow.steps[index];
let step_id = step.step_id.clone();
let wanted = step.workflow_id.clone().unwrap_or_default();
if wanted.starts_with('$') {
return Err(ExecutionError::Unsupported(format!(
"step `{step_id}` calls `{wanted}`, and this executor runs only workflows of the description it was given"
)));
}
let workflow = self
.description
.workflows
.iter()
.find(|workflow| workflow.workflow_id == wanted)
.ok_or_else(|| ExecutionError::UnknownWorkflow(wanted.clone()))?;
let mut arguments = frame.workflow.parameters.clone();
arguments.extend(step.parameters.iter().cloned());
let inputs = self.arguments(&arguments)?;
self.enter(workflow, inputs, Some((step_id, Then::Advance)))
}
fn arguments(&self, arguments: &[ReusableOr<Parameter>]) -> Result<Value, ExecutionError> {
let frame = self.frames.last().expect("a frame to pass arguments from");
let scope = scope(frame, &frame.steps, None, &self.finished, &self.ambient);
let mut inputs = Map::new();
for parameter in parameters(arguments, self.description, &scope)? {
inputs.insert(parameter.name, parameter.value);
}
Ok(Value::Object(inputs))
}
fn build(&self, index: usize) -> Result<(HttpRequest, Exchange), ExecutionError> {
let frame = self.frames.last().expect("a frame to build in");
let step = &frame.workflow.steps[index];
let endpoint = operation::resolve(step, &self.options.sources, &self.options.base_urls)?;
let scope = scope(frame, &frame.steps, None, &self.finished, &self.ambient);
let mut resolved = parameters(&frame.workflow.parameters, self.description, &scope)?;
for parameter in parameters(&step.parameters, self.description, &scope)? {
resolved.retain(|existing| {
!(existing.name == parameter.name && existing.location == parameter.location)
});
resolved.push(parameter);
}
let mut path = BTreeMap::new();
let mut query = BTreeMap::new();
let mut headers: Vec<(String, String)> = Vec::new();
let mut cookies = Vec::new();
let mut querystring = None;
for parameter in resolved {
match parameter.location {
ParameterLocation::Path => {
path.insert(parameter.name, parameter.value);
}
ParameterLocation::Query => {
query.insert(parameter.name, parameter.value);
}
ParameterLocation::Querystring => {
querystring = Some(text(¶meter.value));
}
ParameterLocation::Header => {
headers.push((parameter.name, text(¶meter.value)));
}
ParameterLocation::Cookie => {
cookies.push(format!("{}={}", parameter.name, text(¶meter.value)));
}
ParameterLocation::Channel => {
return Err(ExecutionError::Unsupported(format!(
"step `{}` has a `channel` parameter, which belongs to an AsyncAPI step",
step.step_id
)));
}
}
}
if !cookies.is_empty() {
headers.push(("Cookie".to_owned(), cookies.join("; ")));
}
for (name, value) in &self.options.headers {
if !headers
.iter()
.any(|(existing, _)| existing.eq_ignore_ascii_case(name))
{
headers.push((name.clone(), value.clone()));
}
}
let body = body(step, &scope)?;
if let Some(body) = &body
&& !headers
.iter()
.any(|(name, _)| name.eq_ignore_ascii_case("content-type"))
{
headers.push(("Content-Type".to_owned(), body.content_type.clone()));
}
let url = url(&endpoint, &path, &query, querystring.as_deref()).map_err(|reason| {
ExecutionError::BadRequest {
step: step.step_id.clone(),
reason,
}
})?;
let request = HttpRequest {
method: endpoint.method,
url,
headers,
body: body.as_ref().map(|body| body.bytes.clone()),
timeout: step
.timeout
.and_then(|timeout| u64::try_from(timeout).ok())
.map(Duration::from_millis),
};
Ok((
request.clone(),
Exchange {
request,
path,
query,
body: body.map(|body| body.value),
response: None,
response_body: None,
},
))
}
fn decide(
&mut self,
index: usize,
passed: bool,
exchange: Option<&Exchange>,
) -> Result<Action, ExecutionError> {
let frame = self.frames.last().expect("a frame to decide in");
let step = &frame.workflow.steps[index];
let scope = scope(frame, &frame.steps, exchange, &self.finished, &self.ambient);
if passed {
let actions = step
.on_success
.iter()
.chain(frame.workflow.success_actions.iter());
for action in actions {
let action = success_action(action, self.description)?;
if !holds(&action.criteria, &scope)? {
continue;
}
return Ok(match action.type_ {
SuccessActionType::End => Action::End(Outcome::Ended),
SuccessActionType::Goto => Action::Goto {
step: action.step_id.clone(),
workflow: action.workflow_id.clone(),
parameters: action.parameters.clone(),
},
});
}
return Ok(Action::Advance);
}
let actions = step
.on_failure
.iter()
.chain(frame.workflow.failure_actions.iter());
for (at, action) in actions.enumerate() {
let action = failure_action(action, self.description)?;
if !holds(&action.criteria, &scope)? {
continue;
}
return Ok(match action.type_ {
FailureActionType::End => Action::End(Outcome::Failed),
FailureActionType::Goto => Action::Goto {
step: action.step_id.clone(),
workflow: action.workflow_id.clone(),
parameters: action.parameters.clone(),
},
FailureActionType::Retry => {
let allowed = action
.retry_limit
.map_or(1, |limit| u32::try_from(limit).unwrap_or(u32::MAX));
let spent = frame
.retries
.get(&(step.step_id.clone(), at))
.copied()
.unwrap_or(0);
let taken = frame.attempts.get(&step.step_id).copied().unwrap_or(0);
if spent >= allowed || taken >= self.options.limits.retries {
continue;
}
Action::Retry {
at,
after: action.retry_after,
step: action.step_id.clone(),
workflow: action.workflow_id.clone(),
parameters: action.parameters.clone(),
}
}
});
}
Ok(Action::End(Outcome::Failed))
}
fn apply(&mut self, action: Action) -> Result<(), ExecutionError> {
let frame = self.frames.last_mut().expect("a frame to act in");
match action {
Action::Advance => {
match frame.detour.take() {
Some(back) => frame.at = back,
None => frame.at += 1,
}
Ok(())
}
Action::End(outcome) => {
frame.outcome = outcome;
frame.at = frame.order.len();
Ok(())
}
Action::Retry {
at,
after,
step: target,
workflow,
parameters: arguments,
} => {
let index = frame.order[frame.at];
let step_id = frame.workflow.steps[index].step_id.clone();
*frame.attempts.entry(step_id.clone()).or_insert(0) += 1;
*frame.retries.entry((step_id.clone(), at)).or_insert(0) += 1;
if let Some(after) = after.filter(|after| *after > 0.0) {
self.wait = Some(Duration::from_secs_f64(after));
}
match (target, workflow) {
(Some(target), _) => {
let at = position_of(frame, &target)?;
frame.detour = Some(frame.at);
frame.at = at;
Ok(())
}
(None, Some(workflow_id)) => {
let workflow = self
.description
.workflows
.iter()
.find(|workflow| workflow.workflow_id == workflow_id)
.ok_or_else(|| ExecutionError::UnknownWorkflow(workflow_id.clone()))?;
let inputs = self.arguments(&arguments)?;
self.enter(workflow, inputs, Some((step_id, Then::Retry)))
}
(None, None) => Ok(()),
}
}
Action::Goto {
step: Some(step_id),
..
} => {
frame.at = position_of(frame, &step_id)?;
Ok(())
}
Action::Goto {
workflow: Some(workflow_id),
parameters: arguments,
..
} => {
let index = frame.order[frame.at];
let step_id = frame.workflow.steps[index].step_id.clone();
let workflow = self
.description
.workflows
.iter()
.find(|workflow| workflow.workflow_id == workflow_id)
.ok_or_else(|| ExecutionError::UnknownWorkflow(workflow_id.clone()))?;
let inputs = self.arguments(&arguments)?;
self.enter(workflow, inputs, Some((step_id, Then::EndCaller)))
}
Action::Goto { .. } => Ok(()),
}
}
fn finish(&mut self) -> ExecutionReport {
ExecutionReport {
workflow_id: self
.options
.workflow
.clone()
.or_else(|| {
self.description
.workflows
.first()
.map(|workflow| workflow.workflow_id.clone())
})
.unwrap_or_default(),
outcome: Outcome::Succeeded,
outputs: BTreeMap::new(),
steps: std::mem::take(&mut self.records),
}
}
}
fn position_of(frame: &Frame<'_>, step_id: &str) -> Result<usize, ExecutionError> {
let index = frame
.workflow
.steps
.iter()
.position(|step| step.step_id == step_id)
.ok_or_else(|| ExecutionError::UnknownStep {
workflow: frame.workflow.workflow_id.clone(),
step: step_id.to_owned(),
})?;
Ok(frame
.order
.iter()
.position(|&candidate| candidate == index)
.unwrap_or(frame.order.len()))
}
#[derive(Clone, Debug)]
enum Action {
Advance,
End(Outcome),
Retry {
at: usize,
after: Option<f64>,
step: Option<String>,
workflow: Option<String>,
parameters: Vec<ReusableOr<Parameter>>,
},
Goto {
step: Option<String>,
workflow: Option<String>,
parameters: Vec<ReusableOr<Parameter>>,
},
}
fn describe(action: &Action) -> Option<String> {
match action {
Action::Advance => None,
Action::End(Outcome::Failed) => Some("ended, failed".to_owned()),
Action::End(_) => Some("ended".to_owned()),
Action::Retry {
step: Some(step), ..
} => Some(format!("retry via step `{step}`")),
Action::Retry {
workflow: Some(workflow),
..
} => Some(format!("retry via workflow `{workflow}`")),
Action::Retry { .. } => Some("retry".to_owned()),
Action::Goto {
step: Some(step), ..
} => Some(format!("goto step `{step}`")),
Action::Goto {
workflow: Some(workflow),
..
} => Some(format!("goto workflow `{workflow}`")),
Action::Goto { .. } => None,
}
}
struct Resolved {
name: String,
location: ParameterLocation,
value: Value,
}
fn parameters(
list: &[ReusableOr<Parameter>],
description: &Description,
scope: &Scope<'_>,
) -> Result<Vec<Resolved>, ExecutionError> {
let mut resolved = Vec::with_capacity(list.len());
for entry in list {
let (parameter, overridden) = match entry {
ReusableOr::Item(parameter) => (parameter.clone(), None),
ReusableOr::Reusable(reusable) => {
let name = reusable
.reference
.strip_prefix("$components.parameters.")
.ok_or_else(|| {
ExecutionError::Unsupported(format!(
"`{}` is not a component this executor can follow",
reusable.reference
))
})?;
let parameter = description
.components
.as_ref()
.and_then(|components| components.parameters.get(name))
.ok_or_else(|| {
ExecutionError::Unsupported(format!(
"`{}` names a component the description has not got",
reusable.reference
))
})?;
(parameter.clone(), reusable.value.clone())
}
};
let value = match overridden {
Some(value) => select::resolve(&value, scope)?,
None => select::value_of(¶meter.value, scope)?,
};
resolved.push(Resolved {
name: parameter.name,
location: parameter.in_.unwrap_or(ParameterLocation::Query),
value,
});
}
Ok(resolved)
}
fn success_action(
entry: &ReusableOr<roas_arazzo::v1_1::SuccessAction>,
description: &Description,
) -> Result<roas_arazzo::v1_1::SuccessAction, ExecutionError> {
match entry {
ReusableOr::Item(action) => Ok(action.clone()),
ReusableOr::Reusable(reusable) => reusable
.reference
.strip_prefix("$components.successActions.")
.and_then(|name| {
description
.components
.as_ref()
.and_then(|components| components.success_actions.get(name))
})
.cloned()
.ok_or_else(|| {
ExecutionError::Unsupported(format!(
"`{}` names a component the description has not got",
reusable.reference
))
}),
}
}
fn failure_action(
entry: &ReusableOr<roas_arazzo::v1_1::FailureAction>,
description: &Description,
) -> Result<roas_arazzo::v1_1::FailureAction, ExecutionError> {
match entry {
ReusableOr::Item(action) => Ok(action.clone()),
ReusableOr::Reusable(reusable) => reusable
.reference
.strip_prefix("$components.failureActions.")
.and_then(|name| {
description
.components
.as_ref()
.and_then(|components| components.failure_actions.get(name))
})
.cloned()
.ok_or_else(|| {
ExecutionError::Unsupported(format!(
"`{}` names a component the description has not got",
reusable.reference
))
}),
}
}
fn holds(criteria: &[Criterion], scope: &Scope<'_>) -> Result<bool, ExecutionError> {
for criterion in criteria {
if !criterion::passes(criterion, scope)? {
return Ok(false);
}
}
Ok(true)
}
fn evaluate_what_ran(
outputs: &BTreeMap<String, ValueOrSelector>,
scope: &Scope<'_>,
) -> Result<BTreeMap<String, Value>, ExecutionError> {
let mut named = BTreeMap::new();
for (name, value) in outputs {
match select::value_of(value, scope) {
Ok(value) => {
named.insert(name.clone(), value);
}
Err(SelectError::Expression(ExpressionError::NotRun { .. })) => {}
Err(error) => return Err(error.into()),
}
}
Ok(named)
}
fn evaluate_outputs(
outputs: &BTreeMap<String, ValueOrSelector>,
scope: &Scope<'_>,
) -> Result<BTreeMap<String, Value>, ExecutionError> {
let mut resolved = BTreeMap::new();
for (name, value) in outputs {
resolved.insert(name.clone(), select::value_of(value, scope)?);
}
Ok(resolved)
}
struct Body {
bytes: Vec<u8>,
value: Value,
content_type: String,
}
fn body(step: &Step, scope: &Scope<'_>) -> Result<Option<Body>, ExecutionError> {
let Some(request_body) = &step.request_body else {
return Ok(None);
};
let mut payload = match &request_body.payload {
Some(payload) => select::resolve(payload, scope)?,
None => Value::Null,
};
for replacement in &request_body.replacements {
let value = select::value_of(&replacement.value, scope)?;
let language = match &replacement.target_selector_type {
Some(type_) => select::kind_of(type_)?,
None => select::Language::Pointer,
};
select::place(language, &replacement.target, &mut payload, value).map_err(|reason| {
ExecutionError::BadRequest {
step: step.step_id.clone(),
reason,
}
})?;
}
let content_type = request_body
.content_type
.clone()
.unwrap_or_else(|| "application/json".to_owned());
let bytes = match (&payload, content_type.contains("json")) {
(Value::String(text), false) => text.clone().into_bytes(),
_ => payload.to_string().into_bytes(),
};
Ok(Some(Body {
bytes,
value: payload,
content_type,
}))
}
fn url(
endpoint: &operation::Endpoint,
path: &BTreeMap<String, Value>,
query: &BTreeMap<String, Value>,
querystring: Option<&str>,
) -> Result<String, String> {
let mut filled = endpoint.path.clone();
for (name, value) in path {
filled = filled.replace(&format!("{{{name}}}"), &encode(&text(value)));
}
if let Some(start) = filled.find('{') {
return Err(format!(
"`{}` still has `{}` in it, which no parameter filled in",
endpoint.path,
&filled[start
..filled[start..]
.find('}')
.map_or(filled.len(), |end| start + end + 1)]
));
}
let mut url = format!("{}{filled}", endpoint.base);
let pairs: Vec<String> = query
.iter()
.map(|(name, value)| format!("{}={}", encode(name), encode(&text(value))))
.collect();
let query = match (pairs.is_empty(), querystring) {
(true, None) => String::new(),
(true, Some(raw)) => raw.to_owned(),
(false, None) => pairs.join("&"),
(false, Some(raw)) => format!("{}&{raw}", pairs.join("&")),
};
if !query.is_empty() {
url.push('?');
url.push_str(&query);
}
url::Url::parse(&url).map_err(|error| format!("`{url}` is not a URL: {error}"))?;
Ok(url)
}
fn text(value: &Value) -> String {
match value {
Value::String(text) => text.clone(),
other => other.to_string(),
}
}
fn encode(text: &str) -> String {
let mut encoded = String::with_capacity(text.len());
for byte in text.bytes() {
match byte {
b'A'..=b'Z' | b'a'..=b'z' | b'0'..=b'9' | b'-' | b'.' | b'_' | b'~' => {
encoded.push(byte as char);
}
other => encoded.push_str(&format!("%{other:02X}")),
}
}
encoded
}
fn ordered_workflows<'d>(
description: &'d Description,
wanted: &'d Workflow,
) -> Result<Vec<&'d Workflow>, ExecutionError> {
let mut ordered = Vec::new();
let mut visiting = BTreeSet::new();
let mut done = BTreeSet::new();
visit(description, wanted, &mut ordered, &mut visiting, &mut done)?;
Ok(ordered)
}
fn visit<'d>(
description: &'d Description,
workflow: &'d Workflow,
ordered: &mut Vec<&'d Workflow>,
visiting: &mut BTreeSet<String>,
done: &mut BTreeSet<String>,
) -> Result<(), ExecutionError> {
if done.contains(&workflow.workflow_id) {
return Ok(());
}
if !visiting.insert(workflow.workflow_id.clone()) {
return Err(ExecutionError::Circular(workflow.workflow_id.clone()));
}
for id in &workflow.depends_on {
let dependency = description
.workflows
.iter()
.find(|candidate| &candidate.workflow_id == id)
.ok_or_else(|| ExecutionError::UnknownWorkflow(id.clone()))?;
visit(description, dependency, ordered, visiting, done)?;
}
visiting.remove(&workflow.workflow_id);
done.insert(workflow.workflow_id.clone());
ordered.push(workflow);
Ok(())
}
fn steps_named_by(step: &Step, description: &Description) -> BTreeSet<String> {
let mut found = BTreeSet::new();
fn named_in(expression: &str) -> Option<String> {
let rest = expression.strip_prefix("$steps.")?;
let (id, _) = rest.split_once('.')?;
Some(id.to_owned())
}
fn read(text: &str, found: &mut BTreeSet<String>) {
found.extend(
expression::references(text)
.into_iter()
.filter_map(named_in),
);
}
fn read_condition(criterion: &Criterion, found: &mut BTreeSet<String>) {
let simple = matches!(
criterion.type_,
None | Some(CriterionType::Simple(CriterionKind::Simple))
);
if simple {
found.extend(
criterion::expressions_in(&criterion.condition)
.iter()
.filter_map(|expression| named_in(expression)),
);
} else {
found.extend(
expression::interpolations(&criterion.condition)
.into_iter()
.filter_map(named_in),
);
}
}
fn read_value(value: &ValueOrSelector, found: &mut BTreeSet<String>) {
match value {
ValueOrSelector::Literal(literal) => read_literal(literal, found),
ValueOrSelector::Selector(selector) => read(&selector.context, found),
}
}
fn read_literal(value: &Value, found: &mut BTreeSet<String>) {
match value {
Value::String(text) => read(text, found),
Value::Array(items) => items.iter().for_each(|item| read_literal(item, found)),
Value::Object(members) => members
.values()
.for_each(|member| read_literal(member, found)),
_ => {}
}
}
fn read_parameters(
list: &[ReusableOr<Parameter>],
description: &Description,
found: &mut BTreeSet<String>,
) {
for entry in list {
match entry {
ReusableOr::Item(parameter) => read_value(¶meter.value, found),
ReusableOr::Reusable(reusable) => {
if let Some(overridden) = &reusable.value {
read_literal(overridden, found);
} else if let Some(parameter) = reusable
.reference
.strip_prefix("$components.parameters.")
.and_then(|name| {
description
.components
.as_ref()
.and_then(|components| components.parameters.get(name))
})
{
read_value(¶meter.value, found);
}
}
}
}
}
fn read_criteria(list: &[Criterion], found: &mut BTreeSet<String>) {
for criterion in list {
if let Some(context) = &criterion.context {
read(context, found);
}
read_condition(criterion, found);
}
}
read_parameters(&step.parameters, description, &mut found);
read_criteria(&step.success_criteria, &mut found);
for output in step.outputs.values() {
read_value(output, &mut found);
}
if let Some(body) = &step.request_body {
if let Some(payload) = &body.payload {
read_literal(payload, &mut found);
}
for replacement in &body.replacements {
read_value(&replacement.value, &mut found);
}
}
for entry in &step.on_success {
if let Ok(action) = success_action(entry, description) {
read_criteria(&action.criteria, &mut found);
read_parameters(&action.parameters, description, &mut found);
}
}
for entry in &step.on_failure {
if let Ok(action) = failure_action(entry, description) {
read_criteria(&action.criteria, &mut found);
read_parameters(&action.parameters, description, &mut found);
}
}
found.remove(&step.step_id);
found
}
fn ordered_steps(
workflow: &Workflow,
description: &Description,
) -> Result<Vec<usize>, ExecutionError> {
let index: BTreeMap<&str, usize> = workflow
.steps
.iter()
.enumerate()
.map(|(at, step)| (step.step_id.as_str(), at))
.collect();
let mut ordered = Vec::with_capacity(workflow.steps.len());
let mut visiting = BTreeSet::new();
let mut done = BTreeSet::new();
for step in &workflow.steps {
visit_step(
workflow,
description,
&index,
step,
&mut ordered,
&mut visiting,
&mut done,
)?;
}
Ok(ordered)
}
fn visit_step(
workflow: &Workflow,
description: &Description,
index: &BTreeMap<&str, usize>,
step: &Step,
ordered: &mut Vec<usize>,
visiting: &mut BTreeSet<String>,
done: &mut BTreeSet<String>,
) -> Result<(), ExecutionError> {
if done.contains(&step.step_id) {
return Ok(());
}
if !visiting.insert(step.step_id.clone()) {
return Err(ExecutionError::Circular(step.step_id.clone()));
}
for id in &step.depends_on {
let at = index
.get(id.as_str())
.ok_or_else(|| ExecutionError::UnknownStep {
workflow: workflow.workflow_id.clone(),
step: id.clone(),
})?;
visit_step(
workflow,
description,
index,
&workflow.steps[*at],
ordered,
visiting,
done,
)?;
}
for id in steps_named_by(step, description) {
let Some(at) = index.get(id.as_str()) else {
continue;
};
visit_step(
workflow,
description,
index,
&workflow.steps[*at],
ordered,
visiting,
done,
)?;
}
visiting.remove(&step.step_id);
done.insert(step.step_id.clone());
ordered.push(index[step.step_id.as_str()]);
Ok(())
}