use std::{net::IpAddr, time::Duration}; use anyhow::{Result, bail}; use serde::Deserialize; use serde_json::json; use tracing::warn; use super::Sink; use crate::event::{Event, RoundReport, describe}; /// Sends one message per round that has events. pub struct TelegramSink { http: reqwest::Client, token: String, chat_id: i64, } #[derive(Deserialize)] struct Reply { ok: bool, description: Option, parameters: Option, } #[derive(Deserialize)] struct ReplyParameters { retry_after: Option, } impl TelegramSink { pub fn new(token: String, chat_id: i64) -> Result { let http = reqwest::Client::builder().timeout(Duration::from_secs(15)).build()?; Ok(Self { http, token, chat_id }) } pub async fn send(&self, text: &str) -> Result<()> { const ATTEMPTS: u64 = 3; let url = format!("https://api.telegram.org/bot{}/sendMessage", self.token); let body = json!({ "chat_id": self.chat_id, "text": text, "disable_web_page_preview": true }); let mut last_err = String::new(); for attempt in 1..=ATTEMPTS { let mut wait = Duration::from_secs(2 * attempt); // reqwest errors embed the request URL, which contains the bot token: strip it. match self.http.post(&url).json(&body).send().await { Err(e) => last_err = e.without_url().to_string(), Ok(resp) => { let status = resp.status(); match resp.json::().await { Err(e) => last_err = format!("HTTP {status}: {}", e.without_url()), Ok(reply) if reply.ok => return Ok(()), Ok(reply) => { last_err = format!("HTTP {status}: {}", reply.description.unwrap_or_default()); match reply.parameters.and_then(|p| p.retry_after) { Some(secs) => wait = Duration::from_secs(secs), // Bad token, bot not in the chat, ...: retrying won't help. None if status.is_client_error() => break, None => {} } } } } } if attempt < ATTEMPTS { tokio::time::sleep(wait).await; } } bail!("Telegram sendMessage failed: {last_err}") } } impl Sink for TelegramSink { async fn handle(&mut self, report: &RoundReport) { let Some(text) = render(report) else { return }; if let Err(e) = self.send(&text).await { warn!("{e:#}"); } } } /// The message for one round, or `None` if nothing happened worth telling. pub fn render(r: &RoundReport) -> Option { if r.events.is_empty() { return None; } let mut lines = vec![r.group.clone()]; for event in &r.events { match event { Event::Started { missing } => { lines.push("🚀 dns-monitor 已启动".to_string()); for s in &r.ips { let mark = if s.up { "✅" } else { "❌" }; lines.push(format!("{mark} {} ({}) {}", s.ip, s.members.join(", "), describe(&s.probe))); } for m in missing { lines.push(format!("⚠️ {m} 没有 {} 记录", r.record_type)); } } Event::MemberChanged { member, old, new } => { lines.push(format!("🔄 {member}: {} → {}", join_ips(old), join_ips(new))) } Event::HealthChanged { ip, members, up: true, probe } => { lines.push(format!("🟢 {ip} ({}) 恢复: {}", members.join(", "), describe(probe))) } Event::HealthChanged { ip, members, up: false, probe } => { lines.push(format!("🔴 {ip} ({}) 不通: {}", members.join(", "), describe(probe))) } Event::DnsUpdated { added, removed } => { if !added.is_empty() { lines.push(format!("➕ DNS 添加 {}", join_ips(added))); } if !removed.is_empty() { lines.push(format!("➖ DNS 移除 {}", join_ips(removed))); } } Event::NoHealthyMembers => lines.push("🆘 所有成员都不通,DNS 保持不变".to_string()), Event::CheckFailed { error } => lines.push(format!("⚠️ 检查失败: {error}")), Event::CheckRecovered => lines.push("✅ 检查恢复正常".to_string()), } } if let Some(now) = &r.dns_now { lines.push(format!("当前解析: {}", join_ips(now))); } Some(lines.join("\n")) } fn join_ips(ips: &[IpAddr]) -> String { if ips.is_empty() { return "(无)".to_string(); } ips.iter().map(ToString::to_string).collect::>().join(", ") } #[cfg(test)] mod tests { use super::*; use crate::event::IpStatus; fn ip(s: &str) -> IpAddr { s.parse().unwrap() } fn report(events: Vec) -> RoundReport { RoundReport { group: "jp.hoshino.app".to_string(), record_type: "A", dry_run: false, ips: vec![ IpStatus { ip: ip("43.207.42.149"), members: vec!["jp6.hoshino.app".to_string()], up: false, probe: Err("ping: 3 lost".to_string()), }, IpStatus { ip: ip("52.195.154.53"), members: vec!["jp5.hoshino.app".to_string()], up: true, probe: Ok(Duration::from_millis(49)), }, ], events, dns_now: Some(vec![ip("52.195.154.53")]), error: None, } } #[test] fn quiet_round_sends_nothing() { assert_eq!(render(&report(vec![])), None); } #[test] fn startup_lists_every_ip() { let text = render(&report(vec![Event::Started { missing: vec!["jp7.hoshino.app".to_string()] }])); assert_eq!( text.unwrap(), "jp.hoshino.app\n\ 🚀 dns-monitor 已启动\n\ ❌ 43.207.42.149 (jp6.hoshino.app) ping: 3 lost\n\ ✅ 52.195.154.53 (jp5.hoshino.app) 49ms\n\ ⚠️ jp7.hoshino.app 没有 A 记录\n\ 当前解析: 52.195.154.53" ); } #[test] fn node_down_and_member_change_in_one_message() { let text = render(&report(vec![ Event::MemberChanged { member: "jp5.hoshino.app".to_string(), old: vec![ip("52.195.154.53")], new: vec![], }, Event::HealthChanged { ip: ip("43.207.42.149"), members: vec!["jp6.hoshino.app".to_string()], up: false, probe: Err("ping: 3 lost".to_string()), }, Event::DnsUpdated { added: vec![], removed: vec![ip("43.207.42.149")] }, ])); assert_eq!( text.unwrap(), "jp.hoshino.app\n\ 🔄 jp5.hoshino.app: 52.195.154.53 → (无)\n\ 🔴 43.207.42.149 (jp6.hoshino.app) 不通: ping: 3 lost\n\ ➖ DNS 移除 43.207.42.149\n\ 当前解析: 52.195.154.53" ); } #[test] fn outage_and_errors() { let mut r = report(vec![Event::NoHealthyMembers, Event::CheckFailed { error: "API down".to_string() }]); r.dns_now = None; assert_eq!( render(&r).unwrap(), "jp.hoshino.app\n🆘 所有成员都不通,DNS 保持不变\n⚠️ 检查失败: API down" ); } }