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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
71 changes: 44 additions & 27 deletions lib/durable_server/lifecycle_manager.ex
Original file line number Diff line number Diff line change
Expand Up @@ -2116,7 +2116,7 @@ defmodule DurableServer.LifecycleManager do
:skip ->
# permanently non-restartable — cache without meta
now = System.monotonic_time(:millisecond)
:ets.insert(state.discovery_skip_table, {key, etag, :skip, now})
skip_cache_put(state, {key, etag, :skip, now})
clear_restart_gate_state(state, key)
:noop

Expand All @@ -2125,7 +2125,7 @@ defmodule DurableServer.LifecycleManager do
# health check) — cache trimmed meta for re-evaluation next round
now = System.monotonic_time(:millisecond)
trimmed = trim_meta_for_cache(meta)
:ets.insert(state.discovery_skip_table, {key, etag, trimmed, now})
skip_cache_put(state, {key, etag, trimmed, now})
clear_restart_gate_state(state, key)
:noop

Expand Down Expand Up @@ -2165,6 +2165,15 @@ defmodule DurableServer.LifecycleManager do
:ok
end

defp skip_cache_put(%LifecycleManager{} = state, entry) do
case :ets.whereis(state.discovery_skip_table) do
:undefined -> :ok
_ -> :ets.insert(state.discovery_skip_table, entry)
end
rescue
ArgumentError -> :ok
end

defp discovery_diag_snapshot(%LifecycleManager{} = state) do
case :ets.whereis(state.discovery_diag_table) do
:undefined -> %{}
Expand Down Expand Up @@ -2562,38 +2571,46 @@ defmodule DurableServer.LifecycleManager do
end

defp discovery_skip?(%LifecycleManager{} = state, key, etag) do
case :ets.lookup(state.discovery_skip_table, key) do
[{^key, ^etag, :skip, _ts}] ->
true
case :ets.whereis(state.discovery_skip_table) do
:undefined ->
false

[{^key, ^etag, %Meta{} = meta, _ts}] ->
# Etag unchanged so stored state is identical, but time-dependent checks
# (placement gates, circuit breaker, node health) may have changed.
case CircuitBreaker.check_module_circuit_breaker(state.circuit_breaker, meta.module) do
{:circuit_open, _} ->
_ ->
case :ets.lookup(state.discovery_skip_table, key) do
[{^key, ^etag, :skip, _ts}] ->
true

:ok ->
case appears_restartable?(state, meta) do
{:restartable, _claim_context} ->
# Now restartable — remove from cache, let async task do fresh GET
:ets.delete(state.discovery_skip_table, key)
false

:transient ->
# Heartbeat/read uncertainty is not a stable "non-restartable" state.
# Drop the cache entry so the next async pass does a fresh read.
:ets.delete(state.discovery_skip_table, key)
false

false ->
[{^key, ^etag, %Meta{} = meta, _ts}] ->
# Etag unchanged so stored state is identical, but time-dependent checks
# (placement gates, circuit breaker, node health) may have changed.
case CircuitBreaker.check_module_circuit_breaker(state.circuit_breaker, meta.module) do
{:circuit_open, _} ->
true

:ok ->
case appears_restartable?(state, meta) do
{:restartable, _claim_context} ->
# Now restartable — remove from cache, let async task do fresh GET
:ets.delete(state.discovery_skip_table, key)
false

:transient ->
# Heartbeat/read uncertainty is not a stable "non-restartable" state.
# Drop the cache entry so the next async pass does a fresh read.
:ets.delete(state.discovery_skip_table, key)
false

false ->
true
end
end
end

_ ->
false
_ ->
false
end
end
rescue
ArgumentError -> false
end

# Strip fields not needed for re-evaluation to reduce ETS memory.
Expand Down
222 changes: 222 additions & 0 deletions test/durable_server/lifecycle_test.exs
Original file line number Diff line number Diff line change
Expand Up @@ -3301,4 +3301,226 @@ defmodule DurableServer.LifecycleTest do
GenServer.stop(manager_pid)
end
end

defmodule PauseOnFetchBackend do
@behaviour DurableServer.StorageBackend

alias DurableServer.StorageBackend

@impl true
def init_backend(opts) do
delegate = Keyword.fetch!(opts, :delegate)

{:ok,
%{
state: %{
delegate: delegate,
notify_pid: Keyword.fetch!(opts, :notify_pid),
pause_key: Keyword.fetch!(opts, :pause_key),
pause_list: Keyword.get(opts, :pause_list, false)
},
defaults: StorageBackend.defaults(delegate),
features: StorageBackend.features(delegate)
}}
end

@impl true
def ensure_ready(%{delegate: delegate}), do: StorageBackend.ensure_ready(delegate)

@impl true
def get_object(%{delegate: delegate, notify_pid: notify_pid, pause_key: pause_key}, key, opts) do
if String.ends_with?(key, pause_key) do
send(notify_pid, {:paused_in_discovery_worker, self()})

receive do
:continue -> :ok
after
10_000 -> :ok
end
end

StorageBackend.get_object(delegate, key, opts)
end

@impl true
def list_all_objects_stream(%{delegate: delegate} = state, prefix, opts) do
# Heartbeat refresh also lists "#{prefix}__nodes/" during init. Only the
# discovery listing should block, otherwise the manager never finishes starting.
if state.pause_list and not String.contains?(prefix, "__nodes/") do
send(state.notify_pid, {:paused_in_discovery_list, self()})

