Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
14 changes: 14 additions & 0 deletions test/jepsen/harness_modules.exs
Original file line number Diff line number Diff line change
@@ -0,0 +1,14 @@
# Require once so independent harness regressions share the same real modules.
# Keep top-level executable startup, transport adapters, and unrelated modules
# out of the test VM. Log paths are configured at runtime, never rewritten here.
path = Path.join(__DIR__, "node.exs")
{:__block__, metadata, forms} = path |> File.read!() |> Code.string_to_quoted!()
modules = [Group.Jepsen.Transport.Stats, Group.Jepsen.Owner, Group.Jepsen.Driver]

forms =
Enum.filter(forms, fn
{:defmodule, _, [{:__aliases__, _, parts}, _]} -> Module.concat(parts) in modules
_ -> false
end)

Code.compile_quoted({:__block__, metadata, forms}, path)
100 changes: 49 additions & 51 deletions test/jepsen/node.exs
Original file line number Diff line number Diff line change
Expand Up @@ -19,7 +19,7 @@ defmodule Group.Jepsen.Transport.Stats do

def increment_persistent(event, amount \\ 1) do
increment(event, amount)
File.write!(@persistent_event_log, "#{event}\t#{amount}\n", [:append])
File.write!(persistent_event_log(), "#{event}\t#{amount}\n", [:append])
:ok
end

Expand Down Expand Up @@ -52,8 +52,11 @@ defmodule Group.Jepsen.Transport.Stats do
{:ok, %{}}
end

defp persistent_event_log,
do: System.get_env("GROUP_JEPSEN_PERSISTENT_EVENT_LOG", @persistent_event_log)

defp persistent_events do
case File.read(@persistent_event_log) do
case File.read(persistent_event_log()) do
{:ok, contents} ->
contents
|> String.split("\n", trim: true)
Expand Down Expand Up @@ -508,7 +511,7 @@ defmodule Group.Jepsen.Owner do
registrations = drop_cluster(state.registrations, cluster)
memberships = drop_cluster(state.memberships, cluster)
state = %{state | registrations: registrations, memberships: memberships}
{:reply, :ok, state}
{:reply, snapshot(state), state}
end

def handle_call(:snapshot, _from, state), do: {:reply, snapshot(state), state}
Expand Down Expand Up @@ -582,7 +585,13 @@ defmodule Group.Jepsen.Driver do
def kill(logical_owner), do: GenServer.call(driver(logical_owner), {:kill, logical_owner})

def drop_cluster(cluster) do
Enum.each(names(), &GenServer.call(&1, {:drop_cluster, cluster}, 10_000))
Enum.each(names(), fn driver ->
case GenServer.call(driver, {:drop_cluster, cluster}, 30_000) do
:ok -> :ok
{:error, reason} -> raise "owner cluster cleanup failed: #{inspect(reason)}"
end
end)

:ok
end

Expand Down Expand Up @@ -661,61 +670,47 @@ defmodule Group.Jepsen.Driver do
end

def handle_call({:drop_cluster, cluster}, _from, state) do
Enum.each(state.owners, fn {_logical_owner, {pid, _token, _monitor, _owner_state}} ->
if Process.alive?(pid), do: GenServer.call(pid, {:drop_cluster, cluster}, 10_000)
end)

owners =
Map.new(state.owners, fn {logical_owner, {pid, token, monitor, _owner_state}} ->
owner_state = if Process.alive?(pid), do: GenServer.call(pid, :snapshot), else: nil
{logical_owner, {pid, token, monitor, owner_state}}
end)

{:reply, :ok, %{state | owners: owners}}
end

def handle_call(:owner_snapshots, _from, state) do
result =
Enum.reduce_while(state.owners, {[], %{}}, fn
{logical_owner, {pid, token, monitor_ref, cached}}, {snapshots, acc} ->
case live_owner_snapshot(pid) do
{:ok, owner_state} ->
{:cont,
{
[owner_state | snapshots],
Map.put(acc, logical_owner, {pid, token, monitor_ref, owner_state})
}}

{:error, reason} ->
{:halt, {:error, {logical_owner, token, reason, cached}}}
end
end)

case result do
{:error, reason} ->
{:reply, {:error, reason}, state}

