diff --git a/lib/durable_server.ex b/lib/durable_server.ex index 34bfbc8..f4fc530 100644 --- a/lib/durable_server.ex +++ b/lib/durable_server.ex @@ -3595,6 +3595,13 @@ defmodule DurableServer do report_lock_diagnostic(sup_name, :check_lock_calls) cond do + # A partial heartbeat snapshot must never be used to justify taking a lock + # from another owner. The lifecycle manager clears this gate only after a + # complete cache reconciliation. + LifecycleManager.discovery_degraded?(sup_name) -> + report_lock_diagnostic(sup_name, :check_lock_discovery_degraded) + {:error, :discovery_degraded} + # delete tombstones are not live process locks and may not have owner fields Meta.deleting?(meta) -> report_lock_diagnostic(sup_name, :check_lock_deleting) diff --git a/lib/durable_server/lifecycle_manager.ex b/lib/durable_server/lifecycle_manager.ex index 6d48585..d6a5bf2 100644 --- a/lib/durable_server/lifecycle_manager.ex +++ b/lib/durable_server/lifecycle_manager.ex @@ -43,7 +43,7 @@ defmodule DurableServer.LifecycleManager do - `:crashed` - Server crashed or failed, always restart ### Discovery Process - 1. Write node heartbeat and refresh heartbeat cache + 1. Write node heartbeats and independently reconcile the heartbeat cache 2. List all DurableServer objects from ObjectStore (paginated with continuation tokens) 3. Check health using check_server_health/2 (uses group + heartbeat cache) 4. Apply consistent hashing to determine which servers this node should handle @@ -94,6 +94,8 @@ defmodule DurableServer.LifecycleManager do - Discovery continues if individual object reads fail (logged but skipped) - List operations retried at the GenServer level - Atomic claim failures are treated as "already claimed" rather than errors + - Partial heartbeat reconciliation retains the last complete cache snapshot and + disables takeovers and remote placement until a complete refresh succeeds """ use GenServer @@ -115,6 +117,7 @@ defmodule DurableServer.LifecycleManager do node_module: nil, current_discovery_task: nil, current_heartbeat_task: nil, + current_heartbeat_reconcile_task: nil, heartbeat_table: nil, discovery_interval_ms: nil, initial_discovery_delay_ms: nil, @@ -133,6 +136,7 @@ defmodule DurableServer.LifecycleManager do discovery_skip_table: nil, restart_gate_table: nil, discovery_stopped: false, + discovery_degraded: false, discovery_burst_remaining: 0, discovery_shuffle_batch_size: nil, parallel_restart_batch_size: nil, @@ -187,6 +191,7 @@ defmodule DurableServer.LifecycleManager do Returns a map with: - `:last_heartbeat_timing` - The timing from the last successful heartbeat (put_ms, cache_ms, total_ms) - `:last_successful_heartbeat_at` - Timestamp of last successful heartbeat + - `:discovery_degraded` - Whether an incomplete heartbeat snapshot is gating discovery - `:node` - This node's name This is used by the admin dashboard to monitor heartbeat health across the cluster. @@ -410,6 +415,7 @@ defmodule DurableServer.LifecycleManager do node_module: Keyword.get(opts, :node_module, Node), current_discovery_task: nil, current_heartbeat_task: nil, + current_heartbeat_reconcile_task: nil, heartbeat_table: hearbeat_tab, discovery_interval_ms: config.discovery_interval_ms, initial_discovery_delay_ms: config.initial_discovery_delay_ms, @@ -432,6 +438,10 @@ defmodule DurableServer.LifecycleManager do restart_claim_gate_disable_after_ms: config.restart_claim_gate_disable_after_ms } + # Fail closed while the initial heartbeat snapshot is being established. + # A complete reconciliation below clears this shared gate before init returns. + set_discovery_degraded_marker(state, true) + :ok = :pg.join( DurableServer.Supervisor.presence_pg_scope(supervisor_name), @@ -458,14 +468,17 @@ defmodule DurableServer.LifecycleManager do state = maybe_start_heartbeat_subscription(state) - if state.heartbeat_tracking_mode == :subscribe do - Process.send_after(self(), :heartbeat_reconcile, state.heartbeat_reconcile_interval_ms) - end + Process.send_after(self(), :heartbeat_reconcile, heartbeat_reconcile_delay_ms(state)) # we MUST start with a populated node heartbeat cache - # perform_heartbeat writes our heartbeat and refreshes the node health cache - {timing, heartbeat_entry, heartbeat_monotonic_at} = - perform_heartbeat(state, refresh_cache?: true) + {timing, heartbeat_entry, heartbeat_monotonic_at} = perform_heartbeat(state) + {cache_duration, cache_status} = refresh_heartbeat_cache_with_timing(state) + + timing = %{ + timing + | cache_ms: cache_duration, + total_ms: timing.total_ms + cache_duration + } # Join Group with heartbeat data so other nodes see us instantly via peer_connect. # S3 is the source of truth for liveness; Group is the fast path for discovery. @@ -479,13 +492,16 @@ defmodule DurableServer.LifecycleManager do heartbeat_hard_deadline_ms(state) ) - {:ok, - %{ - state - | last_successful_heartbeat_at: heartbeat_entry_timestamp(heartbeat_entry), - last_successful_heartbeat_monotonic_at: heartbeat_monotonic_at, - last_heartbeat_timing: timing - }} + state = + %{ + state + | last_successful_heartbeat_at: heartbeat_entry_timestamp(heartbeat_entry), + last_successful_heartbeat_monotonic_at: heartbeat_monotonic_at, + last_heartbeat_timing: timing + } + |> apply_heartbeat_cache_status(cache_status) + + {:ok, state} end @impl true @@ -508,9 +524,9 @@ defmodule DurableServer.LifecycleManager do task = Task.Supervisor.async(state.task_sup, fn -> :ok = HeartbeatWatchdog.track_heartbeat_task(heartbeat_watchdog, owner, self()) - {timing, heartbeat_entry, heartbeat_monotonic_at} = perform_heartbeat(state) - HeartbeatWatchdog.renew(heartbeat_watchdog, owner, heartbeat_monotonic_at) + {timing, heartbeat_entry, heartbeat_monotonic_at} = + perform_heartbeat(state, watchdog_owner: owner) {:heartbeat, {timing, heartbeat_entry, heartbeat_monotonic_at}} end) @@ -557,25 +573,20 @@ defmodule DurableServer.LifecycleManager do def handle_info( :heartbeat_reconcile, - %LifecycleManager{heartbeat_tracking_mode: :subscribe} = state + %LifecycleManager{current_heartbeat_reconcile_task: nil} = state ) do - Task.Supervisor.start_child(state.task_sup, fn -> - case refresh_node_heartbeat_cache(state) do - {:ok, _count, _cleaned_count, error_count} when error_count > 0 -> - log(state, :warning, fn -> - "Heartbeat reconcile observed #{error_count} heartbeat fetch error(s)" - end) - - {:ok, _count, _cleaned_count, _error_count} -> - :ok - end - end) + task = + Task.Supervisor.async_nolink(state.task_sup, fn -> + {_cache_duration, cache_status} = refresh_heartbeat_cache_with_timing(state) + {:heartbeat_reconcile, cache_status} + end) - Process.send_after(self(), :heartbeat_reconcile, state.heartbeat_reconcile_interval_ms) - {:noreply, state} + Process.send_after(self(), :heartbeat_reconcile, heartbeat_reconcile_delay_ms(state)) + {:noreply, %{state | current_heartbeat_reconcile_task: task}} end def handle_info(:heartbeat_reconcile, %LifecycleManager{} = state) do + Process.send_after(self(), :heartbeat_reconcile, heartbeat_reconcile_delay_ms(state)) {:noreply, state} end @@ -589,6 +600,11 @@ defmodule DurableServer.LifecycleManager do {:noreply, state} end + def handle_info(:discover_and_restart, %LifecycleManager{discovery_degraded: true} = state) do + Process.send_after(self(), :discover_and_restart, state.discovery_interval_ms) + {:noreply, state} + end + def handle_info(:discover_and_restart, %LifecycleManager{} = state) do with %Task{ref: ref} <- state.current_discovery_task do Process.demonitor(ref, [:flush]) @@ -647,14 +663,30 @@ defmodule DurableServer.LifecycleManager do end) end - {:noreply, - %{ - state - | current_heartbeat_task: nil, - last_successful_heartbeat_at: heartbeat_entry_timestamp(heartbeat_entry), - last_successful_heartbeat_monotonic_at: heartbeat_monotonic_at, - last_heartbeat_timing: timing - }} + state = %{ + state + | current_heartbeat_task: nil, + last_successful_heartbeat_at: heartbeat_entry_timestamp(heartbeat_entry), + last_successful_heartbeat_monotonic_at: heartbeat_monotonic_at, + last_heartbeat_timing: timing + } + + {:noreply, state} + end + + def handle_info( + {ref, {:heartbeat_reconcile, cache_status}}, + %LifecycleManager{current_heartbeat_reconcile_task: %Task{ref: ref}} = state + ) do + Process.demonitor(ref, [:flush]) + + state = + apply_heartbeat_cache_status( + %{state | current_heartbeat_reconcile_task: nil}, + cache_status + ) + + {:noreply, state} end def handle_info( @@ -686,6 +718,23 @@ defmodule DurableServer.LifecycleManager do {:stop, {:heartbeat_failed, reason}, state} end + def handle_info( + {:DOWN, ref, :process, _pid, reason}, + %LifecycleManager{current_heartbeat_reconcile_task: %Task{ref: ref}} = state + ) do + log(state, :error, fn -> + "heartbeat cache reconciliation task failed: #{inspect(reason)}" + end) + + state = + apply_heartbeat_cache_status( + %{state | current_heartbeat_reconcile_task: nil}, + {:partial, 1} + ) + + {:noreply, state} + end + @impl true def handle_call(:stop_discovery, _from, %LifecycleManager{} = state) do state = @@ -750,7 +799,8 @@ defmodule DurableServer.LifecycleManager do deadline_ms: heartbeat_hard_deadline_ms(state), resources: resources, capacity: capacity, - heartbeat_meta: heartbeat_meta + heartbeat_meta: heartbeat_meta, + discovery_degraded: state.discovery_degraded or discovery_degraded?(state.supervisor_name) } {:reply, metrics, state} @@ -765,8 +815,8 @@ defmodule DurableServer.LifecycleManager do end defp perform_heartbeat(%LifecycleManager{} = state, opts \\ []) do - opts = Keyword.validate!(opts, [:refresh_cache?]) - refresh_cache? = Keyword.get(opts, :refresh_cache?, state.heartbeat_tracking_mode == :poll) + opts = Keyword.validate!(opts, [:watchdog_owner]) + watchdog_owner = Keyword.get(opts, :watchdog_owner) start_time = System.monotonic_time(:millisecond) # Do the critical heartbeat PUT inline @@ -775,6 +825,8 @@ defmodule DurableServer.LifecycleManager do {heartbeat_entry, heartbeat_monotonic_at} = case write_node_heartbeat(state) do {:ok, {entry, monotonic_at}} -> + renew_heartbeat_watchdog(state, watchdog_owner, monotonic_at) + log(state, :debug, fn -> put_duration = System.monotonic_time(:millisecond) - put_start "Node heartbeat written successfully in #{put_duration}ms" @@ -793,49 +845,65 @@ defmodule DurableServer.LifecycleManager do put_duration = System.monotonic_time(:millisecond) - put_start - cache_duration = - if refresh_cache? do - refresh_heartbeat_cache_with_timing!(state) - else - 0 - end - :ok = CircuitBreaker.prune_stale_entries(state.circuit_breaker) total_duration = System.monotonic_time(:millisecond) - start_time { - %{put_ms: put_duration, cache_ms: cache_duration, total_ms: total_duration}, + %{put_ms: put_duration, cache_ms: 0, total_ms: total_duration}, heartbeat_entry, heartbeat_monotonic_at } end - defp refresh_heartbeat_cache_with_timing!(%LifecycleManager{} = state) do + defp renew_heartbeat_watchdog(%LifecycleManager{} = state, owner, heartbeat_monotonic_at) + when is_pid(owner) do + # Cache reconciliation is not part of our liveness contract, so renewing the + # watchdog must not wait for the cluster read to finish. + HeartbeatWatchdog.renew( + state.heartbeat_watchdog, + owner, + heartbeat_monotonic_at + ) + end + + defp renew_heartbeat_watchdog( + %LifecycleManager{}, + _owner, + _heartbeat_monotonic_at + ), + do: :ok + + defp refresh_heartbeat_cache_with_timing(%LifecycleManager{} = state) do cache_start = System.monotonic_time(:millisecond) cache_result = refresh_node_heartbeat_cache(state) cache_duration = System.monotonic_time(:millisecond) - cache_start case cache_result do - {:ok, _count, _cleaned_count, error_count} when error_count > 0 -> - # A partial view of the cluster is dangerous - we could incorrectly treat - # healthy nodes as expired and steal their locks - raise RuntimeError, - "failed to refresh heartbeat cache: #{error_count} heartbeat fetch errors" + {:partial, count, _cleaned_count, error_count} -> + # Close the shared takeover/placement gate from the reconciliation task + # before its result waits in the lifecycle manager mailbox. + set_discovery_degraded_marker(state, true) + + log(state, :warning, fn -> + "Heartbeat cache reconciliation was partial (#{error_count} error(s), #{count} node(s) read); retaining the last complete snapshot" + end) + + {cache_duration, {:partial, error_count}} {:ok, count, cleaned_count, _error_count} when cleaned_count > 0 -> log(state, :debug, fn -> "Refreshed heartbeat cache with #{count} nodes, cleaned up #{cleaned_count} dead nodes in #{cache_duration}ms" end) - cache_duration + {cache_duration, :complete} {:ok, count, _cleaned_count, _error_count} -> log(state, :debug, fn -> "Refreshed heartbeat cache with #{count} nodes in #{cache_duration}ms" end) - cache_duration + {cache_duration, :complete} end end @@ -1139,6 +1207,83 @@ defmodule DurableServer.LifecycleManager do _ -> false end + defp heartbeat_reconcile_delay_ms( + %LifecycleManager{heartbeat_tracking_mode: :subscribe} = state + ), + do: state.heartbeat_reconcile_interval_ms + + defp heartbeat_reconcile_delay_ms(%LifecycleManager{} = state), + do: state.heartbeat_interval_ms + + defp apply_heartbeat_cache_status( + %LifecycleManager{discovery_degraded: true} = state, + :complete + ) do + set_discovery_degraded_marker(state, false) + + log(state, :info, fn -> + "Heartbeat cache reconciliation recovered; resuming discovery and remote placement" + end) + + %{state | discovery_degraded: false} + end + + defp apply_heartbeat_cache_status(%LifecycleManager{} = state, :complete) do + set_discovery_degraded_marker(state, false) + state + end + + defp apply_heartbeat_cache_status( + %LifecycleManager{discovery_degraded: true} = state, + {:partial, _error_count} + ), + do: state + + defp apply_heartbeat_cache_status( + %LifecycleManager{} = state, + {:partial, error_count} + ) do + # Publish the gate before the GenServer state transition. Discovery runs in a + # separate task, and placement/lock checks can run in arbitrary callers. + set_discovery_degraded_marker(state, true) + :ets.delete_all_objects(state.restart_gate_table) + + log(state, :warning, fn -> + "Entering discovery_degraded after #{error_count} heartbeat cache reconciliation error(s); " <> + "orphan detection, lock stealing, and remote placement are disabled" + end) + + %{state | discovery_degraded: true} + end + + defp set_discovery_degraded_marker(%LifecycleManager{} = state, value) + when is_boolean(value) do + table_name = + case state.config do + %{ets_table: table_name} -> table_name + _ -> DurableServer.Supervisor.__ets_table_name__(state.supervisor_name) + end + + :ets.insert(table_name, {:discovery_degraded, value}) + :ok + end + + @doc false + def discovery_degraded?(supervisor_name) when is_atom(supervisor_name) do + table_name = DurableServer.Supervisor.__ets_table_name__(supervisor_name) + match?([{:discovery_degraded, true}], :ets.lookup(table_name, :discovery_degraded)) + rescue + # If the supervisor's coordination table is unavailable, no caller has a + # complete heartbeat view from which it can safely take over a lock. + ArgumentError -> true + end + + defp discovery_takeovers_allowed?(%LifecycleManager{discovery_degraded: true}), do: false + + defp discovery_takeovers_allowed?(%LifecycleManager{} = state) do + not discovery_degraded?(state.supervisor_name) + end + # this gets run async inside a task defp refresh_node_heartbeat_cache(%LifecycleManager{} = state) do dead_node_threshold_ms = state.config.dead_node_threshold_ms @@ -1228,8 +1373,9 @@ defmodule DurableServer.LifecycleManager do end) end - error_count = length(errors) - listing_complete? = errors == [] and :counters.get(warnings_counter, 1) == 0 + listing_error_count = :counters.get(warnings_counter, 1) + error_count = length(errors) + listing_error_count + listing_complete? = error_count == 0 # extract heartbeat data for alive nodes live_heartbeats = @@ -1238,37 +1384,46 @@ defmodule DurableServer.LifecycleManager do {node, node_ref, timestamp, capacity, resources, env_vars, heartbeat_meta} end) - # attempt to clean up dead nodes (race conditions are OK, delete might fail) - cleaned_count = - dead_nodes - |> Enum.map(fn {:dead, key, node, _node_ref, _timestamp} -> - :ets.delete(state.heartbeat_table, node) + if listing_complete? do + # Do not mutate the cache until every listed heartbeat has been read. + # This makes the existing ETS contents the last complete snapshot. + cleaned_count = cleanup_dead_node_heartbeats(state, dead_nodes) + :ets.insert(state.heartbeat_table, live_heartbeats) + prune_orphaned_heartbeat_entries(state, seen_nodes, current_time, dead_node_threshold_ms) - case StorageBackend.delete_object(state.heartbeat_store, key) do - :ok -> - log(state, :info, fn -> "Cleaned up dead node heartbeat: #{node}" end) - 1 + {:ok, length(live_heartbeats), cleaned_count, 0} + else + {:partial, length(live_heartbeats), 0, error_count} + end + end - {:error, reason} -> - # race condition or other error - another node might have cleaned it up - log(state, :info, fn -> - "Failed to clean up dead node heartbeat #{inspect(node)}: #{inspect(reason)}" - end) + defp cleanup_dead_node_heartbeats(%LifecycleManager{} = state, dead_nodes) do + dead_nodes + |> Enum.map(fn {:dead, key, node, _node_ref, _timestamp} -> + :ets.delete(state.heartbeat_table, node) - 0 + delete_result = + try do + StorageBackend.delete_object(state.heartbeat_store, key) + catch + kind, reason -> {:error, {kind, reason}} end - end) - |> Enum.sum() - :ets.insert(state.heartbeat_table, live_heartbeats) + case delete_result do + :ok -> + log(state, :info, fn -> "Cleaned up dead node heartbeat: #{node}" end) + 1 - # The listing must be complete, otherwise nodes missing from - # `live_heartbeats` would look prunable - if listing_complete? do - prune_orphaned_heartbeat_entries(state, seen_nodes, current_time, dead_node_threshold_ms) - end + {:error, reason} -> + # race condition or other error - another node might have cleaned it up + log(state, :info, fn -> + "Failed to clean up dead node heartbeat #{inspect(node)}: #{inspect(reason)}" + end) - {:ok, length(live_heartbeats), cleaned_count, error_count} + 0 + end + end) + |> Enum.sum() end defp process_heartbeat_list_entry( @@ -2028,6 +2183,18 @@ defmodule DurableServer.LifecycleManager do end defp discover_and_restart_servers(%LifecycleManager{} = state) do + if discovery_takeovers_allowed?(state) do + do_discover_and_restart_servers(state) + else + log(state, :debug, fn -> + "Skipping discovery while heartbeat cache reconciliation is degraded" + end) + + :ok + end + end + + defp do_discover_and_restart_servers(%LifecycleManager{} = state) do diagnostics_before = discovery_diag_snapshot(state) restart_gate_config = restart_claim_gate_config(state) skip_count = :ets.info(state.discovery_skip_table, :size) @@ -2403,19 +2570,23 @@ defmodule DurableServer.LifecycleManager do end defp appears_restartable?(%LifecycleManager{} = state, %Meta{} = meta) do - case check_server_health(state, meta) do - :healthy -> - false + if discovery_takeovers_allowed?(state) do + case check_server_health(state, meta) do + :healthy -> + false - {:orphaned, claim_context} -> - # server is orphaned, check if this node should handle it - if orphan_claimable?(meta), do: {:restartable, claim_context}, else: false + {:orphaned, claim_context} -> + # server is orphaned, check if this node should handle it + if orphan_claimable?(meta), do: {:restartable, claim_context}, else: false - :orphaned -> - if orphan_claimable?(meta), do: {:restartable, :check_lock}, else: false + :orphaned -> + if orphan_claimable?(meta), do: {:restartable, :check_lock}, else: false - :transient -> - :transient + :transient -> + :transient + end + else + :transient end end @@ -2631,6 +2802,20 @@ defmodule DurableServer.LifecycleManager do %StoredState{meta: %Meta{} = meta} = stored_state, claim_context ) do + if discovery_takeovers_allowed?(state) do + do_attempt_restart(state, meta, stored_state, claim_context) + else + report_diagnostic(state.supervisor_name, :restart_claim_discovery_degraded) + :ok + end + end + + defp do_attempt_restart( + %LifecycleManager{} = state, + %Meta{} = meta, + %StoredState{} = stored_state, + claim_context + ) do claim_opts = [ttl: restart_claim_ttl_ms(state)] claim_result = @@ -2818,6 +3003,14 @@ defmodule DurableServer.LifecycleManager do end defp fetch_orphaned_slow_path(%Meta{} = meta) do + if discovery_degraded?(meta.supervisor) do + :transient + else + do_fetch_orphaned_slow_path(meta) + end + end + + defp do_fetch_orphaned_slow_path(%Meta{} = meta) do report_diagnostic(meta.supervisor, :slow_path_lock_checks) node_health = lookup_node_health(meta) lock_result = DurableServer.check_lock_status(meta) @@ -2854,6 +3047,7 @@ defmodule DurableServer.LifecycleManager do defp transient_lock_check_error?({:error, reason}), do: conflict_consistent_read_error?(reason) defp transient_lock_check_error?(_), do: false + defp conflict_consistent_read_error?(:discovery_degraded), do: true defp conflict_consistent_read_error?(:conflict), do: true defp conflict_consistent_read_error?(":conflict"), do: true @@ -3341,6 +3535,14 @@ defmodule DurableServer.LifecycleManager do """ def find_eligible_nodes(supervisor_name, module, opts \\ []) when is_atom(supervisor_name) do + if discovery_degraded?(supervisor_name) do + [] + else + find_eligible_nodes_from_cache(supervisor_name, module, opts) + end + end + + defp find_eligible_nodes_from_cache(supervisor_name, module, opts) do heartbeat_table = heartbeat_table_name(supervisor_name) # Handle case where table doesn't exist yet during startup case :ets.whereis(heartbeat_table) do diff --git a/lib/durable_server/supervisor.ex b/lib/durable_server/supervisor.ex index bd2b425..6f4c31a 100644 --- a/lib/durable_server/supervisor.ex +++ b/lib/durable_server/supervisor.ex @@ -1837,13 +1837,19 @@ defmodule DurableServer.Supervisor do end end - defp try_nodes(supervisor, child_spec, nodes, placement_opts \\ []) + defp try_nodes(supervisor, child_spec, nodes, placement_opts \\ []) do + if LifecycleManager.discovery_degraded?(supervisor) do + {:error, {:capacity_limit, :discovery_degraded}} + else + do_try_nodes(supervisor, child_spec, nodes, placement_opts) + end + end - defp try_nodes(_supervisor, _child_spec, [], _placement_opts) do + defp do_try_nodes(_supervisor, _child_spec, [], _placement_opts) do {:error, {:capacity_limit, :all_placement_attempts_failed}} end - defp try_nodes( + defp do_try_nodes( supervisor, {module, _init_arg, _boot_info} = child_spec, [node | rest], diff --git a/test/durable_server/lifecycle_test.exs b/test/durable_server/lifecycle_test.exs index 18686e8..091b07f 100644 --- a/test/durable_server/lifecycle_test.exs +++ b/test/durable_server/lifecycle_test.exs @@ -250,29 +250,56 @@ defmodule DurableServer.LifecycleTest do 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) + def get_object(%{delegate: delegate} = state, key, opts) do + if current_mode(state) == :fetch_error and key == state.inject_key do + {:error, :simulated_fetch_failure} + else + StorageBackend.get_object(delegate, key, opts) + end + end @impl true def list_all_objects_stream(%{delegate: delegate} = state, prefix, opts) do stream = StorageBackend.list_all_objects_stream(delegate, prefix, opts) if prefix == state.heartbeat_prefix do - inject_list_fault(stream, state, opts) + inject_list_fault(stream, current_mode(state), state, opts) else stream end end - defp inject_list_fault(stream, %{mode: :truncate}, opts) do + defp inject_list_fault(stream, :truncate, _state, opts) do error_handler = Keyword.get(opts, :error_handler, fn _reason -> :continue end) Stream.concat(stream, reported_error_stream(error_handler, :simulated_page_failure)) end - defp inject_list_fault(stream, %{mode: :inject_key} = state, _opts) do + defp inject_list_fault(stream, :inject_key, state, _opts) do Stream.concat(stream, [%{key: state.inject_key}]) end + defp inject_list_fault(stream, :fetch_error, state, _opts) do + Stream.concat(stream, [%{key: state.inject_key}]) + end + + defp inject_list_fault(stream, {:block, notify_pid}, _state, _opts) do + send(notify_pid, {:heartbeat_cache_refresh_blocked, self()}) + + receive do + :continue_heartbeat_cache_refresh -> stream + after + 5_000 -> raise "timed out waiting to continue heartbeat cache refresh" + end + end + + defp inject_list_fault(stream, :ok, _state, _opts), do: stream + + defp current_mode(%{mode: {:controlled, table}}) do + :ets.lookup_element(table, :mode, 2) + end + + defp current_mode(%{mode: mode}), do: mode + defp reported_error_stream(error_handler, reason) do Stream.flat_map([reason], fn reason -> _ = error_handler.({:list_failed, reason}) @@ -1945,6 +1972,18 @@ defmodule DurableServer.LifecycleTest do end) end + defp await_heartbeat_reconciliation(manager_pid) do + send(manager_pid, :heartbeat_reconcile) + + # This call reaches the manager after the reconciliation message, proving + # that the task was either started or completed. + _ = :sys.get_state(manager_pid) + + assert_eventually(fn -> + :sys.get_state(manager_pid).current_heartbeat_reconcile_task == nil + end) + end + defp seed_peer_heartbeat(config, prefix, node, current_time) do ObjectStore.put_object( config.object_store, @@ -2844,6 +2883,85 @@ defmodule DurableServer.LifecycleTest do GenServer.stop(manager_pid) end + test "continues heartbeat PUTs while cache reconciliation is blocked", %{ + supervisor_name: supervisor_name, + prefix: prefix, + config: config + } do + control = :ets.new(__MODULE__.HeartbeatListBackend, [:set, :public]) + :ets.insert(control, {:mode, :ok}) + + heartbeat_store = + heartbeat_list_store(config, prefix, + mode: {:controlled, control}, + inject_key: "#{prefix}__nodes/injected@test" + ) + + {:ok, manager_pid} = + start_standalone_lifecycle_manager(supervisor_name, config, + heartbeat_store: heartbeat_store + ) + + manager_state = :sys.get_state(manager_pid) + watchdog = manager_state.heartbeat_watchdog + previous_watchdog_heartbeat = :sys.get_state(watchdog).last_heartbeat_at + + :ets.insert(control, {:mode, {:block, self()}}) + send(manager_pid, :heartbeat_reconcile) + + assert_receive {:heartbeat_cache_refresh_blocked, reconcile_task_pid}, 1_000 + reconcile_task = :sys.get_state(manager_pid).current_heartbeat_reconcile_task + assert reconcile_task.pid == reconcile_task_pid + + Process.sleep(2) + await_heartbeat_cycle(manager_pid) + + first_heartbeat = :sys.get_state(manager_pid).last_successful_heartbeat_monotonic_at + assert :sys.get_state(watchdog).last_heartbeat_at > previous_watchdog_heartbeat + assert :sys.get_state(manager_pid).current_heartbeat_reconcile_task == reconcile_task + + Process.sleep(2) + await_heartbeat_cycle(manager_pid) + + assert :sys.get_state(manager_pid).last_successful_heartbeat_monotonic_at > + first_heartbeat + + assert :sys.get_state(manager_pid).current_heartbeat_reconcile_task == reconcile_task + + send(reconcile_task_pid, :continue_heartbeat_cache_refresh) + + assert_eventually(fn -> + :sys.get_state(manager_pid).current_heartbeat_reconcile_task == nil + end) + + GenServer.stop(manager_pid) + end + + test "starts in degraded mode when the initial heartbeat snapshot is partial", %{ + supervisor_name: supervisor_name, + prefix: prefix, + config: config + } do + heartbeat_store = + heartbeat_list_store(config, prefix, + mode: :fetch_error, + inject_key: "#{prefix}__nodes/failed-peer@test" + ) + + assert {:ok, manager_pid} = + start_standalone_lifecycle_manager(supervisor_name, config, + heartbeat_store: heartbeat_store + ) + + manager_state = :sys.get_state(manager_pid) + assert manager_state.discovery_degraded + assert is_integer(manager_state.last_successful_heartbeat_monotonic_at) + assert LifecycleManager.discovery_degraded?(manager_state.supervisor_name) + assert Process.alive?(manager_pid) + + GenServer.stop(manager_pid) + end + test "retries retryable heartbeat write failures until success within deadline", %{ supervisor_name: supervisor_name, config: config @@ -3043,11 +3161,7 @@ defmodule DurableServer.LifecycleTest do {:ok, manager_pid} = start_standalone_lifecycle_manager(supervisor_name, test_config) - # Trigger heartbeat cycle which should clean up dead nodes - send(manager_pid, :heartbeat) - - # Give some time for the heartbeat task to complete - Process.sleep(200) + await_heartbeat_reconciliation(manager_pid) # Verify alive node still exists {:ok, _} = ObjectStore.get_object(config.object_store, "#{prefix}__nodes/#{alive_node}") @@ -3077,7 +3191,7 @@ defmodule DurableServer.LifecycleTest do heartbeat_table = :sys.get_state(manager_pid).heartbeat_table insert_stale_heartbeat(heartbeat_table, orphan_node, current_time) - await_heartbeat_cycle(manager_pid) + await_heartbeat_reconciliation(manager_pid) assert :ets.lookup(heartbeat_table, orphan_node) == [] assert :ets.lookup(heartbeat_table, peer_node) != [] @@ -3085,7 +3199,7 @@ defmodule DurableServer.LifecycleTest do GenServer.stop(manager_pid) end - test "does not prune cached heartbeat entries when the listing was truncated by an error", %{ + test "retains the last complete heartbeat snapshot when a listing is truncated", %{ supervisor_name: supervisor_name, prefix: prefix, config: config @@ -3095,19 +3209,110 @@ defmodule DurableServer.LifecycleTest do orphan_node = "orphan-truncated@test" seed_peer_heartbeat(config, prefix, peer_node, current_time) + control = :ets.new(__MODULE__.HeartbeatListBackend, [:set, :public]) + :ets.insert(control, {:mode, :ok}) {:ok, manager_pid} = start_standalone_lifecycle_manager(supervisor_name, config, - heartbeat_store: heartbeat_list_store(config, prefix, mode: :truncate) + heartbeat_store: + heartbeat_list_store(config, prefix, + mode: {:controlled, control}, + inject_key: "#{prefix}__nodes/injected@test" + ) ) heartbeat_table = :sys.get_state(manager_pid).heartbeat_table + assert :ets.lookup(heartbeat_table, peer_node) != [] insert_stale_heartbeat(heartbeat_table, orphan_node, current_time) - await_heartbeat_cycle(manager_pid) + :ets.insert(control, {:mode, :truncate}) + await_heartbeat_reconciliation(manager_pid) assert :ets.lookup(heartbeat_table, peer_node) != [] assert :ets.lookup(heartbeat_table, orphan_node) != [] + assert :sys.get_state(manager_pid).discovery_degraded + assert LifecycleManager.discovery_degraded?(:sys.get_state(manager_pid).supervisor_name) + + GenServer.stop(manager_pid) + end + + test "a heartbeat GET failure degrades discovery without changing the snapshot, then recovers", + %{ + supervisor_name: supervisor_name, + prefix: prefix, + config: config + } do + current_time = System.system_time(:millisecond) + known_node = "known-peer@test" + new_node = "new-peer@test" + failed_node = "failed-peer@test" + + seed_peer_heartbeat(config, prefix, known_node, current_time) + + control = :ets.new(__MODULE__.HeartbeatListBackend, [:set, :public]) + :ets.insert(control, {:mode, :ok}) + + heartbeat_store = + heartbeat_list_store(config, prefix, + mode: {:controlled, control}, + inject_key: "#{prefix}__nodes/#{failed_node}" + ) + + {:ok, manager_pid} = + start_standalone_lifecycle_manager(supervisor_name, config, + heartbeat_store: heartbeat_store + ) + + initial_state = :sys.get_state(manager_pid) + heartbeat_table = initial_state.heartbeat_table + manager_supervisor = initial_state.supervisor_name + known_snapshot = :ets.lookup(heartbeat_table, known_node) + assert known_snapshot != [] + + seed_peer_heartbeat(config, prefix, new_node, System.system_time(:millisecond)) + :ets.insert(control, {:mode, :fetch_error}) + + await_heartbeat_reconciliation(manager_pid) + + degraded_state = :sys.get_state(manager_pid) + assert Process.alive?(manager_pid) + assert degraded_state.discovery_degraded + assert LifecycleManager.discovery_degraded?(manager_supervisor) + assert :ets.lookup(heartbeat_table, known_node) == known_snapshot + assert :ets.lookup(heartbeat_table, new_node) == [] + + # All takeover and outbound placement decisions fail closed while the + # previous complete snapshot remains available for diagnostics. + assert {:error, :discovery_degraded} = + DurableServer.check_lock_status(%Meta{supervisor: manager_supervisor}) + + assert LifecycleManager.find_eligible_nodes(manager_supervisor, TestServer) == [] + + send(manager_pid, :discover_and_restart) + Process.sleep(25) + assert :sys.get_state(manager_pid).current_discovery_task == nil + + # Degraded discovery does not stop this node from renewing its own + # heartbeat or retrying the failed reconciliation. + Process.sleep(2) + await_heartbeat_cycle(manager_pid) + await_heartbeat_reconciliation(manager_pid) + still_degraded_state = :sys.get_state(manager_pid) + assert still_degraded_state.discovery_degraded + + assert still_degraded_state.last_successful_heartbeat_monotonic_at > + degraded_state.last_successful_heartbeat_monotonic_at + + assert :ets.lookup(heartbeat_table, known_node) == known_snapshot + assert :ets.lookup(heartbeat_table, new_node) == [] + + :ets.insert(control, {:mode, :ok}) + await_heartbeat_reconciliation(manager_pid) + + recovered_state = :sys.get_state(manager_pid) + refute recovered_state.discovery_degraded + refute LifecycleManager.discovery_degraded?(manager_supervisor) + assert :ets.lookup(heartbeat_table, new_node) != [] GenServer.stop(manager_pid) end @@ -3137,7 +3342,7 @@ defmodule DurableServer.LifecycleTest do heartbeat_table = :sys.get_state(manager_pid).heartbeat_table insert_stale_heartbeat(heartbeat_table, orphan_node, current_time) - await_heartbeat_cycle(manager_pid) + await_heartbeat_reconciliation(manager_pid) assert :ets.lookup(heartbeat_table, peer_node) != [] assert :ets.lookup(heartbeat_table, orphan_node) != [] diff --git a/test/durable_server/remote_placement_test.exs b/test/durable_server/remote_placement_test.exs index ff8660b..df69095 100644 --- a/test/durable_server/remote_placement_test.exs +++ b/test/durable_server/remote_placement_test.exs @@ -127,6 +127,42 @@ defmodule DurableServer.RemotePlacementTest do # Empty since no remote nodes available in test assert length(nodes) <= 1 end + + test "returns no candidates while discovery is degraded", %{ + supervisor_name: supervisor_name, + prefix: prefix + } do + start_supervised!( + {DurableServer.Supervisor, + name: supervisor_name, + prefix: prefix, + object_store: test_object_store_opts(), + max_children: %{:total => 100}} + ) + + config = DurableServer.Supervisor.__get_config__(supervisor_name) + heartbeat_table = :"durable_server_heartbeats_#{supervisor_name}" + remote_node = :"remote_#{:erlang.unique_integer([:positive])}@host" + + :ets.insert( + heartbeat_table, + {to_string(remote_node), 1, System.system_time(:millisecond), nil, nil, %{}, nil} + ) + + assert [remote_node] == + DurableServer.LifecycleManager.find_eligible_nodes( + supervisor_name, + RemotePlacementTestServer + ) + + :ets.insert(config.ets_table, {:discovery_degraded, true}) + + assert [] == + DurableServer.LifecycleManager.find_eligible_nodes( + supervisor_name, + RemotePlacementTestServer + ) + end end describe "start_child with max_placement_retries" do