Packages
phoenix_kit
1.7.216
1.7.216
1.7.215
1.7.214
1.7.213
1.7.212
1.7.211
1.7.210
1.7.209
1.7.208
1.7.207
1.7.206
1.7.205
1.7.204
1.7.203
1.7.202
1.7.201
1.7.200
1.7.199
1.7.198
1.7.197
1.7.196
1.7.194
1.7.193
1.7.192
1.7.191
1.7.190
1.7.189
1.7.187
1.7.186
1.7.185
1.7.184
1.7.183
1.7.182
1.7.181
1.7.180
1.7.179
1.7.178
1.7.177
1.7.176
1.7.175
1.7.174
1.7.173
1.7.172
1.7.171
1.7.170
1.7.169
1.7.168
1.7.167
1.7.166
1.7.165
1.7.164
1.7.162
1.7.161
1.7.160
1.7.159
1.7.157
1.7.156
1.7.155
1.7.154
1.7.153
1.7.152
1.7.151
1.7.150
1.7.149
1.7.146
1.7.145
1.7.144
1.7.143
1.7.138
1.7.133
1.7.132
1.7.131
1.7.130
1.7.128
1.7.126
1.7.125
1.7.121
1.7.120
1.7.119
1.7.118
1.7.117
1.7.116
1.7.115
1.7.114
1.7.113
1.7.112
1.7.111
1.7.110
1.7.109
1.7.108
1.7.107
1.7.106
1.7.105
1.7.104
1.7.103
1.7.102
1.7.101
1.7.100
1.7.99
1.7.98
1.7.97
1.7.96
1.7.95
1.7.94
1.7.93
1.7.92
1.7.91
1.7.90
1.7.89
1.7.88
1.7.87
1.7.86
1.7.85
1.7.84
1.7.83
1.7.82
1.7.81
1.7.80
1.7.79
1.7.78
1.7.77
1.7.76
1.7.75
1.7.74
1.7.71
1.7.70
1.7.69
1.7.66
1.7.65
1.7.64
1.7.63
1.7.62
1.7.61
1.7.59
1.7.58
1.7.57
1.7.56
1.7.55
1.7.54
1.7.53
1.7.52
1.7.51
1.7.49
1.7.44
1.7.43
1.7.42
1.7.41
1.7.39
1.7.38
1.7.37
1.7.36
1.7.34
1.7.33
1.7.31
1.7.30
1.7.29
1.7.28
1.7.27
1.7.26
1.7.25
1.7.24
1.7.23
1.7.22
1.7.21
1.7.20
1.7.19
1.7.18
1.7.17
1.7.16
1.7.15
1.7.14
1.7.13
1.7.12
1.7.11
1.7.10
1.7.9
1.7.8
1.7.7
1.7.6
1.7.5
1.7.4
1.7.3
1.7.2
1.7.1
1.7.0
1.6.20
1.6.19
1.6.18
1.6.17
1.6.16
1.6.15
1.6.14
1.6.13
1.6.12
1.6.11
1.6.10
1.6.9
1.6.8
1.6.7
1.6.6
1.6.5
1.6.4
1.6.3
1.5.2
1.5.1
1.5.0
1.4.9
1.4.8
1.4.7
1.4.6
1.4.5
1.4.4
1.4.3
1.4.2
1.4.1
1.4.0
1.3.2
1.3.1
1.3.0
1.2.10
1.2.9
1.2.8
1.2.7
1.2.5
1.2.4
1.2.2
1.2.1
1.2.0
1.1.0
1.0.0
A foundation for building Elixir Phoenix apps — SaaS, social networks, ERP systems, marketplaces, and more
Current section
Files
Jump to
Current section
Files
lib/phoenix_kit/notifications/delivery_worker.ex
defmodule PhoenixKit.Notifications.DeliveryWorker do
@moduledoc """
Delivers one notification to one external channel (Telegram, …), off the hot
path. Enqueued transactionally from the notification creation choke point —
one job per `{source, channel}` — so a slow bot API call never blocks the
business action that logged the activity.
## Job args
%{"channel" => "telegram",
"recipient_uuid" => "<uuid>",
"type_key" => "posts.likes",
# exactly ONE source:
"activity_uuid" => "<uuid>" # activity-driven (renders from the Entry)
# or
"notification_uuid" => "<uuid>"} # standalone (renders from the row)
## Retry / failure policy
The channel returns a permanent/transient verdict:
* `:ok` — done.
* `{:error, {:transient, _}}` / `{:error, {:transient, _}, retry_ms}` — Oban
retries (with a snooze on `retry_ms`), up to `max_attempts`.
* `{:error, {:permanent, reason}}` — **soft-disable** the channel in the
user's config (`enabled: false`, `status`, `disabled_reason`) and stop.
A bot the user blocked shouldn't retry-storm.
A channel/user/source that's gone by run time is discarded (`:ok`), not
retried. The channel resolves its own credentials owner-scoped to the
recipient, so this worker never handles a token.
"""
use Oban.Worker,
queue: :notifications,
max_attempts: 5,
# One job per source+channel — dedupes a double-logged activity or a
# re-run of the creation path. (Retries of THIS job reuse the same row.)
unique: [keys: [:source_uuid, :channel], period: 300]
require Logger
alias PhoenixKit.Activity
alias PhoenixKit.Notifications
alias PhoenixKit.Notifications.ChannelConfig
alias PhoenixKit.Notifications.Channels
alias PhoenixKit.Notifications.Notification
alias PhoenixKit.Notifications.Render
alias PhoenixKit.Users.Auth
alias PhoenixKit.Utils.Routes
@doc """
Builds the Oban job changeset for a delivery, stamping the `source_uuid` used
by the unique key (whichever of activity/notification uuid is present). String
keys throughout (JSONB args). Callers pass this to `Oban.insert/*`.
"""
@spec build(map()) :: Ecto.Changeset.t()
def build(args) when is_map(args) do
args = Map.new(args, fn {k, v} -> {to_string(k), v} end)
source_uuid = args["activity_uuid"] || args["notification_uuid"]
args
|> Map.put("source_uuid", source_uuid)
|> new()
end
@impl Oban.Worker
def perform(%Oban.Job{args: %{"channel" => channel_key, "recipient_uuid" => recipient} = args}) do
type_key = args["type_key"]
with {:ok, channel} <- fetch_channel(channel_key),
{:ok, user} <- fetch_user(recipient),
config = ChannelConfig.for_channel(user, channel_key),
true <- deliverable?(channel, user, config, type_key),
{:ok, notification} <- load_source(args) do
envelope = build_envelope(notification, user, recipient, type_key)
handle(channel.deliver(envelope, config), user, channel_key)
else
# Channel/user/source gone, or the user turned the channel/type off
# between enqueue and run — discard, don't retry.
_ -> :ok
end
end
def perform(_job), do: :ok
# --- Delivery result handling ---------------------------------------------
defp handle(:ok, _user, _channel_key), do: :ok
defp handle({:error, {:transient, reason}}, _user, _channel_key),
do: {:error, reason}
defp handle({:error, {:transient, _reason}, retry_ms}, _user, _channel_key)
when is_integer(retry_ms),
do: {:snooze, max(1, div(retry_ms, 1000))}
defp handle({:error, {:permanent, reason}}, user, channel_key) do
soft_disable(user, channel_key, reason)
# Discard — a permanent failure must not consume retries.
:ok
end
defp handle(other, _user, _channel_key) do
Logger.warning("[DeliveryWorker] unexpected deliver/2 result: #{inspect(other)}")
:ok
end
# Turn the channel off + record why, so the user sees it in settings and we
# stop hammering a dead destination. `reason` is a channel-supplied
# description/atom (never a raw error carrying a token).
defp soft_disable(user, channel_key, reason) do
ChannelConfig.update(user, channel_key, fn config ->
config
|> Map.put("enabled", false)
|> Map.put("status", "disabled")
|> Map.put("disabled_reason", describe(reason))
end)
end
defp describe(reason) when is_binary(reason), do: reason
defp describe(reason), do: inspect(reason)
# --- Loading + rendering ---------------------------------------------------
defp fetch_channel(key) do
case Channels.get(key) do
nil -> :error
mod -> {:ok, mod}
end
end
defp fetch_user(uuid) do
case Auth.get_user(uuid) do
%{} = user -> {:ok, user}
_ -> :error
end
end
# Re-check the user's CURRENT config: honor a channel/type turned off after
# the job was enqueued (least-surprising — a disabled channel stops sending).
defp deliverable?(channel, user, config, type_key) do
ChannelConfig.enabled?(config) and
(is_nil(type_key) or ChannelConfig.type_enabled?(config, type_key)) and
channel.configured?(user.uuid, config)
end
# Standalone → load the persisted row (recipient-scoped). Activity-driven →
# load the Entry and wrap it in a transient notification so `Render` takes the
# activity path (the inbox row may not exist when in-app was muted).
defp load_source(%{"notification_uuid" => uuid, "recipient_uuid" => recipient})
when is_binary(uuid) do
case Notifications.get_notification(recipient, uuid) do
%Notification{} = n -> {:ok, n}
_ -> :error
end
end
defp load_source(%{"activity_uuid" => uuid, "recipient_uuid" => recipient})
when is_binary(uuid) do
case Activity.get_entry(uuid) do
nil ->
:error
entry ->
{:ok,
%Notification{
activity: entry,
recipient_uuid: recipient,
metadata: entry.metadata || %{}
}}
end
end
defp load_source(_), do: :error
defp build_envelope(notification, user, recipient, type_key) do
locale = recipient_locale(user)
rendered = Render.render(notification, locale)
%{
recipient_uuid: recipient,
type_key: type_key,
notification_uuid: Map.get(notification, :uuid),
locale: locale,
icon: rendered.icon,
title: nil,
text: rendered.text,
url: absolutize(rendered.link)
}
end
# The recipient's own locale for send-time rendering (in-app renders per
# viewer; external must resolve the recipient's). Stored full-dialect
# ("en-GB") in custom_fields; Gettext wants the base ("en").
defp recipient_locale(user) do
case get_in(user.custom_fields || %{}, ["preferred_locale"]) do
loc when is_binary(loc) and loc != "" -> loc |> String.split("-") |> hd()
_ -> nil
end
end
# A notification `link` is already url-/locale-prefixed (or nil). Prepend the
# bare base; never pass it back through `Routes.url/1` (double-prefix).
defp absolutize(nil), do: nil
defp absolutize(link) when is_binary(link), do: Routes.base_url() <> link
end