diff --git a/lib/trinity/scheduler.ex b/lib/trinity/scheduler.ex index da997ac..2f358aa 100644 --- a/lib/trinity/scheduler.ex +++ b/lib/trinity/scheduler.ex @@ -6,6 +6,8 @@ defmodule Trinity.Scheduler do queue: :ets.table, proc_queue_keys: :ets.table, proc_links: :ets.table, + proc_aliases: :ets.table, + now: :atomics.atomics_ref, supervisor_pid: pid, } @@ -14,6 +16,7 @@ defmodule Trinity.Scheduler do :proc_queue_keys, :proc_links, :proc_aliases, + :now, :supervisor_pid, ] @@ -22,18 +25,37 @@ defmodule Trinity.Scheduler do defmacro simulation_key, do: :_trinity_simulation + @spec get_sim :: Simulation.t defp get_sim, do: Process.get(simulation_key()) def run_simulation(fun) do SimulationSupervisor.run_simulation(fun) end + @spec now :: integer + def now do + %{now: now} = get_sim() + read_now(now) + end + @spec yield(non_neg_integer) :: :ok def yield(delay \\ 0) do %Simulation{queue: queue, proc_queue_keys: proc_queue_keys, now: now} = get_sim() - yield_ref = enqueue_self(queue, proc_queue_keys, now, delay) + yield_ref = enqueue_resume(queue, proc_queue_keys, now, delay) + + perform_next(queue, proc_queue_keys, now) + + receive do + ^yield_ref -> :ok + end + end + + @spec yield_until_message(non_neg_integer) :: :ok + def yield_until_message(timeout) do + %Simulation{queue: queue, proc_queue_keys: proc_queue_keys, now: now} = get_sim() + yield_ref = enqueue_timeout(queue, proc_queue_keys, now, timeout) - perform_next(queue, now) + perform_next(queue, proc_queue_keys, now) receive do ^yield_ref -> :ok @@ -50,7 +72,7 @@ defmodule Trinity.Scheduler do # 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_self(sim.queue, sim.proc_queue_keys, sim.now, 0) + yield_ref = enqueue_resume(sim.queue, sim.proc_queue_keys, sim.now, 0) send parent_pid, spawn_ref receive do ^yield_ref -> :noop @@ -94,7 +116,12 @@ defmodule Trinity.Scheduler do end end - case Trinity.Scheduler.receive_loop(receive_fun, ref, unquote(timeout)) do + start = Trinity.Scheduler.now() + # This will always perform one receive() check synchronously before yielding, which is + # needed to avoid a race where the message is already in the queue so we never unsuspend + # + # TODO: we could perform a yield(0) here to ensure all calls to `receive_yield()` actually yield + case Trinity.Scheduler.receive_loop(receive_fun, ref, unquote(timeout), start) do # We re-use the same ref to detect a timeout ^ref -> unquote(timeout_value) value -> value @@ -103,17 +130,29 @@ defmodule Trinity.Scheduler do end @doc false - def receive_loop(receive_fun, ref, timeout) do + def receive_loop(receive_fun, ref, timeout, start) do case receive_fun.() do ^ref -> # TODO: timeout - yield(10) - receive_loop(receive_fun, ref, timeout) + now = Trinity.Scheduler.now() + case (now - start) >= timeout do + true -> + ref + + false -> + remaining = timeout_remaining(timeout, now - start) + yield_until_message(remaining) + receive_loop(receive_fun, ref, timeout, start) + end + value -> value end end + defp timeout_remaining(:infinity, _elapsed), do: :infinity + defp timeout_remaining(timeout, elapsed), do: timeout - elapsed + @spec add_link(pid) :: :ok def add_link(to_pid) do # TODO: noproc if `to_pid` has exited @@ -143,12 +182,13 @@ defmodule Trinity.Scheduler do end def handle_down(%Simulation{} = sim, _pid, :normal) do - perform_next(sim.queue, sim.now) + %{queue: queue, proc_queue_keys: proc_queue_keys} = sim + perform_next(queue, proc_queue_keys, sim.now) [] end def handle_down(%Simulation{} = sim, pid, reason) do - %Simulation{proc_links: proc_links} = sim + %{queue: queue, proc_queue_keys: proc_queue_keys, proc_links: proc_links} = sim linked = gather_linked(pid, proc_links) Enum.each(linked, fn pid -> @@ -156,44 +196,100 @@ defmodule Trinity.Scheduler do Process.exit(pid, reason) end) - perform_next(sim.queue, sim.now) + perform_next(queue, proc_queue_keys, sim.now) linked end - defp perform_next(queue, now) do - case pop_next(queue) do - {time, {:resume, pid, ref}} -> + 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() + + case :ets.lookup(proc_queue_keys, dest_pid) do + [{_key, {:timeout_infinity, ref}}] -> + # This is an infinite timeout so we add the dest process to the queue + queue_key = enqueue_event(queue, now, 0, {:resume, dest_pid, ref}) + :ets.insert(proc_queue_keys, {dest_pid, {:resume, queue_key}}) + + [{_key, {:timeout, ref, existing_queue_key}}] -> + # This is a finite timeout so we remove the existing (:timeout) queue entry + :ets.delete(queue, existing_queue_key) + # Then add the dest process back to the queue at now() + queue_key = enqueue_event(queue, now, 0, {:resume, dest_pid, ref}) + :ets.insert(proc_queue_keys, {dest_pid, {:resume, queue_key}}) + + _ -> + # Process is already queued for :resume or unknown (dead) + :noop + end + end + + defp perform_next(queue, proc_queue_keys, now) do + case pop_next(queue, proc_queue_keys) do + {time, {_event, pid, ref}} -> set_now(now, time) send pid, ref end end - defp pop_next(queue) do + defp pop_next(queue, proc_queue_keys) do case :ets.next_lookup(queue, {-1, 0}) do {_, [{{time, _i} = key, entry}]} -> :ets.delete(queue, key) + {_event, pid, _ref} = entry + :ets.delete(proc_queue_keys, pid) + {time, entry} _ -> raise "Queue empty!" end end - defp enqueue_self(queue, proc_queue_keys, now, delay) do + defp enqueue_timeout(_queue, proc_queue_keys, _now, :infinity) do ref = make_ref() - time = read_now(now) + delay + pid = self() + :ets.insert(proc_queue_keys, {pid, {:timeout_infinity, ref}}) + + ref + end + + defp enqueue_timeout(queue, proc_queue_keys, now, timeout) do + ref = make_ref() + pid = self() + queue_key = enqueue_event(queue, now, timeout, {:timeout, pid, ref}) + :ets.insert(proc_queue_keys, {pid, {:timeout, ref, queue_key}}) + + ref + end + defp enqueue_resume(queue, proc_queue_keys, now, delay) do + ref = make_ref() pid = self() - entry = {:resume, pid, ref} + queue_key = enqueue_event(queue, now, delay, {:resume, pid, ref}) + :ets.insert(proc_queue_keys, {pid, {:resume, queue_key}}) + + ref + end + + defp enqueue_event(queue, now, delay, event) do + time = read_now(now) + delay i = case :ets.prev(queue, {time, :infinity}) do {^time, prev} -> prev + 1 _ -> 0 end queue_key = {time, i} - :ets.insert(queue, {{time, i}, entry}) - :ets.insert(proc_queue_keys, {pid, queue_key}) + :ets.insert(queue, {{time, i}, event}) - ref + queue_key end defp destroy_process(%Simulation{} = sim, pid) do @@ -206,7 +302,7 @@ defmodule Trinity.Scheduler do destroy_links(proc_links, pid) case :ets.lookup(proc_queue_keys, pid) do - [{^pid, qk}] -> + [{_event, ^pid, qk}] -> :ets.delete(queue, qk) :ets.delete(proc_queue_keys, pid) _ -> :noop diff --git a/lib/trinity/scheduler/simulation_supervisor.ex b/lib/trinity/scheduler/simulation_supervisor.ex index c2bcd88..46ee3d0 100644 --- a/lib/trinity/scheduler/simulation_supervisor.ex +++ b/lib/trinity/scheduler/simulation_supervisor.ex @@ -15,6 +15,7 @@ defmodule Trinity.Scheduler.SimulationSupervisor do defstruct @enforce_keys end + @spec run_simulation(function) :: :ok def run_simulation(fun) do ref = make_ref() GenServer.start_link(__MODULE__, %{ @@ -24,7 +25,10 @@ defmodule Trinity.Scheduler.SimulationSupervisor do }) receive do - ^ref -> :ok + {^ref, :normal} -> :ok + {^ref, :shutdown} -> :ok + # Propagate root exit to parent process so tests fail + {^ref, reason} -> exit(reason) end end @@ -41,6 +45,7 @@ defmodule Trinity.Scheduler.SimulationSupervisor do proc_queue_keys: :ets.new(__MODULE__, [:set, :public]), proc_links: :ets.new(__MODULE__, [:set, :public]), proc_aliases: :ets.new(__MODULE__, [:set, :public]), + now: :atomics.new(1, signed: false), supervisor_pid: self(), } @@ -78,7 +83,7 @@ defmodule Trinity.Scheduler.SimulationSupervisor do ^root_pid -> # The root has died, end the simulation - send parent_pid, parent_ref + send parent_pid, {parent_ref, reason} {:noreply, state} _ -> @@ -92,7 +97,7 @@ defmodule Trinity.Scheduler.SimulationSupervisor do # If the root was killed, end the simulation case pid == root_pid do - true -> send parent_pid, parent_ref + true -> send parent_pid, {parent_ref, reason} _ -> :noop end end) diff --git a/lib/trinity/sim_process.ex b/lib/trinity/sim_process.ex index dc0e00e..8abba8b 100644 --- a/lib/trinity/sim_process.ex +++ b/lib/trinity/sim_process.ex @@ -1,4 +1,5 @@ defmodule Trinity.SimProcess do + alias Trinity.Scheduler alias Trinity.Scheduler.Simulation import Trinity.Scheduler, only: [simulation_key: 0] @@ -7,6 +8,7 @@ defmodule Trinity.SimProcess do defp get_sim, do: Process.get(simulation_key()) def send(dest, message) do + Scheduler.handle_sent(dest) Kernel.send(dest, message) end