Expose bounded delivery and rate-limit outcomes

This commit is contained in:
SimpleTest 2026-08-20 22:51:23 +03:00
parent 1f9c5d74d5
commit 5efb2c319c
11 changed files with 289 additions and 2 deletions

View File

@ -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

View File

@ -16,6 +16,7 @@ defmodule WhoNeedHelp.Push.DeliveryWorker do
] ]
alias WhoNeedHelp.Push alias WhoNeedHelp.Push
alias WhoNeedHelp.Push.DeliveryTelemetry
@impl Oban.Worker @impl Oban.Worker
def perform(%Oban.Job{ def perform(%Oban.Job{
@ -37,21 +38,27 @@ defmodule WhoNeedHelp.Push.DeliveryWorker do
case Push.deliver(notification, Push.delivery_options()) do case Push.deliver(notification, Push.delivery_options()) do
{:ok, _receipt} -> {:ok, _receipt} ->
DeliveryTelemetry.record(:gateway, :delivered)
:ok :ok
{:error, :disabled} -> {:error, :disabled} ->
DeliveryTelemetry.record(:gateway, :disabled)
{:cancel, :push_disabled} {:cancel, :push_disabled}
{:error, {:invalid_notification, _field} = reason} -> {:error, {:invalid_notification, _field} = reason} ->
DeliveryTelemetry.record(:gateway, :invalid_notification)
{:cancel, reason} {:cancel, reason}
{:error, {:invalid_option, _option} = reason} -> {:error, {:invalid_option, _option} = reason} ->
DeliveryTelemetry.record(:gateway, :configuration_error)
{:cancel, reason} {:cancel, reason}
{:error, {:rejected, _status, _body} = reason} -> {:error, {:rejected, _status, _body} = reason} ->
DeliveryTelemetry.record(:gateway, :provider_rejected)
{:cancel, reason} {:cancel, reason}
{:error, reason} -> {:error, reason} ->
DeliveryTelemetry.record(:gateway, :retryable_error)
{:error, reason} {:error, reason}
end end
end end

View File

@ -13,6 +13,7 @@ defmodule WhoNeedHelp.Push.DeviceDeliveryWorker do
alias WhoNeedHelp.Notifications alias WhoNeedHelp.Notifications
alias WhoNeedHelp.Notifications.{Notification, PushDevice, Text} alias WhoNeedHelp.Notifications.{Notification, PushDevice, Text}
alias WhoNeedHelp.Push alias WhoNeedHelp.Push
alias WhoNeedHelp.Push.DeliveryTelemetry
alias WhoNeedHelp.Repo alias WhoNeedHelp.Repo
@impl Oban.Worker @impl Oban.Worker
@ -51,24 +52,30 @@ defmodule WhoNeedHelp.Push.DeviceDeliveryWorker do
Push.device_delivery_options(device.provider) Push.device_delivery_options(device.provider)
) do ) do
{:ok, _receipt} -> {:ok, _receipt} ->
DeliveryTelemetry.record(device.provider, :delivered)
:ok :ok
{:error, :expired} -> {:error, :expired} ->
DeliveryTelemetry.record(device.provider, :expired_device)
_ = Notifications.disable_invalid_device(device) _ = Notifications.disable_invalid_device(device)
:ok :ok
{:error, {:rejected, status, body}} -> {:error, {:rejected, status, body}} ->
if invalid_device_rejection?(device.provider, status, body) do if invalid_device_rejection?(device.provider, status, body) do
DeliveryTelemetry.record(device.provider, :invalid_device)
_ = Notifications.disable_invalid_device(device) _ = Notifications.disable_invalid_device(device)
{:cancel, :device_rejected} {:cancel, :device_rejected}
else else
DeliveryTelemetry.record(device.provider, :provider_rejected)
{:cancel, {:provider_rejected, status}} {:cancel, {:provider_rejected, status}}
end end
{:error, {:invalid_configuration, _field} = reason} -> {:error, {:invalid_configuration, _field} = reason} ->
DeliveryTelemetry.record(device.provider, :configuration_error)
{:cancel, reason} {:cancel, reason}
{:error, reason} -> {:error, reason} ->
DeliveryTelemetry.record(device.provider, :retryable_error)
{:error, reason} {:error, reason}
end end
end end

View File

@ -134,6 +134,8 @@ defmodule WhoNeedHelp.Trust.RateLimiter do
} }
end) end)
emit_telemetry(results)
if Enum.all?(results, &(&1.count <= &1.limit)) do if Enum.all?(results, &(&1.count <= &1.limit)) do
{:ok, results} {:ok, results}
else 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_atom(action), do: Atom.to_string(action)
defp normalize_action(action) when is_binary(action), do: 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 end

