diff --git a/test/durable_server/lifecycle_test.exs b/test/durable_server/lifecycle_test.exs index 18686e8..8011add 100644 --- a/test/durable_server/lifecycle_test.exs +++ b/test/durable_server/lifecycle_test.exs @@ -2037,27 +2037,6 @@ defmodule DurableServer.LifecycleTest do end describe "lifecycle manager edge cases" do - test "handles object store failures during discovery gracefully", %{ - supervisor_name: supervisor_name, - prefix: _prefix, - config: config, - circuit_breaker: _circuit_breaker - } do - # This test would be better with ObjectStore mocking, but we can at least verify - # the LifecycleManager doesn't crash when encountering errors - {: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) - - assert_process_alive(pid) - GenServer.stop(pid) - end - test "restart claimer gate uses LM-local gate age, not object age" do supervisor_name = :"restart_gate_#{DurableServer.UUID.uuid4()}" setup_restart_gate_tables(supervisor_name) @@ -2585,57 +2564,69 @@ defmodule DurableServer.LifecycleTest do GenServer.stop(manager_pid) end - test "restart claim TTL expiration during process", %{ - supervisor_name: supervisor_name, - prefix: prefix, - config: config, - circuit_breaker: _circuit_breaker - } do - key = "ttl-expire-test-#{DurableServer.UUID.uuid4()}" - - # Create object with expired restart attempt from another node - current_time = System.system_time(:millisecond) - # Already expired - expired_ttl = current_time - 1000 - - meta_attrs = %{ - status: :running, - node_str: "other_node@test", - node_ref: "other-ref", - pid: self(), - last_heartbeat_at: current_time, - restart_attempt_node: "other_node@test", - # 31 seconds ago - restart_attempt_time: current_time - 31_000, - restart_attempt_ttl: expired_ttl, - module: TestServer - } - - create_test_object(config.object_store, "#{prefix}#{key}", %{count: 0}, meta_attrs) + for {claim_state, ttl_offset} <- [expired: -1_000, active: 60_000] do + test "discovery respects a #{claim_state} foreign restart claim", %{ + supervisor_name: supervisor_name, + prefix: prefix, + config: config + } do + key = "ttl-#{DurableServer.UUID.uuid4()}" + pid = start_test_server(supervisor_name, key) + assert GenServer.call(pid, :increment) == 1 + assert :ok = DurableServer.Supervisor.terminate_child(supervisor_name, pid) + + {:ok, stored} = + DurableServer.fetch_stored_state(config.object_store, %{key: key, prefix: prefix}) - {:ok, manager_pid} = - start_standalone_lifecycle_manager(supervisor_name, config, - node_module: DurableServer.LifecycleTest.MockNodeModule - ) + now = System.system_time(:millisecond) + + claimed_meta = %{ + stored.meta + | status: :crashed, + permanent: true, + node_str: "dead_node@test", + node_ref: "dead-ref", + last_heartbeat_at: now - 600_000, + restart_attempt_node: "other_node@test", + restart_attempt_time: now - 31_000, + restart_attempt_ttl: now + unquote(ttl_offset) + } - send(manager_pid, :discover_and_restart) + {:ok, _} = + ObjectStore.put_object( + config.object_store, + prefix <> key, + encode_legacy_stored_state(%{stored | meta: Meta.encode_to_binary(claimed_meta)}) + ) - # Wait for discovery to complete or manager to crash - wait_for_discovery_completion(manager_pid, 100) + {:ok, manager} = get_supervisor_lifecycle_manager(supervisor_name) + send(manager, :discover_and_restart) + assert_eventually(fn -> :sys.get_state(manager).current_discovery_task == nil end) - # Should be able to claim since TTL expired - {:ok, data} = - DurableServer.fetch_stored_state(config.object_store, %{key: key, prefix: prefix}) + if unquote(claim_state) == :expired do + assert {new_pid, _} = DurableServer.Supervisor.lookup(supervisor_name, key) + refute new_pid == pid + assert GenServer.call(new_pid, :get_count) == 1 - # The object might not actually be processed if it's not considered orphaned - # Since the node is "other_node@test" and MockNodeModule.ping returns :pong for everything, - # the orphan check might not trigger as expected. This tests the TTL logic exists - # but the specific scenario may not result in processing due to other conditions. + assert {:ok, %{meta: %Meta{status: :running, pid: ^new_pid} = meta}} = + DurableServer.fetch_stored_state( + config.object_store, + %{key: key, prefix: prefix} + ) - # Verify the test setup ran without errors - assert is_map(data) + assert meta.restart_attempt_node == nil + assert meta.restart_attempt_time == nil + assert meta.restart_attempt_ttl == nil + else + assert DurableServer.Supervisor.lookup(supervisor_name, key) == nil - GenServer.stop(manager_pid) + assert {:ok, %{meta: ^claimed_meta}} = + DurableServer.fetch_stored_state( + config.object_store, + %{key: key, prefix: prefix} + ) + end + end end test "node reachable but lock validation fails", %{ @@ -2762,32 +2753,55 @@ defmodule DurableServer.LifecycleTest do } do table = :ets.new(__MODULE__.DiscoveryCrashOnceBackend, [:set, :public]) delegate = DurableServer.Supervisor.__get_config__(supervisor_name).storage_backend + recovery_supervisor = :"recovery_#{DurableServer.UUID.uuid4()}" + prefix = "recovery_#{DurableServer.UUID.uuid4()}/" + key = "recover-after-list-failure" - {:ok, storage_backend} = - DurableServer.StorageBackend.init_backend(DiscoveryCrashOnceBackend, - delegate: delegate, - table: table, - notify_pid: self() - ) - - test_config = - config - |> Map.put(:object_store, storage_backend) - |> Map.put(:storage_backend, storage_backend) - |> Map.put(:initial_discovery_delay_ms, 60_000) - |> Map.put(:discovery_interval_ms, 10) + start_supervised!( + {DurableServer.Supervisor, + name: recovery_supervisor, + prefix: prefix, + backend: + {DiscoveryCrashOnceBackend, delegate: delegate, table: table, notify_pid: self()}, + initial_discovery_delay_ms: 60_000, + discovery_interval_ms: 10} + ) - {:ok, manager_pid} = start_standalone_lifecycle_manager(supervisor_name, test_config) + assert {:ok, _} = + create_test_object(config.object_store, prefix <> key, %{count: 7}, %{ + status: :crashed, + permanent: true, + module: TestServer, + supervisor: recovery_supervisor, + node_str: "dead_node@test", + node_ref: "dead-ref", + pid: self(), + last_heartbeat_at: System.system_time(:millisecond) - 600_000 + }) + + {:ok, manager_pid} = get_supervisor_lifecycle_manager(recovery_supervisor) manager_ref = Process.monitor(manager_pid) send(manager_pid, :discover_and_restart) assert_receive {:discovery_list_attempt, 1}, 1_000 assert_receive {:discovery_list_attempt, 2}, 1_000 - refute_receive {:DOWN, ^manager_ref, :process, ^manager_pid, _reason}, 50 - assert Process.alive?(manager_pid) - GenServer.stop(manager_pid) + assert_eventually(fn -> + match?( + {pid, _} when is_pid(pid), + DurableServer.Supervisor.lookup(recovery_supervisor, key) + ) + end) + + {recovered_pid, _} = DurableServer.Supervisor.lookup(recovery_supervisor, key) + assert GenServer.call(recovered_pid, :get_count) == 7 + + assert {:ok, %{meta: %Meta{status: :running, pid: ^recovered_pid}}} = + DurableServer.fetch_stored_state(config.object_store, %{key: key, prefix: prefix}) + + refute_receive {:DOWN, ^manager_ref, :process, ^manager_pid, _reason}, 0 + Process.demonitor(manager_ref, [:flush]) end test "task supervision completes discovery cycle properly", %{ diff --git a/test/durable_server_test.exs b/test/durable_server_test.exs index 3337697..277fbf1 100644 --- a/test/durable_server_test.exs +++ b/test/durable_server_test.exs @@ -1916,34 +1916,29 @@ defmodule DurableServerTest do %{pid: pid, key: key} end - test "explicit sync with :sync action", %{pid: pid} do - # Increment with sync - assert GenServer.call(pid, :increment_and_sync) == 1 - - # Verify state persisted by checking if we can start another server with same key - # This would require modifying the test to use a known key - assert GenServer.call(pid, :get_count) == 1 - end - - test "sync with cast and :sync action", %{pid: pid} do - GenServer.cast(pid, :increment_and_sync) - # Allow cast and sync to complete - :sys.get_state(pid) - assert GenServer.call(pid, :get_count) == 1 - end + for callback <- [:call, :cast, :info] do + test "#{callback} with :sync persists state before the next request", %{ + pid: pid, + key: key, + prefix: prefix, + supervisor_name: supervisor_name + } do + case unquote(callback) do + :call -> assert GenServer.call(pid, :increment_and_sync) == 1 + :cast -> GenServer.cast(pid, :increment_and_sync) + :info -> send(pid, :increment_and_sync) + end - test "sync with info and :sync action", %{pid: pid} do - send(pid, :increment_and_sync) - # Allow info and sync to complete - :sys.get_state(pid) - assert GenServer.call(pid, :get_count) == 1 - end + # Same-sender ordering makes this a callback barrier, not a timed wait. + assert GenServer.call(pid, :get_count) == 1 - test "handles sync errors gracefully", %{pid: pid} do - # This test would require mocking ObjectStore to return errors - # For now, verify the server continues operating - assert GenServer.call(pid, :increment_and_sync) == 1 - assert GenServer.call(pid, :get_count) == 1 + assert {:ok, %StoredState{state: %{"count" => 1}, meta: %Meta{pid: ^pid}}} = + DurableServer.fetch_stored_state( + supervisor_name, + %{key: key, prefix: prefix}, + consistent: true + ) + end end end