Packages
step_flow
1.7.2-rc0
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"
]
def handle(
%Plug.Conn{assigns: %{current_user: user}} = conn,
%{"workflow_id" => workflow_id} = params
) 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}, _} ->
retry(conn, workflow, user, job_id)
{%{"event" => "stop"}, true} ->
stop(conn, workflow, user)
{%{"event" => "delete"}, _} ->
delete(conn, workflow, user)
{_, _} ->
send_resp(conn, 422, "event is not supported")
end
end
def handle(conn, _) do
conn
|> put_status(:forbidden)
|> json(%{status: "error", message: "Forbidden to handle workflow with this identifier"})
end
defp internal_retry_handle(conn, 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)
conn
|> put_status(:ok)
|> json(%{status: "ok"})
_ ->
conn
|> put_status(:ok)
|> json(%{status: "error", message: "unable to publish message"})
end
end
defp internal_retry_handle(conn, _workflow, _job, _job_name, _last_status_state, _user_uuid) do
send_resp(conn, :forbidden, "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")
CommonEmitter.publish_json(
"job_progression",
0,
%{
job_id: job.id,
datetime: datetime,
docker_container_id: "workflow",
progression: 0
},
"job_response"
)
end
end
defp update(conn, workflow, user, job_id, parameters) do
if has_right?("workflow::" <> workflow.identifier, user, "update") &&
workflow.parent_id == nil do
Logger.warn("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 this job"})
end
else
conn
|> put_status(:forbidden)
|> json(%{status: "error", message: "Forbidden to update this workflow"})
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 this workflow"})
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 this workflow"})
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 this workflow"})
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.warn("User #{user.uuid} is retrying job #{job_id}")
last_status = StepFlow.Controllers.Jobs.get_last_status(job.status)
internal_retry_handle(conn, workflow, job, job.name, last_status.state, user.uuid)
else
conn
|> put_status(:forbidden)
|> json(%{status: "error", message: "Forbidden to retry this workflow"})
end
end
defp stop(conn, workflow, user) do
if has_right?("workflow::" <> workflow.identifier, user, "abort") do
workflow_jobs = Repo.preload(workflow, [:jobs]).jobs
workflow_jobs
|> Live.stop_jobs()
{: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"})
else
conn
|> put_status(:forbidden)
|> json(%{status: "error", message: "Forbidden to stop this workflow"})
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 this workflow"})
end
end
end