diff --git a/lib/who_need_help/push/delivery_telemetry.ex b/lib/who_need_help/push/delivery_telemetry.ex new file mode 100644 index 0000000..97c2ef9 --- /dev/null +++ b/lib/who_need_help/push/delivery_telemetry.ex @@ -0,0 +1,23 @@ +defmodule WhoNeedHelp.Push.DeliveryTelemetry do + @moduledoc false + + @providers [:fcm, :web_push, :gateway] + @outcomes [ + :configuration_error, + :delivered, + :disabled, + :expired_device, + :invalid_device, + :invalid_notification, + :provider_rejected, + :retryable_error + ] + + def record(provider, outcome) when provider in @providers and outcome in @outcomes do + :telemetry.execute( + [:who_need_help, :push, :delivery], + %{count: 1}, + %{provider: Atom.to_string(provider), outcome: Atom.to_string(outcome)} + ) + end +end diff --git a/lib/who_need_help/push/delivery_worker.ex b/lib/who_need_help/push/delivery_worker.ex index ed6cd08..31a00d1 100644 --- a/lib/who_need_help/push/delivery_worker.ex +++ b/lib/who_need_help/push/delivery_worker.ex @@ -16,6 +16,7 @@ defmodule WhoNeedHelp.Push.DeliveryWorker do ] alias WhoNeedHelp.Push + alias WhoNeedHelp.Push.DeliveryTelemetry @impl Oban.Worker def perform(%Oban.Job{ @@ -37,21 +38,27 @@ defmodule WhoNeedHelp.Push.DeliveryWorker do case Push.deliver(notification, Push.delivery_options()) do {:ok, _receipt} -> + DeliveryTelemetry.record(:gateway, :delivered) :ok {:error, :disabled} -> + DeliveryTelemetry.record(:gateway, :disabled) {:cancel, :push_disabled} {:error, {:invalid_notification, _field} = reason} -> + DeliveryTelemetry.record(:gateway, :invalid_notification) {:cancel, reason} {:error, {:invalid_option, _option} = reason} -> + DeliveryTelemetry.record(:gateway, :configuration_error) {:cancel, reason} {:error, {:rejected, _status, _body} = reason} -> + DeliveryTelemetry.record(:gateway, :provider_rejected) {:cancel, reason} {:error, reason} -> + DeliveryTelemetry.record(:gateway, :retryable_error) {:error, reason} end end diff --git a/lib/who_need_help/push/device_delivery_worker.ex b/lib/who_need_help/push/device_delivery_worker.ex index 0a24ef0..d0dd8fd 100644 --- a/lib/who_need_help/push/device_delivery_worker.ex +++ b/lib/who_need_help/push/device_delivery_worker.ex @@ -13,6 +13,7 @@ defmodule WhoNeedHelp.Push.DeviceDeliveryWorker do alias WhoNeedHelp.Notifications alias WhoNeedHelp.Notifications.{Notification, PushDevice, Text} alias WhoNeedHelp.Push + alias WhoNeedHelp.Push.DeliveryTelemetry alias WhoNeedHelp.Repo @impl Oban.Worker @@ -51,24 +52,30 @@ defmodule WhoNeedHelp.Push.DeviceDeliveryWorker do Push.device_delivery_options(device.provider) ) do {:ok, _receipt} -> + DeliveryTelemetry.record(device.provider, :delivered) :ok {:error, :expired} -> + DeliveryTelemetry.record(device.provider, :expired_device) _ = Notifications.disable_invalid_device(device) :ok {:error, {:rejected, status, body}} -> if invalid_device_rejection?(device.provider, status, body) do + DeliveryTelemetry.record(device.provider, :invalid_device) _ = Notifications.disable_invalid_device(device) {:cancel, :device_rejected} else + DeliveryTelemetry.record(device.provider, :provider_rejected) {:cancel, {:provider_rejected, status}} end {:error, {:invalid_configuration, _field} = reason} -> + DeliveryTelemetry.record(device.provider, :configuration_error) {:cancel, reason} {:error, reason} -> + DeliveryTelemetry.record(device.provider, :retryable_error) {:error, reason} end end diff --git a/lib/who_need_help/trust/rate_limiter.ex b/lib/who_need_help/trust/rate_limiter.ex index c633d87..d5d221a 100644 --- a/lib/who_need_help/trust/rate_limiter.ex +++ b/lib/who_need_help/trust/rate_limiter.ex @@ -134,6 +134,8 @@ defmodule WhoNeedHelp.Trust.RateLimiter do } end) + emit_telemetry(results) + if Enum.all?(results, &(&1.count <= &1.limit)) do {:ok, results} else @@ -143,4 +145,16 @@ defmodule WhoNeedHelp.Trust.RateLimiter do defp normalize_action(action) when is_atom(action), do: Atom.to_string(action) defp normalize_action(action) when is_binary(action), do: action + + defp emit_telemetry(results) do + Enum.each(results, fn result -> + outcome = if result.count <= result.limit, do: "allowed", else: "limited" + + :telemetry.execute( + [:who_need_help, :rate_limit, :check], + %{count: 1}, + %{action: result.action, outcome: outcome} + ) + end) + end end diff --git a/lib/who_need_help_web/telemetry.ex b/lib/who_need_help_web/telemetry.ex index 605466a..71e936a 100644 --- a/lib/who_need_help_web/telemetry.ex +++ b/lib/who_need_help_web/telemetry.ex @@ -157,6 +157,15 @@ defmodule WhoNeedHelpWeb.Telemetry do tags: [:queue], description: "Failed Oban job attempts" ), + counter("who_need_help.oban.jobs.outcomes.total", + event_name: [:oban, :job, :stop], + measurement: :duration, + tags: [:queue, :outcome], + tag_values: fn metadata -> + %{queue: metadata.queue, outcome: to_string(metadata.state)} + end, + description: "Terminal and snoozed Oban job outcomes by queue" + ), counter("who_need_help.email.deliveries.total", event_name: [:swoosh, :deliver, :stop], measurement: :duration, @@ -178,6 +187,20 @@ defmodule WhoNeedHelpWeb.Telemetry do description: "Email delivery attempts by fixed purpose and outcome without personal labels" ), + counter("who_need_help.push.delivery.outcomes.total", + event_name: [:who_need_help, :push, :delivery], + measurement: :count, + tags: [:provider, :outcome], + description: + "Push delivery outcomes by fixed provider and outcome without device or user labels" + ), + counter("who_need_help.rate_limit.checks.total", + event_name: [:who_need_help, :rate_limit, :check], + measurement: :count, + tags: [:action, :outcome], + description: + "Rate-limit bucket checks by configured action and allowed or limited outcome" + ), sum("who_need_help.oban.job.duration.microseconds.total", event_name: [:oban, :job, :stop], measurement: fn measurements -> diff --git a/ops/observability/grafana/dashboards/who-need-help-overview.json b/ops/observability/grafana/dashboards/who-need-help-overview.json index 988cab6..0a682db 100644 --- a/ops/observability/grafana/dashboards/who-need-help-overview.json +++ b/ops/observability/grafana/dashboards/who-need-help-overview.json @@ -293,8 +293,8 @@ "targets": [ { "editorMode": "code", - "expr": "sum by (queue) (rate(who_need_help_oban_jobs_completed_total{job=\"who-need-help-worker\"}[1m]))", - "legendFormat": "completed {{queue}}", + "expr": "sum by (queue, outcome) (rate(who_need_help_oban_jobs_outcomes_total{job=\"who-need-help-worker\"}[1m]))", + "legendFormat": "{{outcome}} {{queue}}", "range": true, "refId": "A" }, @@ -411,6 +411,48 @@ ], "title": "Transactional email attempts by purpose and outcome (1h)", "type": "timeseries" + }, + { + "datasource": { + "type": "prometheus", + "uid": "wnh-prometheus" + }, + "fieldConfig": { + "defaults": { + "decimals": 0, + "unit": "short" + }, + "overrides": [] + }, + "gridPos": { + "h": 8, + "w": 24, + "x": 0, + "y": 48 + }, + "id": 12, + "options": { + "legend": { + "displayMode": "table", + "placement": "bottom", + "showLegend": true + }, + "tooltip": { + "mode": "multi", + "sort": "desc" + } + }, + "targets": [ + { + "editorMode": "code", + "expr": "sum by (action, outcome) (increase(who_need_help_rate_limit_checks_total[1h]))", + "legendFormat": "{{action}} {{outcome}}", + "range": true, + "refId": "A" + } + ], + "title": "Rate-limit bucket outcomes by action (1h)", + "type": "timeseries" } ], "refresh": "5s", diff --git a/scripts/production-external-monitor.py b/scripts/production-external-monitor.py index 0b37888..3821b23 100755 --- a/scripts/production-external-monitor.py +++ b/scripts/production-external-monitor.py @@ -25,6 +25,7 @@ MONITORED_COUNTERS = { "who_need_help_email_deliveries_total", "who_need_help_email_delivery_exceptions_total", "who_need_help_email_by_kind_deliveries_total", + "who_need_help_push_delivery_outcomes_total", } MONITORED_EMAIL_KINDS = { @@ -41,6 +42,14 @@ MONITORED_EMAIL_KINDS = { "support_update", } +MONITORED_PUSH_FAILURE_OUTCOMES = { + "configuration_error", + "invalid_notification", + "provider_rejected", +} + +MONITORED_PUSH_PROVIDERS = {"fcm", "gateway", "web_push"} + PROMETHEUS_SAMPLE = re.compile( r"^(?P[a-zA-Z_:][a-zA-Z0-9_:]*)(?:\{(?P.*)\})?\s+" r"(?P[-+]?(?:[0-9]+(?:\.[0-9]*)?|\.[0-9]+)(?:[eE][-+]?[0-9]+)?)" @@ -128,6 +137,14 @@ def parse_monitored_counters(payload: str) -> dict[str, float]: key = f"{name}|kind={kind}|status={labels['status']}" elif name == "who_need_help_oban_jobs_failed_total": key = f"{name}|queue={labels.get('queue', 'unknown')}" + elif name == "who_need_help_push_delivery_outcomes_total": + provider = labels.get("provider") + outcome = labels.get("outcome") + if provider not in MONITORED_PUSH_PROVIDERS: + continue + if outcome not in MONITORED_PUSH_FAILURE_OUTCOMES: + continue + key = f"{name}|provider={provider}|outcome={outcome}" else: key = name diff --git a/test/scripts/production_external_monitor_test.py b/test/scripts/production_external_monitor_test.py index 847b388..d99dda8 100644 --- a/test/scripts/production_external_monitor_test.py +++ b/test/scripts/production_external_monitor_test.py @@ -120,6 +120,12 @@ who_need_help_email_by_kind_deliveries_total{kind="auth_login",status="error"} 2 who_need_help_email_by_kind_deliveries_total{kind="support_update",status="exception"} 1 who_need_help_email_by_kind_deliveries_total{kind="user@example.test",status="error"} 99 who_need_help_http_requests_total 999 +who_need_help_push_delivery_outcomes_total{provider="fcm",outcome="delivered"} 10 +who_need_help_push_delivery_outcomes_total{provider="fcm",outcome="invalid_device"} 2 +who_need_help_push_delivery_outcomes_total{provider="fcm",outcome="provider_rejected"} 3 +who_need_help_push_delivery_outcomes_total{provider="web_push",outcome="configuration_error"} 4 +who_need_help_push_delivery_outcomes_total{provider="gateway",outcome="invalid_notification"} 5 +who_need_help_push_delivery_outcomes_total{provider="unbounded-user-value",outcome="provider_rejected"} 99 """ self.assertEqual( @@ -132,6 +138,9 @@ who_need_help_http_requests_total 999 "who_need_help_email_delivery_exceptions_total": 1.0, "who_need_help_email_by_kind_deliveries_total|kind=auth_login|status=error": 2.0, "who_need_help_email_by_kind_deliveries_total|kind=support_update|status=exception": 1.0, + "who_need_help_push_delivery_outcomes_total|provider=fcm|outcome=provider_rejected": 3.0, + "who_need_help_push_delivery_outcomes_total|provider=web_push|outcome=configuration_error": 4.0, + "who_need_help_push_delivery_outcomes_total|provider=gateway|outcome=invalid_notification": 5.0, }, ) diff --git a/test/who_need_help/push_product_test.exs b/test/who_need_help/push_product_test.exs index f7d1f7f..082a469 100644 --- a/test/who_need_help/push_product_test.exs +++ b/test/who_need_help/push_product_test.exs @@ -200,6 +200,7 @@ defmodule WhoNeedHelp.PushProductTest do end test "delivery worker reconstructs and sends the persisted notification" do + attach_push_outcomes() Application.put_env(:who_need_help, :push_adapter, RecordingAdapter) Application.put_env(:who_need_help, :push_delivery_options, test_pid: self()) @@ -213,6 +214,8 @@ defmodule WhoNeedHelp.PushProductTest do assert :ok = perform_job(DeliveryWorker, args) + assert_receive {:push_outcome, %{provider: "gateway", outcome: "delivered"}} + assert_received {:push_delivered, %{ idempotency_key: "message-created:event:user", @@ -289,6 +292,7 @@ defmodule WhoNeedHelp.PushProductTest do test "device delivery retries transient failures and only disables a rejected FID", context do + attach_push_outcomes() Application.put_env(:who_need_help, :fcm_adapter, RecordingDeviceAdapter) {:ok, device} = @@ -320,6 +324,7 @@ defmodule WhoNeedHelp.PushProductTest do }) assert_received {:device_delivery, notification_id, device_id} + assert_receive {:push_outcome, %{provider: "fcm", outcome: "retryable_error"}} assert notification_id == notification.id assert device_id == device.id assert is_nil(Repo.reload!(device).disabled_at) @@ -340,6 +345,8 @@ defmodule WhoNeedHelp.PushProductTest do "device_id" => device.id }) + assert_receive {:push_outcome, %{provider: "fcm", outcome: "provider_rejected"}} + assert is_nil(Repo.reload!(device).disabled_at) Application.put_env(:who_need_help, :device_delivery_options, %{ @@ -355,6 +362,8 @@ defmodule WhoNeedHelp.PushProductTest do "device_id" => device.id }) + assert_receive {:push_outcome, %{provider: "fcm", outcome: "provider_rejected"}} + assert is_nil(Repo.reload!(device).disabled_at) Application.put_env(:who_need_help, :device_delivery_options, %{ @@ -372,6 +381,23 @@ defmodule WhoNeedHelp.PushProductTest do "device_id" => device.id }) + assert_receive {:push_outcome, %{provider: "fcm", outcome: "invalid_device"}} + assert Repo.reload!(device).disabled_at end + + defp attach_push_outcomes do + handler_id = "push-outcomes-#{System.unique_integer([:positive])}" + test_pid = self() + + :ok = + :telemetry.attach( + handler_id, + [:who_need_help, :push, :delivery], + fn _event, _measurements, metadata, pid -> send(pid, {:push_outcome, metadata}) end, + test_pid + ) + + on_exit(fn -> :telemetry.detach(handler_id) end) + end end diff --git a/test/who_need_help/trust_safety_test.exs b/test/who_need_help/trust_safety_test.exs index 1c62943..fbb598c 100644 --- a/test/who_need_help/trust_safety_test.exs +++ b/test/who_need_help/trust_safety_test.exs @@ -470,6 +470,42 @@ defmodule WhoNeedHelp.TrustSafetyTest do ) == 4 end + test "rate-limit telemetry contains only the action and bounded outcome", context do + old = Application.get_env(:who_need_help, :rate_limit_policies) + + Application.put_env(:who_need_help, :rate_limit_policies, %{ + "telemetry_test" => %{"limit" => 1, "window_seconds" => 60} + }) + + on_exit(fn -> Application.put_env(:who_need_help, :rate_limit_policies, old) end) + + handler_id = "rate-limiter-telemetry-#{System.unique_integer([:positive])}" + test_process = self() + + :ok = + :telemetry.attach( + handler_id, + [:who_need_help, :rate_limit, :check], + fn _event, measurements, metadata, recipient -> + send(recipient, {:rate_limit_telemetry, measurements, metadata}) + end, + test_process + ) + + on_exit(fn -> :telemetry.detach(handler_id) end) + + assert {:ok, _limit} = RateLimiter.check(:telemetry_test, context.helper.email) + + assert_receive {:rate_limit_telemetry, %{count: 1}, + %{action: "telemetry_test", outcome: "allowed"}} + + assert {:error, :rate_limited} = + RateLimiter.check(:telemetry_test, context.helper.email) + + assert_receive {:rate_limit_telemetry, %{count: 1}, + %{action: "telemetry_test", outcome: "limited"}} + end + test "multiple limiter scopes use one PostgreSQL statement", context do old = Application.get_env(:who_need_help, :rate_limit_policies) diff --git a/test/who_need_help_web/controllers/metrics_controller_test.exs b/test/who_need_help_web/controllers/metrics_controller_test.exs index a367507..b14b8a3 100644 --- a/test/who_need_help_web/controllers/metrics_controller_test.exs +++ b/test/who_need_help_web/controllers/metrics_controller_test.exs @@ -110,4 +110,87 @@ defmodule WhoNeedHelpWeb.MetricsControllerTest do assert metric.measurement.(%{queue_time: 100, total_time: 100}) == 0 end + + test "exports Oban stop outcomes without job identifiers" do + reporter_name = :oban_outcome_metrics_controller_test + + metrics = + WhoNeedHelpWeb.Telemetry.prometheus_metrics() + |> Enum.filter(fn metric -> + metric.name == [:who_need_help, :oban, :jobs, :outcomes, :total] + end) + + start_supervised!( + {TelemetryMetricsPrometheus.Core, name: reporter_name, metrics: metrics, start_async: false} + ) + + :telemetry.execute( + [:oban, :job, :stop], + %{duration: 100, queue_time: 10}, + %{queue: "mail", state: :cancelled, job: %{id: 123, args: %{"recipient" => "private"}}} + ) + + body = TelemetryMetricsPrometheus.Core.scrape(reporter_name) + + assert body =~ + ~s(who_need_help_oban_jobs_outcomes_total{outcome="cancelled",queue="mail"} 1) + + refute body =~ "recipient" + refute body =~ "job_id" + refute body =~ "private" + end + + test "exports push outcomes with fixed provider and outcome labels only" do + reporter_name = :push_delivery_metrics_controller_test + + metrics = + WhoNeedHelpWeb.Telemetry.prometheus_metrics() + |> Enum.filter(fn metric -> + metric.name == [:who_need_help, :push, :delivery, :outcomes, :total] + end) + + start_supervised!( + {TelemetryMetricsPrometheus.Core, name: reporter_name, metrics: metrics, start_async: false} + ) + + WhoNeedHelp.Push.DeliveryTelemetry.record(:fcm, :provider_rejected) + body = TelemetryMetricsPrometheus.Core.scrape(reporter_name) + + assert body =~ + ~s(who_need_help_push_delivery_outcomes_total{outcome="provider_rejected",provider="fcm"} 1) + + refute body =~ "device_id" + refute body =~ "notification_id" + refute body =~ "user_id" + end + + test "exports aggregate rate-limit outcomes without scope labels" do + reporter_name = :rate_limit_metrics_controller_test + + metrics = + WhoNeedHelpWeb.Telemetry.prometheus_metrics() + |> Enum.filter(fn metric -> + metric.name == [:who_need_help, :rate_limit, :checks, :total] + end) + + start_supervised!( + {TelemetryMetricsPrometheus.Core, name: reporter_name, metrics: metrics, start_async: false} + ) + + :telemetry.execute( + [:who_need_help, :rate_limit, :check], + %{count: 1}, + %{action: "magic_link_email", outcome: "limited"} + ) + + body = TelemetryMetricsPrometheus.Core.scrape(reporter_name) + + assert body =~ + ~s(who_need_help_rate_limit_checks_total{action="magic_link_email",outcome="limited"} 1) + + refute body =~ "scope" + refute body =~ "email_address" + refute body =~ "ip_address" + refute body =~ "user_id" + end end