diff --git a/lib/trinity/scheduler.ex b/lib/trinity/scheduler.ex index 2f358aa..1549dd3 100644 --- a/lib/trinity/scheduler.ex +++ b/lib/trinity/scheduler.ex @@ -8,6 +8,9 @@ defmodule Trinity.Scheduler do proc_links: :ets.table, proc_aliases: :ets.table, + proc_nodes: :ets.table, + node_procs: :ets.table, + now: :atomics.atomics_ref, supervisor_pid: pid, } @@ -17,6 +20,9 @@ defmodule Trinity.Scheduler do :proc_links, :proc_aliases, + :proc_nodes, + :node_procs, + :now, :supervisor_pid, ] @@ -64,15 +70,32 @@ defmodule Trinity.Scheduler do @spec spawn_and_yield(function) :: pid def spawn_and_yield(fun, link? \\ false) do - sim = get_sim() + spawn_node_and_yield(fun, nil, link?) + end + + @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) + parent_pid = self() spawn_ref = make_ref() spawned_pid = SimulationSupervisor.spawn_child(sim.supervisor_pid, fn -> + %{ + queue: queue, + proc_queue_keys: proc_queue_keys, + proc_nodes: proc_nodes, + node_procs: node_procs, + now: now, + } = sim + + # Update the process/node mapping in both directions + put_proc_node(proc_nodes, node_procs, node) # Enqueue the new process and unsuspend the parent # so that it can register any links before yielding back # to us (eventually, if there are others in queue at `time=now`) - yield_ref = enqueue_resume(sim.queue, sim.proc_queue_keys, sim.now, 0) + yield_ref = enqueue_resume(queue, proc_queue_keys, now, 0) send parent_pid, spawn_ref receive do ^yield_ref -> :noop @@ -181,6 +204,13 @@ defmodule Trinity.Scheduler do :ok end + @spec set_up_root(atom) :: :ok + def set_up_root(root_node) do + %{proc_nodes: proc_nodes, node_procs: node_procs} = get_sim() + put_proc_node(proc_nodes, node_procs, root_node) + :ok + end + def handle_down(%Simulation{} = sim, _pid, :normal) do %{queue: queue, proc_queue_keys: proc_queue_keys} = sim perform_next(queue, proc_queue_keys, sim.now) @@ -297,9 +327,13 @@ defmodule Trinity.Scheduler do queue: queue, proc_links: proc_links, proc_queue_keys: proc_queue_keys, + proc_nodes: proc_nodes, + node_procs: node_procs, } = sim destroy_links(proc_links, pid) + destroy_proc_node(proc_nodes, node_procs, pid) + # Note that we do not currently bother to destroy aliases here case :ets.lookup(proc_queue_keys, pid) do [{_event, ^pid, qk}] -> @@ -352,6 +386,29 @@ defmodule Trinity.Scheduler do end end + defp get_proc_node(proc_nodes) do + [{_pid, node}] = :ets.lookup(proc_nodes, self()) + node + end + + defp put_proc_node(proc_nodes, node_procs, node) do + pid = self() + :ets.insert(proc_nodes, {pid, node}) + existing_procs = case :ets.lookup(node_procs, node) do + [{_node, procs}] -> procs + [] -> [] + end + :ets.insert(node_procs, {node, [pid | existing_procs]}) + end + + defp destroy_proc_node(proc_nodes, node_procs, pid) do + [{_pid, node}] = :ets.lookup(proc_nodes, pid) + :ets.delete(proc_nodes, pid) + + [{_node, existing_procs}] = :ets.lookup(node_procs, node) + :ets.insert(node_procs, {node, List.delete(existing_procs, pid)}) + end + defp read_now(now), do: :atomics.get(now, 1) defp set_now(now, time), do: :atomics.put(now, 1, time) @@ -361,6 +418,10 @@ defmodule Trinity.Scheduler do %{ queue: :ets.tab2list(sim.queue), proc_links: :ets.tab2list(sim.proc_links), + + proc_nodes: :ets.tab2list(sim.proc_nodes), + node_procs: :ets.tab2list(sim.node_procs), + now: read_now(sim.now), supervisor_pid: sim.supervisor_pid, } diff --git a/lib/trinity/scheduler/simulation_supervisor.ex b/lib/trinity/scheduler/simulation_supervisor.ex index 46ee3d0..4b0008f 100644 --- a/lib/trinity/scheduler/simulation_supervisor.ex +++ b/lib/trinity/scheduler/simulation_supervisor.ex @@ -46,11 +46,19 @@ defmodule Trinity.Scheduler.SimulationSupervisor do proc_links: :ets.new(__MODULE__, [:set, :public]), proc_aliases: :ets.new(__MODULE__, [:set, :public]), + proc_nodes: :ets.new(__MODULE__, [:set, :public]), + node_procs: :ets.new(__MODULE__, [:set, :public]), + now: :atomics.new(1, signed: false), supervisor_pid: self(), } - root_pid = spawn_sim_child(sim, fun) + # TODO: configurable? + root_node = :nonode + root_pid = spawn_sim_child(sim, fn -> + Trinity.Scheduler.set_up_root(root_node) + fun.() + end) state = %State{ sim: sim, diff --git a/lib/trinity/sim_process.ex b/lib/trinity/sim_process.ex index 8abba8b..ea12e31 100644 --- a/lib/trinity/sim_process.ex +++ b/lib/trinity/sim_process.ex @@ -12,6 +12,14 @@ defmodule Trinity.SimProcess do Kernel.send(dest, message) end + @spec spawn_node(atom, function) :: pid + def spawn_node(node, fun) do + case get_sim() do + nil -> raise "spawn_node() only works in simulation" + _sim -> Scheduler.spawn_node_and_yield(fun, node, false) + end + end + @spec alias :: reference def alias do case get_sim() do diff --git a/lib/trinity/sim_server.ex b/lib/trinity/sim_server.ex index bfe954a..7a451f5 100644 --- a/lib/trinity/sim_server.ex +++ b/lib/trinity/sim_server.ex @@ -11,6 +11,10 @@ defmodule Trinity.SimServer do SimGen.start_link(module, arg, options) end + def start(module, arg, options \\ []) do + SimGen.start(module, arg, options) + end + @spec call(GenServer.server, term, timeout) :: term def call(server, request, timeout \\ 5000) do case get_sim() do diff --git a/test/trinity_test.exs b/test/trinity_test.exs index 2b1669a..639bfa1 100644 --- a/test/trinity_test.exs +++ b/test/trinity_test.exs @@ -1,7 +1,7 @@ defmodule TrinityTest do use ExUnit.Case - alias Trinity.Scheduler + alias Trinity.{SimProcess, Scheduler} defmodule Counter do use GenServer @@ -22,11 +22,16 @@ defmodule TrinityTest do test "scheduler" do Scheduler.run_simulation(fn -> + nodes = [:n1, :n2, :n3] pids = - Enum.map(1..10, fn i -> - {:ok, pid} = Counter.start_link(i) - pid + Enum.map(nodes, fn i -> + Enum.map(1..10, fn i -> + {:ok, pid} = Counter.start_link(i) + pid + end) end) + |> Enum.concat() + dbg pids Enum.each(pids, fn pid ->