use std::collections::HashMap;
use std::sync::{Mutex, OnceLock};
use std::time::{Duration, Instant};
#[cfg(test)]
use std::sync::atomic::{AtomicU64, Ordering};
use teloxide::Bot;
use teloxide::prelude::Requester;
use teloxide::types::{ChatId, MessageId, ParseMode};
use crate::config::Config;
#[cfg(test)]
static CLOCK_OFFSET_MS: AtomicU64 = AtomicU64::new(0);
#[cfg(not(test))]
fn gate_now() -> Instant {
Instant::now()
}
#[cfg(test)]
fn gate_now() -> Instant {
let off = CLOCK_OFFSET_MS.load(Ordering::Relaxed);
Instant::now()
.checked_add(Duration::from_millis(off))
.expect("virtual clock offset overflow")
}
const SEND_MAX_HOLD: Duration = Duration::from_secs(30);
const FINAL_MAX_ATTEMPTS: u32 = 8;
const DRAIN_TICK: Duration = Duration::from_millis(400);
struct Limits {
enabled: bool,
typing_interval: Duration,
typing_burst: u32,
typing_max_hold: Duration,
edit_rate_per_sec: f64,
edit_burst: u32,
send_interval: Duration,
send_minute_ceiling: u32,
send_burst: u32,
rich_rate_per_sec: f64,
rich_burst: u32,
summary_log_period: Duration,
}
impl Limits {
fn from_config() -> Self {
let rl = &Config::current().channels.telegram.rate_limiter;
let ceiling = rl.sends_ceiling_per_minute.max(1);
Self {
enabled: rl.enabled,
typing_interval: Duration::from_secs(rl.typing_min_interval_secs.max(1)),
typing_burst: rl.typing_burst.max(1),
typing_max_hold: Duration::from_secs(rl.typing_max_hold_secs),
edit_rate_per_sec: (rl.edits_per_minute.max(1) as f64) / 60.0,
edit_burst: rl.edit_burst.max(1),
send_interval: Duration::from_millis(rl.send_min_interval_millis.max(50)),
send_minute_ceiling: ceiling,
send_burst: rl.sends_burst.max(1),
rich_rate_per_sec: (rl.rich_per_minute.max(1) as f64) / 60.0,
rich_burst: rl.rich_burst.max(1),
summary_log_period: Duration::from_secs(rl.summary_log_secs.max(30)),
}
}
}
pub(crate) struct Bucket {
pub(crate) tokens: f64,
capacity: f64,
refill_per_sec: f64,
last_refill: Instant,
}
impl Bucket {
pub(crate) fn new(capacity: u32, refill_per_sec: f64) -> Self {
Self {
tokens: f64::from(capacity),
capacity: f64::from(capacity),
refill_per_sec,
last_refill: Instant::now(),
}
}
fn refill(&mut self, now: Instant) {
let elapsed = now.saturating_duration_since(self.last_refill);
self.last_refill = now;
let gained = elapsed.as_secs_f64() * self.refill_per_sec;
self.tokens = (self.tokens + gained).min(self.capacity);
}
pub(crate) fn take(&mut self, now: Instant) -> Result<(), Duration> {
let wait = self.next_token_in(now);
if wait.is_zero() {
self.tokens -= 1.0;
Ok(())
} else {
Err(wait)
}
}
fn next_token_in(&mut self, now: Instant) -> Duration {
self.refill(now);
if self.tokens >= 1.0 {
Duration::ZERO
} else {
Duration::from_secs_f64((1.0 - self.tokens) / self.refill_per_sec)
}
}
}
pub(crate) fn ensure_bucket(
slot: &mut Option<Bucket>,
capacity: u32,
refill_per_sec: f64,
) -> &mut Bucket {
let stale = match slot {
Some(b) => {
b.capacity != f64::from(capacity)
|| (b.refill_per_sec - refill_per_sec).abs() > f64::EPSILON
}
None => true,
};
if stale {
*slot = Some(Bucket::new(capacity, refill_per_sec));
}
slot.as_mut().expect("bucket was just built")
}
#[derive(Clone)]
struct PendingFinal {
bot: Bot,
html: String,
rich: bool,
attempts: u32,
}
#[derive(Default)]
pub(crate) struct Counters {
pub(crate) admitted_typing: u64,
pub(crate) admitted_edits: u64,
pub(crate) admitted_sends: u64,
pub(crate) admitted_rich: u64,
pub(crate) dropped_typing: u64,
pub(crate) dropped_clock: u64,
pub(crate) dropped_brain_preview: u64,
pub(crate) dropped_intermediary: u64,
pub(crate) dropped_status: u64,
pub(crate) queued_finals: u64,
pub(crate) superseded_finals: u64,
pub(crate) delivered_finals: u64,
pub(crate) failed_finals: u64,
pub(crate) throttled_typing_ms: u64,
pub(crate) throttled_send_ms: u64,
pub(crate) throttled_rich_ms: u64,
}
impl Counters {
pub(crate) fn note_drop(&mut self, class: EditClass) {
match class {
EditClass::Clock => self.dropped_clock += 1,
EditClass::BrainPreview => self.dropped_brain_preview += 1,
EditClass::Intermediary => self.dropped_intermediary += 1,
EditClass::Status => self.dropped_status += 1,
EditClass::Final | EditClass::Interactive => {}
}
}
fn all_zero(&self) -> bool {
self.admitted_typing == 0
&& self.admitted_edits == 0
&& self.admitted_sends == 0
&& self.dropped_typing == 0
&& self.dropped_clock == 0
&& self.dropped_brain_preview == 0
&& self.dropped_intermediary == 0
&& self.dropped_status == 0
&& self.queued_finals == 0
&& self.superseded_finals == 0
&& self.delivered_finals == 0
&& self.failed_finals == 0
&& self.throttled_typing_ms == 0
&& self.throttled_send_ms == 0
}
}
#[derive(Default)]
struct Peer {
forum_seen: bool,
typing: Option<Bucket>,
edits: Option<Bucket>,
sends_sec: Option<Bucket>,
sends_min: Option<Bucket>,
rich: Option<Bucket>,
finals: HashMap<i32, PendingFinal>,
draining: bool,
counters: Counters,
}
fn peers() -> &'static Mutex<HashMap<i64, Peer>> {
static PEERS: OnceLock<Mutex<HashMap<i64, Peer>>> = OnceLock::new();
PEERS.get_or_init(|| Mutex::new(HashMap::new()))
}
pub(crate) fn format_summary(chat_id: i64, c: &Counters, finals_pending: usize) -> Option<String> {
if c.all_zero() && finals_pending == 0 {
return None;
}
Some(format!(
"Telegram rate-limiter chat={chat_id}: \
admitted{{typing={},edits={},sends={},rich={}}} \
dropped{{clock={},brain_preview={},intermediary={},status={},typing={}}} \
finals{{queued={},superseded={},delivered={},failed={},pending={}}} \
throttled_ms{{typing={},send={},rich={}}}",
c.admitted_typing,
c.admitted_edits,
c.admitted_sends,
c.admitted_rich,
c.dropped_clock,
c.dropped_brain_preview,
c.dropped_intermediary,
c.dropped_status,
c.dropped_typing,
c.queued_finals,
c.superseded_finals,
c.delivered_finals,
c.failed_finals,
finals_pending,
c.throttled_typing_ms,
c.throttled_send_ms,
c.throttled_rich_ms,
))
}
#[cfg(not(test))]
async fn summary_loop() {
loop {
let period = Limits::from_config().summary_log_period;
tokio::time::sleep(period).await;
let lines: Vec<String> = {
let map = peers().lock().unwrap_or_else(|e| e.into_inner());
map.iter()
.filter_map(|(chat_id, peer)| {
format_summary(*chat_id, &peer.counters, peer.finals.len())
})
.collect()
};
for line in lines {
tracing::info!("{line}");
}
}
}
fn ensure_summary_task() {
#[cfg(test)]
{
let _ = Limits::from_config().summary_log_period;
}
#[cfg(not(test))]
{
static STARTED: OnceLock<()> = OnceLock::new();
STARTED.get_or_init(|| {
tokio::spawn(summary_loop());
});
}
}
enum Decision {
Admit,
Drop,
Hold(Duration),
}
pub(crate) async fn admit_chat_action(chat: ChatId, thread_id: Option<i32>) -> bool {
let chat_id = chat.0;
if chat_id >= 0 {
return true;
}
let lim = Limits::from_config();
if !lim.enabled {
return true;
}
ensure_summary_task();
let mut waited = Duration::ZERO;
loop {
let decision = {
let mut map = peers().lock().unwrap_or_else(|e| e.into_inner());
let peer = map.entry(chat_id).or_default();
if thread_id.is_some() {
peer.forum_seen = true;
}
if !peer.forum_seen {
return true;
}
let bucket = ensure_bucket(
&mut peer.typing,
lim.typing_burst,
1.0 / lim.typing_interval.as_secs_f64(),
);
match bucket.take(gate_now()) {
Ok(()) => {
peer.counters.admitted_typing += 1;
Decision::Admit
}
Err(wait) => {
if waited + wait > lim.typing_max_hold {
peer.counters.dropped_typing += 1;
Decision::Drop
} else {
Decision::Hold(wait)
}
}
}
};
match decision {
Decision::Admit => {
fold_throttle_ms(chat_id, waited);
return true;
}
Decision::Drop => {
tracing::debug!(
"Telegram rate-limiter: typing refresh dropped for chat={chat_id} after \
holding {waited:?} (hold cap {}s)",
lim.typing_max_hold.as_secs()
);
fold_throttle_ms(chat_id, waited);
return false;
}
Decision::Hold(wait) => {
let start = gate_now();
tokio::time::sleep(wait).await;
#[cfg(test)]
test_support::advance(wait.as_millis() as u64);
waited += gate_now().duration_since(start);
}
}
}
}
fn fold_throttle_ms(chat_id: i64, waited: Duration) {
if waited.is_zero() {
return;
}
let ms = waited.as_millis() as u64;
if ms == 0 {
return;
}
if let Some(peer) = peers()
.lock()
.unwrap_or_else(|e| e.into_inner())
.get_mut(&chat_id)
{
peer.counters.throttled_typing_ms += ms;
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum EditClass {
Clock,
BrainPreview,
Intermediary,
Status,
Final,
Interactive,
}
impl EditClass {
pub(crate) fn drop_rank(self) -> u8 {
match self {
EditClass::Clock => 0,
EditClass::BrainPreview => 1,
EditClass::Intermediary => 2,
EditClass::Status => 3,
EditClass::Final => 4,
EditClass::Interactive => 5,
}
}
}
enum Admission {
Now,
Dropped(EditClass),
Queued,
}
pub(crate) async fn edit_admission(
bot: &Bot,
chat_id: ChatId,
msg_id: MessageId,
class: EditClass,
html: String,
rich: bool,
) -> bool {
if chat_id.0 >= 0 {
return true;
}
let lim = Limits::from_config();
if !lim.enabled {
return true;
}
ensure_summary_task();
let now = gate_now();
let admission = {
let mut map = peers().lock().unwrap_or_else(|e| e.into_inner());
let peer = map.entry(chat_id.0).or_default();
if !peer.forum_seen {
return true;
}
let bucket = ensure_bucket(&mut peer.edits, lim.edit_burst, lim.edit_rate_per_sec);
if bucket.take(now).is_ok() {
peer.counters.admitted_edits += 1;
Admission::Now
} else if class == EditClass::Interactive {
Admission::Now
} else if class == EditClass::Final {
let superseded = peer
.finals
.insert(
msg_id.0,
PendingFinal {
bot: bot.clone(),
html,
rich,
attempts: 0,
},
)
.is_some();
if superseded {
peer.counters.superseded_finals += 1;
}
peer.counters.queued_finals += 1;
Admission::Queued
} else {
peer.counters.note_drop(class);
Admission::Dropped(class)
}
};
match admission {
Admission::Now => true,
Admission::Dropped(dropped) => {
tracing::debug!(
"Telegram rate-limiter: {:?} edit (rank {}) dropped for chat={} msg={} — \
next full-state refresh carries it",
dropped,
dropped.drop_rank(),
chat_id.0,
msg_id.0
);
false
}
Admission::Queued => {
tracing::debug!(
"Telegram rate-limiter: final edit queued latest-wins for chat={} msg={}",
chat_id.0,
msg_id.0
);
ensure_drain(chat_id.0);
false
}
}
}
fn take_due_final(chat_id: i64, lim: &Limits) -> Option<(i32, PendingFinal)> {
let mut map = peers().lock().unwrap_or_else(|e| e.into_inner());
let peer = map.get_mut(&chat_id)?;
if peer.finals.is_empty() {
peer.draining = false;
return None;
}
let bucket = ensure_bucket(&mut peer.edits, lim.edit_burst, lim.edit_rate_per_sec);
if bucket.take(gate_now()).is_err() {
return None;
}
let key = peer.finals.keys().next().copied()?;
peer.finals.remove(&key).map(|pending| (key, pending))
}
async fn drain_finals(chat_id: i64) {
loop {
tokio::time::sleep(DRAIN_TICK).await;
#[cfg(test)]
test_support::advance(DRAIN_TICK.as_millis() as u64);
let lim = Limits::from_config();
let Some(job) = take_due_final(chat_id, &lim) else {
let empty = peers()
.lock()
.unwrap_or_else(|e| e.into_inner())
.get(&chat_id)
.is_none_or(|p| p.finals.is_empty());
if empty {
return;
}
continue;
};
deliver_final(chat_id, job.0, job.1).await;
}
}
fn ensure_drain(chat_id: i64) {
let should_spawn = {
let mut map = peers().lock().unwrap_or_else(|e| e.into_inner());
let peer = map.entry(chat_id).or_default();
if peer.draining || peer.finals.is_empty() {
false
} else {
peer.draining = true;
true
}
};
if should_spawn {
tokio::spawn(drain_finals(chat_id));
}
}
pub(crate) fn is_permanent_edit_error(error: &str) -> bool {
error.contains("message to edit not found")
|| error.contains("message is not modified")
|| error.contains("message can't be edited")
|| error.contains("MESSAGE_ID_INVALID")
|| error.contains("chat not found")
}
enum Verdict {
Delivered,
Abandoned(String),
Retried,
}
async fn deliver_final(chat_id: i64, msg_id: i32, mut pending: PendingFinal) {
let result = run_final_edit(&pending.bot, chat_id, msg_id, &pending.html, pending.rich).await;
let retry_after = match &result {
Err(e) => super::rate_limit::parse_retry_after(e),
Ok(()) => None,
};
let verdict = {
let mut map = peers().lock().unwrap_or_else(|e| e.into_inner());
let Some(peer) = map.get_mut(&chat_id) else {
#[cfg(test)]
test_support::notify_settled();
return;
};
match result {
Ok(()) => {
peer.counters.delivered_finals += 1;
Verdict::Delivered
}
Err(e) => {
pending.attempts += 1;
if pending.attempts >= FINAL_MAX_ATTEMPTS || is_permanent_edit_error(&e) {
peer.counters.failed_finals += 1;
Verdict::Abandoned(e)
} else {
peer.finals.insert(msg_id, pending);
Verdict::Retried
}
}
}
};
match verdict {
Verdict::Delivered => {
tracing::debug!(
"Telegram rate-limiter: queued final delivered chat={chat_id} msg={msg_id}"
);
}
Verdict::Abandoned(e) => {
tracing::warn!(
"Telegram rate-limiter: queued final abandoned after {FINAL_MAX_ATTEMPTS} \
attempts chat={chat_id} msg={msg_id}: {e}"
);
}
Verdict::Retried => {
if let Some(window) = retry_after {
let (wait, _) = super::rate_limit::clamp_inline_wait(window);
tokio::time::sleep(wait).await;
#[cfg(test)]
test_support::advance(wait.as_millis() as u64);
}
}
}
#[cfg(test)]
test_support::notify_settled();
}
async fn run_final_edit(
bot: &Bot,
chat_id: i64,
msg_id: i32,
html: &str,
rich: bool,
) -> Result<(), String> {
if rich {
super::rich::api::edit_rich_html(
bot.api_url().as_str(),
bot.token(),
chat_id,
msg_id,
html,
None,
"turn",
"-",
)
.await
.map_err(|e| e.to_string())
} else {
use teloxide::payloads::EditMessageTextSetters;
bot.edit_message_text(ChatId(chat_id), MessageId(msg_id), html)
.parse_mode(ParseMode::Html)
.await
.map(|_| ())
.map_err(|e| e.to_string())
}
}
enum PaceVerdict {
Go,
FailOpen(Duration),
Wait(Duration),
}
pub(crate) async fn pace_send(chat: ChatId) {
let chat_id = chat.0;
if chat_id >= 0 {
return;
}
let lim = Limits::from_config();
if !lim.enabled {
return;
}
let mut waited = Duration::ZERO;
loop {
let verdict = {
let mut map = peers().lock().unwrap_or_else(|e| e.into_inner());
let peer = map.entry(chat_id).or_default();
if !peer.forum_seen {
return;
}
let sec = ensure_bucket(
&mut peer.sends_sec,
lim.send_burst,
1.0 / lim.send_interval.as_secs_f64(),
);
let min = ensure_bucket(
&mut peer.sends_min,
lim.send_minute_ceiling,
f64::from(lim.send_minute_ceiling) / 60.0,
);
let now = gate_now();
let need = sec.next_token_in(now).max(min.next_token_in(now));
if need.is_zero() {
let _ = sec.take(now);
let _ = min.take(now);
peer.counters.admitted_sends += 1;
PaceVerdict::Go
} else if waited + need > SEND_MAX_HOLD {
peer.counters.admitted_sends += 1;
PaceVerdict::FailOpen(need)
} else {
PaceVerdict::Wait(need)
}
};
match verdict {
PaceVerdict::Go => {
fold_send_ms(chat_id, waited);
return;
}
PaceVerdict::FailOpen(next) => {
tracing::warn!(
"Telegram rate-limiter: send pacing held {waited:?} for chat={chat_id}, \
still {} from a token — failing open to the reactive backstop",
next.as_secs_f64()
);
fold_send_ms(chat_id, waited);
return;
}
PaceVerdict::Wait(delay) => {
let start = gate_now();
tokio::time::sleep(delay).await;
#[cfg(test)]
test_support::advance(delay.as_millis() as u64);
waited += gate_now().duration_since(start);
}
}
}
}
pub(crate) async fn pace_rich(chat: ChatId, thread_id: Option<i32>) {
let chat_id = chat.0;
if chat_id >= 0 {
return;
}
let lim = Limits::from_config();
if !lim.enabled {
return;
}
ensure_summary_task();
let mut waited = Duration::ZERO;
loop {
let verdict = {
let mut map = peers().lock().unwrap_or_else(|e| e.into_inner());
let peer = map.entry(chat_id).or_default();
if thread_id.is_some() {
peer.forum_seen = true;
}
if !peer.forum_seen {
return;
}
let bucket = ensure_bucket(&mut peer.rich, lim.rich_burst, lim.rich_rate_per_sec);
let now = gate_now();
let need = bucket.next_token_in(now);
if need.is_zero() {
let _ = bucket.take(now);
peer.counters.admitted_rich += 1;
PaceVerdict::Go
} else if waited + need > SEND_MAX_HOLD {
peer.counters.admitted_rich += 1;
PaceVerdict::FailOpen(need)
} else {
PaceVerdict::Wait(need)
}
};
match verdict {
PaceVerdict::Go => {
fold_rich_ms(chat_id, waited);
return;
}
PaceVerdict::FailOpen(next) => {
tracing::warn!(
"Telegram rate-limiter: rich pacing held {waited:?} for chat={chat_id}, \
still {} from a token — failing open to the reactive backstop",
next.as_secs_f64()
);
fold_rich_ms(chat_id, waited);
return;
}
PaceVerdict::Wait(delay) => {
let start = gate_now();
tokio::time::sleep(delay).await;
#[cfg(test)]
test_support::advance(delay.as_millis() as u64);
waited += gate_now().duration_since(start);
}
}
}
}
fn fold_rich_ms(chat_id: i64, waited: Duration) {
if waited.is_zero() {
return;
}
let ms = waited.as_millis() as u64;
if ms == 0 {
return;
}
if let Some(peer) = peers()
.lock()
.unwrap_or_else(|e| e.into_inner())
.get_mut(&chat_id)
{
peer.counters.throttled_rich_ms += ms;
}
}
fn fold_send_ms(chat_id: i64, waited: Duration) {
if waited.is_zero() {
return;
}
let ms = waited.as_millis() as u64;
if ms == 0 {
return;
}
if let Some(peer) = peers()
.lock()
.unwrap_or_else(|e| e.into_inner())
.get_mut(&chat_id)
{
peer.counters.throttled_send_ms += ms;
}
}
#[cfg(test)]
pub(crate) mod test_support {
use super::*;
use std::sync::atomic::Ordering;
pub(crate) async fn registry_guard() -> tokio::sync::MutexGuard<'static, ()> {
static GUARD: tokio::sync::Mutex<()> = tokio::sync::Mutex::const_new(());
GUARD.lock().await
}
pub(crate) fn reset(offset_ms: u64) {
peers().lock().unwrap_or_else(|e| e.into_inner()).clear();
CLOCK_OFFSET_MS.store(offset_ms, Ordering::Relaxed);
}
pub(crate) fn advance(ms: u64) {
CLOCK_OFFSET_MS.fetch_add(ms, Ordering::Relaxed);
}
pub(crate) fn settle_notifier() -> &'static tokio::sync::Notify {
static SETTLED: std::sync::OnceLock<tokio::sync::Notify> = std::sync::OnceLock::new();
SETTLED.get_or_init(tokio::sync::Notify::new)
}
pub(crate) fn notify_settled() {
settle_notifier().notify_one();
}
pub(crate) fn mark_forum(chat: ChatId) {
let chat_id = chat.0;
peers()
.lock()
.unwrap_or_else(|e| e.into_inner())
.entry(chat_id)
.or_default()
.forum_seen = true;
}
pub(crate) fn burn_bucket(chat: ChatId, kind: BucketKind, capacity: u32, rate_per_sec: f64) {
let chat_id = chat.0;
let mut map = peers().lock().unwrap_or_else(|e| e.into_inner());
let peer = map.entry(chat_id).or_default();
peer.forum_seen = true;
let slot = match kind {
BucketKind::Typing => &mut peer.typing,
BucketKind::Edits => &mut peer.edits,
BucketKind::SendsSec => &mut peer.sends_sec,
BucketKind::SendsMin => &mut peer.sends_min,
BucketKind::Rich => &mut peer.rich,
};
let bucket = ensure_bucket(slot, capacity, rate_per_sec);
for _ in 0..capacity {
let _ = bucket.take(gate_now());
}
}
#[derive(Debug, Clone, Copy)]
pub(crate) enum BucketKind {
Typing,
Edits,
SendsSec,
SendsMin,
Rich,
}
#[derive(Debug, Default, PartialEq, Eq)]
pub(crate) struct Snap {
pub forum_seen: bool,
pub typing_admitted: u64,
pub typing_dropped: u64,
pub edits_admitted: u64,
pub dropped_clock: u64,
pub dropped_brain_preview: u64,
pub dropped_intermediary: u64,
pub dropped_status: u64,
pub queued_finals: u64,
pub superseded_finals: u64,
pub delivered_finals: u64,
pub failed_finals: u64,
pub throttled_typing_ms: u64,
pub admitted_rich: u64,
pub throttled_rich_ms: u64,
pub throttled_send_ms: u64,
pub finals_pending: usize,
}
pub(crate) fn snapshot(chat: ChatId) -> Option<Snap> {
let chat_id = chat.0;
let map = peers().lock().unwrap_or_else(|e| e.into_inner());
map.get(&chat_id).map(|p| Snap {
forum_seen: p.forum_seen,
typing_admitted: p.counters.admitted_typing,
typing_dropped: p.counters.dropped_typing,
edits_admitted: p.counters.admitted_edits,
dropped_clock: p.counters.dropped_clock,
dropped_brain_preview: p.counters.dropped_brain_preview,
dropped_intermediary: p.counters.dropped_intermediary,
dropped_status: p.counters.dropped_status,
queued_finals: p.counters.queued_finals,
superseded_finals: p.counters.superseded_finals,
delivered_finals: p.counters.delivered_finals,
failed_finals: p.counters.failed_finals,
throttled_typing_ms: p.counters.throttled_typing_ms,
admitted_rich: p.counters.admitted_rich,
throttled_rich_ms: p.counters.throttled_rich_ms,
throttled_send_ms: p.counters.throttled_send_ms,
finals_pending: p.finals.len(),
})
}
}