View File

@ -157,6 +157,15 @@ defmodule WhoNeedHelpWeb.Telemetry do
tags: [:queue], tags: [:queue],
description: "Failed Oban job attempts" 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", counter("who_need_help.email.deliveries.total",
event_name: [:swoosh, :deliver, :stop], event_name: [:swoosh, :deliver, :stop],
measurement: :duration, measurement: :duration,
@ -178,6 +187,20 @@ defmodule WhoNeedHelpWeb.Telemetry do
description: description:
"Email delivery attempts by fixed purpose and outcome without personal labels" "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", sum("who_need_help.oban.job.duration.microseconds.total",
event_name: [:oban, :job, :stop], event_name: [:oban, :job, :stop],
measurement: fn measurements -> measurement: fn measurements ->

View File

@ -293,8 +293,8 @@
"targets": [ "targets": [
{ {
"editorMode": "code", "editorMode": "code",
"expr": "sum by (queue) (rate(who_need_help_oban_jobs_completed_total{job=\"who-need-help-worker\"}[1m]))", "expr": "sum by (queue, outcome) (rate(who_need_help_oban_jobs_outcomes_total{job=\"who-need-help-worker\"}[1m]))",
"legendFormat": "completed {{queue}}", "legendFormat": "{{outcome}} {{queue}}",
"range": true, "range": true,
"refId": "A" "refId": "A"
}, },
@ -411,6 +411,48 @@
], ],
"title": "Transactional email attempts by purpose and outcome (1h)", "title": "Transactional email attempts by purpose and outcome (1h)",
"type": "timeseries" "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", "refresh": "5s",

View File

@ -25,6 +25,7 @@ MONITORED_COUNTERS = {
"who_need_help_email_deliveries_total", "who_need_help_email_deliveries_total",
"who_need_help_email_delivery_exceptions_total", "who_need_help_email_delivery_exceptions_total",
"who_need_help_email_by_kind_deliveries_total", "who_need_help_email_by_kind_deliveries_total",
"who_need_help_push_delivery_outcomes_total",
} }
MONITORED_EMAIL_KINDS = { MONITORED_EMAIL_KINDS = {
@ -41,6 +42,14 @@ MONITORED_EMAIL_KINDS = {
"support_update", "support_update",
} }
MONITORED_PUSH_FAILURE_OUTCOMES = {
"configuration_error",
"invalid_notification",
"provider_rejected",
}
MONITORED_PUSH_PROVIDERS = {"fcm", "gateway", "web_push"}
PROMETHEUS_SAMPLE = re.compile( PROMETHEUS_SAMPLE = re.compile(
r"^(?P<name>[a-zA-Z_:][a-zA-Z0-9_:]*)(?:\{(?P<labels>.*)\})?\s+" r"^(?P<name>[a-zA-Z_:][a-zA-Z0-9_:]*)(?:\{(?P<labels>.*)\})?\s+"
r"(?P<value>[-+]?(?:[0-9]+(?:\.[0-9]*)?|\.[0-9]+)(?:[eE][-+]?[0-9]+)?)" r"(?P<value>[-+]?(?:[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']}" key = f"{name}|kind={kind}|status={labels['status']}"
elif name == "who_need_help_oban_jobs_failed_total": elif name == "who_need_help_oban_jobs_failed_total":
key = f"{name}|queue={labels.get('queue', 'unknown')}" 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: else:
key = name key = name

View File

@ -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="support_update",status="exception"} 1
who_need_help_email_by_kind_deliveries_total{kind="user@example.test",status="error"} 99 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_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( 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_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=auth_login|status=error": 2.0,
"who_need_help_email_by_kind_deliveries_total|kind=support_update|status=exception": 1.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,
}, },
) )