{owners, refreshed} ->
{:reply, {:ok, Enum.reverse(owners)}, %{state | owners: refreshed}}
case refresh_owners(state, {:drop_cluster, cluster}) do
{:ok, _snapshots, state} -> {:reply, :ok, state}
{:error, reason, state} -> {:reply, {:error, reason}, state}
end
end

defp live_owner_snapshot(pid) do
if Process.alive?(pid) do
try do
{:ok, GenServer.call(pid, :snapshot, 10_000)}
catch
:exit, reason -> {:error, reason}
end
else
{:error, :not_alive}
def handle_call(:owner_snapshots, _from, state) do
case refresh_owners(state, :snapshot) do
{:ok, snapshots, state} -> {:reply, {:ok, snapshots}, state}
{:error, reason, state} -> {:reply, {:error, reason}, state}
end
end

def handle_call(:unexpected_deaths, _from, state) do
{:reply, state.unexpected_deaths, state}
end

# Cleanup returns its snapshot in the same Owner turn. There is no second
# unguarded call, and every confirmed death uses the normal monitor path.
defp refresh_owners(state, request) do
Enum.reduce_while(state.owners, {:ok, [], state}, fn
{logical_owner, {pid, token, monitor_ref, cached}}, {:ok, snapshots, state} ->
try do
owner_state = GenServer.call(pid, request, 10_000)
state = put_owner_state(state, logical_owner, pid, owner_state)
{:cont, {:ok, [owner_state | snapshots], state}}
catch
:exit, reason ->
# A timeout is not proof of death. Only the tracked monitor can
# remove this owner, preserving every unrelated owner and monitor.
receive do
{:DOWN, ^monitor_ref, :process, ^pid, _death_reason} = down ->
{:noreply, state} = handle_info(down, state)
{:cont, {:ok, snapshots, state}}
after
0 -> {:halt, {:error, {logical_owner, token, reason, cached}, state}}
end
end
end)
end

@impl true
def handle_info({:DOWN, monitor_ref, :process, _pid, reason}, state) do
case Map.pop(state.monitors, monitor_ref) do
Expand Down Expand Up @@ -789,11 +784,14 @@ defmodule Group.Jepsen.Driver do
defp name(index), do: :"group_jepsen_driver_#{index}"

defp persist_unexpected_death(%{token: token, reason: reason}) do
File.write(@unexpected_death_log, token <> "\t" <> reason <> "\n", [:append])
File.write(unexpected_death_log(), token <> "\t" <> reason <> "\n", [:append])
end

defp unexpected_death_log,
do: System.get_env("GROUP_JEPSEN_UNEXPECTED_DEATH_LOG", @unexpected_death_log)

defp persisted_unexpected_deaths do
case File.read(@unexpected_death_log) do
case File.read(unexpected_death_log()) do
{:ok, contents} ->
contents
|> String.split("\n", trim: true)
Expand Down
139 changes: 139 additions & 0 deletions test/jepsen_cleanup_deaths_test.exs
Original file line number Diff line number Diff line change
@@ -0,0 +1,139 @@
defmodule Group.JepsenCleanupDeathsTest do
use ExUnit.Case, async: false
@moduletag :local
@moduletag :tmp_dir
Code.require_file("jepsen/harness_modules.exs", __DIR__)

setup %{tmp_dir: tmp_dir} do
for variable <- ["GROUP_JEPSEN_UNEXPECTED_DEATH_LOG", "GROUP_JEPSEN_PERSISTENT_EVENT_LOG"] do
previous = System.get_env(variable)
System.put_env(variable, Path.join(tmp_dir, variable))

on_exit(fn ->
if previous, do: System.put_env(variable, previous), else: System.delete_env(variable)
end)
end

start_supervised!({Group, name: :jepsen_group, shards: 1, log: false})

driver =
start_supervised!({Group.Jepsen.Driver, [index: 0, node_id: "n1", boot_id: "boot"]})

for logical <- ["loser", "survivor"] do
assert %{status: :ok} =
GenServer.call(driver, {:mutate, :join, logical, nil, 0, 1})
end

