Created
October 24, 2015 18:05
-
-
Save siscia/f8f563aa2ad2f3517738 to your computer and use it in GitHub Desktop.
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| defmodule Numerino.Supervisor do | |
| use Supervisor | |
| def start_link(_) do | |
| {:ok, sup} = Supervisor.start_link(__MODULE__, [], name: :supervisor) end | |
| def init(_) do | |
| processes = [ | |
| worker(Numerino.QueueAddress, []), | |
| #worker(Numerino.Restarted, []), | |
| supervisor(Numerino.QueueManager, []) | |
| ] | |
| {:ok, {{:one_for_one, 10, 10}, processes}} | |
| end | |
| def kill_queue_manager(pid, reason) do | |
| IO.inspect self | |
| Process.exit(pid, reason) | |
| end | |
| end | |
| defmodule Numerino.Restarted do | |
| def start_link do | |
| Process.spawn loop, [name: Numerino.Restarted] | |
| end | |
| def loop do | |
| receive do | |
| :restart -> restart_all(); loop() | |
| end | |
| end | |
| defp restart_all do | |
| :ets.foldl(&do_restart_one/2, 0, Numerino.QueueAddress) | |
| end | |
| defp do_restart_one {name, _pid, priorities}, acc do | |
| Numerino.QueueManager.new_queue(name, priorities) | |
| acc + 1 | |
| end | |
| end | |
| defmodule Numerino.QueueManager do | |
| use Supervisor | |
| @name QueueManager | |
| def start_link (opts \\ [] ) do | |
| #t = restart_old_queues() | |
| Supervisor.start_link(__MODULE__, :ok, [name: @name]) | |
| #Task.await(t) | |
| end | |
| def init :ok do | |
| process = [ | |
| worker(Numerino, []) | |
| ] | |
| supervise(process, strategy: :simple_one_for_one) | |
| end | |
| defp generate_callback(name, priorities) do | |
| fn p -> | |
| Numerino.QueueAddress.follow(p, name, priorities) | |
| end | |
| end | |
| def new_queue(name, priorities) do | |
| callback = generate_callback(name, priorities) | |
| {:ok, _p} = Supervisor.start_child(@name, [priorities, callback]) | |
| end | |
| defp do_restart_old_queue {name, _pid, priorities}, acc do | |
| new_queue(name, priorities) | |
| acc + 1 | |
| end | |
| defp restart_old_queues do | |
| restart = Task.async(fn -> :ets.foldl(&do_restart_old_queue/2, 0, Numerino.QueueAddress) end) | |
| end | |
| end | |
| defmodule Numerino.QueueAddress do | |
| use GenServer | |
| @table_name Numerino.QueueAddress | |
| @server_name QueueAddressServer | |
| def start_link opts \\ [name: @server_name] do | |
| GenServer.start_link(__MODULE__, :ok, opts) | |
| end | |
| def follow pid, name, priorities do | |
| GenServer.call(@server_name, {:follow, pid, name, priorities}) | |
| end | |
| def init :ok do | |
| @table_name = :ets.new(@table_name, | |
| [:named_table, {:read_concurrency, true}]) | |
| {:ok, @table_name} | |
| end | |
| def handle_call {:follow, pid, name, priorities}, _from, state do | |
| :ets.insert(@table_name, {name, pid, priorities}) | |
| {:reply, name ,state} | |
| end | |
| end |
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment