Skip to content

Commit 5467582

Browse files
committed
Solution: Simplify TaskRegistry and make tests deterministic
1 parent 35905ad commit 5467582

5 files changed

Lines changed: 67 additions & 144 deletions

File tree

CHANGELOG.md

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -6,6 +6,7 @@ This project adheres to [Semantic Versioning](http://semver.org/).
66

77
### Fixed
88
- Properly override jobs with duplicate name (#392)
9+
- Simplify `TaskRegistry` and make tests deterministic
910

1011
Diff for [unreleased]
1112

lib/quantum/task_registry.ex

Lines changed: 26 additions & 97 deletions
Original file line numberDiff line numberDiff line change
@@ -3,17 +3,16 @@ defmodule Quantum.TaskRegistry do
33

44
# Registry to check if a task is already running on a node.
55

6-
use GenServer
7-
86
require Logger
97

10-
alias __MODULE__.{InitOpts, StartOpts, State}
8+
alias __MODULE__.StartOpts
9+
alias Quantum.Job
1110

1211
# Start the registry
1312
@spec start_link(StartOpts.t()) :: GenServer.on_start()
14-
def start_link(%StartOpts{name: name}) do
15-
__MODULE__
16-
|> GenServer.start_link(%InitOpts{}, name: name)
13+
def start_link(%StartOpts{name: name, listeners: listeners}) do
14+
[keys: :unique, name: name, listeners: listeners]
15+
|> Registry.start_link()
1716
|> case do
1817
{:ok, pid} ->
1918
{:ok, pid}
@@ -27,114 +26,44 @@ defmodule Quantum.TaskRegistry do
2726
end
2827
end
2928

29+
@spec child_spec(options :: StartOpts.t()) :: Supervisor.child_spec()
30+
def child_spec(options),
31+
do:
32+
[]
33+
|> Registry.child_spec()
34+
|> Map.put(:start, {__MODULE__, :start_link, [options]})
35+
3036
# Mark a task as Running
3137
#
3238
# ### Examples
3339
#
34-
# iex> Quantum.TaskRegistry.mark_running(server, running_job.name, self())
40+
# iex> Quantum.TaskRegistry.mark_running(server, running_job.name, Node.self())
3541
# :already_running
3642
#
37-
# iex> Quantum.TaskRegistry.mark_running(server, not_running_job.name, self())
43+
# iex> Quantum.TaskRegistry.mark_running(server, not_running_job.name, Node.self())
3844
# :marked_running
45+
@spec mark_running(server :: atom, task :: Job.name(), node :: Node.t()) ::
46+
:already_running | :marked_running
3947
def mark_running(server, task, node) do
40-
GenServer.call(server, {:running, task, node})
48+
server
49+
|> Registry.register({task, node}, true)
50+
|> case do
51+
{:ok, _pid} -> :marked_running
52+
{:error, {:already_registered, _other_pid}} -> :already_running
53+
end
4154
end
4255

4356
# Mark a task as Finished
4457
#
4558
# ### Examples
4659
#
47-
# iex> Quantum.TaskRegistry.mark_running(server, running_job.name, self())
60+
# iex> Quantum.TaskRegistry.mark_running(server, running_job.name, Node.self())
4861
# :ok
4962
#
50-
# iex> Quantum.TaskRegistry.mark_running(server, not_running_job.name, self())
63+
# iex> Quantum.TaskRegistry.mark_running(server, not_running_job.name, Node.self())
5164
# :ok
65+
@spec mark_finished(server :: atom, task :: Job.name(), node :: Node.t()) :: :ok
5266
def mark_finished(server, task, node) do
53-
GenServer.cast(server, {:finished, task, node})
54-
end
55-
56-
# Query if a task with given name is running
57-
#
58-
# ### Examples
59-
#
60-
# iex> Quantum.TaskRegistry.is_running?(server, running_job.name)
61-
# true
62-
#
63-
# iex> Quantum.TaskRegistry.is_running?(server, not_running_job.name)
64-
# false
65-
def is_running?(server, task) do
66-
GenServer.call(server, {:is_running?, task})
67-
end
68-
69-
# Query if any tasks are running in the cluster
70-
#
71-
# ### Examples
72-
#
73-
# iex> Quantum.TaskRegistry.any_running?(server_with_running_tasks)
74-
# true
75-
#
76-
# iex> Quantum.TaskRegistry.any_running?(server_without_running_tasks)
77-
# false
78-
def any_running?(server) do
79-
GenServer.call(server, :any_running?)
80-
end
81-
82-
@impl GenServer
83-
def init(%InitOpts{}) do
84-
{:ok, %State{running_tasks: %{}}}
85-
end
86-
87-
@impl GenServer
88-
def handle_call({:running, task, node}, _caller, %State{running_tasks: running_tasks} = state) do
89-
if Enum.member?(Map.get(running_tasks, task, []), node) do
90-
{:reply, :already_running, state}
91-
else
92-
{:reply, :marked_running,
93-
%{state | running_tasks: Map.update(running_tasks, task, [node], &[node | &1])}}
94-
end
95-
end
96-
97-
@impl GenServer
98-
def handle_call({:is_running?, task}, _caller, %State{running_tasks: running_tasks} = state) do
99-
case running_tasks do
100-
%{^task => [_ | _]} ->
101-
{:reply, true, state}
102-
103-
%{^task => []} ->
104-
{:reply, false, state}
105-
106-
%{} ->
107-
{:reply, false, state}
108-
end
109-
end
110-
111-
@impl GenServer
112-
def handle_call(:any_running?, _caller, %State{running_tasks: running_tasks} = state) do
113-
if Enum.empty?(running_tasks) do
114-
{:reply, false, state}
115-
else
116-
{:reply, true, state}
117-
end
118-
end
119-
120-
@impl GenServer
121-
def handle_cast({:finished, task, node}, %State{running_tasks: running_tasks} = state) do
122-
running_tasks =
123-
running_tasks
124-
|> Map.update(task, [], &(&1 -- [node]))
125-
|> case do
126-
%{^task => []} = still_running_tasks ->
127-
Map.delete(still_running_tasks, task)
128-
129-
still_running_tasks ->
130-
still_running_tasks
131-
end
132-
133-
{:noreply, %{state | running_tasks: running_tasks}}
134-
end
135-
136-
@impl GenServer
137-
def handle_info(_message, state) do
138-
{:noreply, [], state}
67+
Registry.unregister(server, {task, node})
13968
end
14069
end

lib/quantum/task_registry/start_opts.ex

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -4,9 +4,10 @@ defmodule Quantum.TaskRegistry.StartOpts do
44
# Start Options for Quantum.TaskRegistry
55

66
@type t :: %__MODULE__{
7-
name: GenServer.server()
7+
name: GenServer.server(),
8+
listeners: [atom]
89
}
910

1011
@enforce_keys [:name]
11-
defstruct @enforce_keys
12+
defstruct @enforce_keys ++ [listeners: []]
1213
end

test/quantum/executor_test.exs

Lines changed: 34 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -17,20 +17,28 @@ defmodule Quantum.ExecutorTest do
1717
use Quantum, otp_app: :job_broadcaster_test
1818
end
1919

20-
setup do
20+
setup tags do
2121
{:ok, _task_supervisor} =
2222
start_supervised({Task.Supervisor, [name: Module.concat(__MODULE__, TaskSupervisor)]})
2323

24-
{:ok, task_registry} =
24+
process_name = Module.concat(__MODULE__, tags.test)
25+
26+
Process.register(self(), process_name)
27+
28+
{:ok, _task_registry} =
2529
start_supervised(
26-
{TaskRegistry, %TaskRegistryStartOpts{name: Module.concat(__MODULE__, TaskRegistry)}}
30+
{TaskRegistry,
31+
%TaskRegistryStartOpts{
32+
name: Module.concat(__MODULE__, TaskRegistry),
33+
listeners: [process_name]
34+
}}
2735
)
2836

2937
{
3038
:ok,
3139
%{
3240
task_supervisor: Module.concat(__MODULE__, TaskSupervisor),
33-
task_registry: task_registry,
41+
task_registry: Module.concat(__MODULE__, TaskRegistry),
3442
debug_logging: true
3543
}
3644
}
@@ -97,8 +105,8 @@ defmodule Quantum.ExecutorTest do
97105
job =
98106
TestScheduler.new_job()
99107
|> Job.set_task(fn ->
100-
Process.sleep(50)
101108
send(caller, :executed)
109+
Process.sleep(500)
102110
end)
103111
|> Job.set_overlap(false)
104112

@@ -136,28 +144,38 @@ defmodule Quantum.ExecutorTest do
136144
job =
137145
TestScheduler.new_job()
138146
|> Job.set_task(fn ->
139-
Process.sleep(50)
140-
send(caller, :executed)
147+
send(caller, {:executing, self()})
148+
149+
receive do
150+
:continue -> nil
151+
end
152+
153+
send(caller, :execution_end)
141154
end)
142155
|> Job.set_overlap(false)
143156

157+
job_name = job.name
158+
node = Node.self()
159+
144160
capture_log(fn ->
145161
Executor.start_link(
146162
%StartOpts{
147163
task_supervisor_reference: task_supervisor,
148164
task_registry_reference: task_registry,
149165
debug_logging: debug_logging
150166
},
151-
%Event{job: job, node: Node.self()}
167+
%Event{job: job, node: node}
152168
)
153169

154-
# Wait until running
155-
Process.sleep(25)
170+
assert_receive {:executing, job_pid}
156171

157172
assert :already_running = TaskRegistry.mark_running(task_registry, job.name, Node.self())
158173

159-
assert_receive :executed
160-
refute_receive :executed
174+
send(job_pid, :continue)
175+
176+
assert_receive :execution_end
177+
178+
assert_receive {:unregister, _, {^job_name, ^node}, _pid}
161179

162180
assert :marked_running = TaskRegistry.mark_running(task_registry, job.name, Node.self())
163181
end)
@@ -173,6 +191,9 @@ defmodule Quantum.ExecutorTest do
173191
|> Job.set_task(fn -> raise "failed" end)
174192
|> Job.set_overlap(false)
175193

194+
job_name = job.name
195+
node = Node.self()
196+
176197
# Mute Error
177198
capture_log(fn ->
178199
Executor.start_link(
@@ -184,7 +205,7 @@ defmodule Quantum.ExecutorTest do
184205
%Event{job: job, node: Node.self()}
185206
)
186207

187-
Process.sleep(150)
208+
assert_receive {:unregister, _, {^job_name, ^node}, _pid}
188209
end)
189210

190211
assert :marked_running = TaskRegistry.mark_running(task_registry, job.name, Node.self())

test/quantum/task_registry_test.exs

Lines changed: 3 additions & 32 deletions
Original file line numberDiff line numberDiff line change
@@ -7,12 +7,12 @@ defmodule Quantum.TaskRegistryTest do
77
alias Quantum.TaskRegistry.StartOpts
88

99
doctest TaskRegistry,
10-
except: [mark_running: 3, mark_finished: 3, is_running?: 2, any_running?: 1]
10+
except: [mark_running: 3, mark_finished: 3]
1111

1212
setup do
13-
{:ok, registry} = start_supervised({TaskRegistry, %StartOpts{name: __MODULE__}})
13+
{:ok, _registry} = start_supervised({TaskRegistry, %StartOpts{name: __MODULE__}})
1414

15-
{:ok, %{registry: registry}}
15+
{:ok, %{registry: __MODULE__}}
1616
end
1717

1818
describe "running" do
@@ -46,33 +46,4 @@ defmodule Quantum.TaskRegistryTest do
4646
assert :ok = TaskRegistry.mark_finished(registry, task, self())
4747
end
4848
end
49-
50-
describe "is_running?" do
51-
test "not running", %{registry: registry} do
52-
task = make_ref()
53-
assert false == TaskRegistry.is_running?(registry, task)
54-
end
55-
56-
test "running", %{registry: registry} do
57-
task = make_ref()
58-
59-
TaskRegistry.mark_running(registry, task, self())
60-
61-
assert true == TaskRegistry.is_running?(registry, task)
62-
end
63-
end
64-
65-
describe "any_running?" do
66-
test "not running", %{registry: registry} do
67-
assert false == TaskRegistry.any_running?(registry)
68-
end
69-
70-
test "running", %{registry: registry} do
71-
task = make_ref()
72-
73-
TaskRegistry.mark_running(registry, task, self())
74-
75-
assert true == TaskRegistry.any_running?(registry)
76-
end
77-
end
7849
end

0 commit comments

Comments
 (0)