owners = :sys.get_state(driver).owners
{loser, _, _, _} = owners["loser"]
{survivor, survivor_token, _, _} = owners["survivor"]

on_exit(fn ->
for pid <- [loser, survivor], Process.alive?(pid), do: Process.exit(pid, :kill)
end)

%{driver: driver, loser: loser, survivor: survivor, token: survivor_token}
end

for request <- [{:drop_cluster, "red"}, :owner_snapshots] do
@request request
test "monitored conflict death during #{inspect(request)} preserves other owners", ctx do
:ok = :sys.suspend(ctx.loser)
task = Task.async(fn -> GenServer.call(ctx.driver, @request, 15_000) end)
expected = if @request == :owner_snapshots, do: :snapshot, else: @request

Group.LocalCase.wait_until(fn ->
{:messages, messages} = Process.info(ctx.loser, :messages)
Enum.any?(messages, &match?({:"$gen_call", _, ^expected}, &1))
end)

Process.exit(
ctx.loser,
{:group_registry_conflict, "jepsen/registry/0", %{token: "winner", revision: 2}}
)

response = Task.await(task, 15_000)

if @request == :owner_snapshots,
do: assert(match?({:ok, [_]}, response)),
else: assert(response == :ok)

assert Process.alive?(ctx.driver)
assert Process.alive?(ctx.survivor)

assert {:ok, [%{token: token, memberships: [_]}]} =
GenServer.call(ctx.driver, :owner_snapshots)

assert token == ctx.token
state = :sys.get_state(ctx.driver)
assert Map.keys(state.owners) == ["survivor"]
assert map_size(state.monitors) == 1
assert state.unexpected_deaths == []

Group.LocalCase.wait_until(fn ->
Group.members(:jepsen_group, "jepsen/pg/0") ==
[{ctx.survivor, %{token: ctx.token, revision: 1}}]
end)
end
end

test "death queued before snapshot is consumed without replacing the Driver", ctx do
:ok = :sys.suspend(ctx.driver)
task = Task.async(fn -> GenServer.call(ctx.driver, :owner_snapshots) end)

Group.LocalCase.wait_until(fn ->
{:messages, messages} = Process.info(ctx.driver, :messages)
Enum.any?(messages, &match?({:"$gen_call", _, :owner_snapshots}, &1))
end)

ref = Process.monitor(ctx.loser)

Process.exit(
ctx.loser,
{:group_registry_conflict, "jepsen/registry/0", %{token: "winner", revision: 2}}
)

assert_receive {:DOWN, ^ref, :process, _, _}
:ok = :sys.resume(ctx.driver)
assert {:ok, [%{token: token}]} = Task.await(task)
assert token == ctx.token
assert Process.alive?(ctx.driver)
end

test "unexpected owner death still leaves failure evidence", ctx do
:ok = :sys.suspend(ctx.loser)
task = Task.async(fn -> GenServer.call(ctx.driver, {:drop_cluster, "red"}, 15_000) end)

Group.LocalCase.wait_until(fn ->
{:messages, messages} = Process.info(ctx.loser, :messages)
Enum.any?(messages, &match?({:"$gen_call", _, {:drop_cluster, "red"}}, &1))
end)

Process.exit(ctx.loser, :implementation_bug)
assert :ok = Task.await(task, 15_000)

assert [%{reason: ":implementation_bug"}] =
GenServer.call(ctx.driver, :unexpected_deaths)

assert Process.alive?(ctx.survivor)
assert Process.alive?(ctx.driver)
end

test "a live owner's timeout is reported without deleting any owner state", ctx do
:ok = :sys.suspend(ctx.loser)

assert {:error, {"loser", _, {:timeout, _}, _}} =
GenServer.call(ctx.driver, {:drop_cluster, "red"}, 15_000)

assert Process.alive?(ctx.driver)
assert Process.alive?(ctx.loser)
assert Process.alive?(ctx.survivor)
state = :sys.get_state(ctx.driver)
assert map_size(state.owners) == 2
assert map_size(state.monitors) == 2
assert state.unexpected_deaths == []
:ok = :sys.resume(ctx.loser)
assert {:ok, snapshots} = GenServer.call(ctx.driver, :owner_snapshots)
assert length(snapshots) == 2
end
end