Packages

phoenix_kit

1.7.71
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
phoenix_kit lib modules newsletters broadcaster.ex
Raw

lib/modules/newsletters/broadcaster.ex

defmodule PhoenixKit.Modules.Newsletters.Broadcaster do
@moduledoc """
Orchestrates broadcast sending: paginates list members, creates Delivery
records and Oban jobs in batches.
"""
require Logger
import Ecto.Query
alias PhoenixKit.Modules.Newsletters
alias PhoenixKit.Modules.Newsletters.{Broadcast, Delivery, ListMember}
alias PhoenixKit.Modules.Newsletters.Workers.DeliveryWorker
alias PhoenixKit.Utils.Date, as: UtilsDate
@batch_size 500
@doc """
Starts sending a broadcast. Transitions status to `sending`,
creates delivery records, and enqueues Oban jobs.
"""
def send(%Broadcast{status: "draft"} = broadcast) do
do_send(broadcast)
end
def send(%Broadcast{status: "scheduled"} = broadcast) do
do_send(broadcast)
end
def send(%Broadcast{status: status}) do
{:error, {:invalid_status, status}}
end
defp do_send(broadcast) do
repo = repo()
# Render markdown to HTML before sending
html =
case Earmark.as_html(broadcast.markdown_body || "") do
{:ok, html, _warnings} -> html
{:error, html, _errors} -> html
end
text = strip_html(html)
{:ok, broadcast} =
Newsletters.update_broadcast(broadcast, %{
status: "sending",
html_body: html,
text_body: text,
sent_at: UtilsDate.utc_now()
})
# Count total active members
total = Newsletters.count_active_members(broadcast.list_uuid)
{:ok, broadcast} = Newsletters.update_broadcast(broadcast, %{total_recipients: total})
# Process in batches using transaction-wrapped stream
repo.transaction(fn ->
stream_active_members(broadcast.list_uuid)
|> Stream.chunk_every(@batch_size)
|> Enum.each(fn batch ->
process_batch(broadcast, batch, repo)
end)
end)
Logger.info("Broadcaster: Enqueued #{total} deliveries for broadcast #{broadcast.uuid}")
{:ok, broadcast}
end
defp stream_active_members(list_uuid) do
ListMember
|> where([m], m.list_uuid == ^list_uuid and m.status == "active")
|> select([m], m.user_uuid)
|> repo().stream()
end
defp process_batch(broadcast, user_uuids, repo) do
now = UtilsDate.utc_now()
deliveries =
Enum.map(user_uuids, fn user_uuid ->
%{
uuid: UUIDv7.generate(),
broadcast_uuid: broadcast.uuid,
user_uuid: user_uuid,
status: "pending",
inserted_at: now,
updated_at: now
}
end)
{_count, inserted} = repo.insert_all(Delivery, deliveries, returning: [:uuid])
jobs =
Enum.map(inserted, fn %{uuid: delivery_uuid} ->
DeliveryWorker.new(%{
delivery_uuid: delivery_uuid,
broadcast_uuid: broadcast.uuid
})
end)
Oban.insert_all(jobs)
end
defp strip_html(html) do
html
|> String.replace(~r/<br\s*\/?>/, "\n")
|> String.replace(~r/<\/p>/, "\n\n")
|> String.replace(~r/<[^>]+>/, "")
|> String.trim()
end
defp repo, do: PhoenixKit.RepoHelper.repo()
end