diff --git a/README.md b/README.md index 789c27b..db6ff98 100644 --- a/README.md +++ b/README.md @@ -301,18 +301,51 @@ These are independent - joining does not monitor events, and monitoring does not ## Running Tests -### Unit Tests (with LocalStack) +### Storage-free tests -Start LocalStack for S3-compatible storage: +The default suite needs neither Docker nor cloud credentials: ```bash -docker run -d --name localstack -p 4566:4566 localstack/localstack +mix deps.get +mix test ``` -Run the tests: +### LocalStack tests + +Start the pinned LocalStack version, wait for its health endpoint to respond, +then include the storage-backed server, lifecycle, placement, and mirror tests: ```bash -mix test +docker run -d --name localstack -p 127.0.0.1:4566:4566 localstack/localstack:4.14.0 +curl http://localhost:4566/_localstack/health +mix test --include localstack +``` + +Each suite invocation creates a unique bucket only if a LocalStack test runs. +Tests use unique prefixes within it, and the bucket is emptied and deleted after +test supervisors stop. No shared bucket is cleared at startup, so concurrent +suite invocations do not delete each other's data. + +New LocalStack-backed test modules should use `DurableServer.LocalStackCase`. +Selecting a storage-free test file never connects to LocalStack. + +### Local EKV tests + +EKV tests, including the two-node tests, run locally without cloud credentials. +The Erlang port mapper must be running for the peer nodes: + +```bash +epmd -daemon +mix test --include ekv +``` + +EKV data directories are unique to each test allocation, live under the +gitignored `tmp/` directory, and are removed on exit. + +Run all local tests (including LocalStack/EKV migration tests) with: + +```bash +mix test --include localstack --include ekv ``` ### Integration Tests (with Tigris) @@ -330,8 +363,12 @@ export DURABLE_AWS_REGION= export DURABLE_BUCKET= ``` -Run integration tests (which hit t3.storage.dev directly): +The `integration` tag is reserved for credentialed cloud tests. These hit +Tigris directly and create cloud resources: ```bash mix test --include integration ``` + +For repeatability, pass `--seed `. To replay a failed storage-backed +test, retain its inclusion flag, for example `mix test --failed --include localstack`. diff --git a/test/durable_server/lifecycle_helpers_test.exs b/test/durable_server/lifecycle_helpers_test.exs new file mode 100644 index 0000000..efe7414 --- /dev/null +++ b/test/durable_server/lifecycle_helpers_test.exs @@ -0,0 +1,99 @@ +defmodule DurableServer.LifecycleHelpersTest do + use ExUnit.Case, async: true + import DurableServer.LifecycleHelpers + + @moduletag :capture_log + + defmodule Manager do + use GenServer + + def start_link(opts), do: GenServer.start_link(__MODULE__, opts) + + def init(opts) do + {:ok, Map.merge(Map.new(opts), %{current_discovery_task: nil})} + end + + def handle_info(:discover_and_restart, %{ignore?: true} = state), do: {:noreply, state} + + def handle_info(:discover_and_restart, state) do + owner = state.owner + + task = + Task.Supervisor.async_nolink(state.task_supervisor, fn -> + send(owner, {:discovery_started, self()}) + + receive do + :finish -> {:discover, :ok} + :crash -> exit(:injected_failure) + end + end) + + {:noreply, %{state | current_discovery_task: task}} + end + + def handle_info({ref, {:discover, :ok}}, %{current_discovery_task: %Task{ref: ref}} = state) do + Process.demonitor(ref, [:flush]) + {:noreply, %{state | current_discovery_task: nil}} + end + + def handle_info({:DOWN, ref, :process, _pid, _reason}, state) do + %Task{ref: ^ref} = state.current_discovery_task + {:noreply, %{state | current_discovery_task: nil}} + end + end + + setup do + task_supervisor = start_supervised!(Task.Supervisor) + + manager = + start_supervised!( + {Manager, owner: self(), task_supervisor: task_supervisor, ignore?: false} + ) + + %{manager: manager} + end + + test "waits until the started discovery task completes and its result is processed", context do + waiter = Task.async(fn -> discover_and_wait(context.manager) end) + assert_receive {:discovery_started, task_pid} + assert Task.yield(waiter, 0) == nil + + send(task_pid, :finish) + assert Task.await(waiter) == :ok + assert :sys.get_state(context.manager).current_discovery_task == nil + end + + test "an idle manager is not mistaken for a completed discovery", context do + :sys.replace_state(context.manager, &%{&1 | ignore?: true}) + + assert_raise ExUnit.AssertionError, ~r/task_not_started/, fn -> + discover_and_wait(context.manager) + end + end + + test "a failed discovery task is not mistaken for a completed discovery", context do + waiter = + Task.async(fn -> + assert_raise ExUnit.AssertionError, ~r/task_exit.*injected_failure/, fn -> + discover_and_wait(context.manager) + end + end) + + assert_receive {:discovery_started, task_pid} + send(task_pid, :crash) + Task.await(waiter) + end + + test "a timed out wait does not leave a debug hook sending late results", context do + assert_raise ExUnit.AssertionError, ~r/did not complete/, fn -> + discover_and_wait(context.manager, 20) + end + + assert_receive {:discovery_started, task_pid} + ref = Process.monitor(task_pid) + send(task_pid, :finish) + assert_receive {:DOWN, ^ref, :process, ^task_pid, :normal} + :sys.get_state(context.manager) + refute_receive {:discovery, _ref, _result}, 0 + end +end diff --git a/test/durable_server/lifecycle_test.exs b/test/durable_server/lifecycle_test.exs index 18686e8..c7d99bb 100644 --- a/test/durable_server/lifecycle_test.exs +++ b/test/durable_server/lifecycle_test.exs @@ -1,6 +1,7 @@ defmodule DurableServer.LifecycleTest do - use ExUnit.Case, async: true + use DurableServer.LocalStackCase, async: true import DurableServer.TestHelper + import DurableServer.LifecycleHelpers alias DurableServer alias DurableServer.{CircuitBreaker, LifecycleManager, Meta} @@ -52,6 +53,69 @@ defmodule DurableServer.LifecycleTest do end end + defmodule FailingLoadServer do + use DurableServer, vsn: 1 + + def dump_state(state), do: state + def load_state(_vsn, _state), do: raise("injected state loading failure") + def init(state), do: {:ok, state} + end + + defmodule RestartClaimBarrierBackend do + @behaviour DurableServer.StorageBackend + alias DurableServer.StorageBackend + + @impl true + def init_backend(opts) do + {:ok, delegate} = + StorageBackend.init_backend( + DurableServer.Backends.ObjectStore, + Keyword.fetch!(opts, :store) + ) + + {:ok, %{state: %{delegate: delegate, owner: Keyword.fetch!(opts, :owner)}}} + end + + @impl true + def ensure_ready(%{delegate: delegate}), do: StorageBackend.ensure_ready(delegate) + + @impl true + def get_object(%{delegate: delegate}, key, opts), + do: StorageBackend.get_object(delegate, key, opts) + + @impl true + def list_all_objects_stream(%{delegate: delegate}, prefix, opts), + do: StorageBackend.list_all_objects_stream(delegate, prefix, opts) + + @impl true + def put_object(%{delegate: delegate, owner: owner}, key, body, opts) do + send(owner, {:restart_claim_ready, self()}) + + receive do + :release_restart_claim -> StorageBackend.put_object(delegate, key, body, opts) + after + 5_000 -> raise "restart claim barrier was not released" + end + end + + @impl true + def delete_object(%{delegate: delegate}, key), do: StorageBackend.delete_object(delegate, key) + + @impl true + def try_claim(%{delegate: delegate}, key, body), + do: StorageBackend.try_claim(delegate, key, body) + + @impl true + def update_object(%{delegate: delegate}, key, update_fn, opts), + do: StorageBackend.update_object(delegate, key, update_fn, opts) + + @impl true + def encode(%{delegate: delegate}, data), do: StorageBackend.encode(delegate, data) + + @impl true + def decode(%{delegate: delegate}, data), do: StorageBackend.decode(delegate, data) + end + defmodule DelayedTerminateServer do use DurableServer, vsn: 1 @@ -702,56 +766,46 @@ defmodule DurableServer.LifecycleTest do supervisor_name = :"test_supervisor_#{DurableServer.UUID.uuid4()}" prefix = "test_#{DurableServer.UUID.uuid4()}/" + # Discovery tests explicitly drive each cycle, without startup bursts racing + # their storage assertions or contributing extra diagnostic counts. _supervisor_pid = start_supervised!({ DurableServer.Supervisor, - name: supervisor_name, prefix: prefix, object_store: test_object_store_opts() + name: supervisor_name, + prefix: prefix, + object_store: test_object_store_opts(), + initial_discovery_delay_ms: 60_000, + discovery_burst_count: 0 }) object_store = test_object_store() - test_bucket_name = "durable-test-lifecycle-#{DurableServer.UUID.uuid4()}" - - case ObjectStore.create_bucket_with_credentials(object_store, test_bucket_name) do - {:ok, %ObjectStore{} = store} -> - on_exit(fn -> - try do - ObjectStore.delete_bucket(store, test_bucket_name) - catch - _, _ -> :ok - end - end) + supervisor_config = DurableServer.Supervisor.__get_config__(supervisor_name) + circuit_breaker = supervisor_config.circuit_breaker - supervisor_config = DurableServer.Supervisor.__get_config__(supervisor_name) - circuit_breaker = supervisor_config.circuit_breaker - - # Create test config that mimics what supervisor provides - test_config = %{ - name: supervisor_name, - prefix: prefix, - object_store: object_store, - discovery_interval_ms: 60_000, - heartbeat_interval_ms: 10_000, - graceful_shutdown_timeout_ms: 30_000, - dead_node_threshold_ms: 24 * 60 * 60 * 1000, - crash_threshold_count: 5, - crash_threshold_window_ms: 60 * 60 * 1000, - module_circuit_breaker_count: 50, - module_circuit_breaker_window_ms: 5 * 60 * 1000, - module_circuit_breaker_cooldown_ms: 30 * 60 * 1000, - ets_table: supervisor_config.ets_table - } - - {:ok, - test_bucket: test_bucket_name, - store: store, - supervisor_name: supervisor_name, - prefix: prefix, - config: test_config, - circuit_breaker: circuit_breaker} + # Standalone managers use the same storage as the real test supervisor. + test_config = %{ + name: supervisor_name, + prefix: prefix, + object_store: object_store, + initial_discovery_delay_ms: 60_000, + discovery_burst_count: 0, + discovery_interval_ms: 60_000, + heartbeat_interval_ms: 10_000, + graceful_shutdown_timeout_ms: 30_000, + dead_node_threshold_ms: 24 * 60 * 60 * 1000, + crash_threshold_count: 5, + crash_threshold_window_ms: 60 * 60 * 1000, + module_circuit_breaker_count: 50, + module_circuit_breaker_window_ms: 5 * 60 * 1000, + module_circuit_breaker_cooldown_ms: 30 * 60 * 1000, + ets_table: supervisor_config.ets_table + } - {:error, reason} -> - {:skip, "Failed to create test bucket: #{inspect(reason)}"} - end + {:ok, + supervisor_name: supervisor_name, + prefix: prefix, + config: test_config, + circuit_breaker: circuit_breaker} end describe "stop modes" do @@ -1014,9 +1068,7 @@ defmodule DurableServer.LifecycleTest do # Use the real supervisor's LifecycleManager for actual restart testing {:ok, manager_pid} = get_supervisor_lifecycle_manager(supervisor_name) - # Manually trigger discovery - send(manager_pid, :discover_and_restart) - wait_for_discovery_completion(manager_pid, 5000) + discover_and_wait(manager_pid) assert {restarted_pid, _meta} = DurableServer.Supervisor.lookup(supervisor_name, key) @@ -1232,8 +1284,7 @@ defmodule DurableServer.LifecycleTest do {key, etag, stored_state.meta, System.monotonic_time(:millisecond)} ) - send(manager_pid, :discover_and_restart) - wait_for_discovery_completion(manager_pid, 1000) + discover_and_wait(manager_pid) assert [] = :ets.lookup(manager_state.discovery_skip_table, key) assert nil == DurableServer.Supervisor.lookup(supervisor_name, key) @@ -1294,8 +1345,7 @@ defmodule DurableServer.LifecycleTest do {:ok, manager_pid} = get_supervisor_lifecycle_manager(supervisor_name) - send(manager_pid, :discover_and_restart) - wait_for_discovery_completion(manager_pid, 1000) + discover_and_wait(manager_pid) diagnostics = LifecycleManager.get_discovery_diagnostics(supervisor_name) assert Map.get(diagnostics, :restart_claim_ok, 0) >= 1 @@ -1383,19 +1433,17 @@ defmodule DurableServer.LifecycleTest do # Use the real supervisor's LifecycleManager for actual restart testing {:ok, manager_pid} = get_supervisor_lifecycle_manager(supervisor_name) - send(manager_pid, :discover_and_restart) - wait_for_discovery_completion(manager_pid, 3000) + discover_and_wait(manager_pid) - # Check if restart was attempted - after failure, metadata gets cleaned up - {:ok, final_data} = DurableServer.fetch_stored_state(store, %{key: key, prefix: prefix}) + assert {restarted_pid, _meta} = DurableServer.Supervisor.lookup(supervisor_name, key) + assert restarted_pid != pid + assert GenServer.call(restarted_pid, :get_count) == original_count - # Restart should have been attempted and failed (since TestServer interface mismatch) - # After failure, restart metadata is cleaned up, but state should be preserved - assert final_data.state["count"] == original_count, - "State was not preserved during restart attempt" + {:ok, final_data} = DurableServer.fetch_stored_state(store, %{key: key, prefix: prefix}) - # Status should no longer be crashed since restart succeeded - refute final_data.meta.status == :crashed + assert final_data.state["count"] == original_count + assert final_data.meta.status == :running + assert final_data.meta.pid == restarted_pid GenServer.stop(manager_pid) end @@ -1461,185 +1509,139 @@ defmodule DurableServer.LifecycleTest do # Use the real supervisor's LifecycleManager for actual restart testing {:ok, manager_pid} = get_supervisor_lifecycle_manager(supervisor_name) - send(manager_pid, :discover_and_restart) - wait_for_discovery_completion(manager_pid, 2000) + discover_and_wait(manager_pid) + + assert {restarted_pid, _meta} = DurableServer.Supervisor.lookup(supervisor_name, key) + assert restarted_pid != pid + assert GenServer.call(restarted_pid, :get_count) == 5 - # Verify restart was attempted (metadata cleaned up after failure) and state preserved {:ok, post_restart_data} = DurableServer.fetch_stored_state(config.object_store, %{key: key, prefix: prefix}) - # After restart failure, metadata is cleaned up but state should be preserved - assert post_restart_data.state["count"] == 5, "Should use updated state from storage" - - refute post_restart_data.meta.status == :crashed, - "Status should change from crashed after successful restart" + assert post_restart_data.state["count"] == 5 + assert post_restart_data.meta.status == :running + assert post_restart_data.meta.pid == restarted_pid GenServer.stop(manager_pid) end end describe "real-world restart scenarios" do - test "handles restart timing conflicts between multiple managers", %{ + test "concurrent restart claims have exactly one winner and recover its state", %{ supervisor_name: supervisor_name, prefix: prefix, - config: config, - circuit_breaker: _circuit_breaker + config: config } do key = "timing-conflict-test-#{DurableServer.UUID.uuid4()}" + stored = seed_restartable_object(config, key, 3) - # Create crashed server state - meta_attrs = %{ - status: :crashed, - node_str: "unreachable@test", - node_ref: "dead-ref", - pid: self(), - last_heartbeat_at: System.system_time(:millisecond), - module: TestServer - } - - create_test_object(config.object_store, "#{prefix}#{key}", %{count: 3}, meta_attrs) - - # Start two managers that will compete for restart (use spawn to avoid name conflicts) - {:ok, manager1} = - start_standalone_lifecycle_manager(supervisor_name, config, - node_module: DurableServer.LifecycleTest.MockNodeModule + {:ok, backend} = + DurableServer.StorageBackend.init_backend(RestartClaimBarrierBackend, + store: config.object_store, + owner: self() ) - {:ok, manager2} = - start_standalone_lifecycle_manager(supervisor_name, config, - node_module: DurableServer.LifecycleTest.MockNodeModule - ) + # Different manager timeouts make these distinct claims, not the identical + # same-owner write that ObjectStore intentionally treats as idempotent. + attempts = + for ttl <- [30_000, 60_000] do + Task.async(fn -> DurableServer.claim_restart_attempt(backend, stored, ttl: ttl) end) + end - # Trigger discovery on both simultaneously - send(manager1, :discover_and_restart) - send(manager2, :discover_and_restart) + # Both claimers must read the unclaimed object before either CAS + # proceeds. This exercises the losing write, not just a serialized lookup. + assert_receive {:restart_claim_ready, first} + assert_receive {:restart_claim_ready, second} + assert first != second + send(first, :release_restart_claim) + send(second, :release_restart_claim) + + results = Enum.map(attempts, &Task.await/1) + assert [{:ok, claim}] = Enum.filter(results, &match?({:ok, _}, &1)) + assert [{:error, :not_eligible}] = Enum.filter(results, &match?({:error, _}, &1)) + + # Use the same preloaded startup path as LifecycleManager after a claim. + assert {:ok, {pid, _meta}} = + DurableServer.Supervisor.__start_child__( + supervisor_name, + {TestServer, [key: key, initial_state: %{}], + %{ + preloaded: claim, + restart_attempt_node: to_string(node()), + is_sticky_local: false + }} + ) - # Give time for both to process - Process.sleep(2000) + assert {^pid, _meta} = DurableServer.Supervisor.lookup(supervisor_name, key) + assert GenServer.call(pid, :get_count) == 3 - # Check final state - should have atomic claim (only one winner) - {:ok, data} = + {:ok, recovered} = DurableServer.fetch_stored_state(config.object_store, %{key: key, prefix: prefix}) - restart_node = Map.get(data.meta, :restart_attempt_node) - - # Should have exactly one restart attempt node or completion - if restart_node do - assert restart_node == to_string(Node.self()) - else - # Or restart completed successfully with no remaining metadata - assert Map.get(data.meta, :restart_attempt_time) == nil - end + assert recovered.meta.status == :running + assert recovered.meta.pid == pid + assert recovered.meta.restart_attempt_node == nil end test "restart coordination across discovery cycles", %{ supervisor_name: supervisor_name, prefix: prefix, - config: config, - circuit_breaker: _circuit_breaker + config: config } do key = "coordination-test-#{DurableServer.UUID.uuid4()}" + seed_restartable_object(config, key, 7) + {:ok, manager_pid} = get_supervisor_lifecycle_manager(supervisor_name) - meta_attrs = %{ - status: :crashed, - # Use unreachable node - node_str: "unreachable@test", - node_ref: "dead-ref", - pid: self(), - last_heartbeat_at: System.system_time(:millisecond), - module: TestServer - } - - create_test_object(config.object_store, "#{prefix}#{key}", %{count: 0}, meta_attrs) - - {:ok, manager_pid} = - start_standalone_lifecycle_manager(supervisor_name, config, - node_module: DurableServer.LifecycleTest.MockNodeModule - ) - - # Multiple discovery cycles - for i <- 1..3 do - send(manager_pid, :discover_and_restart) - wait_for_discovery_completion(manager_pid, 1500) + recovered_pids = + for _ <- 1..3 do + discover_and_wait(manager_pid) + assert {pid, _meta} = DurableServer.Supervisor.lookup(supervisor_name, key) + assert GenServer.call(pid, :get_count) == 7 - # Check state after each cycle - {:ok, data} = - DurableServer.fetch_stored_state(config.object_store, %{key: key, prefix: prefix}) + {:ok, stored} = + DurableServer.fetch_stored_state(config.object_store, %{key: key, prefix: prefix}) - case i do - 1 -> - # First cycle should claim restart since node is unreachable - restart_attempt = Map.get(data.meta, :restart_attempt_node) - - if restart_attempt == nil do - # If no restart attempt, server might have been restarted successfully - # or the discovery logic didn't find it as orphaned - # This is acceptable - just log for debugging - :ok - end - - 2 -> - # Second cycle - restart might be in progress or cleaned up - :ok - - 3 -> - # Final cycle should have completed or cleaned up - restart_time = Map.get(data.meta, :restart_attempt_time) - - if restart_time do - # If still present, should be recent - assert System.system_time(:millisecond) - restart_time < 30_000 - end + assert stored.meta.status == :running + assert stored.meta.pid == pid + assert stored.meta.restart_attempt_node == nil + pid end - end - GenServer.stop(manager_pid) + assert [_pid] = Enum.uniq(recovered_pids) + diagnostics = LifecycleManager.get_discovery_diagnostics(supervisor_name) + assert diagnostics.restart_claim_ok == 1 + assert diagnostics.restart_start_ok == 1 end end describe "complete error recovery testing" do - test "cleanup prevents future restart attempts after repeated failures", %{ + test "state loading failures clear their claim so a later discovery can retry", %{ supervisor_name: supervisor_name, prefix: prefix, - config: config, - circuit_breaker: _circuit_breaker + config: config } do key = "repeated-failure-test-#{DurableServer.UUID.uuid4()}" - meta_attrs = %{ - status: :crashed, - node_str: "dead_node@test", - node_ref: "dead-ref", - pid: self(), - last_heartbeat_at: System.system_time(:millisecond), - # This will always fail to restart - module: NonExistentModule - } - - create_test_object(config.object_store, "#{prefix}#{key}", %{count: 0}, meta_attrs) - - {:ok, manager_pid} = - start_standalone_lifecycle_manager(supervisor_name, config) - - # Trigger multiple discovery rounds to simulate repeated failures - for _ <- 1..3 do - send(manager_pid, :discover_and_restart) - wait_for_discovery_completion(manager_pid, 1500) - end + seed_restartable_object(config, key, 0, module: FailingLoadServer) + {:ok, manager_pid} = get_supervisor_lifecycle_manager(supervisor_name) - # Check final state - should have cleanup after failure - {:ok, data} = - DurableServer.fetch_stored_state(config.object_store, %{key: key, prefix: prefix}) + for attempt <- 1..3 do + discover_and_wait(manager_pid) - # Should have no restart attempt metadata (cleaned up after failure) - assert Map.get(data.meta, :restart_attempt_ttl) == nil - assert Map.get(data.meta, :restart_attempt_node) == nil - assert Map.get(data.meta, :restart_attempt_time) == nil + diagnostics = LifecycleManager.get_discovery_diagnostics(supervisor_name) + assert diagnostics.restart_claim_ok == attempt + assert diagnostics.restart_start_error == attempt + assert DurableServer.Supervisor.lookup(supervisor_name, key) == nil - # Status should remain crashed since restart failed - assert data.meta.status == :crashed + {:ok, data} = + DurableServer.fetch_stored_state(config.object_store, %{key: key, prefix: prefix}) - GenServer.stop(manager_pid) + assert data.meta.restart_attempt_ttl == nil + assert data.meta.restart_attempt_node == nil + assert data.meta.restart_attempt_time == nil + assert data.meta.status == :crashed + assert data.state["count"] == 0 + end end test "recovery from corrupted restart attempt metadata", %{ @@ -1670,9 +1672,7 @@ defmodule DurableServer.LifecycleTest do {:ok, manager_pid} = start_standalone_lifecycle_manager(supervisor_name, config) - # Should handle corrupted metadata gracefully - send(manager_pid, :discover_and_restart) - wait_for_discovery_completion(manager_pid, 1500) + discover_and_wait(manager_pid) # Manager should still be alive despite corrupted data assert_process_alive(manager_pid) @@ -1689,87 +1689,52 @@ defmodule DurableServer.LifecycleTest do test "error recovery during discovery with mixed valid/invalid objects", %{ supervisor_name: supervisor_name, prefix: prefix, - config: config, - circuit_breaker: _circuit_breaker + config: config } do - # Create multiple objects with mixed validity - keys = for i <- 1..5, do: "mixed-test-#{i}-#{DurableServer.UUID.uuid4()}" - - # Create mix of valid and invalid objects - for {key, i} <- Enum.with_index(keys, 1) do - case rem(i, 3) do - 0 -> - # Valid crashed server - meta_attrs = %{ - status: :crashed, - node_str: "dead_node@test", - node_ref: "dead-ref", - pid: self(), - last_heartbeat_at: System.system_time(:millisecond), - module: TestServer - } - - create_test_object(config.object_store, "#{prefix}#{key}", %{count: i}, meta_attrs) - - 1 -> - # Corrupted metadata - corrupted_data = %{ - vsn: 1, - state: %{count: i}, - meta: "invalid_base64_#{i}" - } - - ObjectStore.put_object( - config.object_store, - "#{prefix}#{key}", - JSON.encode!(corrupted_data) - ) + crashed = for count <- 1..2, do: {"crashed-#{DurableServer.UUID.uuid4()}", count} + invalid = for _ <- 1..2, do: "invalid-#{DurableServer.UUID.uuid4()}" + live_key = "live-#{DurableServer.UUID.uuid4()}" + live_pid = start_test_server(supervisor_name, live_key) - 2 -> - # Valid running server (should be ignored) - meta_attrs = %{ - status: :running, - node_str: to_string(Node.self()), - node_ref: "current-ref", - pid: self(), - last_heartbeat_at: System.system_time(:millisecond), - module: TestServer - } - - create_test_object(config.object_store, "#{prefix}#{key}", %{count: i}, meta_attrs) - end + for {key, count} <- crashed do + seed_restartable_object(config, key, count) end - {:ok, manager_pid} = - start_standalone_lifecycle_manager(supervisor_name, config) + for key <- invalid do + assert {:ok, _} = + ObjectStore.put_object( + config.object_store, + prefix <> key, + JSON.encode!(%{vsn: 1, state: %{}, meta: "invalid_base64"}) + ) + end - # Should handle mixed object states gracefully - send(manager_pid, :discover_and_restart) - wait_for_discovery_completion(manager_pid, 2000) + {:ok, manager_pid} = get_supervisor_lifecycle_manager(supervisor_name) + discover_and_wait(manager_pid) - # Manager should survive processing mixed objects - assert_process_alive(manager_pid) + for {key, count} <- crashed do + assert {pid, _meta} = DurableServer.Supervisor.lookup(supervisor_name, key) + assert GenServer.call(pid, :get_count) == count - # Should have attempted restart only on valid crashed servers - # Every 3rd key (0-indexed) - valid_crashed_keys = Enum.take_every(keys, 3) + assert {:ok, data} = + DurableServer.fetch_stored_state(config.object_store, %{key: key, prefix: prefix}) - for key <- valid_crashed_keys do - case DurableServer.fetch_stored_state(config.object_store, %{key: key, prefix: prefix}) do - {:ok, data} -> - # Should have restart attempt or completion - restart_attempt = Map.get(data.meta, :restart_attempt_node) + assert data.meta.status == :running + assert data.meta.pid == pid + end - assert restart_attempt != nil or data.meta.status != :crashed, - "Valid crashed server #{key} should have restart attempt or be recovered" + for key <- invalid do + assert DurableServer.Supervisor.lookup(supervisor_name, key) == nil - {:error, _} -> - # Object might have been cleaned up during restart - acceptable - :ok - end + assert {:error, %ArgumentError{}} = + DurableServer.fetch_stored_state(config.object_store, %{key: key, prefix: prefix}) end - GenServer.stop(manager_pid) + assert {^live_pid, _meta} = DurableServer.Supervisor.lookup(supervisor_name, live_key) + assert GenServer.call(live_pid, :get_count) == 0 + diagnostics = LifecycleManager.get_discovery_diagnostics(supervisor_name) + assert diagnostics.restart_claim_ok == length(crashed) + assert diagnostics.restart_start_ok == length(crashed) end end @@ -1819,7 +1784,29 @@ defmodule DurableServer.LifecycleTest do meta: encoded_meta } - ObjectStore.put_object(store, key, JSON.encode!(test_data)) + assert {:ok, object} = ObjectStore.put_object(store, key, JSON.encode!(test_data)) + {:ok, object} + end + + defp seed_restartable_object(config, key, count, overrides \\ []) do + %{object_store: store, prefix: prefix, name: supervisor_name} = config + + meta = %{ + status: :crashed, + permanent: true, + supervisor: supervisor_name, + key: key, + prefix: prefix, + node_str: "unreachable@test", + node_ref: "dead-ref", + pid: self(), + last_heartbeat_at: System.system_time(:millisecond), + module: TestServer + } + + create_test_object(store, prefix <> key, %{count: count}, Map.merge(meta, Map.new(overrides))) + assert {:ok, stored} = DurableServer.fetch_stored_state(store, %{key: key, prefix: prefix}) + stored end defp encode_legacy_stored_state(%DurableServer.StoredState{ @@ -1918,20 +1905,6 @@ defmodule DurableServer.LifecycleTest do ) end - # Helper function to wait for discovery task completion or manager crash - defp wait_for_discovery_completion(manager_pid, timeout) do - ref = Process.monitor(manager_pid) - - receive do - {:DOWN, ^ref, :process, ^manager_pid, reason} -> - flunk("LifecycleManager crashed: #{inspect(reason)}") - after - timeout -> - Process.demonitor(ref, [:flush]) - :ok - end - end - @stale_heartbeat_age_ms 48 * 60 * 60 * 1000 defp await_heartbeat_cycle(manager_pid) do @@ -1981,24 +1954,6 @@ defmodule DurableServer.LifecycleTest do store end - defp assert_eventually(fun, timeout \\ 5_000, interval \\ 25) when is_function(fun, 0) do - deadline = System.monotonic_time(:millisecond) + timeout - do_assert_eventually(fun, deadline, interval) - end - - defp do_assert_eventually(fun, deadline, interval) do - if fun.() do - :ok - else - if System.monotonic_time(:millisecond) >= deadline do - flunk("condition was not met within timeout") - else - Process.sleep(interval) - do_assert_eventually(fun, deadline, interval) - end - end - end - defp setup_restart_gate_tables(supervisor_name) do config_table = :"durable_supervisor_#{supervisor_name}" heartbeat_table = :"durable_server_heartbeats_#{supervisor_name}" @@ -2048,11 +2003,7 @@ defmodule DurableServer.LifecycleTest do {:ok, pid} = start_standalone_lifecycle_manager(supervisor_name, config) - # Send a discovery message manually to trigger discovery - send(pid, :discover_and_restart) - - # Wait for discovery to complete or manager to crash - wait_for_discovery_completion(pid, 500) + discover_and_wait(pid) assert_process_alive(pid) GenServer.stop(pid) @@ -2292,10 +2243,7 @@ defmodule DurableServer.LifecycleTest do {:ok, manager_pid} = start_standalone_lifecycle_manager(supervisor_name, config) - send(manager_pid, :discover_and_restart) - - # Wait for discovery to complete or manager to crash - wait_for_discovery_completion(manager_pid, 500) + discover_and_wait(manager_pid) assert_process_alive(manager_pid) GenServer.stop(manager_pid) @@ -2338,10 +2286,7 @@ defmodule DurableServer.LifecycleTest do create_test_object(config.object_store, "#{prefix}#{key}", %{count: 0}, meta_attrs) end - send(manager_pid, :discover_and_restart) - - # Wait for discovery to complete or manager to crash - wait_for_discovery_completion(manager_pid, 500) + discover_and_wait(manager_pid) assert_process_alive(manager_pid) GenServer.stop(manager_pid) @@ -2373,10 +2318,7 @@ defmodule DurableServer.LifecycleTest do create_test_object(config.object_store, "#{prefix}#{key}", %{count: 0}, meta_attrs) - send(manager_pid, :discover_and_restart) - - # Wait for discovery to complete or manager to crash - wait_for_discovery_completion(manager_pid, 100) + discover_and_wait(manager_pid) assert_process_alive(manager_pid) GenServer.stop(manager_pid) @@ -2407,10 +2349,7 @@ defmodule DurableServer.LifecycleTest do create_test_object(config.object_store, "#{prefix}#{key}", %{count: 0}, meta_attrs) - send(manager_pid, :discover_and_restart) - - # Wait for discovery to complete or manager to crash - wait_for_discovery_completion(manager_pid, 100) + discover_and_wait(manager_pid) assert_process_alive(manager_pid) GenServer.stop(manager_pid) @@ -2439,10 +2378,7 @@ defmodule DurableServer.LifecycleTest do create_test_object(config.object_store, "#{prefix}#{key}", %{count: 0}, meta_attrs) - send(manager_pid, :discover_and_restart) - - # Wait for discovery to complete or manager to crash - wait_for_discovery_completion(manager_pid, 100) + discover_and_wait(manager_pid) assert_process_alive(manager_pid) GenServer.stop(manager_pid) @@ -2473,10 +2409,7 @@ defmodule DurableServer.LifecycleTest do create_test_object(config.object_store, "#{prefix}#{key}", %{count: 0}, meta_attrs) - send(manager_pid, :discover_and_restart) - - # Wait for discovery to complete or manager to crash - wait_for_discovery_completion(manager_pid, 100) + discover_and_wait(manager_pid) assert_process_alive(manager_pid) GenServer.stop(manager_pid) @@ -2506,10 +2439,7 @@ defmodule DurableServer.LifecycleTest do create_test_object(config.object_store, "#{prefix}#{key}", %{count: 0}, meta_attrs) - send(manager_pid, :discover_and_restart) - - # Wait for discovery to complete or manager to crash - wait_for_discovery_completion(manager_pid, 100) + discover_and_wait(manager_pid) assert_process_alive(manager_pid) GenServer.stop(manager_pid) @@ -2541,10 +2471,7 @@ defmodule DurableServer.LifecycleTest do create_test_object(config.object_store, "#{prefix}#{key}", %{count: 0}, meta_attrs) - send(manager_pid, :discover_and_restart) - - # Wait for discovery to complete or manager to crash - wait_for_discovery_completion(manager_pid, 100) + discover_and_wait(manager_pid) assert_process_alive(manager_pid) GenServer.stop(manager_pid) @@ -2576,10 +2503,7 @@ defmodule DurableServer.LifecycleTest do create_test_object(config.object_store, "#{prefix}#{key}", %{count: 0}, meta_attrs) - send(manager_pid, :discover_and_restart) - - # Wait for discovery to complete or manager to crash - wait_for_discovery_completion(manager_pid, 100) + discover_and_wait(manager_pid) assert_process_alive(manager_pid) GenServer.stop(manager_pid) @@ -2618,10 +2542,7 @@ defmodule DurableServer.LifecycleTest do node_module: DurableServer.LifecycleTest.MockNodeModule ) - send(manager_pid, :discover_and_restart) - - # Wait for discovery to complete or manager to crash - wait_for_discovery_completion(manager_pid, 100) + discover_and_wait(manager_pid) # Should be able to claim since TTL expired {:ok, data} = @@ -2667,10 +2588,7 @@ defmodule DurableServer.LifecycleTest do create_test_object(config.object_store, "#{prefix}#{key}", %{count: 0}, meta_attrs) - send(manager_pid, :discover_and_restart) - - # Wait for discovery to complete or manager to crash - wait_for_discovery_completion(manager_pid, 100) + discover_and_wait(manager_pid) assert_process_alive(manager_pid) GenServer.stop(manager_pid) @@ -2700,10 +2618,7 @@ defmodule DurableServer.LifecycleTest do {:ok, manager_pid} = start_standalone_lifecycle_manager(supervisor_name, config) - send(manager_pid, :discover_and_restart) - - # Wait longer for restart attempt and cleanup to complete - wait_for_discovery_completion(manager_pid, 200) + discover_and_wait(manager_pid) # Check that restart attempt metadata was cleaned up after failure {:ok, data} = @@ -2740,8 +2655,7 @@ defmodule DurableServer.LifecycleTest do create_test_object(config.object_store, "#{prefix}#{key}", %{count: 0}, meta_attrs) # First discovery - should not be orphaned - send(manager_pid, :discover_and_restart) - wait_for_discovery_completion(manager_pid, 100) + discover_and_wait(manager_pid) # Update object to be orphaned (old heartbeat) old_time = System.system_time(:millisecond) - 5 * 60 * 1000 @@ -2749,8 +2663,7 @@ defmodule DurableServer.LifecycleTest do create_test_object(config.object_store, "#{prefix}#{key}", %{count: 0}, meta_attrs) # Second discovery - should now be orphaned - send(manager_pid, :discover_and_restart) - wait_for_discovery_completion(manager_pid, 100) + discover_and_wait(manager_pid) assert_process_alive(manager_pid) GenServer.stop(manager_pid) @@ -2799,11 +2712,7 @@ defmodule DurableServer.LifecycleTest do {:ok, manager_pid} = start_standalone_lifecycle_manager(supervisor_name, config) - # Trigger discovery and wait for completion - send(manager_pid, :discover_and_restart) - - # Wait for discovery task to complete - wait_for_discovery_completion(manager_pid, 200) + discover_and_wait(manager_pid) # Manager should still be alive and ready for next cycle assert_process_alive(manager_pid) @@ -3172,9 +3081,7 @@ defmodule DurableServer.LifecycleTest do {:ok, manager_pid} = start_standalone_lifecycle_manager(supervisor_name, config) - # Trigger discovery - send(manager_pid, :discover_and_restart) - wait_for_discovery_completion(manager_pid, 1000) + discover_and_wait(manager_pid) # Verify server was NOT restarted (should remain crashed) {:ok, data} = @@ -3212,9 +3119,7 @@ defmodule DurableServer.LifecycleTest do {:ok, manager_pid} = start_standalone_lifecycle_manager(supervisor_name, config) - # Trigger discovery - send(manager_pid, :discover_and_restart) - wait_for_discovery_completion(manager_pid, 1000) + discover_and_wait(manager_pid) # Verify server was NOT restarted (should remain permanently crashed) {:ok, data} = @@ -3247,8 +3152,7 @@ defmodule DurableServer.LifecycleTest do {:ok, manager_pid} = start_standalone_lifecycle_manager(supervisor_name, config) - send(manager_pid, :discover_and_restart) - wait_for_discovery_completion(manager_pid, 1000) + discover_and_wait(manager_pid) {:ok, data} = DurableServer.fetch_stored_state(config.object_store, %{key: key, prefix: prefix}) diff --git a/test/durable_server/object_store_retry_test.exs b/test/durable_server/object_store_retry_test.exs index d61ef64..58aa44d 100644 --- a/test/durable_server/object_store_retry_test.exs +++ b/test/durable_server/object_store_retry_test.exs @@ -3,7 +3,9 @@ defmodule DurableServer.ObjectStoreRetryTest do alias DurableServer.ObjectStore - test "put retries transient Req failures within its deadline" do + # Retry outcomes are scripted, so these cases don't impose a wall-clock + # deadline while other test modules compile. Deadline behavior is covered below. + test "put retries transient Req failures" do adapter = adapter([ %Req.Response{status: 503}, @@ -14,7 +16,7 @@ defmodule DurableServer.ObjectStoreRetryTest do assert {:ok, %{etag: "retry-etag", body: "heartbeat"}} = ObjectStore.put_object(store(adapter), "__nodes/test@localhost", "heartbeat", max_retries: 10, - timeout: 1_000 + timeout: :infinity ) end @@ -45,7 +47,7 @@ defmodule DurableServer.ObjectStoreRetryTest do assert {:ok, %{etag: "retry-etag"}} = ObjectStore.put_object(store(adapter), "__nodes/test@localhost", "heartbeat", max_retries: 1, - timeout: 1_000 + timeout: :infinity ) for _attempt <- 1..2 do @@ -66,7 +68,7 @@ defmodule DurableServer.ObjectStoreRetryTest do "__nodes/test@localhost", "heartbeat", max_retries: 10, - timeout: 1_000 + timeout: :infinity ) assert Process.get(responses_key) == [:unexpected_retry] @@ -87,12 +89,34 @@ defmodule DurableServer.ObjectStoreRetryTest do "__nodes/test@localhost", "heartbeat", max_retries: 1, - timeout: 1_000 + timeout: :infinity ) assert Process.get(responses_key) == [:unexpected_retry] end + test "an expired operation deadline prevents otherwise retryable failures from retrying" do + for failure <- [ + %Req.Response{status: 503}, + %Req.TransportError{reason: :timeout}, + %Req.HTTPError{protocol: :http2, reason: :unprocessed} + ] do + responses_key = make_ref() + Process.put(responses_key, [failure, :unexpected_retry]) + + assert {:error, ^failure} = + ObjectStore.put_object( + store(adapter(responses_key)), + "__nodes/test@localhost", + "heartbeat", + max_retries: 10, + timeout: 0 + ) + + assert Process.get(responses_key) == [:unexpected_retry] + end + end + test "finite operation deadlines cap each HTTP receive attempt" do parent = self() diff --git a/test/durable_server/remote_placement_test.exs b/test/durable_server/remote_placement_test.exs index ff8660b..4769a5c 100644 --- a/test/durable_server/remote_placement_test.exs +++ b/test/durable_server/remote_placement_test.exs @@ -1,5 +1,5 @@ defmodule DurableServer.RemotePlacementTest do - use ExUnit.Case, async: false + use DurableServer.LocalStackCase, async: false import DurableServer.TestHelper alias DurableServer diff --git a/test/durable_server/sticky_placement_test.exs b/test/durable_server/sticky_placement_test.exs index ce491ce..98fda72 100644 --- a/test/durable_server/sticky_placement_test.exs +++ b/test/durable_server/sticky_placement_test.exs @@ -1,5 +1,5 @@ defmodule DurableServer.StickyPlacementTest do - use ExUnit.Case, async: false + use DurableServer.LocalStackCase, async: false import DurableServer.TestHelper alias DurableServer @@ -33,6 +33,14 @@ defmodule DurableServer.StickyPlacementTest do end setup do + previous_env = Map.new(["FLY_MACHINE_ID", "FLY_REGION"], &{&1, System.get_env(&1)}) + + on_exit(fn -> + for {name, value} <- previous_env do + if is_nil(value), do: System.delete_env(name), else: System.put_env(name, value) + end + end) + supervisor_name = :"test_supervisor_#{:erlang.unique_integer([:positive])}" prefix = "sticky_placement_test_#{:erlang.unique_integer([:positive])}/" @@ -265,10 +273,6 @@ defmodule DurableServer.StickyPlacementTest do %{env_var: "FLY_REGION", value: "ord"}, %{env_var: :any, value: :any} ] - - # Cleanup - System.delete_env("FLY_MACHINE_ID") - System.delete_env("FLY_REGION") end test "builds sticky_placement when starting a child", %{ @@ -312,10 +316,6 @@ defmodule DurableServer.StickyPlacementTest do %{env_var: "FLY_REGION", value: "sjc"}, %{env_var: :any, value: :any} ] - - # Cleanup - System.delete_env("FLY_MACHINE_ID") - System.delete_env("FLY_REGION") end test "handles nil env var values", %{supervisor_name: supervisor_name, prefix: prefix} do @@ -354,16 +354,6 @@ defmodule DurableServer.StickyPlacementTest do supervisor_name: supervisor_name, prefix: prefix } do - previous_region = System.get_env("FLY_REGION") - - on_exit(fn -> - if previous_region do - System.put_env("FLY_REGION", previous_region) - else - System.delete_env("FLY_REGION") - end - end) - System.put_env("FLY_REGION", "ord") start_supervised!( @@ -441,10 +431,6 @@ defmodule DurableServer.StickyPlacementTest do [] -> flunk("Expected heartbeat entry to exist") end - - # Cleanup - System.delete_env("FLY_MACHINE_ID") - System.delete_env("FLY_REGION") end test "env_vars is empty map when no sticky placement", %{ @@ -635,8 +621,6 @@ defmodule DurableServer.StickyPlacementTest do # With no sticky config, augmented should be nil assert augmented == nil - - System.delete_env("FLY_REGION") end test "sticky config with :any allows all nodes to match", %{ @@ -676,8 +660,6 @@ defmodule DurableServer.StickyPlacementTest do assert length(augmented) == 2 assert Enum.any?(augmented, fn p -> p.env_var == "FLY_REGION" and p.value == "ord" end) assert Enum.any?(augmented, fn p -> p.env_var == :any and p.value == :any end) - - System.delete_env("FLY_REGION") end test "sticky config WITHOUT :any does not include :any fallback", %{ @@ -715,8 +697,6 @@ defmodule DurableServer.StickyPlacementTest do assert length(augmented) == 1 assert Enum.any?(augmented, fn p -> p.env_var == "FLY_REGION" and p.value == "ord" end) refute Enum.any?(augmented, fn p -> p.env_var == :any end) - - System.delete_env("FLY_REGION") end test "non-matching node returns nil level when :any not configured", %{ @@ -782,8 +762,6 @@ defmodule DurableServer.StickyPlacementTest do # Non-matching node with no :any should get nil assert matching_level == nil - - System.delete_env("FLY_REGION") end test "matching node returns level 0 for sticky placement", %{ @@ -836,8 +814,6 @@ defmodule DurableServer.StickyPlacementTest do # Matching node should get level 0 assert matching_level == 0 - - System.delete_env("FLY_REGION") end end end diff --git a/test/durable_server/test_helper_test.exs b/test/durable_server/test_helper_test.exs new file mode 100644 index 0000000..9c367d3 --- /dev/null +++ b/test/durable_server/test_helper_test.exs @@ -0,0 +1,57 @@ +defmodule DurableServer.TestHelperTest do + use ExUnit.Case, async: true + import DurableServer.TestHelper + + test "eventual assertions retry the predicate and return as soon as it succeeds" do + attempts = make_ref() + + assert :ok = + assert_eventually( + fn -> + count = Process.get(attempts, 0) + Process.put(attempts, count + 1) + count == 1 + end, + 5_000, + 0 + ) + + assert Process.get(attempts) == 2 + end + + test "eventual assertions fail when the predicate stays false until the deadline" do + assert_raise ExUnit.AssertionError, ~r/condition was not met within timeout/, fn -> + assert_eventually(fn -> false end, 0) + end + end + + test "object-store configuration uses the suite bucket without connecting to storage" do + opts = test_object_store_opts() + assert opts[:bucket] == Application.fetch_env!(:durable_server, :test_object_store_bucket) + assert opts[:bucket] =~ ~r/^durable-test-[0-9a-f-]{36}$/ + assert test_object_store().bucket == opts[:bucket] + assert test_object_store_opts(bucket: "explicit-fixture")[:bucket] == "explicit-fixture" + end + + test "data directories are unique, worktree-local, and cleaned up on exit" do + # on_exit callbacks run in reverse registration order. Keep the allocated + # paths alive past this test process so the last callback can verify cleanup. + {:ok, paths} = Agent.start(fn -> [] end) + + on_exit(fn -> + allocated = Agent.get(paths, & &1) + Agent.stop(paths) + for path <- allocated, do: refute(File.exists?(path)) + end) + + first = test_data_dir("isolation") + second = test_data_dir("isolation") + Agent.update(paths, fn _ -> [first, second] end) + + assert first != second + assert Path.dirname(first) == Path.expand("tmp") + assert Path.dirname(second) == Path.expand("tmp") + assert File.dir?(first) + assert File.dir?(second) + end +end diff --git a/test/durable_server/watermark_test.exs b/test/durable_server/watermark_test.exs index a8a6f7a..1266e10 100644 --- a/test/durable_server/watermark_test.exs +++ b/test/durable_server/watermark_test.exs @@ -1,5 +1,5 @@ defmodule DurableServer.WatermarkTest do - use ExUnit.Case, async: false + use DurableServer.LocalStackCase, async: false import DurableServer.TestHelper alias DurableServer @@ -201,16 +201,27 @@ defmodule DurableServer.WatermarkTest do ) # Terminate one - Process.monitor(pid1) - DurableServer.Supervisor.terminate_child(supervisor_name, pid1) - assert_receive {:DOWN, _ref, :process, ^pid1, :normal} - - # Now third should succeed + ref = Process.monitor(pid1) + assert :ok = DurableServer.Supervisor.terminate_child(supervisor_name, pid1) + assert_receive {:DOWN, ^ref, :process, ^pid1, :normal} + + # The reservation keeper processes its own DOWN asynchronously. Observing + # the child exit does not mean that keeper has released its capacity yet. + assert_eventually(fn -> + DurableServer.Supervisor.current_capacity(supervisor_name) == + %{total: %{current: 1, limit: 2}} + end) + + # Now the third child should fit locally, without remote placement fallback. assert {:ok, {_pid3, _}} = DurableServer.Supervisor.start_child( supervisor_name, - {WatermarkTestServer, key: "key3", initial_state: %{}} + {WatermarkTestServer, key: "key3", initial_state: %{}}, + max_placement_retries: 0 ) + + assert DurableServer.Supervisor.current_capacity(supervisor_name) == + %{total: %{current: 2, limit: 2}} end test "enforces a map limit atomically across concurrent starts", %{ diff --git a/test/durable_server_test.exs b/test/durable_server_test.exs index 3337697..e3beea5 100644 --- a/test/durable_server_test.exs +++ b/test/durable_server_test.exs @@ -1,5 +1,5 @@ defmodule DurableServerTest do - use ExUnit.Case, async: true + use DurableServer.LocalStackCase, async: true import ExUnit.CaptureLog import DurableServer.TestHelper @@ -935,34 +935,9 @@ defmodule DurableServerTest do end setup do - # Create a test bucket - test_bucket_name = - "durable-test-durable-#{DurableServer.UUID.uuid4()}" - - case ObjectStore.create_bucket_with_credentials(test_object_store(), test_bucket_name) do - {:ok, %ObjectStore{} = store} -> - on_exit(fn -> - # Clean up bucket on test completion - try do - ObjectStore.delete_bucket(store, test_bucket_name) - catch - _, _ -> :ok - end - end) + {supervisor_name, supervisor_pid, prefix} = start_test_supervisor() - # Start a DurableServer.Supervisor for tests that need one - {supervisor_name, supervisor_pid, prefix} = start_test_supervisor() - - {:ok, - test_bucket: test_bucket_name, - store: store, - supervisor_name: supervisor_name, - supervisor_pid: supervisor_pid, - prefix: prefix} - - {:error, reason} -> - {:skip, "Failed to create test bucket: #{inspect(reason)}"} - end + {:ok, supervisor_name: supervisor_name, supervisor_pid: supervisor_pid, prefix: prefix} end describe "init/1" do diff --git a/test/ekv_integration_test.exs b/test/ekv_integration_test.exs index a950d2d..2dafec7 100644 --- a/test/ekv_integration_test.exs +++ b/test/ekv_integration_test.exs @@ -1,5 +1,6 @@ defmodule DurableServer.EKVIntegrationTest do use ExUnit.Case, async: false + import DurableServer.TestHelper alias DurableServer.Backends.EKVStore alias DurableServer.{LifecycleManager, Meta, StoredState} @@ -7,7 +8,7 @@ defmodule DurableServer.EKVIntegrationTest do alias DurableServer.TestCounterServer, as: CounterServer alias DurableServer.TestTemporalServer - @moduletag :integration + @moduletag :ekv @moduletag :capture_log setup do @@ -16,9 +17,7 @@ defmodule DurableServer.EKVIntegrationTest do ekv_name = :"durable_ekv_integration_#{unique_id}" supervisor_name = :"durable_ekv_supervisor_#{unique_id}" prefix = "ekv_integration/#{unique_id}/" - data_dir = Path.join(System.tmp_dir!(), "durable_server_ekv_integration_#{unique_id}") - - File.rm_rf(data_dir) + data_dir = test_data_dir("ekv_integration") start_supervised!( {ekv_mod(), @@ -41,10 +40,6 @@ defmodule DurableServer.EKVIntegrationTest do ]} ) - on_exit(fn -> - File.rm_rf(data_dir) - end) - {:ok, supervisor_name: supervisor_name, prefix: prefix, ekv_name: ekv_name, data_dir: data_dir} end @@ -185,16 +180,12 @@ defmodule DurableServer.EKVIntegrationTest do ensure_distributed_node!() unique_id = System.unique_integer([:positive, :monotonic]) - peer_name = :"durable_ekv_client_peer_#{unique_id}" + peer_name = :"durable_ekv_client_peer_#{DurableServer.UUID.uuid4()}" ekv_name = :"durable_ekv_client_cluster_#{unique_id}" - remote_data_dir = - Path.join(System.tmp_dir!(), "durable_server_ekv_client_remote_#{unique_id}") - + remote_data_dir = test_data_dir("ekv_client_remote") key = "client-existing-key" - File.rm_rf(remote_data_dir) - {:ok, peer, peer_node} = :peer.start_link(%{name: peer_name}) on_exit(fn -> @@ -203,8 +194,6 @@ defmodule DurableServer.EKVIntegrationTest do catch :exit, _ -> :ok end - - File.rm_rf(remote_data_dir) end) assert Node.connect(peer_node) @@ -503,17 +492,14 @@ defmodule DurableServer.EKVIntegrationTest do ensure_distributed_node!() unique_id = System.unique_integer([:positive, :monotonic]) - peer_name = :"durable_ekv_peer_#{unique_id}" + peer_name = :"durable_ekv_peer_#{DurableServer.UUID.uuid4()}" ekv_name = :"durable_ekv_cluster_#{unique_id}" supervisor_name = :"durable_ekv_cluster_sup_#{unique_id}" prefix = "ekv_cluster/#{unique_id}/" - local_data_dir = Path.join(System.tmp_dir!(), "durable_server_ekv_local_#{unique_id}") - remote_data_dir = Path.join(System.tmp_dir!(), "durable_server_ekv_remote_#{unique_id}") + local_data_dir = test_data_dir("ekv_local") + remote_data_dir = test_data_dir("ekv_remote") key = "seeded-restart" - File.rm_rf(local_data_dir) - File.rm_rf(remote_data_dir) - {:ok, peer, peer_node} = :peer.start_link(%{name: peer_name}) on_exit(fn -> @@ -522,9 +508,6 @@ defmodule DurableServer.EKVIntegrationTest do catch :exit, _ -> :ok end - - File.rm_rf(local_data_dir) - File.rm_rf(remote_data_dir) end) assert Node.connect(peer_node) @@ -701,28 +684,9 @@ defmodule DurableServer.EKVIntegrationTest do if Node.alive?() do :ok else - name = :"durable_server_test_#{System.unique_integer([:positive, :monotonic])}" + name = :"durable_server_test_#{DurableServer.UUID.uuid4()}" {:ok, _} = Node.start(name, :shortnames) :ok end end - - defp assert_eventually(fun, timeout \\ 2_000, interval \\ 25) - when is_function(fun, 0) and is_integer(timeout) and timeout > 0 do - deadline = System.monotonic_time(:millisecond) + timeout - do_assert_eventually(fun, deadline, interval) - end - - defp do_assert_eventually(fun, deadline, interval) do - if fun.() do - :ok - else - if System.monotonic_time(:millisecond) >= deadline do - flunk("eventual assertion timed out") - else - Process.sleep(interval) - do_assert_eventually(fun, deadline, interval) - end - end - end end diff --git a/test/ekv_store_test.exs b/test/ekv_store_test.exs index 3c8877b..82a73ff 100644 --- a/test/ekv_store_test.exs +++ b/test/ekv_store_test.exs @@ -58,7 +58,9 @@ defmodule DurableServer.EKVStoreTest do name: name, cas_retries: 2, backoff: {0, 0}, - timeout: 50, + # These tests script retry outcomes, not elapsed time. Allow scheduler + # delays while other test modules are compiling/running concurrently. + timeout: 1_000, ekv_mod: FakeEKV, ekv_supervisor_mod: FakeEKVSupervisor ] diff --git a/test/group_test.exs b/test/group_test.exs index 04d7551..7d9f6be 100644 --- a/test/group_test.exs +++ b/test/group_test.exs @@ -1,5 +1,5 @@ defmodule GroupTest do - use ExUnit.Case, async: true + use DurableServer.LocalStackCase, async: true import DurableServer.TestHelper @moduletag :capture_log diff --git a/test/mirror_backend_e2e_test.exs b/test/mirror_backend_e2e_test.exs index 4eb56a1..68909db 100644 --- a/test/mirror_backend_e2e_test.exs +++ b/test/mirror_backend_e2e_test.exs @@ -1,5 +1,5 @@ defmodule DurableServer.MirrorBackendE2ETest do - use ExUnit.Case, async: false + use DurableServer.LocalStackCase, async: false import DurableServer.TestHelper @@ -8,16 +8,13 @@ defmodule DurableServer.MirrorBackendE2ETest do alias DurableServer.StorageBackend alias DurableServer.TestCounterServer, as: CounterServer - @moduletag :integration @moduletag :capture_log setup do unique_id = System.unique_integer([:positive, :monotonic]) prefix = "mirror_e2e/#{unique_id}/" ekv_name = :"durable_mirror_e2e_#{unique_id}" - data_dir = Path.join(System.tmp_dir!(), "durable_server_mirror_e2e_#{unique_id}") - - File.rm_rf(data_dir) + data_dir = test_data_dir("mirror_e2e") object_store = test_object_store() @@ -37,10 +34,6 @@ defmodule DurableServer.MirrorBackendE2ETest do :ok = StorageBackend.ensure_ready(primary_backend) :ok = StorageBackend.ensure_ready(secondary_backend) - on_exit(fn -> - File.rm_rf(data_dir) - end) - {:ok, prefix: prefix, object_store: object_store, @@ -108,9 +101,7 @@ defmodule DurableServer.MirrorBackendE2ETest do supervisor_name = unique_supervisor_name("shadow_write_fail") prefix = "mirror_e2e/write_fail/#{unique_id}/" key = "shadow-write-fail" - data_dir = Path.join(System.tmp_dir!(), "durable_server_mirror_write_fail_#{unique_id}") - - File.rm_rf(data_dir) + data_dir = test_data_dir("mirror_write_fail") ekv_pid = start_supervised!(%{ @@ -129,10 +120,6 @@ defmodule DurableServer.MirrorBackendE2ETest do ]} }) - on_exit(fn -> - File.rm_rf(data_dir) - end) - _supervisor = start_durable_supervisor!( {:shadow_write_fail, supervisor_name}, @@ -452,23 +439,4 @@ defmodule DurableServer.MirrorBackendE2ETest do end defp ekv_mod, do: :"Elixir.EKV" - - defp assert_eventually(fun, timeout \\ 5_000, interval \\ 25) - when is_function(fun, 0) and is_integer(timeout) and timeout > 0 do - deadline = System.monotonic_time(:millisecond) + timeout - do_assert_eventually(fun, deadline, interval) - end - - defp do_assert_eventually(fun, deadline, interval) do - if fun.() do - :ok - else - if System.monotonic_time(:millisecond) >= deadline do - flunk("eventual assertion timed out") - else - Process.sleep(interval) - do_assert_eventually(fun, deadline, interval) - end - end - end end diff --git a/test/mirror_backend_integration_test.exs b/test/mirror_backend_integration_test.exs index d55683f..991ca57 100644 --- a/test/mirror_backend_integration_test.exs +++ b/test/mirror_backend_integration_test.exs @@ -1,5 +1,6 @@ defmodule DurableServer.MirrorBackendIntegrationTest do use ExUnit.Case, async: false + import DurableServer.TestHelper alias DurableServer.StorageBackend alias DurableServer.Backends.EKVStore @@ -55,7 +56,7 @@ defmodule DurableServer.MirrorBackendIntegrationTest do do: StorageBackend.unsubscribe(delegate, subscription_ref) end - @moduletag :integration + @moduletag :ekv @moduletag :capture_log setup do @@ -64,13 +65,8 @@ defmodule DurableServer.MirrorBackendIntegrationTest do primary_name = :"durable_mirror_primary_#{unique_id}" secondary_name = :"durable_mirror_secondary_#{unique_id}" - primary_dir = Path.join(System.tmp_dir!(), "durable_server_mirror_primary_#{unique_id}") - - secondary_dir = - Path.join(System.tmp_dir!(), "durable_server_mirror_secondary_#{unique_id}") - - File.rm_rf(primary_dir) - File.rm_rf(secondary_dir) + primary_dir = test_data_dir("mirror_primary") + secondary_dir = test_data_dir("mirror_secondary") start_supervised!( {ekv_mod(), @@ -98,11 +94,6 @@ defmodule DurableServer.MirrorBackendIntegrationTest do secondary = StorageBackend.new(EKVStore, EKVStore.normalize_opts(name: secondary_name)) mirror = mirror_backend(primary, secondary) - on_exit(fn -> - File.rm_rf(primary_dir) - File.rm_rf(secondary_dir) - end) - {:ok, primary: primary, secondary: secondary, mirror: mirror} end diff --git a/test/support/lifecycle_helpers.ex b/test/support/lifecycle_helpers.ex new file mode 100644 index 0000000..4bf9112 --- /dev/null +++ b/test/support/lifecycle_helpers.ex @@ -0,0 +1,96 @@ +defmodule DurableServer.LifecycleHelpers do + @moduledoc """ + Synchronization helpers for LifecycleManager tests. + """ + import ExUnit.Assertions + + @doc """ + Triggers discovery and waits for that task's successful result to be processed. + + The GenServer debug hook observes the real message loop without adding a + production test API. An idle manager, a failed task, and a completed cycle are + distinct outcomes. The hook and monitor are removed even when an assertion fails. + + Fixtures asserting individual cycles should disable startup discovery bursts + and use long automatic discovery intervals so those sweeps do not race the test. + """ + def discover_and_wait(manager_pid, timeout \\ 5_000) do + owner = self() + ref = Process.monitor(manager_pid) + handler = fn stage, event, _ -> observe_discovery(stage, event, owner, ref) end + + try do + :ok = :sys.install(manager_pid, {ref, handler, :awaiting_request}) + send(manager_pid, :discover_and_restart) + + receive do + {:discovery, ^ref, :ok} -> + :ok + + {:discovery, ^ref, {:error, reason}} -> + flunk("Discovery failed: #{inspect(reason)}") + + {:DOWN, ^ref, :process, ^manager_pid, reason} -> + flunk("LifecycleManager crashed: #{inspect(reason)}") + after + timeout -> + flunk("Discovery did not complete within #{timeout}ms") + end + after + try do + :sys.remove(manager_pid, ref) + catch + :exit, _ -> :ok + end + + Process.demonitor(ref, [:flush]) + + receive do + {:discovery, ^ref, _result} -> :ok + after + 0 -> :ok + end + end + end + + defp observe_discovery(:awaiting_request, {:in, :discover_and_restart}, _owner, _ref), + do: :starting + + defp observe_discovery( + :starting, + {:noreply, %{current_discovery_task: %Task{ref: task_ref}}}, + _owner, + _ref + ), + do: {:running, task_ref} + + defp observe_discovery(:starting, {:noreply, _state}, owner, ref) do + send(owner, {:discovery, ref, {:error, :task_not_started}}) + :done + end + + defp observe_discovery({:running, task_ref}, {:in, {task_ref, {:discover, :ok}}}, _owner, _ref), + do: :finishing + + defp observe_discovery( + :finishing, + {:noreply, %{current_discovery_task: nil}}, + owner, + ref + ) do + send(owner, {:discovery, ref, :ok}) + :done + end + + defp observe_discovery( + {:running, task_ref}, + {:in, {:DOWN, task_ref, :process, _pid, reason}}, + owner, + ref + ) do + send(owner, {:discovery, ref, {:error, {:task_exit, reason}}}) + :done + end + + defp observe_discovery(stage, _event, _owner, _ref), do: stage +end diff --git a/test/support/local_stack_case.ex b/test/support/local_stack_case.ex new file mode 100644 index 0000000..3e9eb41 --- /dev/null +++ b/test/support/local_stack_case.ex @@ -0,0 +1,22 @@ +defmodule DurableServer.LocalStackCase do + @moduledoc """ + Opt-in LocalStack tests sharing a bucket unique to the current suite invocation. + + Tests still use unique prefixes within that bucket. Cleanup happens after all + test supervisors have stopped, never at the beginning of another suite run. + """ + use ExUnit.CaseTemplate + + using do + quote do + @moduletag :localstack + end + end + + setup_all do + store = DurableServer.TestHelper.test_object_store() + :ok = DurableServer.ObjectStore.ensure_bucket_exists(store) + Application.put_env(:durable_server, :test_object_store_used, true) + :ok + end +end diff --git a/test/support/test_helper.ex b/test/support/test_helper.ex index 4054343..b95d90b 100644 --- a/test/support/test_helper.ex +++ b/test/support/test_helper.ex @@ -3,10 +3,37 @@ defmodule DurableServer.TestHelper do Test helpers for DurableServer tests. """ + import ExUnit.Assertions, only: [flunk: 1] + alias DurableServer.ObjectStore @doc """ - Returns the default object store config for testing as a keyword list. + Waits for a predicate to become truthy, failing if it stays false until the deadline. + """ + def assert_eventually(fun, timeout \\ 5_000, interval \\ 25) + when is_function(fun, 0) and is_integer(timeout) and timeout >= 0 do + deadline = System.monotonic_time(:millisecond) + timeout + do_assert_eventually(fun, deadline, interval) + end + + defp do_assert_eventually(fun, deadline, interval) do + cond do + fun.() -> + :ok + + System.monotonic_time(:millisecond) >= deadline -> + flunk("condition was not met within timeout") + + true -> + Process.sleep(interval) + do_assert_eventually(fun, deadline, interval) + end + end + + @doc """ + Returns LocalStack options with a bucket unique to this suite invocation. + + This only builds configuration; it does not contact storage. """ def test_object_store_opts(opts \\ []) do Keyword.merge( @@ -16,7 +43,7 @@ defmodule DurableServer.TestHelper do s3_endpoint: "http://localhost:4566", iam_endpoint: "http://localhost:4566", default_region: "us-east-1", - bucket: "durable-test-bucket" + bucket: Application.fetch_env!(:durable_server, :test_object_store_bucket) ], opts ) @@ -24,10 +51,29 @@ defmodule DurableServer.TestHelper do @doc """ Creates an ObjectStore configured for testing. - - Uses environment variables or defaults suitable for LocalStack. """ def test_object_store(opts \\ []) do ObjectStore.new(test_object_store_opts(opts)) end + + @doc false + def clean_test_object_store! do + store = test_object_store() + + for object <- ObjectStore.list_all_objects_stream(store, "") do + :ok = ObjectStore.delete_object(store, object.key) + end + + :ok = ObjectStore.delete_bucket(store, store.bucket) + end + + @doc """ + Allocates an isolated data directory inside the worktree and removes it on exit. + """ + def test_data_dir(label) do + path = Path.expand(Path.join("tmp", "#{label}-#{DurableServer.UUID.uuid4()}")) + File.mkdir_p!(path) + ExUnit.Callbacks.on_exit(fn -> File.rm_rf!(path) end) + path + end end diff --git a/test/test_helper.exs b/test/test_helper.exs index c438299..ee96ae4 100644 --- a/test/test_helper.exs +++ b/test/test_helper.exs @@ -14,18 +14,23 @@ case File.read(".env") do :noop end -# Exclude integration tests by default (they require real credentials) -ExUnit.configure(exclude: [:integration, :stress]) +# Selecting storage-free tests must not contact LocalStack. Each suite invocation +# gets its own bucket; LocalStackCase creates it only when one of its tests runs. +Application.put_env( + :durable_server, + :test_object_store_bucket, + "durable-test-#{DurableServer.UUID.uuid4()}" +) -alias DurableServer.ObjectStore -import DurableServer.TestHelper +Application.put_env(:durable_server, :test_object_store_used, false) -# Clear object store (local stack) for this run -store = test_object_store() -:ok = ObjectStore.ensure_bucket_exists(store) +ExUnit.start( + exclude: [:localstack, :ekv, :integration, :stress], + assert_receive_timeout: 1_000 +) -for obj <- ObjectStore.list_all_objects_stream(store, "") do - :ok = ObjectStore.delete_object(store, obj.key) -end - -ExUnit.start() +ExUnit.after_suite(fn _results -> + if Application.fetch_env!(:durable_server, :test_object_store_used) do + DurableServer.TestHelper.clean_test_object_store!() + end +end)