Skip to content

Commit d247ee6

Browse files
fix: keep ETS storage tables supervised
Route `Jido.Storage.ETS` table creation through a supervised owner process and give fallback-created tables a supervisor heir so short-lived callers do not drop storage tables. Fixes #302.
1 parent 790901a commit d247ee6

4 files changed

Lines changed: 125 additions & 4 deletions

File tree

lib/jido/application.ex

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -6,7 +6,9 @@ defmodule Jido.Application do
66
def start(_type, _args) do
77
Jido.Telemetry.setup()
88

9-
children = []
9+
children = [
10+
Jido.Storage.ETS.Owner
11+
]
1012

1113
register_signal_extensions()
1214
Jido.Discovery.init_async()

lib/jido/storage/ets.ex

Lines changed: 27 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -218,16 +218,31 @@ defmodule Jido.Storage.ETS do
218218
:"#{base}_thread_meta"
219219
end
220220

221-
defp ensure_tables(opts) do
221+
@doc false
222+
@spec create_tables(opts()) :: :ok
223+
def create_tables(opts) do
222224
ensure_table(checkpoint_table(opts), [:set])
223225
ensure_table(threads_table(opts), [:ordered_set])
224226
ensure_table(meta_table(opts), [:set])
225227
end
226228

229+
defp ensure_tables(opts) do
230+
case Jido.Storage.ETS.Owner.ensure_tables(opts) do
231+
:ok -> :ok
232+
{:error, _reason} -> create_tables(opts)
233+
end
234+
end
235+
227236
defp ensure_table(name, extra_opts) do
228237
case :ets.whereis(name) do
229238
:undefined ->
230-
:ets.new(name, [:named_table, :public, read_concurrency: true] ++ extra_opts)
239+
_ =
240+
:ets.new(
241+
name,
242+
[:named_table, :public, read_concurrency: true] ++ heir_opts(name) ++ extra_opts
243+
)
244+
245+
:ok
231246

232247
_ref ->
233248
:ok
@@ -236,6 +251,16 @@ defmodule Jido.Storage.ETS do
236251
ArgumentError -> :ok
237252
end
238253

254+
defp heir_opts(name) do
255+
case Process.whereis(Jido.Supervisor) do
256+
pid when is_pid(pid) and pid != self() ->
257+
[{:heir, pid, {:jido_storage_ets, name}}]
258+
259+
_other ->
260+
[]
261+
end
262+
end
263+
239264
defp get_current_rev(table, thread_id) do
240265
case :ets.select_reverse(table, [{{{thread_id, :"$1"}, :_}, [], [:"$1"]}], 1) do
241266
{[seq], _cont} -> seq + 1

lib/jido/storage/ets/owner.ex

Lines changed: 47 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,47 @@
1+
defmodule Jido.Storage.ETS.Owner do
2+
@moduledoc false
3+
4+
use GenServer
5+
6+
@name __MODULE__
7+
@call_timeout 5_000
8+
9+
@doc false
10+
@spec child_spec(keyword()) :: Supervisor.child_spec()
11+
def child_spec(opts) do
12+
%{
13+
id: @name,
14+
start: {@name, :start_link, [opts]},
15+
type: :worker,
16+
restart: :permanent,
17+
shutdown: 5_000
18+
}
19+
end
20+
21+
@doc false
22+
@spec start_link(keyword()) :: GenServer.on_start()
23+
def start_link(opts \\ []) do
24+
GenServer.start_link(__MODULE__, opts, name: @name)
25+
end
26+
27+
@doc false
28+
@spec ensure_tables(keyword()) :: :ok | {:error, term()}
29+
def ensure_tables(opts) when is_list(opts) do
30+
case GenServer.whereis(@name) do
31+
nil -> {:error, :not_started}
32+
pid -> GenServer.call(pid, {:ensure_tables, opts}, @call_timeout)
33+
end
34+
catch
35+
:exit, reason -> {:error, {:owner_unavailable, reason}}
36+
end
37+
38+
@impl true
39+
def init(_opts) do
40+
{:ok, %{}}
41+
end
42+
43+
@impl true
44+
def handle_call({:ensure_tables, opts}, _from, state) do
45+
{:reply, Jido.Storage.ETS.create_tables(opts), state}
46+
end
47+
end

test/jido/storage/ets_test.exs

Lines changed: 48 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -320,7 +320,7 @@ defmodule JidoTest.Storage.ETSTest do
320320
thread_id = "thread_#{System.unique_integer([:positive])}"
321321
total_appends = 40
322322

323-
# Ensure tables are created/owned by the test process.
323+
# Ensure tables are created before concurrent appends start.
324324
assert {:ok, _} = ETS.append_thread(thread_id, [%{kind: :note, payload: %{n: 0}}], opts)
325325

326326
results =
@@ -400,4 +400,51 @@ defmodule JidoTest.Storage.ETSTest do
400400
assert {:ok, %Thread{}} = ETS.load_thread("key1", opts)
401401
end
402402
end
403+
404+
describe "table ownership" do
405+
test "tables survive when first used by a short-lived process" do
406+
opts = [table: unique_table(:ephemeral_owner)]
407+
table = :"#{Keyword.fetch!(opts, :table)}_checkpoints"
408+
owner = Process.whereis(Jido.Storage.ETS.Owner)
409+
410+
assert is_pid(owner)
411+
412+
parent = self()
413+
414+
{caller, ref} =
415+
spawn_monitor(fn ->
416+
assert :ok = ETS.put_checkpoint(:ephemeral_key, :value, opts)
417+
send(parent, {:owner, :ets.info(table, :owner), self()})
418+
end)
419+
420+
assert_receive {:owner, ^owner, ^caller}
421+
assert_receive {:DOWN, ^ref, :process, ^caller, _reason}
422+
423+
assert :ets.whereis(table) != :undefined
424+
assert {:ok, :value} = ETS.get_checkpoint(:ephemeral_key, opts)
425+
end
426+
427+
test "fallback-created tables survive a short-lived creator" do
428+
opts = [table: unique_table(:fallback_ephemeral_owner)]
429+
table = :"#{Keyword.fetch!(opts, :table)}_checkpoints"
430+
supervisor = Process.whereis(Jido.Supervisor)
431+
432+
assert is_pid(supervisor)
433+
434+
parent = self()
435+
436+
{caller, ref} =
437+
spawn_monitor(fn ->
438+
assert :ok = ETS.create_tables(opts)
439+
true = :ets.insert(table, {:ephemeral_key, :value})
440+
send(parent, {:owner, :ets.info(table, :owner), self()})
441+
end)
442+
443+
assert_receive {:owner, ^caller, ^caller}
444+
assert_receive {:DOWN, ^ref, :process, ^caller, _reason}
445+
446+
assert :ets.info(table, :owner) == supervisor
447+
assert {:ok, :value} = ETS.get_checkpoint(:ephemeral_key, opts)
448+
end
449+
end
403450
end

0 commit comments

Comments
 (0)