Packages
step_flow
1.7.3-rc1
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.ex
defmodule StepFlow.WebController.Workflow do
use StepFlow, :controller
use OpenApiSpex.ControllerSpecs
require Logger
alias OpenApiSpex.Schema
alias StepFlow.Configuration
alias StepFlow.Repo
alias StepFlow.Step
alias StepFlow.WebController.Helpers
alias StepFlow.WebController.OpenApiSchemas
alias StepFlow.Workflows
alias StepFlow.Workflows.Workflow
@moduledoc false
tags ["Workflows"]
security [%{"authorization" => %OpenApiSpex.SecurityScheme{type: "http", scheme: "bearer"}}]
action_fallback(StepFlow.WebController.Fallback)
operation :index,
summary: "List workflows",
description: "List workflows",
type: :object,
parameters: [
identifier: [
in: :query,
description: "Workflows identifier",
type: :string,
example: "my_workflow"
],
status: [
in: :query,
description: "Workflows status",
type: :string,
example: "completed"
],
start_date: [
in: :query,
description: "After date",
type: :string,
example: "2022-10-21T16:46:17"
],
end_date: [
in: :query,
description: "Before date",
type: :string,
example: "2022-10-21T16:46:17"
],
mode: [
in: :query,
description: "Output mode",
schema: %Schema{
type: :array,
items: %Schema{
type: :string,
enum: [:full, :simple]
}
},
example: ["full"]
]
],
responses: [
ok: {"Workflows", "application/json", OpenApiSchemas.Workflows.Workflows},
forbidden: "Forbidden"
]
def index(%Plug.Conn{assigns: %{current_user: user}} = conn, %{"job_id" => job_id} = params) do
mode = StepFlow.Map.get_by_key_or_atom(params, :mode, "full")
workflow =
Workflows.get_workflow_for_job!(job_id)
|> Repo.preload([:status, jobs: :child_workflow])
if Helpers.has_right?("workflow::" <> workflow.identifier, user, "view") do
conn
|> put_view(StepFlow.WorkflowView)
|> render("show.json", workflow: workflow, mode: mode)
else
conn
|> put_status(:forbidden)
|> put_view(StepFlow.WorkflowDefinitionView)
|> render("error.json",
errors: %{reason: "Forbidden to view workflow with this identifier"}
)
end
end
def index(%Plug.Conn{assigns: %{current_user: user}} = conn, params) do
workflows =
params
|> Map.put("roles", user.roles)
|> Workflows.list_workflows()
conn
|> put_view(StepFlow.WorkflowView)
|> render("index.json", workflows: workflows)
end
def index(conn, _) do
conn
|> put_status(:forbidden)
|> put_view(StepFlow.WorkflowDefinitionView)
|> render("error.json",
errors: %{reason: "Forbidden to view workflows."}
)
end
def create_workflow(%Plug.Conn{assigns: %{current_user: user}} = conn, workflow_params) do
if Map.has_key?(user, :uuid) do
workflow_params = Map.put(workflow_params, :user_uuid, user.uuid)
filenames_list =
workflow_params.parameters
|> Enum.filter(fn x ->
if Map.get(x, "type") == "string" do
has_space(Map.get(x, "value"))
else
false
end
end)
if filenames_list != [] do
conn
|> put_status(:unprocessable_entity)
|> put_view(StepFlow.WorkflowDefinitionView)
|> render("error.json", errors: %{reason: "Whitespace in parameters values"})
else
case Workflows.create_workflow(workflow_params) do
{:ok, %Workflow{} = workflow} ->
Workflows.Status.define_workflow_status(workflow.id, :created_workflow)
Step.start_next(workflow)
StepFlow.Notification.send("new_workflow", %{workflow_id: workflow.id})
conn
|> put_status(:created)
|> put_view(StepFlow.WorkflowView)
|> render("created.json", workflow: workflow)
{:error, changeset} ->
conn
|> put_status(:unprocessable_entity)
|> put_view(StepFlow.ChangesetView)
|> render("error.json", changeset: changeset)
end
end
else
conn
|> put_status(:forbidden)
|> put_view(StepFlow.WorkflowDefinitionView)
|> render("error.json",
errors: %{reason: "Forbidden to create workflow with a user with no uuid"}
)
end
end
def locate_and_create_workflow(conn, user, workflow_definition, workflow_params) do
case workflow_definition do
nil ->
conn
|> put_status(:unprocessable_entity)
|> put_view(StepFlow.WorkflowDefinitionView)
|> render("error.json",
errors: %{reason: "Unable to locate workflow with this identifier"}
)
workflow_definition ->
if Helpers.has_right?("workflow::" <> workflow_definition.identifier, user, "create") do
workflow_description =
workflow_definition
|> Map.put(:reference, Map.get(workflow_params, "reference"))
|> Map.put(
:parameters,
merge_parameters(
StepFlow.Map.get_by_key_or_atom(workflow_definition, :parameters),
Map.get(workflow_params, "parameters", %{})
)
)
|> Map.from_struct()
create_workflow(conn, workflow_description)
else
conn
|> put_status(:forbidden)
|> put_view(StepFlow.WorkflowDefinitionView)
|> render("error.json",
errors: %{reason: "Forbidden to create workflow with this identifier"}
)
end
end
end
operation :create,
summary: "Create workflows",
description: "Create workflows",
type: :object,
request_body: {"WorkflowParams", "application/json", OpenApiSchemas.Workflows.WorkflowParams},
responses: [
created: {"Workflows", "application/json", OpenApiSchemas.Workflows.Workflow},
forbidden: "Forbidden"
]
def create(
%Plug.Conn{assigns: %{current_user: user}} = conn,
%{
"workflow_identifier" => identifier,
"version_major" => version_major,
"version_minor" => version_minor,
"version_micro" => version_micro
} = workflow_params
) do
workflow_definition =
StepFlow.WorkflowDefinitions.get_workflow_definition(
identifier,
version_major,
version_minor,
version_micro
)
locate_and_create_workflow(conn, user, workflow_definition, workflow_params)
end
def create(
%Plug.Conn{assigns: %{current_user: user}} = conn,
%{"workflow_identifier" => identifier} = workflow_params
) do
workflow_definition = StepFlow.WorkflowDefinitions.get_workflow_definition(identifier)
locate_and_create_workflow(conn, user, workflow_definition, workflow_params)
end
def create(conn, _workflow_params) do
conn
|> put_status(:unprocessable_entity)
|> put_view(StepFlow.WorkflowDefinitionView)
|> render("error.json",
errors: %{reason: "Missing Workflow identifier parameter"}
)
end
defp merge_parameters(parameters, request_parameters, result \\ [])
defp merge_parameters([], _request_parameters, result), do: result
defp merge_parameters([parameter | tail], request_parameters, result) do
result =
case Map.get(request_parameters, Map.get(parameter, "id")) do
nil ->
List.insert_at(result, -1, parameter)
parameter_value ->
List.insert_at(result, -1, Map.put(parameter, "value", parameter_value))
end
merge_parameters(tail, request_parameters, result)
end
operation :show,
summary: "Show workflow",
description: "Show workflow by id",
type: :object,
parameters: [
id: [
in: :path,
description: "Workflow ID",
type: :integer,
example: 1
],
mode: [
in: :query,
description: "Output mode",
schema: %Schema{
type: :array,
items: %Schema{
type: :string,
enum: [:full, :simple]
}
},
example: ["full"]
]
],
responses: [
ok: {"Workflow", "application/json", OpenApiSchemas.Workflows.Workflow},
forbidden: "Forbidden",
not_found: "Not Found"
]
def show(%Plug.Conn{assigns: %{current_user: user}} = conn, %{"id" => id} = params) do
mode = StepFlow.Map.get_by_key_or_atom(params, :mode, "full")
workflow =
Workflows.get_workflow!(id)
|> Repo.preload([:status, jobs: :child_workflow])
if Helpers.has_right?("workflow::" <> workflow.identifier, user, "view") do
conn
|> put_view(StepFlow.WorkflowView)
|> render("show.json", workflow: workflow, mode: mode)
else
conn
|> put_status(:forbidden)
|> put_view(StepFlow.WorkflowDefinitionView)
|> render("error.json",
errors: %{reason: "Forbidden to view workflow with this identifier"}
)
end
end
def show(conn, _) do
conn
|> put_status(:forbidden)
|> put_view(StepFlow.WorkflowDefinitionView)
|> render("error.json",
errors: %{reason: "Forbidden to show workflow with this identifier"}
)
end
def get(conn, %{"identifier" => workflow_identifier} = _params) do
workflow =
case workflow_identifier do
_ -> %{}
end
conn
|> json(workflow)
end
def get(conn, _params) do
conn
|> json(%{})
end
operation :statistics,
summary: "Get workflows statistics",
description: "Get workflows statistics",
type: :object,
request_body:
{"StatisticsParams", "application/json", OpenApiSchemas.Workflows.StatisticsParams},
responses: [
ok:
{"WorkflowsStatistics", "application/json", OpenApiSchemas.Statistics.WorkflowStatistics},
forbidden: "Forbidden"
]
def statistics(
%Plug.Conn{assigns: %{current_user: user}} = conn,
params
) do
start_date =
params
|> Map.get(
"start_date",
NaiveDateTime.utc_now()
|> NaiveDateTime.add(-86_400, :second)
|> NaiveDateTime.to_string()
)
|> NaiveDateTime.from_iso8601!()
|> NaiveDateTime.truncate(:second)
end_date =
params
|> Map.get(
"end_date",
NaiveDateTime.utc_now() |> NaiveDateTime.to_string()
)
|> NaiveDateTime.from_iso8601!()
|> NaiveDateTime.truncate(:second)
time_interval =
params
|> Map.get("time_interval", define_time_interval(start_date, end_date))
|> String.to_integer()
|> Kernel.max(1)
identifiers = Map.get(params, "identifiers", [])
workflows_status =
Workflows.Status.list_workflows_status(start_date, end_date, identifiers, user.roles)
conn
|> put_view(StepFlow.WorkflowView)
|> render("statistics.json", %{
workflows_status: workflows_status,
time_interval: time_interval,
end_date: end_date
})
end
defp define_time_interval(start_date, end_date) do
with number_bins <- Configuration.get_var_value(StepFlow.Workflows, :number_bins, 60) do
NaiveDateTime.diff(end_date, start_date, :second)
|> Kernel.div(number_bins)
|> Integer.to_string()
end
end
operation :update, false
def update(conn, _params) do
conn
|> put_status(:forbidden)
|> put_view(StepFlow.WorkflowDefinitionView)
|> render("error.json",
errors: %{reason: "Forbidden to update workflow with this identifier"}
)
end
operation :delete,
summary: "Delete workflow",
description: "Delete workflow by id",
type: :object,
parameters: [
id: [
in: :path,
description: "Workflow ID",
type: :integer,
example: 1
]
],
responses: [
no_content: "No Content",
forbidden: "Forbidden",
not_found: "Not Found"
]
def delete(%Plug.Conn{assigns: %{current_user: user}} = conn, %{"id" => id}) do
workflow = Workflows.get_workflow!(id)
if Helpers.has_right?("workflow::" <> workflow.identifier, user, "delete") do
with {:ok, %Workflow{}} <- Workflows.delete_workflow(workflow) do
send_resp(conn, :no_content, "")
end
else
conn
|> put_status(:forbidden)
|> put_view(StepFlow.WorkflowDefinitionView)
|> render("error.json",
errors: %{reason: "Forbidden to update workflow with this identifier"}
)
end
end
def delete(conn, _) do
conn
|> put_status(:forbidden)
|> put_view(StepFlow.WorkflowDefinitionView)
|> render("error.json",
errors: %{reason: "Forbidden to delete workflow with this identifier"}
)
end
defp has_space(string) when is_binary(string), do: String.contains?(string, " ")
defp has_space(_string), do: false
end