Packages
step_flow
1.8.1-rc3
1.9.0-rc2
1.9.0-rc1
1.9.0-rc0
1.8.2
1.8.1
1.8.1-rc8
1.8.1-rc7
1.8.1-rc6
1.8.1-rc5
1.8.1-rc4
1.8.1-rc3
1.8.1-rc2
1.8.1-rc1
1.8.1-rc0
1.8.0
1.8.0-rc3
1.8.0-rc2
1.8.0-rc1
1.8.0-rc0
1.7.3
1.7.3-rc4
1.7.3-rc3
1.7.3-rc2
1.7.3-rc1
1.7.3-rc0
1.7.2
1.7.2-rc4
1.7.2-rc3
1.7.2-rc2
1.7.2-rc1
1.7.2-rc0
1.7.1
1.7.0
1.7.0-rc1
1.7.0-rc0
1.6.1
1.6.1-rc1
1.6.1-rc0
1.6.0
1.6.0-rc9
1.6.0-rc8
1.6.0-rc7
1.6.0-rc6
1.6.0-rc5
1.6.0-rc4
1.6.0-rc3
1.6.0-rc20
1.6.0-rc2
1.6.0-rc19
1.6.0-rc18
1.6.0-rc17
1.6.0-rc16
1.6.0-rc15
1.6.0-rc14
1.6.0-rc13
1.6.0-rc12
1.6.0-rc11
1.6.0-rc10
1.6.0-rc1
1.5.0
1.5.0-rc1
1.4.2-rc2
1.4.2-rc1
1.4.1
1.4.1-rc1
1.4.0
1.4.0-rc4
1.4.0-rc3
1.4.0-rc2
1.4.0-rc1
1.3.1
1.3.0
1.3.0-rc
1.2.0
1.1.0
1.0.0
1.0.0-rc9
1.0.0-rc8
1.0.0-rc7
1.0.0-rc6
1.0.0-rc5
1.0.0-rc1
0.2.13
0.2.12
0.2.11
0.2.10
0.2.9
0.2.8
0.2.7
0.2.6
0.2.5
0.2.4
0.2.3
0.2.2
0.2.1
0.2.0
0.1.8
0.1.7
0.1.6
0.1.5
0.1.4
0.1.3
0.1.2
0.1.1
0.1.0
0.0.24
0.0.23
0.0.22
0.0.21
0.0.20
0.0.19
0.0.18
0.0.17
0.0.16
0.0.15
0.0.14
0.0.13
0.0.12
0.0.11
0.0.10
0.0.9
0.0.8
0.0.7
0.0.6
0.0.4
0.0.3
0.0.2
0.0.1
Step flow manager for Elixir applications
Current section
Files
Jump to
Current section
Files
lib/step_flow/web_controllers/workflow_events.ex
defmodule StepFlow.WebController.WorkflowEvents do
use StepFlow, :controller
use OpenApiSpex.ControllerSpecs
require Logger
import StepFlow.WebController.Helpers
alias StepFlow.{
Amqp.CommonEmitter,
Jobs,
Repo,
Step.Live,
Updates,
WebController.OpenApiSchemas,
Workflows
}
@moduledoc false
tags(["Workflows"])
security([
%{"authorization" => %OpenApiSpex.SecurityScheme{type: "http", scheme: "bearer"}}
])
action_fallback(StepFlow.WebController.Fallback)
operation(:handle,
summary: "Trigger workflows event",
description: "Trigger workflows event",
type: :string,
parameters: [
workflow_id: [
in: :path,
description: "Workflow ID",
type: :integer,
example: 1
]
],
request_body: {"Eventbody", "application/json", OpenApiSchemas.Workflows.EventBody},
responses: [
ok: "OK",
forbidden: "Forbidden"
]
)
# credo:disable-for-lines:35
def handle(
%Plug.Conn{assigns: %{current_user: user}} = conn,
%{"workflow_id" => workflow_id} = params
) do
handle_workflow(conn, params, user, workflow_id)
end
def handle(conn, _) do
conn
|> put_status(:forbidden)
|> json(%{
status: "error",
message: "Forbidden to handle workflow with this identifier"
})
end
defp handle_workflow(conn, params, user, workflow_id) do
workflow = Workflows.get_workflow!(workflow_id)
case {params, workflow.is_live} do
{%{"event" => "abort"}, false} ->
abort(conn, workflow, user)
{%{"event" => "pause", "post_action" => action, "trigger_at" => trigger_date_time}, _} ->
pause(conn, workflow, user, action, trigger_date_time)
{%{"event" => "resume"}, _} ->
resume(conn, workflow, user)
{%{"event" => "update", "job_id" => job_id, "parameters" => parameters}, _} ->
update(conn, workflow, user, job_id, parameters)
{%{"event" => "retry", "job_id" => job_id}, _} ->
{status, message} = retry(conn, workflow, user, job_id)
conn
|> put_status(status)
|> json(message)
{%{"event" => "retry"}, _} ->
retry(conn, workflow, user)
{%{"event" => "stop"}, true} ->
stop(conn, workflow, user)
{%{"event" => "delete"}, _} ->
delete(conn, workflow, user)
{%{"event" => "duplicate"}, _} ->
duplicate(conn, workflow, user)
{_, _} ->
send_resp(conn, 422, "event is not supported")
end
end
defp internal_retry_handle(workflow, job, _job_name, last_status_state, user_uuid)
when last_status_state in [:error] do
internal_retry_statuses(job, workflow, user_uuid)
params = %{
job_id: job.id,
parameters: job.parameters
}
case CommonEmitter.publish_json(job.name, job.step_id, params) do
:ok ->
StepFlow.Notification.send("retry_job", %{workflow_id: workflow.id, body: params})
Jobs.Status.set_job_status(job.id, :queued)
{:ok, %{status: "ok"}}
# NP: pattern never matches, but should in the future
_ ->
{:ok, %{status: "error", message: "unable to publish message"}}
end
end
defp internal_retry_handle(
_workflow,
_job,
_job_name,
_last_status_state,
_user_uuid
) do
{:forbidden, %{status: "error", message: "Illegal operation"}}
end
defp internal_retry_statuses(job, workflow, user_uuid) do
{:ok, job_status} = Jobs.Status.set_job_status(job.id, :retrying, %{user_uuid: user_uuid})
{:ok, _status} =
Workflows.Status.define_workflow_status(workflow.id, :job_retrying, job_status)
if workflow.parent_id != nil do
job = Jobs.get_job!(workflow.parent_id) |> Repo.preload([:child_workflow])
{:ok, job_status} = Jobs.Status.set_job_status(job.id, :retrying, %{user_uuid: user_uuid})
{:ok, _status} =
Workflows.Status.define_workflow_status(
job.workflow_id,
:job_retrying,
job_status
)
{:ok, _status} = Jobs.Status.set_job_status(job.id, :queued)
{:ok, datetime} = DateTime.now("Etc/UTC")
with :error <-
CommonEmitter.publish_json(
"job_progression",
0,
%{
job_id: job.id,
datetime: datetime,
docker_container_id: "workflow",
progression: 0
},
"job_response"
) do
Logger.error("[#{__MODULE__}] Could not send error message.")
end
end
end
defp retry(_conn, workflow, user, job_id) do
job = Jobs.get_job_with_status!(job_id)
if has_right?("workflow::" <> workflow.identifier, user, "retry") &&
Map.has_key?(user, :uuid) &&
job.child_workflow == nil do
Logger.info("User #{user.uuid} is retrying job #{job_id}")
last_status = StepFlow.Controllers.Jobs.get_last_status(job)
internal_retry_handle(workflow, job, job.name, last_status.state, user.uuid)
else
{:forbidden, %{status: "error", message: "Forbidden to retry Workflow #{workflow.id}"}}
end
end
defp retry(conn, workflow, user) do
jobs =
workflow.jobs
|> Enum.filter(fn job ->
StepFlow.Controllers.Jobs.get_last_status(job).state == :error
end)
{status, message} =
Enum.map(
jobs,
fn job -> retry(conn, workflow, user, job.id) end
)
|> Enum.filter(fn result -> result != {:ok, %{status: "ok"}} end)
|> deduce_response(workflow)
conn
|> put_status(status)
|> json(message)
end
defp deduce_response(results, _) when results == [], do: {:ok, %{status: "ok"}}
defp deduce_response(results, workflow) when results != [],
do: {:forbidden, %{status: "error", message: "Forbidden to retry Workflow #{workflow.id}"}}
defp update(conn, workflow, user, job_id, parameters) do
if has_right?("workflow::" <> workflow.identifier, user, "update") &&
workflow.parent_id == nil do
Logger.info("update job #{job_id}")
job = Jobs.get_job_with_status!(job_id)
if job.is_updatable do
Updates.update_parameters(job, parameters)
conn
|> put_status(:ok)
|> json(%{status: "ok"})
else
Logger.error("Job #{job_id} cannot be updated !")
conn
|> put_status(:error)
|> json(%{status: "error", message: "Forbidden to update Job #{job_id}"})
end
else
conn
|> put_status(:forbidden)
|> json(%{status: "error", message: "Forbidden to update Workflow #{workflow.id}"})
end
end
defp abort(conn, workflow, user) do
if has_right?("workflow::" <> workflow.identifier, user, "abort") &&
workflow.parent_id == nil do
Workflows.abort(workflow)
conn
|> put_status(:ok)
|> json(%{status: "ok"})
else
conn
|> put_status(:forbidden)
|> json(%{status: "error", message: "Forbidden to abort Workflow #{workflow.id}"})
end
end
defp pause(conn, workflow, user, action, trigger_date_time) do
if has_right?("workflow::" <> workflow.identifier, user, "pause") &&
workflow.parent_id == nil do
case Workflows.pause(workflow, action, trigger_date_time) do
{:ok, _} ->
conn
|> put_status(:ok)
|> json(%{status: "ok"})
{:error, message} when is_bitstring(message) ->
Logger.error(
"Something went wrong pausing workflow #{workflow.id} #{workflow.identifier}: #{inspect(message)}"
)
conn
|> put_status(:not_found)
|> json(%{status: "error", message: message})
{:error, message} ->
Logger.error(
"Something went wrong pausing workflow #{workflow.id} #{workflow.identifier}: #{inspect(message)}"
)
conn
|> put_status(:internal_server_error)
|> json(%{status: "error", message: "Internal server error"})
_ ->
conn
|> put_status(:not_found)
|> json(%{status: "error", message: "Unknown error"})
end
else
conn
|> put_status(:forbidden)
|> json(%{status: "error", message: "Forbidden to pause Workflow #{workflow.id}"})
end
end
defp resume(conn, workflow, user) do
if has_right?("workflow::" <> workflow.identifier, user, "resume") &&
workflow.parent_id == nil do
case Workflows.resume(workflow) do
{:ok, _} ->
conn
|> put_status(:ok)
|> json(%{status: "ok"})
{:error, message} ->
Logger.error(
"Something went wrong resuming workflow #{workflow.id} #{workflow.identifier}: #{inspect(message)}"
)
conn
|> put_status(:internal_server_error)
|> json(%{status: "error", message: "Internal server error"})
end
else
conn
|> put_status(:forbidden)
|> json(%{status: "error", message: "Forbidden to resume Workflow #{workflow.id}"})
end
end
defp stop(conn, workflow, user) do
if has_right?("workflow::" <> workflow.identifier, user, "abort") do
stop_workflow(conn, workflow)
else
conn
|> put_status(:forbidden)
|> json(%{status: "error", message: "Forbidden to stop Workflow #{workflow.id}"})
end
end
defp stop_workflow(conn, %Workflows.Workflow{is_live: true} = workflow) do
workflow_jobs = Repo.preload(workflow, [:jobs]).jobs
case stop_live(workflow, workflow_jobs) do
{:error, message} ->
Logger.warning(
"Could not stop live workflow #{workflow.id} #{workflow.identifier}: #{inspect(message)}"
)
conn
|> put_status(:internal_server_error)
|> json(%{status: "error", message: message})
_ ->
{:ok, _status} = Workflows.Status.set_workflow_status(workflow.id, :stopped)
topic = "update_workflow_" <> Integer.to_string(workflow.id)
StepFlow.Notification.send(topic, %{workflow_id: workflow.id})
Logger.warning("Stopped live workflow #{workflow.id} #{workflow.reference}")
conn
|> put_status(:ok)
|> json(%{status: "ok"})
end
end
defp stop_workflow(conn, workflow) do
{:ok, _status} = Workflows.Status.set_workflow_status(workflow.id, :stopped)
topic = "update_workflow_" <> Integer.to_string(workflow.id)
StepFlow.Notification.send(topic, %{workflow_id: workflow.id})
conn
|> put_status(:ok)
|> json(%{status: "ok"})
end
defp stop_live(workflow, jobs) do
case Workflows.Status.get_last_workflow_status(workflow.id).state do
:pending ->
{:error, "Impossible to stop Workflow #{workflow.id} that is initializing"}
_ ->
jobs |> Live.stop_jobs()
end
end
defp delete(conn, workflow, user) do
if has_right?("workflow::" <> workflow.identifier, user, "delete") do
for job <- workflow.jobs do
if job.child_workflow != nil do
child_workflow = Workflows.get_workflow!(job.child_workflow.id)
delete(conn, child_workflow, user)
end
Jobs.delete_job(job)
end
Workflows.delete_workflow(workflow)
StepFlow.Notification.send("delete_workflow", %{workflow_id: workflow.id})
conn
|> put_status(:ok)
|> json(%{status: "ok"})
else
conn
|> put_status(:forbidden)
|> json(%{status: "error", message: "Forbidden to delete Workflow #{workflow.id}"})
end
end
defp duplicate(conn, workflow, user) do
if has_right?("workflow::" <> workflow.identifier, user, "create") do
case Workflows.duplicate_workflow(workflow.id, user.uuid) do
{:ok, duplicated_workflow_id} ->
conn
|> put_status(:ok)
|> json(%{status: "ok", workflow_id: duplicated_workflow_id})
{:error, _} ->
conn
|> put_status(:internal_server_error)
|> json(%{status: "error", message: "Internal server error"})
end
else
conn
|> put_status(:forbidden)
|> json(%{status: "error", message: "Forbidden to duplicate Workflow #{workflow.id}"})
end
end
end