From bd89f69c99e4cfb2965affa508032923ea7d4b67 Mon Sep 17 00:00:00 2001 From: Jason Stiebs Date: Thu, 3 Sep 2026 09:49:06 -0500 Subject: [PATCH 1/2] Improve heartbeat failure observability --- CHANGELOG.md | 3 + README.md | 26 + lib/durable_server/backends/ekv_store.ex | 3 +- lib/durable_server/heartbeat_metrics.ex | 536 ++++++++++++++++++ lib/durable_server/heartbeat_watchdog.ex | 100 +++- lib/durable_server/lifecycle_manager.ex | 198 ++++++- lib/durable_server/object_store.ex | 68 ++- mix.exs | 1 + .../durable_server/heartbeat_metrics_test.exs | 137 +++++ .../heartbeat_watchdog_test.exs | 92 +++ test/durable_server/lifecycle_test.exs | 55 +- .../object_store_retry_test.exs | 64 +++ 12 files changed, 1227 insertions(+), 56 deletions(-) create mode 100644 lib/durable_server/heartbeat_metrics.ex create mode 100644 test/durable_server/heartbeat_metrics_test.exs create mode 100644 test/durable_server/heartbeat_watchdog_test.exs diff --git a/CHANGELOG.md b/CHANGELOG.md index d9a3f1a..780b845 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,3 +1,6 @@ +## Unreleased +- Add bounded heartbeat telemetry and live metrics for attempt outcomes, consecutive failures, last-success age, watchdog budget, cache-degraded duration, watchdog terminations, and fenced child counts. Heartbeat writes now retry every ambiguous `Req.TransportError` variant within the hard deadline while permanent HTTP authentication/configuration failures still fail immediately. + ## 0.1.5 (2026-08-27) - Treat explicit `:sync`, `{:sync, metadata}`, and `sync: true` callback returns as strict durability boundaries. Built-in backends first exhaust their bounded transient retry policy; if the write still fails, the DurableServer exits with a structured `{:sync_failed, reason}` fatal-exit reason before acknowledging the callback. Automatic and periodic sync remain best effort for transient failures, while storage conflicts remain fatal. - Honor the caller-supplied `ensure_started_child/3` timeout while waiting for a live storage owner to finish Group registration, and preserve the caller's remaining overall deadline when sticky placement falls back to a local start instead of applying fresh fixed 5-second waits. diff --git a/README.md b/README.md index 789c27b..4eb060e 100644 --- a/README.md +++ b/README.md @@ -99,6 +99,32 @@ passes it through `dump_state/1`, the configured backend's encode/decode path, and then `load_state/2` before `init/1` or `init/2`. The dumped initial state must therefore be encodable by your configured backend. +## Heartbeat Observability + +`DurableServer.LifecycleManager.get_heartbeat_metrics/1` returns a node-local +snapshot that includes: + +- attempt totals grouped by bounded result, HTTP status, and transport class +- consecutive failed attempts +- monotonic age of the last successful heartbeat write +- remaining watchdog budget +- whether the heartbeat cache is degraded and how long it has been degraded + +DurableServer also emits these telemetry events: + +| Event | Measurements | Bounded metadata | +|---|---|---| +| `[:durable_server, :heartbeat, :attempt]` | `count`, `total_attempts`, `consecutive_failures`, `recovered_after_failures`, `last_success_age_ms`, `remaining_watchdog_budget_ms`, `cache_degraded_duration_ms` | `supervisor`, `result`, `http_status`, `transport_class`, `error_class`, `retryable`, `cache_degraded`, `has_last_success` | +| `[:durable_server, :heartbeat, :cache]` | `count`, `error_count`, `refresh_duration_ms`, `degraded_duration_ms` | `supervisor`, `status`, `transition` | +| `[:durable_server, :heartbeat, :watchdog, :termination]` | `count`, `watchdog_terminations`, `children_fenced`, `consecutive_failures`, `last_success_age_ms`, `remaining_watchdog_budget_ms`, `cache_degraded_duration_ms` | `supervisor`, `child_count_status`, `cache_degraded` | + +Attempt metadata never includes a raw error, request URL, or object key. +`http_status` is limited to valid HTTP statuses plus `:none`/`:other`, and +`transport_class` uses a fixed set of categories. Aggregate the telemetry +`count` and `children_fenced` measurements outside the DurableServer +supervision tree; unlike the telemetry stream, the node-local snapshot resets +when its lifecycle manager restarts. + ## Administrative Cordon Use `terminate_and_cordon_child/3` when you need to stop a DurableServer and diff --git a/lib/durable_server/backends/ekv_store.ex b/lib/durable_server/backends/ekv_store.ex index beae462..d4cabd1 100644 --- a/lib/durable_server/backends/ekv_store.ex +++ b/lib/durable_server/backends/ekv_store.ex @@ -659,7 +659,8 @@ defmodule DurableServer.Backends.EKVStore do :max_results, :continuation_token, :prefix, - :etag + :etag, + :retry_observer ]) end diff --git a/lib/durable_server/heartbeat_metrics.ex b/lib/durable_server/heartbeat_metrics.ex new file mode 100644 index 0000000..e2d80b6 --- /dev/null +++ b/lib/durable_server/heartbeat_metrics.ex @@ -0,0 +1,536 @@ +defmodule DurableServer.HeartbeatMetrics do + @moduledoc """ + Heartbeat observability with bounded telemetry dimensions. + + The lifecycle manager emits the following events: + + * `[:durable_server, :heartbeat, :attempt]` for every backend/HTTP attempt + * `[:durable_server, :heartbeat, :cache]` after every heartbeat-cache refresh + * `[:durable_server, :heartbeat, :watchdog, :termination]` when the watchdog + fences a supervisor tree + + Attempt metadata deliberately contains classifications rather than raw errors, + URLs, or keys. This keeps metric series bounded and avoids leaking request data. + """ + + require Logger + + @attempt_event [:durable_server, :heartbeat, :attempt] + @cache_event [:durable_server, :heartbeat, :cache] + @watchdog_event [:durable_server, :heartbeat, :watchdog, :termination] + + @enforce_keys [:supervisor_name, :table, :started_at] + defstruct [:supervisor_name, :table, :started_at] + + @type t :: %__MODULE__{ + supervisor_name: atom(), + table: :ets.tid(), + started_at: integer() + } + + @type attempt_classification :: %{ + result: :success | :http_error | :transport_error | :protocol_error | :backend_error, + http_status: 100..599 | :none | :other, + transport_class: + :none + | :timeout + | :closed + | :connection_refused + | :dns + | :tls + | :network + | :socket + | :other, + error_class: + :none + | :http + | :authentication + | :configuration + | :transport + | :protocol + | :conflict + | :unavailable + | :timeout + | :other + } + + @doc """ + Returns the telemetry event names emitted for heartbeat operations. + """ + def events do + %{ + attempt: @attempt_event, + cache: @cache_event, + watchdog_termination: @watchdog_event + } + end + + @doc false + @spec new(atom()) :: t() + def new(supervisor_name) when is_atom(supervisor_name) do + table = + :ets.new(__MODULE__, [ + :set, + :public, + read_concurrency: true, + write_concurrency: true + ]) + + :ets.insert(table, [ + {:attempt_count, 0}, + {:consecutive_failures, 0} + ]) + + %__MODULE__{ + supervisor_name: supervisor_name, + table: table, + started_at: System.monotonic_time(:millisecond) + } + end + + @doc false + @spec record_attempt(t(), term(), boolean(), integer()) :: attempt_classification() | :ignored + def record_attempt( + %__MODULE__{} = metrics, + response_or_error, + retryable?, + deadline_at + ) + when is_boolean(retryable?) and is_integer(deadline_at) do + classification = classify_attempt(response_or_error) + total_attempts = increment(metrics, :attempt_count) + increment(metrics, {:attempt_result, classification.result}) + + if classification.http_status != :none do + increment(metrics, {:attempt_http_status, classification.http_status}) + end + + if classification.transport_class != :none do + increment(metrics, {:attempt_transport_class, classification.transport_class}) + end + + previous_failures = counter(metrics, :consecutive_failures) + + {consecutive_failures, recovered_after_failures} = + if classification.result == :success do + :ets.insert(metrics.table, {:consecutive_failures, 0}) + {0, previous_failures} + else + {increment(metrics, :consecutive_failures), 0} + end + + now = System.monotonic_time(:millisecond) + {cache_degraded?, cache_degraded_duration_ms} = cache_degraded(metrics, now) + last_success_age_ms = last_success_age(metrics, now) + + measurements = %{ + count: 1, + total_attempts: total_attempts, + consecutive_failures: consecutive_failures, + recovered_after_failures: recovered_after_failures, + last_success_age_ms: last_success_age_ms || 0, + remaining_watchdog_budget_ms: max(deadline_at - now, 0), + cache_degraded_duration_ms: cache_degraded_duration_ms + } + + metadata = + Map.merge(classification, %{ + supervisor: metrics.supervisor_name, + retryable: retryable?, + cache_degraded: cache_degraded?, + has_last_success: is_integer(last_success_age_ms) + }) + + :telemetry.execute(@attempt_event, measurements, metadata) + maybe_log_attempt(metrics, response_or_error, measurements, metadata) + + classification + rescue + ArgumentError -> + # The owning lifecycle manager may disappear while its heartbeat task is + # being fenced. Observability must never turn that race into a new failure. + :ignored + end + + @doc false + @spec successful_attempt?(term()) :: boolean() + def successful_attempt?(response_or_error) do + classify_attempt(response_or_error).result == :success + end + + @doc false + @spec mark_heartbeat_success(t(), integer()) :: :ok + def mark_heartbeat_success(%__MODULE__{} = metrics, monotonic_at) + when is_integer(monotonic_at) do + :ets.insert(metrics.table, {:last_success_tick, tick(metrics, monotonic_at)}) + :ok + rescue + ArgumentError -> :ok + end + + @doc false + @spec mark_watchdog_renewal(t(), integer()) :: :ok + def mark_watchdog_renewal(%__MODULE__{} = metrics, monotonic_at) + when is_integer(monotonic_at) do + :ets.insert(metrics.table, {:watchdog_renewal_tick, tick(metrics, monotonic_at)}) + :ok + rescue + ArgumentError -> :ok + end + + @doc false + @spec record_cache_refresh(t(), boolean(), non_neg_integer(), non_neg_integer()) :: map() + def record_cache_refresh( + %__MODULE__{} = metrics, + complete?, + error_count, + refresh_duration_ms + ) + when is_boolean(complete?) and is_integer(error_count) and error_count >= 0 and + is_integer(refresh_duration_ms) and refresh_duration_ms >= 0 do + now = System.monotonic_time(:millisecond) + now_tick = tick(metrics, now) + + {status, transition, degraded_duration_ms} = + if complete? do + case :ets.take(metrics.table, :cache_degraded_since_tick) do + [{:cache_degraded_since_tick, degraded_since}] -> + {:healthy, :recovered, max(now_tick - degraded_since, 0)} + + [] -> + {:healthy, :unchanged, 0} + end + else + entered? = + :ets.insert_new( + metrics.table, + {:cache_degraded_since_tick, now_tick} + ) + + degraded_since = + if entered? do + now_tick + else + value(metrics, :cache_degraded_since_tick) + end + + transition = if entered?, do: :entered, else: :continued + {:degraded, transition, max(now_tick - degraded_since, 0)} + end + + measurements = %{ + count: 1, + error_count: error_count, + refresh_duration_ms: refresh_duration_ms, + degraded_duration_ms: degraded_duration_ms + } + + metadata = %{ + supervisor: metrics.supervisor_name, + status: status, + transition: transition + } + + :telemetry.execute(@cache_event, measurements, metadata) + + Map.merge(measurements, metadata) + rescue + ArgumentError -> + %{status: :unknown, transition: :unchanged, degraded_duration_ms: 0} + end + + @doc false + @spec watchdog_termination(atom(), map(), map()) :: :ok + def watchdog_termination(supervisor_name, measurements, metadata) + when is_atom(supervisor_name) and is_map(measurements) and is_map(metadata) do + :telemetry.execute( + @watchdog_event, + Map.put(measurements, :count, 1), + Map.put(metadata, :supervisor, supervisor_name) + ) + + :ok + end + + @doc false + @spec snapshot(t(), pos_integer()) :: map() + def snapshot(%__MODULE__{} = metrics, watchdog_deadline_ms) + when is_integer(watchdog_deadline_ms) and watchdog_deadline_ms > 0 do + now = System.monotonic_time(:millisecond) + entries = :ets.tab2list(metrics.table) + {cache_degraded?, cache_degraded_duration_ms} = cache_degraded(metrics, now) + last_success_age_ms = last_success_age(metrics, now) + watchdog_age_ms = elapsed_since(metrics, :watchdog_renewal_tick, now) + + %{ + attempts: %{ + total: entry_value(entries, :attempt_count, 0), + by_http_status: grouped_attempts(entries, :attempt_http_status), + by_transport_class: grouped_attempts(entries, :attempt_transport_class), + by_result: grouped_attempts(entries, :attempt_result) + }, + consecutive_failures: entry_value(entries, :consecutive_failures, 0), + last_success_age_ms: last_success_age_ms, + remaining_watchdog_budget_ms: + if(is_integer(watchdog_age_ms), + do: max(watchdog_deadline_ms - watchdog_age_ms, 0), + else: nil + ), + cache_degraded?: cache_degraded?, + cache_degraded_duration_ms: cache_degraded_duration_ms + } + rescue + ArgumentError -> + %{ + attempts: %{total: 0, by_http_status: %{}, by_transport_class: %{}, by_result: %{}}, + consecutive_failures: 0, + last_success_age_ms: nil, + remaining_watchdog_budget_ms: nil, + cache_degraded?: false, + cache_degraded_duration_ms: 0 + } + end + + @doc false + @spec classify_attempt(term()) :: attempt_classification() + def classify_attempt({:ok, _result}), do: success_classification() + def classify_attempt(:ok), do: success_classification() + def classify_attempt({:error, reason}), do: classify_attempt(reason) + def classify_attempt({:mirror_failed, reason}), do: classify_attempt(reason) + + def classify_attempt(%Req.Response{status: status}) when status in 200..299 do + %{ + result: :success, + http_status: normalize_http_status(status), + transport_class: :none, + error_class: :none + } + end + + def classify_attempt(%Req.Response{status: status}) do + %{ + result: :http_error, + http_status: normalize_http_status(status), + transport_class: :none, + error_class: http_error_class(status) + } + end + + def classify_attempt(%Req.TransportError{reason: reason}) do + %{ + result: :transport_error, + http_status: :none, + transport_class: transport_class(reason), + error_class: :transport + } + end + + def classify_attempt(%Req.HTTPError{}) do + %{ + result: :protocol_error, + http_status: :none, + transport_class: :none, + error_class: :protocol + } + end + + def classify_attempt(%ArgumentError{}) do + backend_error_classification(:configuration) + end + + def classify_attempt({:raised, :error, %ArgumentError{}}), + do: backend_error_classification(:configuration) + + def classify_attempt({:raised, kind, _reason}) when kind in [:error, :exit, :throw], + do: backend_error_classification(:other) + + def classify_attempt(reason) + when reason in [ + :invalid_credentials, + :unauthorized, + :forbidden, + :invalid_configuration, + :unsupported + ] do + backend_error_classification(:configuration) + end + + def classify_attempt(reason) when reason in [:conflict, :already_exists] do + backend_error_classification(:conflict) + end + + def classify_attempt(reason) + when reason in [ + :no_quorum, + :quorum_timeout, + :unavailable, + :temporarily_unavailable, + :cluster_overflow, + :cluster_not_ready + ] do + backend_error_classification(:unavailable) + end + + def classify_attempt(:timeout), do: backend_error_classification(:timeout) + def classify_attempt(_reason), do: backend_error_classification(:other) + + defp success_classification do + %{ + result: :success, + http_status: :none, + transport_class: :none, + error_class: :none + } + end + + defp backend_error_classification(error_class) do + %{ + result: :backend_error, + http_status: :none, + transport_class: :none, + error_class: error_class + } + end + + defp normalize_http_status(status) when status in 100..599, do: status + defp normalize_http_status(_status), do: :other + + defp http_error_class(status) when status in [401, 403], do: :authentication + defp http_error_class(status) when status in [400, 404], do: :configuration + defp http_error_class(_status), do: :http + + defp transport_class(:timeout), do: :timeout + defp transport_class(:closed), do: :closed + defp transport_class(:econnrefused), do: :connection_refused + + defp transport_class(reason) + when reason in [:nxdomain, :eai_again, :eai_fail, :eai_noname, :host_not_found], + do: :dns + + defp transport_class(reason) + when reason in [ + :enetdown, + :enetreset, + :enetunreach, + :econnaborted, + :econnreset, + :ehostdown, + :ehostunreach + ], + do: :network + + defp transport_class(:protocol_not_negotiated), do: :tls + defp transport_class({:bad_alpn_protocol, _protocol}), do: :tls + defp transport_class({:tls_alert, _alert}), do: :tls + defp transport_class({:bad_cert, _reason}), do: :tls + defp transport_class(reason) when is_atom(reason), do: :socket + defp transport_class(_reason), do: :other + + defp increment(%__MODULE__{table: table}, key) do + :ets.update_counter(table, key, {2, 1}, {key, 0}) + end + + defp counter(%__MODULE__{} = metrics, key), do: value(metrics, key, 0) + + defp value(%__MODULE__{table: table}, key, default \\ nil) do + case :ets.lookup(table, key) do + [{^key, value}] -> value + [] -> default + end + end + + defp tick(%__MODULE__{started_at: started_at}, monotonic_at) do + max(monotonic_at - started_at + 1, 1) + end + + defp elapsed_since(%__MODULE__{} = metrics, key, now) do + case value(metrics, key) do + stored_tick when is_integer(stored_tick) -> max(tick(metrics, now) - stored_tick, 0) + nil -> nil + end + end + + defp last_success_age(%__MODULE__{} = metrics, now) do + elapsed_since(metrics, :last_success_tick, now) + end + + defp cache_degraded(%__MODULE__{} = metrics, now) do + case value(metrics, :cache_degraded_since_tick) do + degraded_since when is_integer(degraded_since) -> + {true, max(tick(metrics, now) - degraded_since, 0)} + + nil -> + {false, 0} + end + end + + defp entry_value(entries, key, default) do + case List.keyfind(entries, key, 0) do + {^key, value} -> value + nil -> default + end + end + + defp grouped_attempts(entries, group) do + for {{^group, classification}, count} <- entries, into: %{} do + {classification, count} + end + end + + defp maybe_log_attempt( + %__MODULE__{supervisor_name: supervisor_name}, + _response_or_error, + measurements, + %{result: :success} + ) + when measurements.recovered_after_failures > 0 do + Logger.info( + "#{inspect(supervisor_name)}: heartbeat request path recovered after " <> + "#{measurements.recovered_after_failures} consecutive failed attempt(s); " <> + "last_success_age_ms=#{measurements.last_success_age_ms} " <> + "remaining_watchdog_budget_ms=#{measurements.remaining_watchdog_budget_ms}", + durable_server_supervisor: supervisor_name, + heartbeat_result: :success + ) + end + + defp maybe_log_attempt( + %__MODULE__{supervisor_name: supervisor_name}, + response_or_error, + measurements, + %{result: result} = metadata + ) + when result != :success do + if not metadata.retryable or log_failure_count?(measurements.consecutive_failures) do + Logger.warning( + "#{inspect(supervisor_name)}: heartbeat attempt failed; " <> + "result=#{metadata.result} http_status=#{metadata.http_status} " <> + "transport_class=#{metadata.transport_class} error_class=#{metadata.error_class} " <> + attempt_detail(response_or_error) <> + "retryable=#{metadata.retryable} " <> + "has_last_success=#{metadata.has_last_success} " <> + "consecutive_failures=#{measurements.consecutive_failures} " <> + "last_success_age_ms=#{measurements.last_success_age_ms} " <> + "remaining_watchdog_budget_ms=#{measurements.remaining_watchdog_budget_ms}", + durable_server_supervisor: supervisor_name, + heartbeat_result: metadata.result, + heartbeat_http_status: metadata.http_status, + heartbeat_transport_class: metadata.transport_class, + heartbeat_retryable: metadata.retryable + ) + end + end + + defp maybe_log_attempt(_metrics, _response_or_error, _measurements, _metadata), do: :ok + + defp attempt_detail(%Req.TransportError{reason: reason}) do + "transport_reason=#{inspect(reason, limit: 20, printable_limit: 200)} " + end + + defp attempt_detail({:error, reason}), do: attempt_detail(reason) + defp attempt_detail({:mirror_failed, reason}), do: attempt_detail(reason) + defp attempt_detail(_response_or_error), do: "" + + defp log_failure_count?(count) when count <= 3, do: true + defp log_failure_count?(count) when count > 0, do: Bitwise.band(count, count - 1) == 0 +end diff --git a/lib/durable_server/heartbeat_watchdog.ex b/lib/durable_server/heartbeat_watchdog.ex index 5101438..73fe0e6 100644 --- a/lib/durable_server/heartbeat_watchdog.ex +++ b/lib/durable_server/heartbeat_watchdog.ex @@ -4,12 +4,15 @@ defmodule DurableServer.HeartbeatWatchdog do use GenServer require Logger + alias DurableServer.HeartbeatMetrics + defstruct supervisor_name: nil, owner: nil, owner_monitor: nil, deadline_ms: nil, last_heartbeat_at: nil, heartbeat_task_pid: nil, + heartbeat_metrics: nil, timer_ref: nil, timer_token: nil @@ -23,7 +26,16 @@ defmodule DurableServer.HeartbeatWatchdog do def arm(name, owner, last_heartbeat_at, deadline_ms) when is_pid(owner) and is_integer(last_heartbeat_at) and is_integer(deadline_ms) and deadline_ms > 0 do - GenServer.call(name, {:arm, owner, last_heartbeat_at, deadline_ms}) + arm(name, owner, last_heartbeat_at, deadline_ms, nil) + end + + def arm(name, owner, last_heartbeat_at, deadline_ms, heartbeat_metrics) + when is_pid(owner) and is_integer(last_heartbeat_at) and is_integer(deadline_ms) and + deadline_ms > 0 do + GenServer.call( + name, + {:arm, owner, last_heartbeat_at, deadline_ms, heartbeat_metrics} + ) end def renew(name, owner, heartbeat_at) @@ -47,7 +59,7 @@ defmodule DurableServer.HeartbeatWatchdog do @impl true def handle_call( - {:arm, owner, last_heartbeat_at, deadline_ms}, + {:arm, owner, last_heartbeat_at, deadline_ms, heartbeat_metrics}, _from, %__MODULE__{} = state ) do @@ -65,7 +77,8 @@ defmodule DurableServer.HeartbeatWatchdog do owner_monitor: owner_monitor, deadline_ms: deadline_ms, last_heartbeat_at: last_heartbeat_at, - heartbeat_task_pid: nil + heartbeat_task_pid: nil, + heartbeat_metrics: heartbeat_metrics } |> schedule_deadline() @@ -123,11 +136,35 @@ defmodule DurableServer.HeartbeatWatchdog do %__MODULE__{timer_token: timer_token} = state ) do elapsed_since_last = System.monotonic_time(:millisecond) - state.last_heartbeat_at + {children_fenced, child_count_status} = count_fenced_children(state.supervisor_name) + metrics = heartbeat_snapshot(state, elapsed_since_last) + + :ok = + HeartbeatMetrics.watchdog_termination( + state.supervisor_name, + %{ + children_fenced: children_fenced, + consecutive_failures: metrics.consecutive_failures, + last_success_age_ms: metrics.last_success_age_ms, + remaining_watchdog_budget_ms: 0, + cache_degraded_duration_ms: metrics.cache_degraded_duration_ms, + watchdog_terminations: 1 + }, + %{ + child_count_status: child_count_status, + cache_degraded: metrics.cache_degraded? + } + ) Logger.error(fn -> "#{inspect(state.supervisor_name)}: heartbeat watchdog deadline exceeded " <> - "(#{elapsed_since_last}ms since last success, deadline #{state.deadline_ms}ms), " <> - "terminating lifecycle manager to prevent orphan conflicts" + "(#{elapsed_since_last}ms since last watchdog renewal, " <> + "#{metrics.last_success_age_ms}ms since last successful write, " <> + "deadline #{state.deadline_ms}ms, remaining budget 0ms, " <> + "consecutive failures #{metrics.consecutive_failures}, " <> + "cache degraded for #{metrics.cache_degraded_duration_ms}ms), " <> + "terminating lifecycle manager to prevent orphan conflicts; " <> + "children_fenced=#{format_child_count(children_fenced, child_count_status)}" end) if is_pid(state.heartbeat_task_pid) do @@ -166,6 +203,7 @@ defmodule DurableServer.HeartbeatWatchdog do deadline_ms: nil, last_heartbeat_at: nil, heartbeat_task_pid: nil, + heartbeat_metrics: nil, timer_ref: nil, timer_token: nil }) @@ -185,4 +223,56 @@ defmodule DurableServer.HeartbeatWatchdog do end defp clear_owner_monitor(%__MODULE__{} = state), do: state + + defp heartbeat_snapshot( + %__MODULE__{heartbeat_metrics: %HeartbeatMetrics{} = metrics, deadline_ms: deadline_ms}, + elapsed_since_last + ) do + snapshot = HeartbeatMetrics.snapshot(metrics, deadline_ms) + + %{ + consecutive_failures: snapshot.consecutive_failures, + last_success_age_ms: snapshot.last_success_age_ms || elapsed_since_last, + cache_degraded?: snapshot.cache_degraded?, + cache_degraded_duration_ms: snapshot.cache_degraded_duration_ms + } + end + + defp heartbeat_snapshot(%__MODULE__{}, elapsed_since_last) do + %{ + consecutive_failures: 0, + last_success_age_ms: elapsed_since_last, + cache_degraded?: false, + cache_degraded_duration_ms: 0 + } + end + + # Querying DynamicSupervisor would be a blocking GenServer call on the process + # responsible for enforcing the hard deadline. Its direct links are the + # supervised children plus its parent, so this gives us a non-blocking count. + defp count_fenced_children(supervisor_name) do + dynamic_supervisor = DurableServer.Supervisor.get_dynamic_supervisor(supervisor_name) + + with pid when is_pid(pid) <- Process.whereis(dynamic_supervisor), + [links: links, dictionary: dictionary] <- + Process.info(pid, [:links, :dictionary]) do + ancestors = + dictionary + |> Keyword.get(:"$ancestors", []) + |> Enum.filter(&is_pid/1) + |> MapSet.new() + + children_fenced = + Enum.count(links, fn linked_pid -> + not MapSet.member?(ancestors, linked_pid) + end) + + {children_fenced, :known} + else + _ -> {0, :unavailable} + end + end + + defp format_child_count(children_fenced, :known), do: children_fenced + defp format_child_count(_children_fenced, :unavailable), do: "unknown" end diff --git a/lib/durable_server/lifecycle_manager.ex b/lib/durable_server/lifecycle_manager.ex index d6a5bf2..2594ff0 100644 --- a/lib/durable_server/lifecycle_manager.ex +++ b/lib/durable_server/lifecycle_manager.ex @@ -103,7 +103,7 @@ defmodule DurableServer.LifecycleManager do alias DurableServer.LifecycleManager alias DurableServer - alias DurableServer.{StoredState, Meta, CircuitBreaker, HeartbeatWatchdog} + alias DurableServer.{StoredState, Meta, CircuitBreaker, HeartbeatMetrics, HeartbeatWatchdog} alias DurableServer.ObjectStore alias DurableServer.StorageBackend @@ -130,6 +130,7 @@ defmodule DurableServer.LifecycleManager do last_successful_heartbeat_at: nil, last_successful_heartbeat_monotonic_at: nil, heartbeat_watchdog: nil, + heartbeat_metrics: nil, # Last successful heartbeat timing for diagnostics last_heartbeat_timing: nil, discovery_diag_table: nil, @@ -192,6 +193,11 @@ defmodule DurableServer.LifecycleManager do - `: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 + - `:attempts` - Bounded attempt counts by result, HTTP status, and transport class + - `:consecutive_failures` - Current run of failed heartbeat attempts + - `:last_success_age_ms` - Monotonic age of the last successful heartbeat write + - `:remaining_watchdog_budget_ms` - Time remaining before the local watchdog fences the tree + - `:cache_degraded_duration_ms` - How long the heartbeat cache has been incomplete - `:node` - This node's name This is used by the admin dashboard to monitor heartbeat health across the cluster. @@ -404,6 +410,8 @@ defmodule DurableServer.LifecycleManager do write_concurrency: true ]) + heartbeat_metrics = HeartbeatMetrics.new(supervisor_name) + state = %LifecycleManager{ supervisor_name: supervisor_name, task_sup: task_supervisor, @@ -425,6 +433,7 @@ defmodule DurableServer.LifecycleManager do capacity_limits: capacity_limits, heartbeat_meta: heartbeat_meta, heartbeat_watchdog: heartbeat_watchdog, + heartbeat_metrics: heartbeat_metrics, discovery_diag_table: diagnostics_tab, discovery_skip_table: skip_tab, restart_gate_table: restart_gate_tab, @@ -489,7 +498,14 @@ defmodule DurableServer.LifecycleManager do state.heartbeat_watchdog, self(), heartbeat_monotonic_at, - heartbeat_hard_deadline_ms(state) + heartbeat_hard_deadline_ms(state), + state.heartbeat_metrics + ) + + :ok = + HeartbeatMetrics.mark_watchdog_renewal( + state.heartbeat_metrics, + heartbeat_monotonic_at ) state = @@ -726,6 +742,8 @@ defmodule DurableServer.LifecycleManager do "heartbeat cache reconciliation task failed: #{inspect(reason)}" end) + HeartbeatMetrics.record_cache_refresh(state.heartbeat_metrics, false, 1, 0) + state = apply_heartbeat_cache_status( %{state | current_heartbeat_reconcile_task: nil}, @@ -766,6 +784,11 @@ defmodule DurableServer.LifecycleManager do heartbeat_monotonic_at ) + HeartbeatMetrics.mark_watchdog_renewal( + state.heartbeat_metrics, + heartbeat_monotonic_at + ) + join_group_heartbeat(state, heartbeat_entry) %{ @@ -790,18 +813,23 @@ defmodule DurableServer.LifecycleManager do capacity = DurableServer.Supervisor.current_capacity(state.supervisor_name) heartbeat_meta = resolve_heartbeat_meta(state.heartbeat_meta) - metrics = %{ - node: Node.self(), - last_heartbeat_timing: state.last_heartbeat_timing, - last_successful_heartbeat_at: state.last_successful_heartbeat_at, - heartbeat_interval_ms: state.heartbeat_interval_ms, - heartbeat_staleness_threshold_ms: heartbeat_staleness_threshold_ms(state), - deadline_ms: heartbeat_hard_deadline_ms(state), - resources: resources, - capacity: capacity, - heartbeat_meta: heartbeat_meta, - discovery_degraded: state.discovery_degraded or discovery_degraded?(state.supervisor_name) - } + heartbeat_metrics = + HeartbeatMetrics.snapshot(state.heartbeat_metrics, heartbeat_hard_deadline_ms(state)) + + metrics = + %{ + node: Node.self(), + last_heartbeat_timing: state.last_heartbeat_timing, + last_successful_heartbeat_at: state.last_successful_heartbeat_at, + heartbeat_interval_ms: state.heartbeat_interval_ms, + heartbeat_staleness_threshold_ms: heartbeat_staleness_threshold_ms(state), + deadline_ms: heartbeat_hard_deadline_ms(state), + resources: resources, + capacity: capacity, + heartbeat_meta: heartbeat_meta, + discovery_degraded: state.discovery_degraded or discovery_degraded?(state.supervisor_name) + } + |> Map.merge(heartbeat_metrics) {:reply, metrics, state} end @@ -865,6 +893,11 @@ defmodule DurableServer.LifecycleManager do owner, heartbeat_monotonic_at ) + + HeartbeatMetrics.mark_watchdog_renewal( + state.heartbeat_metrics, + heartbeat_monotonic_at + ) end defp renew_heartbeat_watchdog( @@ -1073,6 +1106,12 @@ defmodule DurableServer.LifecycleManager do case put_heartbeat_until_deadline(state, key, heartbeat_data, deadline_at) do {:ok, _} -> + :ok = + HeartbeatMetrics.mark_heartbeat_success( + state.heartbeat_metrics, + current_monotonic_time + ) + # update local ets cache with full capacity info :ets.insert(state.heartbeat_table, entry) @@ -1102,9 +1141,7 @@ defmodule DurableServer.LifecycleManager do # Let Req retry transient HTTP and transport failures itself. The timeout keeps # those retries inside the heartbeat budget; the outer loop remains for retryable # errors from non-HTTP backends. - put_opts = [max_retries: @heartbeat_req_max_retries, timeout: remaining_ms] - - case StorageBackend.put_object(state.heartbeat_store, key, heartbeat_data, put_opts) do + case put_heartbeat_once(state, key, heartbeat_data, deadline_at, remaining_ms) do {:ok, _} = ok -> ok @@ -1124,6 +1161,90 @@ defmodule DurableServer.LifecycleManager do end end + defp put_heartbeat_once( + %LifecycleManager{} = state, + key, + heartbeat_data, + deadline_at, + remaining_ms + ) do + observations = :counters.new(2, []) + + retry_observer = fn _request, response_or_error, retryable? -> + :counters.add(observations, 1, 1) + + unless HeartbeatMetrics.successful_attempt?(response_or_error) do + :counters.add(observations, 2, 1) + end + + HeartbeatMetrics.record_attempt( + state.heartbeat_metrics, + response_or_error, + retryable?, + deadline_at + ) + end + + put_opts = [ + max_retries: @heartbeat_req_max_retries, + timeout: remaining_ms, + retry_observer: retry_observer + ] + + try do + result = + StorageBackend.put_object( + state.heartbeat_store, + key, + heartbeat_data, + put_opts + ) + + maybe_record_backend_attempt(state, result, observations, deadline_at) + result + catch + kind, reason -> + if :counters.get(observations, 2) == 0 do + HeartbeatMetrics.record_attempt( + state.heartbeat_metrics, + {:raised, kind, reason}, + false, + deadline_at + ) + end + + :erlang.raise(kind, reason, __STACKTRACE__) + end + end + + defp maybe_record_backend_attempt( + %LifecycleManager{} = state, + result, + observations, + deadline_at + ) do + observed_attempts = :counters.get(observations, 1) + observed_failures = :counters.get(observations, 2) + + if observed_attempts == 0 or + (not HeartbeatMetrics.successful_attempt?(result) and observed_failures == 0) do + retryable? = + case result do + {:error, reason} -> heartbeat_write_retryable?(reason) + _ -> false + end + + HeartbeatMetrics.record_attempt( + state.heartbeat_metrics, + result, + retryable?, + deadline_at + ) + end + + :ok + end + defp heartbeat_deadline_at(%LifecycleManager{} = state, nil), do: System.monotonic_time(:millisecond) + heartbeat_hard_deadline_ms(state) @@ -1142,6 +1263,15 @@ defmodule DurableServer.LifecycleManager do defp heartbeat_write_retryable?({:mirror_failed, reason}), do: heartbeat_write_retryable?(reason) + defp heartbeat_write_retryable?(%Req.TransportError{}), do: true + + defp heartbeat_write_retryable?(%Req.Response{status: status}) + when status in [408, 429, 500, 502, 503, 504], + do: true + + defp heartbeat_write_retryable?(%Req.HTTPError{protocol: :http2, reason: :unprocessed}), + do: true + defp heartbeat_write_retryable?(reason) do reason in [ :no_quorum, @@ -1286,6 +1416,7 @@ defmodule DurableServer.LifecycleManager do # this gets run async inside a task defp refresh_node_heartbeat_cache(%LifecycleManager{} = state) do + cache_refresh_started_at = System.monotonic_time(:millisecond) dead_node_threshold_ms = state.config.dead_node_threshold_ms current_time = System.system_time(:millisecond) heartbeat_prefix = "#{state.prefix}__nodes/" @@ -1384,17 +1515,30 @@ defmodule DurableServer.LifecycleManager do {node, node_ref, timestamp, capacity, resources, env_vars, heartbeat_meta} end) - 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) + result = + 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) - {:ok, length(live_heartbeats), cleaned_count, 0} - else - {:partial, length(live_heartbeats), 0, error_count} - end + {:ok, length(live_heartbeats), cleaned_count, 0} + else + {:partial, length(live_heartbeats), 0, error_count} + end + + cache_refresh_duration_ms = + System.monotonic_time(:millisecond) - cache_refresh_started_at + + HeartbeatMetrics.record_cache_refresh( + state.heartbeat_metrics, + listing_complete?, + error_count, + cache_refresh_duration_ms + ) + + result end defp cleanup_dead_node_heartbeats(%LifecycleManager{} = state, dead_nodes) do diff --git a/lib/durable_server/object_store.ex b/lib/durable_server/object_store.ex index 371b212..5f9ce63 100644 --- a/lib/durable_server/object_store.ex +++ b/lib/durable_server/object_store.ex @@ -909,12 +909,20 @@ defmodule DurableServer.ObjectStore do - `:etag` - The existing etag to match. Conflicts return `{:error, :conflict}` - `:timeout` - Total time in ms for the operation including retries. If exceeded, no further retries will be attempted. Default: no timeout (unlimited retries until max_retries). + - `:retry_observer` - Optional three-arity callback invoked with the request, + response/error, and retryable classification for each HTTP attempt. """ def put_object(%__MODULE__{} = client, key, data, opts \\ []) do opts = validate_opts!(opts) content_type = Keyword.get(opts, :content_type, "application/octet-stream") consistent = Keyword.get(opts, :consistent, true) timeout = Keyword.get(opts, :timeout) + retry_observer = Keyword.get(opts, :retry_observer) + + unless is_nil(retry_observer) or is_function(retry_observer, 3) do + raise ArgumentError, + "expected :retry_observer to be a 3-arity function, got: #{inspect(retry_observer)}" + end # Handle :infinity and nil as no deadline deadline_at = @@ -944,30 +952,22 @@ defmodule DurableServer.ObjectStore do req = new_req(client, consistent: consistent, headers: headers) req = Req.merge(req, receive_timeout: retry_attempt_timeout(timeout)) + retry = fn %Req.Request{} = request, response_or_exception -> + retryable? = put_retryable?(response_or_exception, deadline_at) + notify_retry_observer(retry_observer, request, response_or_exception, retryable?) + retryable? + end + req_with_retries = case Keyword.fetch(opts, :max_retries) do {:ok, retries} when is_integer(retries) and retries >= 0 -> - Req.merge(req, - max_retries: retries, - retry: fn - %Req.Request{}, %Req.Response{status: status} - when status in [408, 429, 500, 502, 503, 504] -> - within_retry_deadline?(deadline_at) - - %Req.Request{}, %Req.TransportError{reason: reason} - when reason in [:timeout, :econnrefused, :closed] -> - within_retry_deadline?(deadline_at) - - %Req.Request{}, %Req.HTTPError{protocol: :http2, reason: :unprocessed} -> - within_retry_deadline?(deadline_at) - - %Req.Request{}, _response_or_exception -> - false - end - ) + Req.merge(req, max_retries: retries, retry: retry) - :error -> + :error when is_nil(retry_observer) -> req + + :error -> + Req.merge(req, max_retries: 0, retry: retry) end case Req.put(req_with_retries, url: key, body: data) do @@ -995,6 +995,33 @@ defmodule DurableServer.ObjectStore do System.monotonic_time(:millisecond) < deadline_at end + defp put_retryable?(%Req.Response{status: status}, deadline_at) + when status in [408, 429, 500, 502, 503, 504] do + within_retry_deadline?(deadline_at) + end + + # A transport failure says that the caller does not have an authoritative + # HTTP response. Mint may add new socket/TLS reasons over time, so all + # variants are ambiguous and safe to retry inside the caller's hard deadline. + defp put_retryable?(%Req.TransportError{}, deadline_at) do + within_retry_deadline?(deadline_at) + end + + defp put_retryable?( + %Req.HTTPError{protocol: :http2, reason: :unprocessed}, + deadline_at + ) do + within_retry_deadline?(deadline_at) + end + + defp put_retryable?(_response_or_exception, _deadline_at), do: false + + defp notify_retry_observer(nil, _request, _response_or_exception, _retryable?), do: :ok + + defp notify_retry_observer(observer, request, response_or_exception, retryable?) do + observer.(request, response_or_exception, retryable?) + end + defp retry_attempt_timeout(timeout) when is_integer(timeout) and timeout >= 0, do: min(timeout, @default_retry_attempt_timeout) @@ -1021,7 +1048,8 @@ defmodule DurableServer.ObjectStore do :max_results, :continuation_token, :prefix, - :etag + :etag, + :retry_observer ]) end diff --git a/mix.exs b/mix.exs index b92ea1e..93e8407 100644 --- a/mix.exs +++ b/mix.exs @@ -44,6 +44,7 @@ defmodule DurableServer.MixProject do {:req, "~> 0.5"}, {:req_s3, "~> 0.2"}, {:finch, "~> 0.18"}, + {:telemetry, "~> 1.0"}, {:sweet_xml, "~> 0.7"}, {:ex_doc, "~> 0.30", only: :dev, runtime: false}, {:ekv, "~> 0.4.0", optional: true} diff --git a/test/durable_server/heartbeat_metrics_test.exs b/test/durable_server/heartbeat_metrics_test.exs new file mode 100644 index 0000000..7ddb826 --- /dev/null +++ b/test/durable_server/heartbeat_metrics_test.exs @@ -0,0 +1,137 @@ +defmodule DurableServer.HeartbeatMetricsTest do + use ExUnit.Case, async: true + + alias DurableServer.HeartbeatMetrics + + @moduletag :capture_log + + test "records attempts with bounded HTTP and transport classifications" do + supervisor_name = unique_name() + metrics = HeartbeatMetrics.new(supervisor_name) + deadline_at = System.monotonic_time(:millisecond) + 5_000 + attach_telemetry(HeartbeatMetrics.events().attempt, supervisor_name) + + HeartbeatMetrics.mark_heartbeat_success(metrics, System.monotonic_time(:millisecond)) + HeartbeatMetrics.mark_watchdog_renewal(metrics, System.monotonic_time(:millisecond)) + + HeartbeatMetrics.record_attempt(metrics, %Req.Response{status: 503}, true, deadline_at) + + HeartbeatMetrics.record_attempt( + metrics, + %Req.TransportError{reason: :enetunreach}, + true, + deadline_at + ) + + HeartbeatMetrics.record_attempt( + metrics, + %Req.TransportError{reason: {:future_transport_reason, make_ref()}}, + true, + deadline_at + ) + + HeartbeatMetrics.record_attempt(metrics, %Req.Response{status: 200}, false, deadline_at) + + snapshot = HeartbeatMetrics.snapshot(metrics, 5_000) + + assert snapshot.attempts == %{ + total: 4, + by_http_status: %{200 => 1, 503 => 1}, + by_transport_class: %{network: 1, other: 1}, + by_result: %{http_error: 1, success: 1, transport_error: 2} + } + + assert snapshot.consecutive_failures == 0 + assert is_integer(snapshot.last_success_age_ms) + assert snapshot.remaining_watchdog_budget_ms in 0..5_000 + + assert_receive {:telemetry, measurements, + %{ + result: :http_error, + http_status: 503, + transport_class: :none, + retryable: true + }} + + assert measurements.consecutive_failures == 1 + + assert_receive {:telemetry, _measurements, + %{result: :transport_error, transport_class: :network}} + + assert_receive {:telemetry, _measurements, + %{result: :transport_error, transport_class: :other} = metadata} + + refute Map.has_key?(metadata, :reason) + refute Map.has_key?(metadata, :key) + + assert_receive {:telemetry, %{recovered_after_failures: 3}, + %{result: :success, http_status: 200}} + end + + test "classifies permanent authentication and configuration failures" do + assert HeartbeatMetrics.classify_attempt(%Req.Response{status: 403}) == %{ + result: :http_error, + http_status: 403, + transport_class: :none, + error_class: :authentication + } + + assert HeartbeatMetrics.classify_attempt(%Req.Response{status: 400}) == %{ + result: :http_error, + http_status: 400, + transport_class: :none, + error_class: :configuration + } + end + + test "tracks cache degradation until a complete refresh recovers" do + supervisor_name = unique_name() + metrics = HeartbeatMetrics.new(supervisor_name) + attach_telemetry(HeartbeatMetrics.events().cache, supervisor_name) + + assert %{status: :degraded, transition: :entered, degraded_duration_ms: 0} = + HeartbeatMetrics.record_cache_refresh(metrics, false, 2, 10) + + Process.sleep(5) + + snapshot = HeartbeatMetrics.snapshot(metrics, 5_000) + assert snapshot.cache_degraded? + assert snapshot.cache_degraded_duration_ms >= 5 + + assert %{status: :healthy, transition: :recovered, degraded_duration_ms: duration_ms} = + HeartbeatMetrics.record_cache_refresh(metrics, true, 0, 4) + + assert duration_ms >= 5 + + refute HeartbeatMetrics.snapshot(metrics, 5_000).cache_degraded? + + assert_receive {:telemetry, %{error_count: 2}, %{status: :degraded, transition: :entered}} + + assert_receive {:telemetry, %{degraded_duration_ms: recovered_duration}, + %{status: :healthy, transition: :recovered}} + + assert recovered_duration >= 5 + end + + defp attach_telemetry(event, supervisor_name) do + handler_id = {__MODULE__, self(), make_ref()} + + :ok = + :telemetry.attach( + handler_id, + event, + fn _event, measurements, metadata, test_pid -> + if metadata.supervisor == supervisor_name do + send(test_pid, {:telemetry, measurements, metadata}) + end + end, + self() + ) + + on_exit(fn -> :telemetry.detach(handler_id) end) + end + + defp unique_name do + :"heartbeat_metrics_test_#{System.unique_integer([:positive])}" + end +end diff --git a/test/durable_server/heartbeat_watchdog_test.exs b/test/durable_server/heartbeat_watchdog_test.exs new file mode 100644 index 0000000..ca63f78 --- /dev/null +++ b/test/durable_server/heartbeat_watchdog_test.exs @@ -0,0 +1,92 @@ +defmodule DurableServer.HeartbeatWatchdogTest do + use ExUnit.Case, async: true + + alias DurableServer.{HeartbeatMetrics, HeartbeatWatchdog} + + @moduletag :capture_log + + test "termination telemetry reports the children that the supervisor tree will fence" do + supervisor_name = unique_name("supervisor") + dynamic_supervisor = DurableServer.Supervisor.get_dynamic_supervisor(supervisor_name) + watchdog_name = unique_name("watchdog") + + start_supervised!({DynamicSupervisor, name: dynamic_supervisor, strategy: :one_for_one}) + + Enum.each(1..2, fn id -> + child_spec = + Supervisor.child_spec( + {Task, fn -> Process.sleep(:infinity) end}, + id: {__MODULE__, id} + ) + + assert {:ok, _pid} = DynamicSupervisor.start_child(dynamic_supervisor, child_spec) + end) + + start_supervised!( + Supervisor.child_spec( + {HeartbeatWatchdog, name: watchdog_name, supervisor_name: supervisor_name}, + id: watchdog_name + ) + ) + + metrics = HeartbeatMetrics.new(supervisor_name) + now = System.monotonic_time(:millisecond) + HeartbeatMetrics.mark_heartbeat_success(metrics, now) + HeartbeatMetrics.mark_watchdog_renewal(metrics, now) + HeartbeatMetrics.record_cache_refresh(metrics, false, 1, 2) + + HeartbeatMetrics.record_attempt( + metrics, + %Req.TransportError{reason: :enetunreach}, + true, + now + 5_000 + ) + + attach_telemetry(HeartbeatMetrics.events().watchdog_termination, supervisor_name) + + owner = spawn(fn -> Process.sleep(:infinity) end) + owner_ref = Process.monitor(owner) + + assert :ok = HeartbeatWatchdog.arm(watchdog_name, owner, now, 50, metrics) + assert_receive {:DOWN, ^owner_ref, :process, ^owner, :killed}, 1_000 + + assert_receive {:telemetry, + %{ + count: 1, + watchdog_terminations: 1, + children_fenced: 2, + consecutive_failures: 1, + remaining_watchdog_budget_ms: 0, + cache_degraded_duration_ms: degraded_duration_ms + }, + %{ + child_count_status: :known, + cache_degraded: true + }}, + 1_000 + + assert degraded_duration_ms >= 0 + end + + defp attach_telemetry(event, supervisor_name) do + handler_id = {__MODULE__, self(), make_ref()} + + :ok = + :telemetry.attach( + handler_id, + event, + fn _event, measurements, metadata, {test_pid, expected_supervisor} -> + if metadata.supervisor == expected_supervisor do + send(test_pid, {:telemetry, measurements, metadata}) + end + end, + {self(), supervisor_name} + ) + + on_exit(fn -> :telemetry.detach(handler_id) end) + end + + defp unique_name(label) do + :"heartbeat_watchdog_test_#{label}_#{System.unique_integer([:positive])}" + end +end diff --git a/test/durable_server/lifecycle_test.exs b/test/durable_server/lifecycle_test.exs index 091b07f..11ebd14 100644 --- a/test/durable_server/lifecycle_test.exs +++ b/test/durable_server/lifecycle_test.exs @@ -111,12 +111,12 @@ defmodule DurableServer.LifecycleTest do def list_all_objects_stream(_state, _prefix, _opts), do: [] @impl true - def put_object(%{table: table, fail_count: fail_count}, _key, data, opts) do + def put_object(%{table: table, fail_count: fail_count} = state, _key, data, opts) do attempt = :ets.update_counter(table, :attempts, {2, 1}, {:attempts, 0}) :ets.insert(table, {:last_put_opts, opts}) if attempt <= fail_count do - {:error, {:mirror_failed, :no_quorum}} + {:error, Map.get(state, :failure_reason, {:mirror_failed, :no_quorum})} else :ets.insert(table, {:last_write, data}) {:ok, %{body: data, etag: Integer.to_string(attempt)}} @@ -2979,7 +2979,7 @@ defmodule DurableServer.LifecycleTest do |> Map.put(:object_store, storage_backend) |> Map.put(:storage_backend, storage_backend) - {:ok, _pid} = start_standalone_lifecycle_manager(supervisor_name, test_config) + {:ok, manager_pid} = start_standalone_lifecycle_manager(supervisor_name, test_config) assert [{:attempts, attempts}] = :ets.lookup(table, :attempts) assert attempts >= 3 @@ -2988,6 +2988,47 @@ defmodule DurableServer.LifecycleTest do assert [{:last_write, heartbeat_data}] = :ets.lookup(table, :last_write) assert is_map(heartbeat_data) assert Map.has_key?(heartbeat_data, "last_heartbeat_at") + + manager_name = :sys.get_state(manager_pid).supervisor_name + metrics = LifecycleManager.get_heartbeat_metrics(manager_name) + + assert metrics.attempts.total == attempts + assert metrics.attempts.by_result == %{backend_error: 2, success: 1} + assert metrics.consecutive_failures == 0 + assert is_integer(metrics.last_success_age_ms) + assert metrics.remaining_watchdog_budget_ms > 0 + refute metrics.cache_degraded? + assert metrics.cache_degraded_duration_ms == 0 + end + + test "retries ambiguous Req transport errors from a heartbeat backend", %{ + supervisor_name: supervisor_name, + config: config + } do + table = :ets.new(__MODULE__.HeartbeatRetryBackend, [:set, :public]) + + storage_backend = + DurableServer.StorageBackend.new(HeartbeatRetryBackend, %{ + table: table, + fail_count: 2, + failure_reason: %Req.TransportError{reason: {:future_transport_reason, make_ref()}} + }) + + test_config = + config + |> Map.put(:object_store, storage_backend) + |> Map.put(:storage_backend, storage_backend) + + {:ok, manager_pid} = start_standalone_lifecycle_manager(supervisor_name, test_config) + + assert [{:attempts, 3}] = :ets.lookup(table, :attempts) + + manager_name = :sys.get_state(manager_pid).supervisor_name + metrics = LifecycleManager.get_heartbeat_metrics(manager_name) + + assert metrics.attempts.by_result == %{success: 1, transport_error: 2} + assert metrics.attempts.by_transport_class == %{other: 2} + assert metrics.consecutive_failures == 0 end test "heartbeat watchdog terminates lifecycle manager even when message handling is suspended", @@ -3306,6 +3347,10 @@ defmodule DurableServer.LifecycleTest do assert :ets.lookup(heartbeat_table, known_node) == known_snapshot assert :ets.lookup(heartbeat_table, new_node) == [] + degraded_metrics = LifecycleManager.get_heartbeat_metrics(manager_supervisor) + assert degraded_metrics.cache_degraded? + assert is_integer(degraded_metrics.cache_degraded_duration_ms) + :ets.insert(control, {:mode, :ok}) await_heartbeat_reconciliation(manager_pid) @@ -3314,6 +3359,10 @@ defmodule DurableServer.LifecycleTest do refute LifecycleManager.discovery_degraded?(manager_supervisor) assert :ets.lookup(heartbeat_table, new_node) != [] + recovered_metrics = LifecycleManager.get_heartbeat_metrics(manager_supervisor) + refute recovered_metrics.cache_degraded? + assert recovered_metrics.cache_degraded_duration_ms == 0 + GenServer.stop(manager_pid) end diff --git a/test/durable_server/object_store_retry_test.exs b/test/durable_server/object_store_retry_test.exs index 568db70..84f976b 100644 --- a/test/durable_server/object_store_retry_test.exs +++ b/test/durable_server/object_store_retry_test.exs @@ -34,6 +34,70 @@ defmodule DurableServer.ObjectStoreRetryTest do assert Process.get(responses_key) == [:unexpected_retry] end + test "put reports authentication responses as non-retryable" do + parent = self() + responses_key = make_ref() + Process.put(responses_key, [%Req.Response{status: 403}, :unexpected_retry]) + + observer = fn _request, response_or_error, retryable? -> + send(parent, {:attempt, response_or_error, retryable?}) + end + + assert {:error, %Req.Response{status: 403}} = + ObjectStore.put_object( + store(adapter(responses_key)), + "__nodes/test@localhost", + "heartbeat", + max_retries: 10, + timeout: 1_000, + retry_observer: observer + ) + + assert_receive {:attempt, %Req.Response{status: 403}, false} + assert Process.get(responses_key) == [:unexpected_retry] + end + + test "put retries every Req transport error variant within its deadline" do + adapter = + adapter([ + %Req.TransportError{reason: :enetunreach}, + %Req.TransportError{reason: {:tls_alert, :unknown_ca}}, + %Req.Response{status: 200, headers: %{"etag" => ["retry-etag"]}} + ]) + + assert {:ok, %{etag: "retry-etag", body: "heartbeat"}} = + ObjectStore.put_object(store(adapter), "__nodes/test@localhost", "heartbeat", + max_retries: 10, + timeout: 1_000 + ) + end + + test "put retry observer receives every classified attempt" do + parent = self() + + adapter = + adapter([ + %Req.Response{status: 503}, + %Req.TransportError{reason: :enetunreach}, + %Req.Response{status: 200, headers: %{"etag" => ["retry-etag"]}} + ]) + + observer = fn _request, response_or_error, retryable? -> + send(parent, {:attempt, response_or_error, retryable?}) + end + + assert {:ok, %{etag: "retry-etag"}} = + ObjectStore.put_object(store(adapter), "__nodes/test@localhost", "heartbeat", + max_retries: 10, + timeout: 1_000, + retry_observer: observer + ) + + assert_receive {:attempt, %Req.Response{status: 503}, true} + assert_receive {:attempt, %Req.TransportError{reason: :enetunreach}, true} + assert_receive {:attempt, %Req.Response{status: 200}, false} + end + test "put stops after the configured transient retry limit" do responses_key = make_ref() From 5b1592ab86cbffa5c197f90c74652ac859ba3196 Mon Sep 17 00:00:00 2001 From: Jason Stiebs Date: Thu, 3 Sep 2026 09:58:44 -0500 Subject: [PATCH 2/2] Block restart claims during degraded discovery --- CHANGELOG.md | 2 +- lib/durable_server.ex | 3 ++ test/durable_server/lifecycle_test.exs | 62 ++++++++++++++++++++++++++ 3 files changed, 66 insertions(+), 1 deletion(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 780b845..c2f6c30 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,5 +1,5 @@ ## Unreleased -- Add bounded heartbeat telemetry and live metrics for attempt outcomes, consecutive failures, last-success age, watchdog budget, cache-degraded duration, watchdog terminations, and fenced child counts. Heartbeat writes now retry every ambiguous `Req.TransportError` variant within the hard deadline while permanent HTTP authentication/configuration failures still fail immediately. +- Add bounded heartbeat telemetry and live metrics for attempt outcomes, consecutive failures, last-success age, watchdog budget, cache-degraded duration, watchdog terminations, and fenced child counts. Heartbeat writes now retry every ambiguous `Req.TransportError` variant within the hard deadline while permanent HTTP authentication/configuration failures still fail immediately. Restart claims recheck the discovery-degraded gate at the storage mutation boundary so a retained stale snapshot cannot authorize an orphan claim. ## 0.1.5 (2026-08-27) - Treat explicit `:sync`, `{:sync, metadata}`, and `sync: true` callback returns as strict durability boundaries. Built-in backends first exhaust their bounded transient retry policy; if the write still fails, the DurableServer exits with a structured `{:sync_failed, reason}` fatal-exit reason before acknowledging the callback. Automatic and periodic sync remain best effort for transient failures, while storage conflicts remain fatal. diff --git a/lib/durable_server.ex b/lib/durable_server.ex index f4fc530..82677bb 100644 --- a/lib/durable_server.ex +++ b/lib/durable_server.ex @@ -1953,6 +1953,9 @@ defmodule DurableServer do storage_key = stored_state.prefix <> stored_state.key cond do + LifecycleManager.discovery_degraded?(meta.supervisor) -> + {:error, :discovery_degraded} + Meta.currently_restarting?(meta) -> {:error, :already_claimed} diff --git a/test/durable_server/lifecycle_test.exs b/test/durable_server/lifecycle_test.exs index 11ebd14..714914b 100644 --- a/test/durable_server/lifecycle_test.exs +++ b/test/durable_server/lifecycle_test.exs @@ -1137,6 +1137,68 @@ defmodule DurableServer.LifecycleTest do DurableServer.claim_restart_attempt(backend, stored_state, ttl: 10_000) end + test "degraded discovery rejects stale orphan claims at the storage mutation boundary", %{ + supervisor_name: supervisor_name, + prefix: prefix + } do + %{ets_table: coordination_table, storage_backend: backend} = + DurableServer.Supervisor.__get_config__(supervisor_name) + + key = "degraded-stale-claim-test-#{DurableServer.UUID.uuid4()}" + now = System.system_time(:millisecond) + + stored_state = %DurableServer.StoredState{ + vsn: 1, + state: %{"count" => 1}, + meta: %Meta{ + key: key, + prefix: prefix, + supervisor: supervisor_name, + module: TestServer, + permanent: true, + status: :running, + node_str: "remote@test", + node_ref: System.unique_integer([:positive]), + pid: self(), + last_heartbeat_at: now - 60_000 + } + } + + assert {:ok, _} = + DurableServer.StorageBackend.put_object( + backend, + "#{prefix}#{key}", + stored_state + ) + + assert {:ok, %DurableServer.StoredState{} = stored_state} = + DurableServer.fetch_stored_state( + backend, + %{key: key, prefix: prefix} + ) + + :ets.insert(coordination_table, {:discovery_degraded, true}) + + assert {:error, :discovery_degraded} = + DurableServer.claim_restart_attempt(backend, stored_state, ttl: 10_000) + + assert {:error, :discovery_degraded} = + DurableServer.claim_restart_attempt_with_verified_expired_lock( + backend, + stored_state, + ttl: 10_000 + ) + + assert {:ok, unchanged_state} = + DurableServer.fetch_stored_state( + backend, + %{key: key, prefix: prefix} + ) + + assert unchanged_state.etag == stored_state.etag + assert unchanged_state.meta.restart_attempt_node == nil + end + test "running restart claims revalidate heartbeat before claiming stale running meta", %{ supervisor_name: supervisor_name, prefix: prefix