Packages
step_flow
1.7.0-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/models/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
require Logger
@moduledoc false
defenum(StateEnum, [
# State can start from: queued, skipped, error
# Unknown, Retrying -->
# --> Processing, Ready_to_init, Paused, Error, Dropped
"queued",
# Processing, Queued -->
# --> Processing, Dropped
"paused",
# Unknown -->
"skipped",
# Queued -->
# --> Completed, Error, Paused, Stopped
"processing",
# Error -->
# --> Queued
"retrying",
# Queued, Processing -->
# --> Retrying, Deleting
"error",
# Processing -->
"completed",
# Queued -->
# --> Initializing
"ready_to_init",
# Initialized, Initializing -->
# --> Starting
"ready_to_start",
# Processing -->
# --> Updating
"update",
# Processing -->
# --> Processing, Deleting
"stopped",
# Ready_to_init -->
# --> Initialized, Ready_to_start
"initializing",
# Ready_to_start -->
# --> Processing
"starting",
# Updating -->
# --> Processing
"updating",
# Stopped -->
# --> Completed
"deleting",
# Queued, Paused -->
"dropped",
# DEPRECATED
"running",
# DEPRECATED
"initialized",
# --> Queued, Skipped
"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,
# DEPRECATED
15 => :running,
# DEPRECATED
16 => :initialized,
17 => :deleting,
18 => :dropped
}
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, :dropped],
1 => [],
2 => [:completed, :error, :paused, :stopped],
3 => [:queued],
4 => [:deleting, :retrying],
5 => [],
6 => [:initializing],
7 => [:starting],
8 => [:updating],
9 => [:deleting, :processing],
10 => [:ready_to_start],
11 => [:processing],
12 => [:processing],
13 => [:queued, :skipped],
14 => [:dropped, :processing],
# DEPRECATED
15 => [],
# DEPRECATED
16 => [],
17 => [:completed],
18 => []
}
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
def convert_to_string(%Status{} = status), do: convert_to_string(status.state)
def convert_to_string(status) when is_binary(status), do: status
def 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 = StepFlow.Controllers.Jobs.get_last_status(job.status)
if convert_to_string(last_status) == string_status do
Logger.warn(
"Transition to the same state #{string_status}, ignoring and fall back to last status."
)
{:ok, last_status}
else
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 for job #{job_id}."}
end
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
end