Packages
phoenix_kit
1.7.214
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/digest_worker.ex
defmodule PhoenixKit.Notifications.DigestWorker do
@moduledoc """
Aggregation engine: sends batched notification **summaries** for types a user
put on a digest cadence (hourly / 12h / daily / weekly) instead of getting
pinged for each event.
One cron entry per cadence enqueues this worker with `%{"cadence" => "hourly"}`
etc. The window is simply the cadence period (the hourly cron counts the last
hour, daily the last day, …) — so there's **no per-user last-sent state** to
track. For each user who routes a type to a channel on this cadence, it counts
that type's events in the window from the activity log and, if any, delivers a
single summary ("You have 1,432 new likes in the last hour.") via the channel.
Only activity-driven types aggregate — standalone notifications are always
immediate. Immediate cadences never reach here (they enqueue `DeliveryWorker`
at creation time).
"""
use Oban.Worker, queue: :notifications, max_attempts: 3
import Ecto.Query
require Logger
alias PhoenixKit.Activity.Entry
alias PhoenixKit.Notifications
alias PhoenixKit.Notifications.ChannelConfig
alias PhoenixKit.Notifications.Channels
alias PhoenixKit.Notifications.Prefs
alias PhoenixKit.Notifications.Types
alias PhoenixKit.RepoHelper
alias PhoenixKit.Users.Auth.User
alias PhoenixKit.Utils.Date, as: UtilsDate
alias PhoenixKit.Utils.Routes
@windows %{"hourly" => 3600, "12h" => 43_200, "daily" => 86_400, "weekly" => 604_800}
@impl Oban.Worker
def perform(%Oban.Job{args: %{"cadence" => cadence}}) when is_map_key(@windows, cadence) do
# Honor the global kill switch — same gate immediate notifications pass.
if Notifications.enabled?() do
since = DateTime.add(UtilsDate.utc_now(), -@windows[cadence], :second)
cadence
|> users_with_cadence()
|> Enum.each(&digest_user(&1, cadence, since))
end
:ok
end
def perform(_job), do: :ok
# --- Per-user sweep -------------------------------------------------------
defp digest_user(user, cadence, since) do
Enum.each(ChannelConfig.all_for(user), fn {channel_key, config} ->
digest_config(channel_key, config, user, cadence, since)
end)
rescue
e -> Logger.warning("[DigestWorker] user #{user.uuid} failed: #{Exception.message(e)}")
end
# In-app is a delivery mode too: post one aggregated inbox row.
defp digest_config("inapp", config, user, cadence, since),
do: digest_inapp(user, config, cadence, since)
defp digest_config(channel_key, config, user, cadence, since) do
channel = Channels.get(channel_key)
if deliverable_channel?(channel, config, user.uuid) do
for type_key <- cadenced_types(config, cadence) do
digest_type(user, channel, config, type_key, cadence, since)
end
end
end
defp deliverable_channel?(nil, _config, _uuid), do: false
defp deliverable_channel?(channel, config, uuid),
do: ChannelConfig.enabled?(config) and channel.configured?(uuid, config)
# In-app digest: for each type the user wants in-app (routing lives in
# `notification_preferences`) that's on this cadence, post ONE summary inbox
# row rather than the per-event rows the creation path suppressed.
defp digest_inapp(user, config, cadence, since) do
for type_key <- inapp_cadenced_types(user, config, cadence) do
count = count_events(user.uuid, type_key, since)
if count > 0 do
Notifications.create_inapp(user.uuid, %{
text: digest_text(count, type_label(type_key), cadence),
icon: "hero-bell",
link: Routes.path("/admin/notifications")
})
end
end
end
# Types routed to this channel on exactly this cadence.
defp cadenced_types(config, cadence) do
config
|> Map.get("types", %{})
|> Enum.filter(fn {type_key, on} ->
on == true and ChannelConfig.cadence(config, type_key) == cadence
end)
|> Enum.map(&elem(&1, 0))
end
# In-app types on this cadence — routing is checked against the user's in-app
# prefs (not a `types` map, which the inapp config doesn't carry).
defp inapp_cadenced_types(user, config, cadence) do
config
|> Map.get("cadences", %{})
|> Enum.filter(fn {type_key, c} ->
c == cadence and Prefs.user_wants_type?(user, type_key)
end)
|> Enum.map(&elem(&1, 0))
end
defp digest_type(user, channel, config, type_key, cadence, since) do
count = count_events(user.uuid, type_key, since)
if count > 0 do
envelope = digest_envelope(user, type_key, count, cadence)
case channel.deliver(envelope, config) do
:ok -> :ok
# A failed digest just waits for the next window — no retry state to keep.
other -> Logger.info("[DigestWorker] #{channel.key()} digest deferred: #{inspect(other)}")
end
end
end
# --- Data -----------------------------------------------------------------
# Users who have any channel config key set — filtered in Elixir to those
# actually using this cadence. (Configs live in the user JSONB; there's no
# separate table to index.)
defp users_with_cadence(cadence) do
# "inapp" is a synthetic mode, not a registered channel, but its cadence
# lives under the same `notification_channel:` prefix.
keys = Enum.map(["inapp" | Channels.keys()], &ChannelConfig.config_key/1)
condition =
Enum.reduce(keys, dynamic(false), fn key, acc ->
dynamic([u], ^acc or fragment("? \\? ?", u.custom_fields, ^key))
end)
from(u in User, where: ^condition)
|> repo().all()
|> Enum.filter(fn user ->
user
|> ChannelConfig.all_for()
|> Enum.any?(fn
{"inapp", config} -> inapp_cadenced_types(user, config, cadence) != []
{_ck, config} -> cadenced_types(config, cadence) != []
end)
end)
end
defp count_events(user_uuid, type_key, since) do
case actions_for_type(type_key) do
[] ->
0
actions ->
# Exclude self-actions (actor == target) to match the immediate path,
# which never notifies a user about their own action.
from(e in Entry,
where:
e.target_uuid == ^user_uuid and e.action in ^actions and e.inserted_at >= ^since and
(is_nil(e.actor_uuid) or e.actor_uuid != e.target_uuid)
)
|> repo().aggregate(:count)
end
end
# The action strings a type (or sub-type) key claims — from the Types registry.
defp actions_for_type(type_key) do
Types.list()
|> Enum.flat_map(fn type ->
[{type.key, type.actions} | Enum.map(type.sub_types, &{&1.key, &1.actions})]
end)
|> Enum.find_value([], fn {key, actions} -> if key == type_key, do: actions end)
end
defp digest_envelope(user, type_key, count, cadence) do
%{
recipient_uuid: user.uuid,
type_key: type_key,
notification_uuid: nil,
locale: recipient_locale(user),
icon: "hero-bell",
title: nil,
text: digest_text(count, type_label(type_key), cadence),
url: Routes.base_url() <> Routes.path("/admin/notifications")
}
end
defp digest_text(count, label, cadence) do
Gettext.dgettext(
PhoenixKitWeb.Gettext,
"default",
"You have %{count} new %{label} %{period}.",
count: count,
label: label,
period: period_phrase(cadence)
)
end
defp period_phrase("hourly"),
do: Gettext.dgettext(PhoenixKitWeb.Gettext, "default", "in the last hour")
defp period_phrase("12h"),
do: Gettext.dgettext(PhoenixKitWeb.Gettext, "default", "in the last 12 hours")
defp period_phrase("daily"),
do: Gettext.dgettext(PhoenixKitWeb.Gettext, "default", "in the last day")
defp period_phrase("weekly"),
do: Gettext.dgettext(PhoenixKitWeb.Gettext, "default", "in the last week")
defp period_phrase(_), do: ""
defp type_label(type_key) do
Types.list()
|> Enum.flat_map(fn type ->
[{type.key, type.label} | Enum.map(type.sub_types, &{&1.key, &1.label})]
end)
|> Enum.find_value(type_key, fn {key, label} ->
if key == type_key, do: String.downcase(label)
end)
end
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
defp repo, do: RepoHelper.repo()
end