diff --git a/lib/microwaveprop/application.ex b/lib/microwaveprop/application.ex index 4e37933c..e87771e2 100644 --- a/lib/microwaveprop/application.ex +++ b/lib/microwaveprop/application.ex @@ -45,6 +45,11 @@ defmodule Microwaveprop.Application do opts = [strategy: :one_for_one, name: Microwaveprop.Supervisor] result = Supervisor.start_link(children, opts) + # Attach the Oban exception telemetry tap so per-worker recovery + # callbacks and structured logs fire on job failure. Must run after + # the Oban supervisor starts (above in `children`). + :ok = Microwaveprop.ObanErrorReporter.attach() + # Run migrations after Repo is started _ = Microwaveprop.Release.migrate() diff --git a/lib/microwaveprop/oban_error_reporter.ex b/lib/microwaveprop/oban_error_reporter.ex new file mode 100644 index 00000000..de87822f --- /dev/null +++ b/lib/microwaveprop/oban_error_reporter.ex @@ -0,0 +1,199 @@ +defmodule Microwaveprop.ObanErrorReporter do + @moduledoc """ + Structured logging + per-worker recovery tap for Oban job exceptions. + + `PromEx.Plugins.Oban` already counts job exceptions as a Prometheus + counter, but a counter alone can't tell us *which* worker failed, why, + or whether the job is about to be retried. This module closes that + gap by attaching a `:telemetry` handler to `[:oban, :job, :exception]` + that: + + * logs a single `Logger.error` line with structured metadata — + `worker`, `queue`, `attempt`, `max_attempts`, `retry_exhausted?`, + `kind`, `reason` class, `job_id`, and a safe-for-log args summary + (keys only, no values that might contain secrets or bodies); and + + * dispatches to opt-in per-worker callbacks for recovery actions: + + def on_permanent_failure(%Oban.Job{} = job), do: ... + def on_transient_failure(%Oban.Job{} = job), do: ... + + Callback failures are caught so a bug in one worker's recovery + path cannot detach the telemetry handler or crash another job's + logging. + + The reporter does not emit a second telemetry event — PromEx counts + come from Oban's own `[:oban, :job, :exception]` event, so emitting + anything here would double-count. + """ + + require Logger + + @handler_id "microwaveprop-oban-error-reporter" + @event [:oban, :job, :exception] + + @doc """ + Attach the reporter to Oban's exception telemetry event. + + Safe to call more than once — a pre-existing attachment is treated + as success. + """ + @spec attach() :: :ok + def attach do + case :telemetry.attach(@handler_id, @event, &__MODULE__.handle_event/4, %{}) do + :ok -> :ok + {:error, :already_exists} -> :ok + end + end + + @doc """ + Detach the reporter. Primarily useful for tests. + """ + @spec detach() :: :ok + def detach do + _ = :telemetry.detach(@handler_id) + :ok + end + + @doc false + def handle_event(@event, measurements, metadata, _config) do + do_handle(measurements, metadata) + rescue + error -> + # A bug in the reporter must never detach the handler. + Logger.error("ObanErrorReporter handler crashed: " <> Exception.format(:error, error, __STACKTRACE__)) + end + + defp do_handle(_measurements, metadata) do + ctx = context(metadata) + log_exception(ctx, metadata) + dispatch_callback(ctx, metadata) + :ok + end + + # ——— logging ———————————————————————————————————————————————————————— + + defp log_exception(ctx, metadata) do + kind = Map.get(metadata, :kind, :error) + reason = Map.get(metadata, :reason) + + Logger.error( + "Oban job exception " <> + "worker=#{ctx.worker} " <> + "queue=#{ctx.queue} " <> + "job_id=#{ctx.job_id} " <> + "attempt=#{ctx.attempt} " <> + "max_attempts=#{ctx.max_attempts} " <> + "retry_exhausted?=#{ctx.retry_exhausted?} " <> + "kind=#{kind} " <> + "reason_class=#{reason_class(reason)} " <> + "reason=#{format_reason(reason)} " <> + "args_keys=#{inspect(ctx.args_summary)}" + ) + end + + defp format_reason(%{__exception__: true} = ex) do + "#{inspect(ex.__struct__)}: #{Exception.message(ex)}" + end + + defp format_reason(other), do: inspect(other) + + defp reason_class(%{__exception__: true} = ex), do: inspect(ex.__struct__) + defp reason_class(other) when is_atom(other), do: inspect(other) + defp reason_class(_), do: "term" + + # ——— context extraction —————————————————————————————————————————————— + + defp context(%{job: job} = metadata) when is_map(job) do + attempt = Map.get(job, :attempt, 0) + max_attempts = Map.get(job, :max_attempts, 0) + worker = Map.get(job, :worker) || Map.get(metadata, :worker) || "unknown" + queue = Map.get(job, :queue) || Map.get(metadata, :queue) || "unknown" + + %{ + worker: worker, + queue: queue, + job_id: Map.get(job, :id, 0), + attempt: attempt, + max_attempts: max_attempts, + retry_exhausted?: retry_exhausted?(attempt, max_attempts), + args_summary: args_summary(Map.get(job, :args)) + } + end + + defp context(metadata) do + %{ + worker: Map.get(metadata, :worker, "unknown"), + queue: Map.get(metadata, :queue, "unknown"), + job_id: 0, + attempt: 0, + max_attempts: 0, + retry_exhausted?: false, + args_summary: [] + } + end + + defp retry_exhausted?(attempt, max) when is_integer(attempt) and is_integer(max) and max > 0, do: attempt >= max + + defp retry_exhausted?(_, _), do: false + + # Never log args *values* — they can contain bodies, tokens, or PII. + # A sorted key list is enough to identify "which shape of payload". + defp args_summary(args) when is_map(args) do + args |> Map.keys() |> Enum.map(&to_string/1) |> Enum.sort() + end + + defp args_summary(_), do: [] + + # ——— per-worker callbacks ——————————————————————————————————————————— + + defp dispatch_callback(%{retry_exhausted?: exhausted?} = ctx, metadata) do + job = Map.get(metadata, :job) + fun = if exhausted?, do: :on_permanent_failure, else: :on_transient_failure + + case worker_module(ctx.worker) do + nil -> + :ok + + mod -> + if function_exported?(mod, fun, 1) do + safely_invoke(mod, fun, job) + else + :ok + end + end + end + + # `Module.safe_concat/1` avoids creating an atom for unknown workers, + # and returns an existing atom otherwise. If the module is not loaded + # in this node (e.g. a test stub name) we get :error and skip. + defp worker_module(worker) when is_binary(worker) do + mod = Module.safe_concat([worker]) + (Code.ensure_loaded?(mod) && mod) || nil + rescue + ArgumentError -> nil + end + + defp worker_module(worker) when is_atom(worker) do + if Code.ensure_loaded?(worker), do: worker + end + + defp worker_module(_), do: nil + + defp safely_invoke(mod, fun, job) do + apply(mod, fun, [job]) + :ok + rescue + error -> + Logger.warning( + "ObanErrorReporter: #{inspect(mod)}.#{fun}/1 raised: " <> Exception.format(:error, error, __STACKTRACE__) + ) + + :ok + catch + kind, reason -> + Logger.warning("ObanErrorReporter: #{inspect(mod)}.#{fun}/1 #{kind}: #{inspect(reason)}") + + :ok + end +end diff --git a/test/microwaveprop/oban_error_reporter_test.exs b/test/microwaveprop/oban_error_reporter_test.exs new file mode 100644 index 00000000..c594b83f --- /dev/null +++ b/test/microwaveprop/oban_error_reporter_test.exs @@ -0,0 +1,264 @@ +defmodule Microwaveprop.ObanErrorReporterTest do + # async: false — test attaches/detaches a named :telemetry handler + # and toggles per-test callback modules in the process dictionary. + use ExUnit.Case, async: false + + import ExUnit.CaptureLog + + alias Microwaveprop.ObanErrorReporter + + # Worker stubs used to assert callback dispatch. Each records a message + # back to the test process so we can `assert_receive`. + + defmodule PermanentWorker do + @moduledoc false + def on_permanent_failure(job) do + send(:oban_error_reporter_test, {:permanent, job.id}) + :ok + end + end + + defmodule TransientWorker do + @moduledoc false + def on_transient_failure(job) do + send(:oban_error_reporter_test, {:transient, job.id}) + :ok + end + end + + defmodule BothWorker do + @moduledoc false + def on_permanent_failure(job) do + send(:oban_error_reporter_test, {:permanent, job.id}) + :ok + end + + def on_transient_failure(job) do + send(:oban_error_reporter_test, {:transient, job.id}) + :ok + end + end + + defmodule RaisingWorker do + @moduledoc false + def on_permanent_failure(_job), do: raise("callback boom") + def on_transient_failure(_job), do: raise("callback boom") + end + + defmodule SilentWorker do + @moduledoc false + # No callbacks — used to confirm the reporter still logs for workers + # that haven't opted in. + end + + setup do + Process.register(self(), :oban_error_reporter_test) + :ok = ObanErrorReporter.attach() + + on_exit(fn -> + ObanErrorReporter.detach() + end) + + :ok + end + + defp execute_exception(metadata, measurements \\ %{duration: 123}) do + :telemetry.execute([:oban, :job, :exception], measurements, metadata) + end + + defp build_job(overrides) do + defaults = %{ + id: 42, + args: %{"foo" => "bar"}, + queue: "propagation", + worker: "Microwaveprop.Workers.SomeWorker", + attempt: 1, + max_attempts: 3, + state: "executing" + } + + # Use a bare map (not Oban.Job struct) so tests don't depend on Oban + # compile-time internals. The reporter only reads fields. + Map.merge(defaults, overrides) + end + + defp build_metadata(overrides) do + job = build_job(Map.get(overrides, :job, %{})) + + Map.merge( + %{ + job: job, + worker: job.worker, + queue: job.queue, + kind: :error, + reason: %RuntimeError{message: "kaboom"}, + stacktrace: [{Mod, :fun, 1, [file: ~c"x.ex", line: 1]}] + }, + Map.delete(overrides, :job) + ) + end + + describe "attach/0" do + test "is idempotent — calling twice does not raise" do + assert :ok = ObanErrorReporter.attach() + end + end + + describe "handle_event/4 on [:oban, :job, :exception]" do + test "logs an error with structured metadata" do + meta = + build_metadata(%{ + job: %{ + id: 7, + worker: "Microwaveprop.Workers.SilentWorker", + queue: "weather", + attempt: 1, + max_attempts: 3 + } + }) + + log = capture_log(fn -> execute_exception(meta) end) + + assert log =~ "Oban job exception" + assert log =~ "worker=" + assert log =~ "queue=weather" + assert log =~ "attempt=1" + assert log =~ "max_attempts=3" + assert log =~ "retry_exhausted?=false" + assert log =~ "kind=error" + assert log =~ "RuntimeError" + assert log =~ "job_id=7" + end + + test "marks retry_exhausted? true when attempt >= max_attempts" do + meta = + build_metadata(%{ + job: %{ + id: 8, + worker: to_string(SilentWorker), + attempt: 3, + max_attempts: 3 + } + }) + + log = capture_log(fn -> execute_exception(meta) end) + assert log =~ "retry_exhausted?=true" + end + + test "invokes on_permanent_failure/1 when retry exhausted and callback defined" do + meta = + build_metadata(%{ + job: %{ + id: 101, + worker: to_string(PermanentWorker), + attempt: 5, + max_attempts: 5 + } + }) + + capture_log(fn -> execute_exception(meta) end) + + assert_receive {:permanent, 101}, 200 + end + + test "invokes on_transient_failure/1 on transient failure when callback defined" do + meta = + build_metadata(%{ + job: %{ + id: 202, + worker: to_string(TransientWorker), + attempt: 1, + max_attempts: 5 + } + }) + + capture_log(fn -> execute_exception(meta) end) + + assert_receive {:transient, 202}, 200 + end + + test "dispatches only the matching callback when both are defined" do + meta_transient = + build_metadata(%{ + job: %{ + id: 303, + worker: to_string(BothWorker), + attempt: 1, + max_attempts: 5 + } + }) + + capture_log(fn -> execute_exception(meta_transient) end) + + assert_receive {:transient, 303}, 200 + refute_receive {:permanent, 303}, 50 + + meta_permanent = + build_metadata(%{ + job: %{ + id: 404, + worker: to_string(BothWorker), + attempt: 5, + max_attempts: 5 + } + }) + + capture_log(fn -> execute_exception(meta_permanent) end) + + assert_receive {:permanent, 404}, 200 + refute_receive {:transient, 404}, 50 + end + + test "a raising callback does not propagate, and the log still lands" do + meta = + build_metadata(%{ + job: %{ + id: 505, + worker: to_string(RaisingWorker), + attempt: 5, + max_attempts: 5 + } + }) + + log = + capture_log(fn -> + # Must not raise. + execute_exception(meta) + end) + + assert log =~ "Oban job exception" + assert log =~ "job_id=505" + # Reporter surfaces the callback failure at warning level. + assert log =~ "on_permanent_failure" + end + + test "does not raise when worker module is unknown / not loaded" do + meta = + build_metadata(%{ + job: %{ + id: 606, + worker: "Nonexistent.Worker.Module.ZZZ", + attempt: 1, + max_attempts: 3 + } + }) + + log = capture_log(fn -> execute_exception(meta) end) + assert log =~ "job_id=606" + end + + test "does not raise when reason is a plain term and stacktrace is missing" do + meta = %{ + job: build_job(%{id: 707, worker: to_string(SilentWorker)}), + worker: to_string(SilentWorker), + queue: "propagation", + kind: :exit, + reason: :shutdown + } + + log = capture_log(fn -> execute_exception(meta) end) + assert log =~ "job_id=707" + assert log =~ "kind=exit" + end + end +end