-module(telega@broadcast). -compile([no_auto_import, nowarn_unused_vars, nowarn_unused_function, nowarn_nomatch, inline]). -define(FILEPATH, "src/telega/broadcast.gleam"). -export([new/3, new_from_iterator/3, send_text/3, with_rate/3, with_on_progress/2, start/1, await/2, cancel/1, progress/1, run/1]). -export_type([broadcast/1, source/0, broadcast_progress/0, broadcast_report/1, broadcast_handle/1, msg/1, state/1]). -if(?OTP_RELEASE >= 27). -define(MODULEDOC(Str), -moduledoc(Str)). -define(DOC(Str), -doc(Str)). -else. -define(MODULEDOC(Str), -compile([])). -define(DOC(Str), -compile([])). -endif. ?MODULEDOC( " Mass messaging with pacing, result classification and reports —\n" " the answer to \"how do I send a message to all my users?\" without\n" " tripping over Telegram's rate limits or losing track of who\n" " actually got the message.\n" "\n" " A broadcast sends to a list (or a stream) of chat ids sequentially,\n" " paced below Telegram's limit, and classifies every result:\n" "\n" " - `sent` — delivered, with the value returned by the send function\n" " - `blocked` — HTTP 403: the user blocked the bot, was deactivated,\n" " or kicked the bot\n" " - `failed` — everything else after retries\n" "\n" " ```gleam\n" " import telega/broadcast\n" "\n" " let assert Ok(report) =\n" " broadcast.send_text(client:, chat_ids:, text: \"Big news!\")\n" " |> broadcast.run\n" "\n" " // 403s are users who blocked the bot — stop sending to them\n" " mark_as_dead(report.blocked)\n" " ```\n" "\n" " ## Telegram's limits\n" "\n" " Telegram allows bots roughly **30 messages per second** across all\n" " chats (and ~20 messages per minute into the same group). Exceeding\n" " it earns HTTP 429 responses and, if you keep pushing, longer and\n" " longer cooldowns.\n" "\n" " The broadcast default is **25 messages per 1000 ms** — a deliberate\n" " safety margin. Tune it with `with_rate`:\n" "\n" " ```gleam\n" " broadcast.send_text(client:, chat_ids:, text:)\n" " |> broadcast.with_rate(rate: 20, window_ms: 1000)\n" " |> broadcast.run\n" " ```\n" "\n" " If the client also has a request queue configured\n" " (`client.new_with_queue` / `client.set_request_queue`), broadcast\n" " calls go through it too, so the effective rate is the **min** of the\n" " two limits. The broadcast's own pacing exists so that mass sends are\n" " throttled even on clients without a queue — and so a broadcast never\n" " starves interactive traffic by monopolizing the queue's default rule.\n" "\n" " On 429: the client itself retries honoring `parameters.retry_after`.\n" " A 429 that still reaches the broadcast means Telegram is pushing back\n" " hard — the broadcast pauses for one full window and retries that\n" " chat id **once**, then reports it as `failed`.\n" "\n" " ## Custom payloads\n" "\n" " `send_text` is a convenience over `api.send_message`. For anything\n" " else — photos, invoices, per-user personalization — pass your own\n" " send function:\n" "\n" " ```gleam\n" " let send_promo = fn(client, chat_id) {\n" " api.send_photo(client, parameters: promo_photo_for(chat_id))\n" " }\n" "\n" " let assert Ok(report) =\n" " broadcast.new(client:, chat_ids:, send: send_promo)\n" " |> broadcast.run\n" " ```\n" "\n" " The function's success value ends up in `report.sent`, so you can\n" " keep the returned `Message` for later edits or deletion.\n" "\n" " Sends are sequential by design: one send at a time, inside the\n" " broadcast actor. Concurrency would break pacing.\n" "\n" " ## Streaming recipients from a database\n" "\n" " For large audiences, don't load every chat id into memory — stream\n" " them in chunks. The broadcast pulls the next chunk when the current\n" " one is exhausted; return `None` (or an empty chunk) to signal the\n" " end:\n" "\n" " ```gleam\n" " let next_page = fn() {\n" " case load_subscriber_page(db) {\n" " [] -> None\n" " chat_ids -> Some(chat_ids)\n" " }\n" " }\n" "\n" " let assert Ok(report) =\n" " broadcast.new_from_iterator(client:, next_chunk: next_page, send: send_promo)\n" " |> broadcast.run\n" " ```\n" "\n" " With an iterator source, `BroadcastProgress.total` is `None` — the\n" " size is unknown upfront.\n" "\n" " ## Background broadcasts: progress and cancellation\n" "\n" " `run` is fine for scripts. In a bot you usually want to start the\n" " broadcast, answer the admin immediately, and check on it later:\n" "\n" " ```gleam\n" " let assert Ok(handle) =\n" " broadcast.send_text(client:, chat_ids:, text:)\n" " |> broadcast.start\n" "\n" " // From any process, at any time:\n" " let progress = broadcast.progress(handle)\n" " broadcast.cancel(handle)\n" " let assert Ok(report) = broadcast.await(handle, timeout: 60_000)\n" " ```\n" "\n" " For live progress messages (\"Sending… 250/1000\"), register a\n" " callback with `with_on_progress`. It runs inside the broadcast actor\n" " after every processed chat id — keep it cheap, a slow callback slows\n" " the whole broadcast down.\n" "\n" " ## Blocked-user hygiene\n" "\n" " A 403 (`Forbidden: bot was blocked by the user` and friends) is\n" " permanent until the user comes back on their own. Every broadcast to\n" " a dead chat id wastes your rate budget, so treat `report.blocked` as\n" " a to-do list: mark those chat ids as inactive in your storage,\n" " exclude them from future broadcasts, and re-activate a user when\n" " they message the bot again (`/start`).\n" "\n" " `failed` is different — those are transient errors (network,\n" " server-side 5xx, a 429 that survived retries). Keep those ids and\n" " retry them in a later broadcast.\n" ). -opaque broadcast(AUEC) :: {broadcast, telega@client:telegram_client(), source(), fun((telega@client:telegram_client(), integer()) -> {ok, AUEC} | {error, telega@error:telega_error()}), integer(), integer(), gleam@option:option(fun((broadcast_progress()) -> nil))}. -type source() :: {chat_id_list, list(integer())} | {chunk_iterator, fun(() -> gleam@option:option(list(integer()))), list(integer())}. -type broadcast_progress() :: {broadcast_progress, gleam@option:option(integer()), integer(), integer(), integer(), integer()}. -type broadcast_report(AUED) :: {broadcast_report, list({integer(), AUED}), list(integer()), list({integer(), telega@error:telega_error()}), integer(), boolean()}. -opaque broadcast_handle(AUEE) :: {broadcast_handle, gleam@erlang@process:subject(msg(AUEE))}. -type msg(AUEF) :: send_next | cancel | {get_progress, gleam@erlang@process:subject(broadcast_progress())} | {await, gleam@erlang@process:subject(broadcast_report(AUEF))}. -type state(AUEG) :: {state, telega@client:telegram_client(), source(), fun((telega@client:telegram_client(), integer()) -> {ok, AUEG} | {error, telega@error:telega_error()}), integer(), integer(), gleam@option:option(fun((broadcast_progress()) -> nil)), gleam@option:option(integer()), gleam@option:option(integer()), list({integer(), AUEG}), list(integer()), list({integer(), telega@error:telega_error()}), integer(), integer(), integer(), list(gleam@erlang@process:subject(broadcast_report(AUEG))), gleam@option:option(broadcast_report(AUEG)), gleam@erlang@process:subject(msg(AUEG))}. -file("src/telega/broadcast.gleam", 208). ?DOC( " Create a broadcast for a known list of chat ids.\n" "\n" " The send function is called once per chat id (twice on a 429 that\n" " survived the client's retries) inside the broadcast actor.\n" ). -spec new( telega@client:telegram_client(), list(integer()), fun((telega@client:telegram_client(), integer()) -> {ok, AUEI} | {error, telega@error:telega_error()}) ) -> broadcast(AUEI). new(Client, Chat_ids, Send) -> {broadcast, Client, {chat_id_list, Chat_ids}, Send, 25, 1000, none}. -file("src/telega/broadcast.gleam", 230). ?DOC( " Create a broadcast that pulls chat ids in chunks — for streaming\n" " millions of recipients from a database without loading them all\n" " into memory.\n" "\n" " The next chunk is requested (inside the broadcast actor) when the\n" " current one is exhausted. Return `None` — or an empty chunk — to\n" " signal the end of the stream.\n" ). -spec new_from_iterator( telega@client:telegram_client(), fun(() -> gleam@option:option(list(integer()))), fun((telega@client:telegram_client(), integer()) -> {ok, AUEO} | {error, telega@error:telega_error()}) ) -> broadcast(AUEO). new_from_iterator(Client, Next_chunk, Send) -> {broadcast, Client, {chunk_iterator, Next_chunk, []}, Send, 25, 1000, none}. -file("src/telega/broadcast.gleam", 247). ?DOC( " Convenience broadcast sending the same text to every chat id\n" " via `sendMessage`.\n" ). -spec send_text(telega@client:telegram_client(), list(integer()), binary()) -> broadcast(telega@model@types:message()). send_text(Client, Chat_ids, Text) -> new( Client, Chat_ids, fun(Client@1, Chat_id) -> telega@api:send_message( Client@1, {send_message_parameters, none, {int, Chat_id}, none, Text, none, none, none, none, none, none, none, none, none} ) end ). -file("src/telega/broadcast.gleam", 276). ?DOC( " Set the pacing: at most `rate` sends per `window_ms` milliseconds.\n" " Default is 25 per 1000 ms. Values below 1 are clamped to 1.\n" ). -spec with_rate(broadcast(AUEU), integer(), integer()) -> broadcast(AUEU). with_rate(Broadcast, Rate, Window_ms) -> {broadcast, erlang:element(2, Broadcast), erlang:element(3, Broadcast), erlang:element(4, Broadcast), gleam@int:max(Rate, 1), gleam@int:max(Window_ms, 1), erlang:element(7, Broadcast)}. -file("src/telega/broadcast.gleam", 290). ?DOC( " Set a progress callback, called from the broadcast actor after every\n" " processed chat id. Keep it cheap — a slow callback slows the broadcast.\n" ). -spec with_on_progress(broadcast(AUEX), fun((broadcast_progress()) -> nil)) -> broadcast(AUEX). with_on_progress(Broadcast, On_progress) -> {broadcast, erlang:element(2, Broadcast), erlang:element(3, Broadcast), erlang:element(4, Broadcast), erlang:element(5, Broadcast), erlang:element(6, Broadcast), {some, On_progress}}. -file("src/telega/broadcast.gleam", 570). -spec progress_of(state(any())) -> broadcast_progress(). progress_of(State) -> Sent = erlang:length(erlang:element(10, State)), Blocked = erlang:length(erlang:element(11, State)), Failed = erlang:length(erlang:element(12, State)), {broadcast_progress, erlang:element(8, State), (Sent + Blocked) + Failed, Sent, Blocked, Failed}. -file("src/telega/broadcast.gleam", 553). -spec finish(state(AUGX), boolean()) -> gleam@otp@actor:next(state(AUGX), msg(AUGX)). finish(State, Cancelled) -> Report = {broadcast_report, lists:reverse(erlang:element(10, State)), lists:reverse(erlang:element(11, State)), lists:reverse(erlang:element(12, State)), telega@internal@utils:current_time_ms() - erlang:element(15, State), Cancelled}, gleam@list:each( erlang:element(16, State), fun(_capture) -> gleam@erlang@process:send(_capture, Report) end ), gleam@otp@actor:continue( {state, erlang:element(2, State), erlang:element(3, State), erlang:element(4, State), erlang:element(5, State), erlang:element(6, State), erlang:element(7, State), erlang:element(8, State), erlang:element(9, State), erlang:element(10, State), erlang:element(11, State), erlang:element(12, State), erlang:element(13, State), erlang:element(14, State), erlang:element(15, State), [], {some, Report}, erlang:element(18, State)} ). -file("src/telega/broadcast.gleam", 543). -spec record(state(AUGR)) -> gleam@otp@actor:next(state(AUGR), msg(AUGR)). record(State) -> case erlang:element(7, State) of {some, On_progress} -> On_progress(progress_of(State)); none -> nil end, gleam@erlang@process:send(erlang:element(18, State), send_next), gleam@otp@actor:continue(State). -file("src/telega/broadcast.gleam", 509). -spec send_to_chat(state(AUGL), integer()) -> gleam@otp@actor:next(state(AUGL), msg(AUGL)). send_to_chat(State, Chat_id) -> Is_retry = erlang:element(9, State) =:= {some, Chat_id}, State@1 = {state, erlang:element(2, State), erlang:element(3, State), erlang:element(4, State), erlang:element(5, State), erlang:element(6, State), erlang:element(7, State), erlang:element(8, State), none, erlang:element(10, State), erlang:element(11, State), erlang:element(12, State), erlang:element(13, State), erlang:element(14, State) + 1, erlang:element(15, State), erlang:element(16, State), erlang:element(17, State), erlang:element(18, State)}, case (erlang:element(4, State@1))(erlang:element(2, State@1), Chat_id) of {ok, Value} -> record( {state, erlang:element(2, State@1), erlang:element(3, State@1), erlang:element(4, State@1), erlang:element(5, State@1), erlang:element(6, State@1), erlang:element(7, State@1), erlang:element(8, State@1), erlang:element(9, State@1), [{Chat_id, Value} | erlang:element(10, State@1)], erlang:element(11, State@1), erlang:element(12, State@1), erlang:element(13, State@1), erlang:element(14, State@1), erlang:element(15, State@1), erlang:element(16, State@1), erlang:element(17, State@1), erlang:element(18, State@1)} ); {error, {telegram_api_error, 403, _}} -> record( {state, erlang:element(2, State@1), erlang:element(3, State@1), erlang:element(4, State@1), erlang:element(5, State@1), erlang:element(6, State@1), erlang:element(7, State@1), erlang:element(8, State@1), erlang:element(9, State@1), erlang:element(10, State@1), [Chat_id | erlang:element(11, State@1)], erlang:element(12, State@1), erlang:element(13, State@1), erlang:element(14, State@1), erlang:element(15, State@1), erlang:element(16, State@1), erlang:element(17, State@1), erlang:element(18, State@1)} ); {error, {telegram_api_error, 429, Description}} -> case Is_retry of false -> State@2 = {state, erlang:element(2, State@1), erlang:element(3, State@1), erlang:element(4, State@1), erlang:element(5, State@1), erlang:element(6, State@1), erlang:element(7, State@1), erlang:element(8, State@1), {some, Chat_id}, erlang:element(10, State@1), erlang:element(11, State@1), erlang:element(12, State@1), erlang:element(13, State@1), erlang:element(14, State@1), erlang:element(15, State@1), erlang:element(16, State@1), erlang:element(17, State@1), erlang:element(18, State@1)}, gleam@erlang@process:send_after( erlang:element(18, State@2), erlang:element(6, State@2), send_next ), gleam@otp@actor:continue(State@2); true -> record( {state, erlang:element(2, State@1), erlang:element(3, State@1), erlang:element(4, State@1), erlang:element(5, State@1), erlang:element(6, State@1), erlang:element(7, State@1), erlang:element(8, State@1), erlang:element(9, State@1), erlang:element(10, State@1), erlang:element(11, State@1), [{Chat_id, {telegram_api_error, 429, Description}} | erlang:element(12, State@1)], erlang:element(13, State@1), erlang:element(14, State@1), erlang:element(15, State@1), erlang:element(16, State@1), erlang:element(17, State@1), erlang:element(18, State@1)} ) end; {error, Reason} -> record( {state, erlang:element(2, State@1), erlang:element(3, State@1), erlang:element(4, State@1), erlang:element(5, State@1), erlang:element(6, State@1), erlang:element(7, State@1), erlang:element(8, State@1), erlang:element(9, State@1), erlang:element(10, State@1), erlang:element(11, State@1), [{Chat_id, Reason} | erlang:element(12, State@1)], erlang:element(13, State@1), erlang:element(14, State@1), erlang:element(15, State@1), erlang:element(16, State@1), erlang:element(17, State@1), erlang:element(18, State@1)} ) end. -file("src/telega/broadcast.gleam", 490). -spec pull(source()) -> {gleam@option:option(integer()), source()}. pull(Source) -> case Source of {chat_id_list, []} -> {none, Source}; {chat_id_list, [Chat_id | Rest]} -> {{some, Chat_id}, {chat_id_list, Rest}}; {chunk_iterator, Next_chunk, [Chat_id@1 | Rest@1]} -> {{some, Chat_id@1}, {chunk_iterator, Next_chunk, Rest@1}}; {chunk_iterator, Next_chunk@1, []} -> case Next_chunk@1() of {some, [Chat_id@2 | Rest@2]} -> {{some, Chat_id@2}, {chunk_iterator, Next_chunk@1, Rest@2}}; {some, []} -> {none, Source}; none -> {none, Source} end end. -file("src/telega/broadcast.gleam", 480). -spec next_chat_id(state(AUGG)) -> {gleam@option:option(integer()), state(AUGG)}. next_chat_id(State) -> case erlang:element(9, State) of {some, Chat_id} -> {{some, Chat_id}, State}; none -> {Chat_id@1, Source} = pull(erlang:element(3, State)), {Chat_id@1, {state, erlang:element(2, State), Source, erlang:element(4, State), erlang:element(5, State), erlang:element(6, State), erlang:element(7, State), erlang:element(8, State), erlang:element(9, State), erlang:element(10, State), erlang:element(11, State), erlang:element(12, State), erlang:element(13, State), erlang:element(14, State), erlang:element(15, State), erlang:element(16, State), erlang:element(17, State), erlang:element(18, State)}} end. -file("src/telega/broadcast.gleam", 453). -spec handle_send_next(state(AUGA)) -> gleam@otp@actor:next(state(AUGA), msg(AUGA)). handle_send_next(State) -> case erlang:element(17, State) of {some, _} -> gleam@otp@actor:continue(State); none -> Now = telega@internal@utils:current_time_ms(), State@1 = case (Now - erlang:element(13, State)) >= erlang:element( 6, State ) of true -> {state, erlang:element(2, State), erlang:element(3, State), erlang:element(4, State), erlang:element(5, State), erlang:element(6, State), erlang:element(7, State), erlang:element(8, State), erlang:element(9, State), erlang:element(10, State), erlang:element(11, State), erlang:element(12, State), Now, 0, erlang:element(15, State), erlang:element(16, State), erlang:element(17, State), erlang:element(18, State)}; false -> State end, case erlang:element(14, State@1) >= erlang:element(5, State@1) of true -> Wait = (erlang:element(13, State@1) + erlang:element( 6, State@1 )) - Now, gleam@erlang@process:send_after( erlang:element(18, State@1), gleam@int:max(Wait, 1), send_next ), gleam@otp@actor:continue(State@1); false -> case next_chat_id(State@1) of {none, State@2} -> finish(State@2, false); {{some, Chat_id}, State@3} -> send_to_chat(State@3, Chat_id) end end end. -file("src/telega/broadcast.gleam", 423). -spec handle_message(state(AUFT), msg(AUFT)) -> gleam@otp@actor:next(state(AUFT), msg(AUFT)). handle_message(State, Message) -> case Message of send_next -> handle_send_next(State); cancel -> case erlang:element(17, State) of {some, _} -> gleam@otp@actor:continue(State); none -> finish(State, true) end; {get_progress, Reply_to} -> gleam@erlang@process:send(Reply_to, progress_of(State)), gleam@otp@actor:continue(State); {await, Reply_to@1} -> case erlang:element(17, State) of {some, Report} -> gleam@erlang@process:send(Reply_to@1, Report), gleam@otp@actor:continue(State); none -> gleam@otp@actor:continue( {state, erlang:element(2, State), erlang:element(3, State), erlang:element(4, State), erlang:element(5, State), erlang:element(6, State), erlang:element(7, State), erlang:element(8, State), erlang:element(9, State), erlang:element(10, State), erlang:element(11, State), erlang:element(12, State), erlang:element(13, State), erlang:element(14, State), erlang:element(15, State), [Reply_to@1 | erlang:element(16, State)], erlang:element(17, State), erlang:element(18, State)} ) end end. -file("src/telega/broadcast.gleam", 298). ?DOC(" Start the broadcast in a background actor and return a handle.\n"). -spec start(broadcast(AUFA)) -> {ok, broadcast_handle(AUFA)} | {error, telega@error:telega_error()}. start(Broadcast) -> {broadcast, Client, Source, Send, Rate, Window_ms, On_progress} = Broadcast, Total = case Source of {chat_id_list, Chat_ids} -> {some, erlang:length(Chat_ids)}; {chunk_iterator, _, _} -> none end, _pipe@2 = gleam@otp@actor:new_with_initialiser( 1000, fun(Self) -> Now = telega@internal@utils:current_time_ms(), Initial_state = {state, Client, Source, Send, Rate, Window_ms, On_progress, Total, none, [], [], [], Now, 0, Now, [], none, Self}, gleam@erlang@process:send(Self, send_next), _pipe = gleam@otp@actor:initialised(Initial_state), _pipe@1 = gleam@otp@actor:returning(_pipe, Self), {ok, _pipe@1} end ), _pipe@3 = gleam@otp@actor:on_message(_pipe@2, fun handle_message/2), _pipe@4 = gleam@otp@actor:start(_pipe@3), _pipe@5 = gleam@result:map( _pipe@4, fun(Started) -> {broadcast_handle, erlang:element(3, Started)} end ), gleam@result:map_error( _pipe@5, fun(Reason) -> {actor_error, <<"Failed to start broadcast: "/utf8, (gleam@string:inspect(Reason))/binary>>} end ). -file("src/telega/broadcast.gleam", 349). ?DOC( " Wait for the broadcast to finish and return the report.\n" " Returns an error if it does not finish within `timeout` milliseconds\n" " (the broadcast itself keeps running).\n" ). -spec await(broadcast_handle(AUFF), integer()) -> {ok, broadcast_report(AUFF)} | {error, telega@error:telega_error()}. await(Handle, Timeout) -> Reply_subject = gleam@erlang@process:new_subject(), gleam@erlang@process:send(erlang:element(2, Handle), {await, Reply_subject}), _pipe = gleam@erlang@process:'receive'(Reply_subject, Timeout), gleam@result:map_error( _pipe, fun(_) -> {actor_error, <<"Broadcast await timed out"/utf8>>} end ). -file("src/telega/broadcast.gleam", 363). ?DOC( " Stop the broadcast. Recipients not yet contacted stay untouched,\n" " the report is finalized with `cancelled: True`. Cancelling a finished\n" " broadcast is a no-op.\n" ). -spec cancel(broadcast_handle(any())) -> nil. cancel(Handle) -> gleam@erlang@process:send(erlang:element(2, Handle), cancel). -file("src/telega/broadcast.gleam", 368). ?DOC(" Get a progress snapshot of the broadcast.\n"). -spec progress(broadcast_handle(any())) -> broadcast_progress(). progress(Handle) -> Reply_subject = gleam@erlang@process:new_subject(), gleam@erlang@process:send( erlang:element(2, Handle), {get_progress, Reply_subject} ), case gleam@erlang@process:'receive'(Reply_subject, 1000) of {ok, Progress} -> Progress; {error, _} -> {broadcast_progress, none, 0, 0, 0, 0} end. -file("src/telega/broadcast.gleam", 381). ?DOC( " Run the broadcast to completion: `start` + `await` forever.\n" " Convenient for scripts and one-off jobs.\n" ). -spec run(broadcast(AUFO)) -> {ok, broadcast_report(AUFO)} | {error, telega@error:telega_error()}. run(Broadcast) -> gleam@result:'try'( start(Broadcast), fun(Handle) -> Reply_subject = gleam@erlang@process:new_subject(), gleam@erlang@process:send( erlang:element(2, Handle), {await, Reply_subject} ), {ok, gleam_erlang_ffi:'receive'(Reply_subject)} end ).