diff --git a/README.md b/README.md index 789c27b..fede217 100644 --- a/README.md +++ b/README.md @@ -315,6 +315,22 @@ Run the tests: mix test ``` +### Distributed ownership tests (local BEAM nodes) + +```bash +mix test --only distributed +``` + +This opt-in lane starts three disposable BEAM VMs with a shared EKV quorum. +It checks concurrent starts, acknowledged writes across owner VM loss, minority +write rejection, recovery after a distribution partition, and stale ETag fencing +after healing. Partitions affect only the test nodes; the controller uses stdio +and does not change host networking. Each case uses isolated names and data +directories and stops its peers on exit. + +These are focused failure scenarios, not a general linearizability checker or +deterministic scheduler. Run them separately from the fast default suite. + ### Integration Tests (with Tigris) Set the required environment variables: diff --git a/test/distributed_ownership_test.exs b/test/distributed_ownership_test.exs new file mode 100644 index 0000000..f68ea7f --- /dev/null +++ b/test/distributed_ownership_test.exs @@ -0,0 +1,205 @@ +defmodule DurableServer.DistributedOwnershipTest do + use ExUnit.Case, async: false + + alias DurableServer.DistributedTestClient, as: Client + alias DurableServer.StoredState + + @moduletag :integration + @moduletag :distributed + @moduletag :capture_log + @moduletag timeout: 90_000 + + setup do + id = DurableServer.UUID.uuid4() + cookie = :"ownership_cookie_#{id}" + store = :"ownership_store_#{id}" + supervisor = :"ownership_supervisor_#{id}" + tmp_dir = Path.expand("tmp/distributed_ownership/#{id}") + on_exit(fn -> File.rm_rf!(tmp_dir) end) + + members = + for index <- 1..3 do + {:ok, peer, node} = + :peer.start(%{ + name: :"ownership_#{id}_#{index}", + connection: :standard_io, + peer_down: :continue, + shutdown: :halt, + args: [~c"-setcookie", Atom.to_charlist(cookie), ~c"-connect_all", ~c"false"] + }) + + # Control uses stdio, so it remains available through a distribution + # partition. Teardown halts only these disposable VMs, not the host. + on_exit(fn -> :peer.stop(peer) end) + :ok = :peer.call(peer, :code, :add_paths, [:code.get_path()]) + {:ok, _} = :peer.call(peer, Application, :ensure_all_started, [:ekv]) + {:ok, _} = :peer.call(peer, Application, :ensure_all_started, [:durable_server]) + :ok = :peer.call(peer, Logger, :configure, [[level: :warning]]) + + {:ok, _} = + :peer.call(peer, Client, :start_store, [store, Path.expand("#{tmp_dir}/#{index}")]) + + %{peer: peer, node: node} + end + + for member <- members, other <- members, member != other do + assert :peer.call(member.peer, Node, :connect, [other.node]) + end + + await_cluster(members, store) + + for member <- members do + assert {:ok, _} = rpc(member, :start_durable, [supervisor, store]) + end + + {:ok, members: members, store: store, supervisor: supervisor, cookie: cookie} + end + + test "concurrent starts and acknowledged increments agree across three nodes", context do + %{members: members, supervisor: supervisor} = context + key = "concurrent" + + results = + for member <- members, _ <- 1..4 do + Task.async(fn -> rpc(member, :ensure_counter, [supervisor, key]) end) + end + |> Enum.map(&Task.await(&1, 15_000)) + + pids = + Enum.map(results, fn result -> + assert {:ok, {pid, _}} = result + pid + end) + + assert [owner] = Enum.uniq(pids) + + acknowledgements = + for member <- members do + Task.async(fn -> rpc(member, :invoke, [owner, :increment_and_sync]) end) + end + |> Enum.map(&Task.await(&1, 15_000)) + + assert Enum.sort(acknowledgements) == [{:ok, 1}, {:ok, 2}, {:ok, 3}] + assert_durable_count(members, supervisor, key, 3) + end + + test "a minority cannot acknowledge writes and stale tokens stay fenced after healing", + context do + %{members: [minority | majority] = members, supervisor: supervisor, cookie: cookie} = context + key = "partition" + assert {:ok, {old_owner, _}} = rpc(minority, :ensure_counter, [supervisor, key]) + assert node(old_owner) == minority.node + assert {:ok, 1} = rpc(minority, :invoke, [old_owner, :increment_and_sync]) + assert {:ok, stale} = rpc(minority, :stored, [supervisor, key]) + + isolate(minority, majority) + + assert {:indeterminate, _} = rpc(minority, :invoke, [old_owner, :increment_and_sync]) + assert :ok = rpc(hd(majority), :discover, [supervisor]) + new_owner = await_owner(majority, supervisor, key, old_owner) + owner_member = Enum.find(majority, &(&1.node == node(new_owner))) + + assert {:ok, 1} = rpc(owner_member, :invoke, [new_owner, :get_count]) + assert {:ok, 2} = rpc(owner_member, :invoke, [new_owner, :increment_and_sync]) + assert_durable_count(majority, supervisor, key, 2) + + for member <- majority do + assert true = :peer.call(minority.peer, :erlang, :set_cookie, [member.node, cookie]) + assert :peer.call(minority.peer, Node, :connect, [member.node]) + end + + await_cluster(members, context.store) + + await(fn -> not :peer.call(minority.peer, Process, :alive?, [old_owner]) end) + await(fn -> :peer.call(minority.peer, DurableServer.Supervisor, :ready?, [supervisor]) end) + + assert {:error, :conflict} = + rpc(minority, :stale_write, [supervisor, key, stale.body, stale.etag]) + + assert_durable_count(members, supervisor, key, 2) + end + + test "acknowledged state survives abrupt owner VM loss and majority recovery", context do + %{members: [owner_member | survivors], supervisor: supervisor} = context + key = "vm-loss" + assert {:ok, {old_owner, _}} = rpc(owner_member, :ensure_counter, [supervisor, key]) + + for expected <- 1..3 do + assert {:ok, ^expected} = rpc(owner_member, :invoke, [old_owner, :increment_and_sync]) + end + + # halt bypasses terminate callbacks and final sync; this is not graceful stop. + :ok = :peer.cast(owner_member.peer, :erlang, :halt, []) + await(fn -> match?({:down, _}, :peer.get_state(owner_member.peer)) end) + assert :ok = rpc(hd(survivors), :discover, [supervisor]) + new_owner = await_owner(survivors, supervisor, key, old_owner) + new_member = Enum.find(survivors, &(&1.node == node(new_owner))) + + assert {:ok, 3} = rpc(new_member, :invoke, [new_owner, :get_count]) + assert {:ok, 4} = rpc(new_member, :invoke, [new_owner, :increment_and_sync]) + assert_durable_count(survivors, supervisor, key, 4) + end + + defp isolate(minority, majority) do + for member <- majority do + # disconnect alone is insufficient: distribution can reconnect on send. + assert true = + :peer.call(minority.peer, :erlang, :set_cookie, [member.node, :partitioned]) + + assert :peer.call(minority.peer, Node, :disconnect, [member.node]) + end + + await(fn -> :peer.call(minority.peer, Node, :list, []) == [] end) + + for member <- majority do + await(fn -> minority.node not in :peer.call(member.peer, Node, :list, []) end) + end + end + + defp await_cluster(members, store) do + for member <- members do + expected = members |> Enum.reject(&(&1 == member)) |> Enum.map(& &1.node) |> Enum.sort() + await(fn -> rpc(member, :connected_store_members, [store]) == expected end) + end + end + + defp await_owner(members, supervisor, key, old_owner) do + await(fn -> + Enum.all?(members, fn member -> + case rpc(member, :lookup, [supervisor, key]) do + {pid, _} -> pid != old_owner and node(pid) in Enum.map(members, & &1.node) + nil -> false + end + end) + end) + + owners = Enum.map(members, fn member -> elem(rpc(member, :lookup, [supervisor, key]), 0) end) + assert [owner] = Enum.uniq(owners) + owner + end + + defp assert_durable_count(members, supervisor, key, expected) do + for member <- members do + assert {:ok, %{body: %StoredState{state: %{count: ^expected}}}} = + rpc(member, :stored, [supervisor, key]) + end + end + + defp rpc(member, function, args), do: :peer.call(member.peer, Client, function, args, 15_000) + + defp await(fun, timeout \\ 20_000) do + deadline = System.monotonic_time(:millisecond) + timeout + await_until(fun, deadline) + end + + defp await_until(fun, deadline) do + unless fun.() do + if System.monotonic_time(:millisecond) >= deadline do + flunk("cluster condition was not met before the deadline") + end + + Process.sleep(25) + await_until(fun, deadline) + end + end +end diff --git a/test/support/distributed_test_client.ex b/test/support/distributed_test_client.ex new file mode 100644 index 0000000..08344b2 --- /dev/null +++ b/test/support/distributed_test_client.ex @@ -0,0 +1,80 @@ +defmodule DurableServer.DistributedTestClient do + @moduledoc false + + alias DurableServer.{LifecycleManager, StorageBackend} + alias DurableServer.Backends.EKVStore + + defmodule Counter do + use DurableServer, vsn: 1 + + def dump_state(state), do: state + def load_state(_vsn, state), do: state + def init(state), do: {:ok, state, permanent: true} + + def handle_call(:get_count, _from, state), do: {:reply, state.count, state} + + def handle_call(:increment_and_sync, _from, state) do + state = %{state | count: state.count + 1} + {:reply, state.count, state, :sync} + end + end + + def start_store(name, data_dir) do + Supervisor.start_child(EKV.AppSupervisor, { + EKV, + name: name, data_dir: data_dir, cluster_size: 3, node_id: to_string(node()), log: false + }) + end + + def start_durable(name, store) do + Supervisor.start_child(DurableServer.AppSupervisor, { + DurableServer.Supervisor, + name: name, + prefix: "ownership/", + backend: {EKVStore, name: store, start: false}, + initial_discovery_delay_ms: 60_000, + discovery_interval_ms: 200, + discovery_burst_count: 0, + heartbeat_interval_ms: 250, + heartbeat_staleness_threshold_ms: 4_000, + graceful_shutdown_timeout_ms: 500 + }) + end + + def ensure_counter(supervisor, key) do + DurableServer.Supervisor.ensure_started_child( + supervisor, + {Counter, key: key, initial_state: %{count: 0}}, + local_only: true, + timeout: 10_000 + ) + end + + def invoke(pid, request) do + {:ok, GenServer.call(pid, request, 10_000)} + catch + # A missing acknowledgement is not evidence that a write did not commit. + :exit, reason -> {:indeterminate, reason} + end + + def stored(supervisor, key) do + %{storage_backend: backend} = DurableServer.Supervisor.__get_config__(supervisor) + StorageBackend.get_object(backend, "ownership/" <> key, consistent: true) + end + + def stale_write(supervisor, key, body, etag) do + %{storage_backend: backend} = DurableServer.Supervisor.__get_config__(supervisor) + StorageBackend.put_object(backend, "ownership/" <> key, body, etag: etag, max_retries: 0) + end + + def discover(supervisor) do + send(LifecycleManager.name(supervisor), :discover_and_restart) + :ok + end + + def lookup(supervisor, key), do: DurableServer.Supervisor.lookup(supervisor, key) + + def connected_store_members(store) do + EKV.info(store).connected_members |> Enum.map(& &1.node) |> Enum.sort() + end +end