View File

@ -200,6 +200,7 @@ defmodule WhoNeedHelp.PushProductTest do
end end
test "delivery worker reconstructs and sends the persisted notification" do 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_adapter, RecordingAdapter)
Application.put_env(:who_need_help, :push_delivery_options, test_pid: self()) 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 :ok = perform_job(DeliveryWorker, args)
assert_receive {:push_outcome, %{provider: "gateway", outcome: "delivered"}}
assert_received {:push_delivered, assert_received {:push_delivered,
%{ %{
idempotency_key: "message-created:event:user", 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", test "device delivery retries transient failures and only disables a rejected FID",
context do context do
attach_push_outcomes()
Application.put_env(:who_need_help, :fcm_adapter, RecordingDeviceAdapter) Application.put_env(:who_need_help, :fcm_adapter, RecordingDeviceAdapter)
{:ok, device} = {:ok, device} =
@ -320,6 +324,7 @@ defmodule WhoNeedHelp.PushProductTest do
}) })
assert_received {:device_delivery, notification_id, device_id} assert_received {:device_delivery, notification_id, device_id}
assert_receive {:push_outcome, %{provider: "fcm", outcome: "retryable_error"}}
assert notification_id == notification.id assert notification_id == notification.id
assert device_id == device.id assert device_id == device.id
assert is_nil(Repo.reload!(device).disabled_at) assert is_nil(Repo.reload!(device).disabled_at)
@ -340,6 +345,8 @@ defmodule WhoNeedHelp.PushProductTest do
"device_id" => device.id "device_id" => device.id
}) })
assert_receive {:push_outcome, %{provider: "fcm", outcome: "provider_rejected"}}
assert is_nil(Repo.reload!(device).disabled_at) assert is_nil(Repo.reload!(device).disabled_at)
Application.put_env(:who_need_help, :device_delivery_options, %{ Application.put_env(:who_need_help, :device_delivery_options, %{
@ -355,6 +362,8 @@ defmodule WhoNeedHelp.PushProductTest do
"device_id" => device.id "device_id" => device.id
}) })
assert_receive {:push_outcome, %{provider: "fcm", outcome: "provider_rejected"}}
assert is_nil(Repo.reload!(device).disabled_at) assert is_nil(Repo.reload!(device).disabled_at)
Application.put_env(:who_need_help, :device_delivery_options, %{ Application.put_env(:who_need_help, :device_delivery_options, %{
@ -372,6 +381,23 @@ defmodule WhoNeedHelp.PushProductTest do
"device_id" => device.id "device_id" => device.id
}) })
assert_receive {:push_outcome, %{provider: "fcm", outcome: "invalid_device"}}
assert Repo.reload!(device).disabled_at assert Repo.reload!(device).disabled_at
end 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 end

View File

@ -470,6 +470,42 @@ defmodule WhoNeedHelp.TrustSafetyTest do
) == 4 ) == 4
end 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 test "multiple limiter scopes use one PostgreSQL statement", context do
old = Application.get_env(:who_need_help, :rate_limit_policies) old = Application.get_env(:who_need_help, :rate_limit_policies)

View File

@ -110,4 +110,87 @@ defmodule WhoNeedHelpWeb.MetricsControllerTest do
assert metric.measurement.(%{queue_time: 100, total_time: 100}) == 0 assert metric.measurement.(%{queue_time: 100, total_time: 100}) == 0
end 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 end