Packages
krug
2.0.12
2.0.36
2.0.35
2.0.34
2.0.33
2.0.32
2.0.31
2.0.30
2.0.29
2.0.27
2.0.26
2.0.25
2.0.24
2.0.23
2.0.22
2.0.20
2.0.19
2.0.18
2.0.17
2.0.16
2.0.15
2.0.14
2.0.13
2.0.12
2.0.11
2.0.10
2.0.9
2.0.8
2.0.7
2.0.6
2.0.5
2.0.4
2.0.1
2.0.0
1.1.53
1.1.52
1.1.50
1.1.49
1.1.48
1.1.47
1.1.46
1.1.45
1.1.44
1.1.43
1.1.42
1.1.41
1.1.40
1.1.39
1.1.38
1.1.37
1.1.36
1.1.35
1.1.34
1.1.33
1.1.32
1.1.31
1.1.30
1.1.29
1.1.28
1.1.27
1.1.26
1.1.25
1.1.24
1.1.23
1.1.22
1.1.21
1.1.20
1.1.19
1.1.18
1.1.17
1.1.16
1.1.15
1.1.14
1.1.12
1.1.10
1.1.9
1.1.8
1.1.7
1.1.6
1.1.5
1.1.4
1.1.3
1.1.2
1.1.1
1.1.0
1.0.9
1.0.8
1.0.7
1.0.6
1.0.5
1.0.4
1.0.3
1.0.2
1.0.1
1.0.0
0.5.5
0.5.4
0.5.3
0.5.2
0.5.1
0.5.0
0.4.31
0.4.30
0.4.29
0.4.28
0.4.27
0.4.26
0.4.25
0.4.24
0.4.23
0.4.22
0.4.20
0.4.19
0.4.18
0.4.17
0.4.16
0.4.15
0.4.14
0.4.13
0.4.12
0.4.11
0.4.10
0.4.9
0.4.8
0.4.7
0.4.6
0.4.5
0.4.4
0.4.3
0.4.2
0.4.1
0.4.0
0.3.9
0.3.8
0.3.7
0.3.6
0.3.5
0.3.4
0.3.3
0.3.2
0.3.1
0.3.0
0.2.1
0.2.0
0.1.0
A Utilitary package functionalities modules for improve a secure performatic development.
Current section
Files
Jump to
Current section
Files
lib/krug/util/cluster.util.ex
defmodule Krug.ClusterUtil do
@moduledoc """
Utilitary module to handle cluster nodes operations
"""
@moduledoc since: "1.1.17"
@connection_node_timeout 100
@doc """
Connect the local node to a list of other nodes by their IP's "cluster_ips".
Return a list (of atom) containing the sucessfully connected nodes.
"""
def connect_nodes(connected_nodes,cluster_name,cluster_ips,connection_timeout) do
connection_timeout = cond do
(nil == connection_timeout)
-> @connection_node_timeout
true
-> connection_timeout
end
cluster_nodes = cluster_ips
|> Enum.map(
fn(cluster_ip) ->
"#{cluster_name}@#{cluster_ip}"
|> String.to_atom()
end
)
|> Enum.chunk_every(25)
connected_nodes
|> connect_nodes2(cluster_nodes,connection_timeout)
end
##########################################
### init functions
##########################################
def connect_nodes2(connected_nodes,cluster_nodes,connection_timeout) do
cond do
(Enum.empty?(cluster_nodes))
-> connected_nodes
true
-> connected_nodes
|> connect_nodes3(cluster_nodes,connection_timeout)
end
end
defp connect_nodes3(connected_nodes,cluster_nodes,connection_timeout) do
cluster_nodes
|> hd()
|> connect_nodes4(connected_nodes,connection_timeout)
|> connect_nodes2(
cluster_nodes |> tl(),
connection_timeout
)
end
defp connect_nodes4(nodes,connected_nodes,connection_timeout) do
task_nodes = nodes
|> enqueue_connection_tasks()
connection_timeout
|> :timer.sleep()
connected_nodes
|> verify_connected_nodes(task_nodes)
end
defp verify_connected_nodes(connected_nodes,task_nodes) do
cond do
(Enum.empty?(task_nodes))
-> connected_nodes
true
-> connected_nodes
|> verify_connected_nodes2(task_nodes)
end
end
defp verify_connected_nodes2(connected_nodes,task_nodes) do
[task,node] = task_nodes
|> hd()
%Task{pid: pid} = task
cond do
(Process.alive?(pid))
-> connected_nodes
|> shutdown_connection_task(task_nodes,task,node)
(Task.await(task))
-> [node | connected_nodes]
|> verify_connected_nodes(task_nodes |> tl())
true
-> connected_nodes
|> verify_connected_nodes(task_nodes |> tl())
end
end
defp shutdown_connection_task(connected_nodes,task_nodes,task,node) do
connected_on_aborting = task
|> Task.shutdown(0)
|> shutdown_connection_task_result()
cond do
(connected_on_aborting)
-> [node | connected_nodes]
|> verify_connected_nodes(task_nodes |> tl())
true
-> connected_nodes
|> verify_connected_nodes(task_nodes |> tl())
end
end
defp shutdown_connection_task_result({:ok, reply}) do
reply == true
end
defp shutdown_connection_task_result(_) do
false
end
defp enqueue_connection_tasks(nodes,tasks \\ []) do
cond do
(Enum.empty?(nodes))
-> tasks
true
-> nodes
|> enqueue_connection_tasks2(tasks)
end
end
defp enqueue_connection_tasks2(nodes,tasks) do
node = nodes
|> hd()
task = Task.async(
fn ->
node
|> :net_kernel.connect_node()
end
)
nodes
|> tl()
|> enqueue_connection_tasks([[task,node] | tasks])
end
end