receive do
:continue_list -> :ok
after
10_000 -> :ok
end
end

StorageBackend.list_all_objects_stream(delegate, prefix, opts)
end

@impl true
def put_object(%{delegate: delegate}, key, data, opts),
do: StorageBackend.put_object(delegate, key, data, opts)

@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

describe "discovery shutdown race" do
test "orphaned discovery worker does not crash when the skip table dies mid-fetch",
%{supervisor_name: supervisor_name, prefix: prefix, config: config} do
key = "skip-race-#{DurableServer.UUID.uuid4()}"

{manager_pid, worker_pid, discovery_pid, skip_table} =
pause_discovery_worker(supervisor_name, prefix, config, key)

worker_ref = Process.monitor(worker_pid)

# the skip table is owned by the manager, so it dies with the manager
assert :ets.info(skip_table, :owner) == manager_pid

# async_nolink means the discovery task is not linked to the manager
{:links, discovery_links} = Process.info(discovery_pid, :links)
refute manager_pid in discovery_links
assert Process.alive?(discovery_pid)

Process.exit(manager_pid, :kill)
assert_eventually(fn -> :ets.whereis(skip_table) == :undefined end)

# the worker outlives the manager and then writes the skip cache
assert Process.alive?(worker_pid)

send(worker_pid, :continue)

# without skip_cache_put/2 this exits with ArgumentError from :ets.insert/2
assert_receive {:DOWN, ^worker_ref, :process, ^worker_pid, :normal}, 10_000
assert :ets.whereis(skip_table) == :undefined
end

test "orphaned discovery task does not crash when discovery_skip? reads a dead skip table",
%{supervisor_name: supervisor_name, prefix: prefix, config: config} do
key = "skip-read-race-#{DurableServer.UUID.uuid4()}"
delegate = DurableServer.Supervisor.__get_config__(supervisor_name).storage_backend
put_non_restartable_object(delegate, prefix, key, supervisor_name)

{:ok, storage_backend} =
DurableServer.StorageBackend.init_backend(PauseOnFetchBackend,
delegate: delegate,
notify_pid: self(),
pause_key: key,
pause_list: true
)

{:ok, manager_pid} =
start_standalone_lifecycle_manager(
supervisor_name,
discovery_race_config(config, storage_backend)
)

send(manager_pid, :discover_and_restart)

assert_receive {:paused_in_discovery_list, discovery_pid}, 10_000
discovery_ref = Process.monitor(discovery_pid)

manager_state = :sys.get_state(manager_pid)
skip_table = manager_state.discovery_skip_table
assert manager_state.current_discovery_task.pid == discovery_pid
assert :ets.info(skip_table, :owner) == manager_pid

# async_nolink means the discovery task is not linked to the manager
{:links, discovery_links} = Process.info(discovery_pid, :links)
refute manager_pid in discovery_links

Process.exit(manager_pid, :kill)
assert_eventually(fn -> :ets.whereis(skip_table) == :undefined end)
assert Process.alive?(discovery_pid)

# listing resumes only after the table is gone, so discovery_skip?/3 runs next
send(discovery_pid, :continue_list)

# without the :ets.whereis/1 guard this task dies with ArgumentError from :ets.lookup/2
# and the worker pause never arrives
assert_receive {:paused_in_discovery_worker, worker_pid}, 10_000
worker_ref = Process.monitor(worker_pid)
assert Process.alive?(discovery_pid)

send(worker_pid, :continue)

assert_receive {:DOWN, ^worker_ref, :process, ^worker_pid, :normal}, 10_000
assert_receive {:DOWN, ^discovery_ref, :process, ^discovery_pid, :normal}, 10_000
assert :ets.whereis(skip_table) == :undefined
end
end

defp pause_discovery_worker(supervisor_name, prefix, config, key) do
delegate = DurableServer.Supervisor.__get_config__(supervisor_name).storage_backend
put_non_restartable_object(delegate, prefix, key, supervisor_name)

{:ok, storage_backend} =
DurableServer.StorageBackend.init_backend(PauseOnFetchBackend,
delegate: delegate,
notify_pid: self(),
pause_key: key
)

{:ok, manager_pid} =
start_standalone_lifecycle_manager(
supervisor_name,
discovery_race_config(config, storage_backend)
)

send(manager_pid, :discover_and_restart)

assert_receive {:paused_in_discovery_worker, worker_pid}, 10_000

manager_state = :sys.get_state(manager_pid)

{manager_pid, worker_pid, manager_state.current_discovery_task.pid,
manager_state.discovery_skip_table}
end

defp discovery_race_config(config, storage_backend) do
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, 60_000)
end

defp put_non_restartable_object(delegate, prefix, key, supervisor_name) do
stored_state = %DurableServer.StoredState{
vsn: 1,
state: %{"count" => 1},
meta: %Meta{
key: key,
prefix: prefix,
supervisor: supervisor_name,
module: TestServer,
permanent: false,
status: :running,
node_str: "unreachable@test",
node_ref: System.unique_integer([:positive]),
pid: self(),
last_heartbeat_at: System.system_time(:millisecond) - 60_000
}
}

assert {:ok, %{etag: _}} =
DurableServer.StorageBackend.put_object(delegate, "#{prefix}#{key}", stored_state)
end
end