Packages
step_flow
1.6.0-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/jobs/status.ex
defmodule StepFlow.Jobs.Status do
use Ecto.Schema
import Ecto.Changeset
import Ecto.Query, warn: false
import EctoEnum
alias StepFlow.Jobs
alias StepFlow.Jobs.Job
alias StepFlow.Jobs.Status
alias StepFlow.Repo
alias StepFlow.Workflows
@moduledoc false
defenum(StateEnum, [
# State can start from: queued, skipped, error
# Retrying -->
# --> Processing, Ready_to_init, Paused, Error
"queued",
# Processing, Queued -->
# --> Processing
"paused",
# Singleton
"skipped",
# Queued -->
# --> Completed, Error, Paused, Stopped
"processing",
# Error -->
# --> Queued
"retrying",
# Queued, Processing, Error -->
# --> Retrying, Error
"error",
# Processing -->
"completed",
# Queued -->
# --> Initializing
"ready_to_init",
# Initialized, Initializing -->
# --> Starting
"ready_to_start",
# Processing -->
# --> Updating
"update",
# Processing -->
# --> Processing, Completed
"stopped",
# Ready_to_init -->
# --> Initialized, Ready_to_start
"initializing",
# Ready_to_start -->
# --> Processing
"starting",
# Updating -->
# --> Processing
"updating",
"unknown"
])
defp state_map_lookup(value, key \\ false) do
state_map = %{
0 => :queued,
1 => :skipped,
2 => :processing,
3 => :retrying,
4 => :error,
5 => :completed,
6 => :ready_to_init,
7 => :ready_to_start,
8 => :update,
9 => :stopped,
10 => :initializing,
11 => :starting,
12 => :updating,
13 => :unknown,
14 => :paused
}
if key do
if is_number(value) do
value
else
state_map
|> Enum.find(fn {_key, val} -> val == value end)
|> elem(0)
end
else
if is_number(value) do
state_map[value]
else
case Map.values(state_map) |> Enum.member?(value) do
true -> value
_ -> nil
end
end
end
end
def state_enum_label(value) do
to_atom_state(value)
|> Atom.to_string()
end
defp to_atom_state(value) do
case state_map_lookup(value) do
nil -> :unknown
value -> value
end
end
defp state_enum_position(value) do
state_map_lookup(value, true)
end
defp transition_map_lookup(value) do
state_map = %{
0 => [:processing, :ready_to_init, :paused, :error],
1 => [],
2 => [:completed, :error, :paused, :stopped, :processing],
3 => [:queued],
4 => [:retrying, :error],
5 => [],
6 => [:initializing],
7 => [:starting],
8 => [:updating],
9 => [:completed, :processing],
10 => [:ready_to_start],
11 => [:processing],
12 => [:processing],
13 => [],
14 => [:processing]
}
if is_number(value) do
state_map[value]
else
case Map.values(state_map) |> Enum.member?(value) do
true -> value
_ -> nil
end
end
end
defp to_atom_transition(value) do
case transition_map_lookup(value) do
nil -> []
value -> value
end
end
schema "step_flow_status" do
field(:state, StepFlow.Jobs.Status.StateEnum)
field(:description, :map, default: %{})
belongs_to(:job, Job, foreign_key: :job_id)
has_many(:workflow_status, Workflows.Status, on_delete: :delete_all)
timestamps()
end
@doc false
def changeset(%Status{} = job, attrs) do
job
|> cast(attrs, [:state, :job_id, :description])
|> foreign_key_constraint(:job_id)
|> validate_required([:state, :job_id])
end
defp convert_to_string(%Status{} = status), do: convert_to_string(status.state)
defp convert_to_string(status) when is_binary(status), do: status
defp convert_to_string(status) do
state_enum_label(status)
end
def set_job_status(job_id, status, description \\ %{}) do
# Ensure status is string
string_status = convert_to_string(status)
job = Jobs.get_job!(job_id) |> Repo.preload([:status])
last_status = get_last_status(job.status)
if check_transition(last_status, string_status) do
%Status{}
|> Status.changeset(%{job_id: job_id, state: status, description: description})
|> Repo.insert()
else
{:error,
"Status transition from #{convert_to_string(last_status)} to #{string_status} not allowed"}
end
end
defp check_transition(last_status, new_status) do
case {last_status, new_status} do
paired when paired in [{nil, "queued"}, {nil, "error"}, {nil, "skipped"}] ->
true
{nil, _} ->
false
{_, _} ->
enum_position = state_enum_position(last_status.state)
allowed_transitions = to_atom_transition(enum_position)
String.to_atom(new_status) in allowed_transitions
end
end
@doc """
Returns the last updated status of a list of status.
"""
def get_last_status(status) when is_list(status) do
status
|> Enum.sort(fn state_1, state_2 ->
case NaiveDateTime.compare(state_1.inserted_at, state_2.inserted_at) do
:lt -> true
:gt -> false
:eq -> state_1.id < state_2.id
end
end)
|> List.last()
end
def get_last_status(%Status{} = status), do: status
def get_last_status(_status), do: nil
@doc """
Returns the last status id of a list of status.
"""
def get_last_status_id(status) when is_list(status) do
status
|> Enum.sort(fn state_1, state_2 ->
state_1.id < state_2.id
end)
|> List.last()
end
def get_last_status_id(%Status{} = status), do: status
def get_last_status_id(_status), do: nil
@doc """
Returns action linked to status
"""
def get_action(status) do
case status.state do
:queued -> "create"
:ready_to_init -> "init_process"
:ready_to_start -> "start_process"
:update -> "update_process"
:stopped -> "delete"
_ -> "none"
end
end
@doc """
Returns action linked to status as parameter
"""
def get_action_parameter(status) do
action = get_action(status)
[%{"id" => "action", "type" => "string", "value" => action}]
end
end