use super::*;
use boatramp_core::project::ProjectRef;
#[cfg(feature = "handlers")]
const MAX_WORKFLOW_INPUT_BYTES: usize = 1024 * 1024;
#[cfg(feature = "handlers")]
#[derive(serde::Deserialize)]
pub(super) struct WorkflowBody {
steps: Vec<boatramp_core::workflow::Step>,
}
#[cfg(feature = "handlers")]
pub(super) async fn define_workflow(
State(deploy): State<DeployStore>,
Extension(project): axum::extract::Extension<ProjectContext>,
Path(name): Path<String>,
Json(body): Json<WorkflowBody>,
) -> Response {
if let Some(resp) = reject_invalid_name("workflow", &name) {
return resp;
}
let workflow = boatramp_core::workflow::Workflow {
name: name.clone(),
steps: body.steps,
};
if let Err(reason) = workflow.validate() {
return (StatusCode::BAD_REQUEST, format!("{reason}\n")).into_response();
}
if let Err(err) = deploy.put_workflow(project.as_ref(), &workflow).await {
return deploy_error_response(err);
}
Json(workflow).into_response()
}
#[cfg(feature = "handlers")]
pub(super) async fn list_workflows_handler(
State(deploy): State<DeployStore>,
Extension(project): axum::extract::Extension<ProjectContext>,
) -> Response {
match deploy.list_workflows(project.as_ref()).await {
Ok(mut list) => {
list.sort_by(|a, b| a.name.cmp(&b.name));
Json(list).into_response()
}
Err(err) => deploy_error_response(err),
}
}
#[cfg(feature = "handlers")]
pub(super) async fn get_workflow_handler(
State(deploy): State<DeployStore>,
Extension(project): axum::extract::Extension<ProjectContext>,
Path(name): Path<String>,
) -> Response {
match deploy.get_workflow(project.as_ref(), &name).await {
Ok(Some(w)) => Json(w).into_response(),
Ok(None) => (StatusCode::NOT_FOUND, format!("no workflow {name:?}\n")).into_response(),
Err(err) => deploy_error_response(err),
}
}
#[cfg(feature = "handlers")]
pub(super) async fn delete_workflow_handler(
State(deploy): State<DeployStore>,
Extension(project): axum::extract::Extension<ProjectContext>,
Path(name): Path<String>,
) -> Response {
match deploy.delete_workflow(project.as_ref(), &name).await {
Ok(_) => StatusCode::NO_CONTENT.into_response(),
Err(err) => deploy_error_response(err),
}
}
#[cfg(feature = "handlers")]
pub(super) async fn start_workflow_run(
State(deploy): State<DeployStore>,
Extension(project): axum::extract::Extension<ProjectContext>,
Path(name): Path<String>,
request: Request,
) -> Response {
let workflow = match deploy.get_workflow(project.as_ref(), &name).await {
Ok(Some(w)) => w,
Ok(None) => {
return (StatusCode::NOT_FOUND, format!("no workflow {name:?}\n")).into_response()
}
Err(err) => return deploy_error_response(err),
};
let body = match axum::body::to_bytes(request.into_body(), MAX_WORKFLOW_INPUT_BYTES).await {
Ok(bytes) => bytes,
Err(_) => {
return (
StatusCode::PAYLOAD_TOO_LARGE,
"workflow input exceeds the cap\n",
)
.into_response()
}
};
let now = now_unix();
let id = new_invocation_id();
let input_b64 = (!body.is_empty()).then(|| b64_encode(&body));
let run = boatramp_core::workflow::WorkflowRun::start(&workflow, id, input_b64, now);
if let Err(err) = deploy.put_workflow_run(project.as_ref(), &run).await {
return deploy_error_response(err);
}
(StatusCode::ACCEPTED, Json(run)).into_response()
}
#[cfg(feature = "handlers")]
pub(super) async fn get_workflow_run_handler(
State(deploy): State<DeployStore>,
Extension(project): axum::extract::Extension<ProjectContext>,
Path((name, id)): Path<(String, String)>,
) -> Response {
match deploy.get_workflow_run(project.as_ref(), &name, &id).await {
Ok(Some(run)) => Json(run).into_response(),
Ok(None) => (
StatusCode::NOT_FOUND,
format!("no run {id:?} for workflow {name:?}\n"),
)
.into_response(),
Err(err) => deploy_error_response(err),
}
}
#[cfg(feature = "handlers")]
pub(super) async fn drain_workflow_runs(
inner: &HandlerRuntimeInner,
deploy: &DeployStore,
project: ProjectRef<'_>,
workflow: &boatramp_core::workflow::Workflow,
) {
let runs = match deploy.list_workflow_runs(project, &workflow.name).await {
Ok(runs) => runs,
Err(err) => {
tracing::warn!(workflow = %workflow.name, %err, "listing workflow runs failed");
return;
}
};
for run in runs {
if run.is_terminal() {
continue;
}
advance_workflow_run(inner, deploy, project, workflow, run).await;
}
}
#[cfg(feature = "handlers")]
async fn advance_workflow_run(
inner: &HandlerRuntimeInner,
deploy: &DeployStore,
project: ProjectRef<'_>,
workflow: &boatramp_core::workflow::Workflow,
mut run: boatramp_core::workflow::WorkflowRun,
) {
use boatramp_core::workflow::StepStatus;
let now = now_unix();
for step_id in run.ready_steps(workflow) {
let Some(step) = workflow.step(&step_id).cloned() else {
continue;
};
let input = build_step_input(&run, &step);
let outcome = run_workflow_step(inner, deploy, project, &step, input).await;
let Some(sr) = run.steps.get_mut(&step_id) else {
continue;
};
sr.attempts = sr.attempts.saturating_add(1);
sr.updated = now;
match outcome {
Some(output) => {
sr.status = StepStatus::Succeeded;
sr.output_b64 = Some(b64_encode(&output));
run.completed_order.push(step_id.clone());
}
None if sr.attempts >= step.retry.max_attempts => sr.status = StepStatus::Failed,
None => sr.status = StepStatus::Pending, }
}
let any_failed = run.steps.values().any(|r| r.status == StepStatus::Failed);
if any_failed {
compensate_run(inner, deploy, project, workflow, &mut run).await;
run.status = boatramp_core::workflow::WorkflowStatus::Failed;
} else if run.all_succeeded() {
run.status = boatramp_core::workflow::WorkflowStatus::Succeeded;
}
run.updated = now_unix();
let _ = deploy.put_workflow_run(project, &run).await;
}
#[cfg(feature = "handlers")]
async fn run_workflow_step(
inner: &HandlerRuntimeInner,
deploy: &DeployStore,
project: ProjectRef<'_>,
step: &boatramp_core::workflow::Step,
input: Vec<u8>,
) -> Option<Vec<u8>> {
let function = match deploy.get_function(project, &step.function).await {
Ok(Some(f)) => f,
_ => return None,
};
let component = function.resolve(&function.active).map(str::to_owned)?;
let request = build_step_request(input);
let (response, _duration) = execute_function(
inner,
deploy,
project,
&function,
&component,
request,
0,
boatramp_handlers::Lane::Async,
)
.await;
let (status, _content_type, body) = capture_response(response).await;
let delivered = status != StatusCode::INTERNAL_SERVER_ERROR
&& status != StatusCode::GATEWAY_TIMEOUT
&& status != StatusCode::SERVICE_UNAVAILABLE;
delivered.then_some(body)
}
#[cfg(feature = "handlers")]
async fn compensate_run(
inner: &HandlerRuntimeInner,
deploy: &DeployStore,
project: ProjectRef<'_>,
workflow: &boatramp_core::workflow::Workflow,
run: &mut boatramp_core::workflow::WorkflowRun,
) {
let completed = run.completed_order.clone();
for step_id in completed.iter().rev() {
let Some(step) = workflow.step(step_id) else {
continue;
};
if let Some(compensate_fn) = &step.compensate {
if let Ok(Some(function)) = deploy.get_function(project, compensate_fn).await {
if let Some(component) = function.resolve(&function.active).map(str::to_owned) {
let request = build_step_request(Vec::new());
let _ = execute_function(
inner,
deploy,
project,
&function,
&component,
request,
0,
boatramp_handlers::Lane::Async,
)
.await;
}
}
}
if let Some(sr) = run.steps.get_mut(step_id) {
sr.status = boatramp_core::workflow::StepStatus::Compensated;
sr.updated = now_unix();
}
}
}
#[cfg(feature = "handlers")]
fn build_step_input(
run: &boatramp_core::workflow::WorkflowRun,
step: &boatramp_core::workflow::Step,
) -> Vec<u8> {
if step.depends_on.is_empty() {
return run.input_b64.as_deref().map(b64_decode).unwrap_or_default();
}
let mut map = serde_json::Map::new();
for dep in &step.depends_on {
let out = run
.steps
.get(dep)
.and_then(|r| r.output_b64.as_deref())
.map(b64_decode)
.unwrap_or_default();
map.insert(
dep.clone(),
serde_json::Value::String(String::from_utf8_lossy(&out).into_owned()),
);
}
serde_json::to_vec(&serde_json::Value::Object(map)).unwrap_or_default()
}
#[cfg(feature = "handlers")]
fn build_step_request(input: Vec<u8>) -> Request {
axum::http::Request::builder()
.method(axum::http::Method::POST)
.uri("/")
.header(header::CONTENT_TYPE, "application/json")
.body(axum::body::Body::from(input))
.unwrap_or_else(|_| Request::new(axum::body::Body::empty()))
}