diff --git a/lib/trinity/scheduler.ex b/lib/trinity/scheduler.ex index 1549dd3..177999e 100644 --- a/lib/trinity/scheduler.ex +++ b/lib/trinity/scheduler.ex @@ -30,6 +30,7 @@ defmodule Trinity.Scheduler do end defmacro simulation_key, do: :_trinity_simulation + defmacro sim_node_key, do: :_trinity_simulation_node @spec get_sim :: Simulation.t defp get_sim, do: Process.get(simulation_key()) @@ -75,8 +76,8 @@ defmodule Trinity.Scheduler do @spec spawn_node_and_yield(function, atom, boolean) :: pid def spawn_node_and_yield(fun, node, link?) do - %{proc_nodes: proc_nodes} = sim = get_sim() - node = node || get_proc_node(proc_nodes) + sim = get_sim() + node = node || get_proc_node() parent_pid = self() spawn_ref = make_ref() @@ -230,16 +231,6 @@ defmodule Trinity.Scheduler do linked end - def handle_sent(dest_ref) when is_reference(dest_ref) do - %{proc_aliases: proc_aliases} = get_sim() - # Look up the alias ref in the aliases table to get its pid - case :ets.lookup(proc_aliases, dest_ref) do - [{_key, dest_pid}] -> handle_sent(dest_pid) - # This alias has no known pid (alias was unaliased), do nothing - [] -> :noop - end - end - def handle_sent(dest_pid) when is_pid(dest_pid) do %{queue: queue, proc_queue_keys: proc_queue_keys, now: now} = get_sim() @@ -386,14 +377,14 @@ defmodule Trinity.Scheduler do end end - defp get_proc_node(proc_nodes) do - [{_pid, node}] = :ets.lookup(proc_nodes, self()) - node - end + defp get_proc_node, do: Process.get(sim_node_key()) defp put_proc_node(proc_nodes, node_procs, node) do pid = self() + + Process.put(sim_node_key(), node) :ets.insert(proc_nodes, {pid, node}) + existing_procs = case :ets.lookup(node_procs, node) do [{_node, procs}] -> procs [] -> [] diff --git a/lib/trinity/sim_process.ex b/lib/trinity/sim_process.ex index ea12e31..66f3a03 100644 --- a/lib/trinity/sim_process.ex +++ b/lib/trinity/sim_process.ex @@ -2,14 +2,28 @@ defmodule Trinity.SimProcess do alias Trinity.Scheduler alias Trinity.Scheduler.Simulation - import Trinity.Scheduler, only: [simulation_key: 0] + import Trinity.Scheduler, only: [simulation_key: 0, sim_node_key: 0] @spec get_sim :: Simulation.t | nil defp get_sim, do: Process.get(simulation_key()) + @spec get_proc_node :: atom + defp get_proc_node, do: Process.get(sim_node_key()) + + @spec send(pid | reference | atom | {atom, node}, term) :: term def send(dest, message) do - Scheduler.handle_sent(dest) - Kernel.send(dest, message) + case get_sim() do + nil -> Kernel.send(dest, message) + _sim -> sim_send(dest, message) + end + end + + @spec node :: node + def node do + case get_sim() do + nil -> Kernel.node() + _sim -> get_proc_node() + end end @spec spawn_node(atom, function) :: pid @@ -28,6 +42,61 @@ defmodule Trinity.SimProcess do end end + @spec unalias(reference) :: boolean + def unalias(alias) do + case get_sim() do + nil -> Process.unalias(alias) + _sim -> sim_unalias(alias) + end + end + + @spec register(pid, atom) :: true + def register(pid, name) do + case get_sim() do + nil -> Process.register(pid, name) + _sim -> sim_register(pid, name) + end + end + + defp sim_send(pid, message) when is_pid(pid) do + Scheduler.handle_sent(pid) + Kernel.send(pid, message) + end + + defp sim_send(alias, message) when is_reference(alias) do + %{proc_aliases: proc_aliases} = get_sim() + case :ets.lookup(proc_aliases, alias) do + [{_key, pid}] -> Scheduler.handle_sent(pid) + # If the alias doesn't exist, do nothing + [] -> :noop + end + + Kernel.send(alias, message) + end + + defp sim_send(name, message) when is_atom(name) do + %{proc_aliases: proc_aliases} = get_sim() + node = get_proc_node() + case :ets.lookup(proc_aliases, {name, node}) do + [{_key, pid}] -> + Scheduler.handle_sent(pid) + Kernel.send(pid, message) + + [] -> raise ArgumentError, "Invalid destination #{inspect(name)}" + end + end + + defp sim_send({_name, _node} = name_node, message) do + %{proc_aliases: proc_aliases} = get_sim() + case :ets.lookup(proc_aliases, name_node) do + [{_key, pid}] -> + Scheduler.handle_sent(pid) + Kernel.send(pid, message) + + [] -> raise ArgumentError, "Invalid destination #{inspect(name_node)}" + end + end + defp sim_alias do %Simulation{proc_aliases: proc_aliases} = get_sim() @@ -37,18 +106,23 @@ defmodule Trinity.SimProcess do alias end - @spec unalias(reference) :: boolean - def unalias(alias) do - case get_sim() do - nil -> Process.unalias(alias) - _sim -> sim_unalias(alias) - end - end - defp sim_unalias(alias) do %Simulation{proc_aliases: proc_aliases} = get_sim() :ets.delete(proc_aliases, alias) Process.unalias(alias) end + + defp sim_register(pid, name) do + %{proc_aliases: proc_aliases} = get_sim() + node = get_proc_node() + + key = {name, node} + case :ets.lookup(proc_aliases, key) do + [] -> :ets.insert(proc_aliases, {key, pid}) + [{_key, _pid}] -> raise ArgumentError, "Could not register #{inspect(pid)} with name #{inspect(name)}" + end + + true + end end diff --git a/test/trinity_test.exs b/test/trinity_test.exs index 5ca53f3..fcb4b39 100644 --- a/test/trinity_test.exs +++ b/test/trinity_test.exs @@ -26,10 +26,17 @@ defmodule TrinityTest do nodes = [:n1, :n2, :n3] pids = Enum.map(nodes, fn node -> + name = String.to_atom(Atom.to_string(node) <> "_proc") parent = self() ref = make_ref() SimProcess.spawn_node(node, fn -> + SimProcess.register(self(), name) + + receive_yield do + :begin -> dbg {"began", name} + end + pids = Enum.map(1..10, fn i -> {:ok, pid} = Counter.start_link(i) pid @@ -38,6 +45,8 @@ defmodule TrinityTest do SimProcess.send parent, {ref, pids} end) + SimProcess.send {name, node}, :begin + receive_yield do {^ref, pids} -> pids end