Packages
journey
0.10.23
0.10.58
0.10.57
0.10.56
0.10.55
0.10.54
0.10.53
0.10.52
0.10.51
0.10.50
0.10.49
0.10.48
0.10.47
0.10.46
0.10.45
0.10.44
0.10.43
0.10.41
0.10.40
0.10.39
0.10.38
0.10.37
0.10.36
0.10.35
0.10.34
0.10.33
0.10.32
0.10.31
0.10.30
0.10.29
0.10.28
0.10.27
0.10.26
0.10.25
0.10.24
0.10.23
0.10.22
0.0.9
retired
0.0.8
0.0.7
0.0.6
0.0.5
0.0.3
0.0.2
Journey is a library for defining and running durable workflows with persistence, reliability, and scalability.
Current section
Files
Jump to
Current section
Files
lib/journey/graph/validations.ex
defmodule Journey.Graph.Validations do
@moduledoc false
def validate(graph) do
graph
|> ensure_no_duplicate_node_names()
|> validate_dependencies()
end
def ensure_known_node_name(execution, node_name) do
# Fetch the graph from catalog to get all node names
graph = Journey.Graph.Catalog.fetch(execution.graph_name, execution.graph_version)
all_node_names = graph.nodes |> Enum.map(& &1.name)
if node_name in all_node_names do
:ok
else
raise "'#{inspect(node_name)}' is not a known node in execution '#{execution.id}' / graph '#{execution.graph_name}'. Valid node names: #{inspect(Enum.sort(all_node_names))}."
end
end
def ensure_known_input_node_name(graph, node_name)
when is_struct(graph, Journey.Graph) and is_atom(node_name) do
all_input_node_names =
graph.nodes
|> Enum.filter(fn n -> n.type == :input end)
|> Enum.map(& &1.name)
if node_name in all_input_node_names do
:ok
else
raise "'#{inspect(node_name)}' is not a valid input node in graph '#{graph.name}'.'#{graph.version}'. Valid input node names: #{inspect(Enum.sort(all_input_node_names))}."
end
end
def ensure_known_input_node_name(execution, node_name)
when is_struct(execution, Journey.Persistence.Schema.Execution) and is_atom(node_name) do
graph = Journey.Graph.Catalog.fetch(execution.graph_name, execution.graph_version)
all_input_node_names =
graph.nodes
|> Enum.filter(fn n -> n.type == :input end)
|> Enum.map(& &1.name)
if node_name in all_input_node_names do
:ok
else
raise "'#{inspect(node_name)}' is not a valid input node in execution '#{execution.id}' / graph '#{execution.graph_name}'. Valid input node names: #{inspect(Enum.sort(all_input_node_names))}."
end
end
defp validate_dependencies(graph) do
all_node_names = Enum.map(graph.nodes, & &1.name)
graph.nodes |> Enum.each(fn node -> validate_node(node, all_node_names) end)
graph
end
defp ensure_no_duplicate_node_names(graph) do
graph.nodes
|> Enum.map(& &1.name)
|> Enum.frequencies()
|> Enum.filter(fn {_, v} -> v > 1 end)
|> Enum.each(fn {k, _v} ->
raise "Duplicate node name in graph definition: #{inspect(k)}"
end)
graph
end
defp validate_node(%Journey.Graph.Step{} = step, all_node_names) do
all_upstream_node_names =
step.gated_by
|> Journey.Node.UpstreamDependencies.Computations.list_all_node_names()
|> MapSet.new()
unknown_deps = MapSet.difference(all_upstream_node_names, MapSet.new(all_node_names))
if Enum.any?(unknown_deps) do
raise "Unknown upstream nodes in input node '#{inspect(step.name)}': #{Enum.join(unknown_deps, ", ")}"
end
if not is_nil(step.mutates) and step.mutates not in all_node_names do
raise "Mutation node '#{inspect(step.name)}' mutates an unknown node '#{inspect(step.mutates)}'"
end
if not is_nil(step.mutates) and step.mutates == step.name do
raise "Mutation node '#{inspect(step.name)}' attempts to mutate itself"
end
step
end
defp validate_node(%Journey.Graph.Input{} = input, _all_node_names) do
input
end
end