-
Notifications
You must be signed in to change notification settings - Fork 11
Expand file tree
/
Copy pathtelegram.js
More file actions
228 lines (198 loc) · 7.69 KB
/
Copy pathtelegram.js
File metadata and controls
228 lines (198 loc) · 7.69 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
const { Telegram } = require('telegraf');
const { HttpsProxyAgent } = require('https-proxy-agent');
const { escapeHtml } = require('./escape');
/**
* Telegram expects message_thread_id to be a positive integer. parseInt on a
* label like 'general' yields NaN, which is serialised as null and answered
* with a 400 that says nothing about the actual mistake.
*/
function parseThreadId(value) {
const id = Number.parseInt(String(value).trim(), 10);
return Number.isSafeInteger(id) && id > 0 ? id : null;
}
// Telegram accepts roughly 20 messages per minute into one group. A
// crash-looping container produces far more than that, so sends are spaced
// out instead of being fired as fast as docker reports them.
const configuredInterval = Number.parseInt(process.env.TELEGRAM_NOTIFIER_SEND_INTERVAL_MS, 10);
const SEND_INTERVAL_MS = Number.isSafeInteger(configuredInterval) && configuredInterval >= 0 ?
configuredInterval : 1000;
// A burst that outruns the queue is dropped rather than kept in memory
// forever: by the time a backlog this long drains, the notifications are of
// no use anyway.
const MAX_QUEUED = 200;
const MAX_RATE_LIMIT_RETRIES = 3;
// Every failed send used to produce another notification, so one bad topic id
// could turn a restart loop into a stream of error messages.
const ERROR_WINDOW_MS = 300000;
const MAX_ERRORS_PER_WINDOW = 5;
function sleep(ms) {
return new Promise(resolve => setTimeout(resolve, ms));
}
/**
* node-fetch, which telegraf sends through, does not read the proxy
* environment variables by itself, so the agent has to be built explicitly.
* Only when one is actually configured: HttpsProxyAgent throws on an empty or
* missing value, which would take the container down for everyone not behind
* a proxy.
*/
function proxyOptions() {
const proxy = process.env.HTTPS_PROXY || process.env.https_proxy;
if (!proxy) return {};
try {
return { agent: new HttpsProxyAgent(proxy) };
} catch (e) {
console.error(`HTTPS_PROXY is not a usable URL: ${proxy}`);
console.error(e.message);
process.exit(100);
}
}
class TelegramClient {
constructor() {
this.telegram = new Telegram(process.env.TELEGRAM_NOTIFIER_BOT_TOKEN, proxyOptions());
this.queue = Promise.resolve();
this.queued = 0;
this.lastSentAt = 0;
this.recentErrors = [];
this.suppressionAnnounced = false;
this.dropped = 0;
this.threadId =
process.env.TELEGRAM_NOTIFIER_TOPIC_ID ||
process.env.TELEGRAM_NOTIFIER_THREAD_ID ||
null;
}
// Runs tasks one at a time, never closer together than SEND_INTERVAL_MS.
enqueue(task) {
if (this.queued >= MAX_QUEUED) {
// One line per drop would bury the log during exactly the burst that
// caused it, so report the first drop and the total once it clears.
if (this.dropped === 0) {
console.error(`Send queue is full (${MAX_QUEUED} waiting), dropping notifications.`);
}
this.dropped++;
return Promise.resolve(null);
}
if (this.dropped > 0) {
console.error(`Send queue has room again, ${this.dropped} notification(s) were dropped.`);
this.dropped = 0;
}
this.queued++;
const run = this.queue.then(async () => {
const wait = SEND_INTERVAL_MS - (Date.now() - this.lastSentAt);
if (wait > 0) await sleep(wait);
try {
return await task();
} finally {
this.lastSentAt = Date.now();
this.queued--;
}
});
// The chain must survive a failed send, or nothing is ever sent again.
this.queue = run.catch(() => {});
return run;
}
// Honours the retry_after Telegram sends with a 429 instead of hammering it.
async withRateLimitRetry(call) {
for (let attempt = 0; ; attempt++) {
try {
return await call();
} catch (e) {
const retryAfter = e.response?.parameters?.retry_after;
if (retryAfter === undefined || attempt >= MAX_RATE_LIMIT_RETRIES) throw e;
console.error(`Rate limited by Telegram, retrying in ${retryAfter}s.`);
await sleep((retryAfter + 1) * 1000);
}
}
}
async send(message, overrides = {}) {
const options = {
parse_mode: 'HTML',
disable_web_page_preview: true
};
// Check if threadId was explicitly provided in overrides (even if empty)
const threadId = 'threadId' in overrides
? overrides.threadId
: this.threadId;
// Only set message_thread_id if threadId has a truthy value
if (threadId) {
const parsedThreadId = parseThreadId(threadId);
if (parsedThreadId === null) {
console.error(
`Ignoring invalid topic/thread id ${JSON.stringify(threadId)}: ` +
`expected a positive integer. Sending to the chat without a topic.`
);
} else {
options.message_thread_id = parsedThreadId;
}
}
const chatId = overrides.chatId || process.env.TELEGRAM_NOTIFIER_CHAT_ID;
return this.enqueue(() => this.withRateLimitRetry(
() => this.telegram.sendMessage(chatId, message, options)
));
}
// True while the error budget for the current window still has room.
mayReportError() {
const now = Date.now();
this.recentErrors = this.recentErrors.filter(at => now - at < ERROR_WINDOW_MS);
if (this.recentErrors.length >= MAX_ERRORS_PER_WINDOW) {
if (!this.suppressionAnnounced) {
console.error(
`More than ${MAX_ERRORS_PER_WINDOW} errors in ${ERROR_WINDOW_MS / 1000}s; ` +
`further error notifications are suppressed until the rate drops. ` +
`They are still written to the log.`
);
this.suppressionAnnounced = true;
}
return false;
}
this.recentErrors.push(now);
this.suppressionAnnounced = false;
return true;
}
async sendError(e, overrides = {}) {
if (!this.mayReportError()) return null;
const options = {
parse_mode: 'HTML',
disable_web_page_preview: true
};
const chatId = overrides.chatId || process.env.TELEGRAM_NOTIFIER_CHAT_ID;
// Extract error details from TelegramError response JSON. Errors that did
// not come from the Telegram API carry neither response nor on, so the
// method and payload are left out rather than mislabelled.
const errorCode = e.response?.error_code || '0';
const errorDescription = e.response?.description || e.message || 'Unknown error';
const failedMethod = e.on?.method;
const failedPayload = e.on?.payload;
let errorMessage = failedMethod
? `[Error ${escapeHtml(failedMethod)}] ${escapeHtml(errorCode)} - ${escapeHtml(errorDescription)}`
: `[Error] ${escapeHtml(errorCode)} - ${escapeHtml(errorDescription)}`;
// The payload contains the original message text, which is our own HTML
// template output. Unescaped it would either be rendered instead of shown,
// or break the error report itself.
if (failedPayload !== undefined) {
errorMessage += `\n<pre>${escapeHtml(JSON.stringify(failedPayload, null, 2))}</pre>`;
}
try {
// Try to send error WITHOUT the topic_id to avoid recursive errors
// This ensures the message reaches the chat even if topic_id is wrong
return await this.enqueue(() => this.withRateLimitRetry(
() => this.telegram.sendMessage(chatId, errorMessage, options)
));
} catch (fallbackError) {
// If even this fails, log to console only
console.error('Failed to send error notification to Telegram:', {
originalError: {
code: errorCode,
description: errorDescription,
chatId: chatId
},
fallbackError: fallbackError.message
});
// Re-throw so it's visible in logs
throw fallbackError;
}
}
check() {
return this.telegram.getMe();
}
}
module.exports = TelegramClient;