Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
24 changes: 22 additions & 2 deletions broadcaster/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -17,12 +17,13 @@ if Telegram or the service goes down, the pool is unaffected.

| | when |
|---|---|
| **Daily digest** | once per UTC day at `DIGEST_UTC_HOUR` (default 12): hashrate 1h/24h, active workers, shares and reject rate, best share, blocks in the last 24h and all time, mode and fee |
| **Daily digest** | once per UTC day at `DIGEST_UTC_HOUR` (default 12): hashrate now (5m), 1h and 24h, active workers, accepted and rejected shares, best share, blocks in the last 24h and all time, mode and fee |
| **Block found** | the first time a block appears in `/api/blocks` as `pending` or `confirmed` |
| **Block lost** | a block it announced turns `orphaned` or `rejected` |
| **Pool quiet / back** | no share accepted for `STALE_SHARES_SEC` (default 15m), and when shares resume |
| **Health failing / recovered** | `/health` changes state, with the failing checks listed |
| **Pinned live message** (`BROADCASTER_LIVE=1`) | one message, pinned, edited every `BROADCASTER_LIVE_MS` (default 5m) with current stats |
| **Pinned live message** (`BROADCASTER_LIVE=1`) | one message, pinned unless `BROADCASTER_LIVE_PIN=0`, edited every `BROADCASTER_LIVE_MS` (default 5m) with current stats and the last block |
| **`/pool_status`** (`BROADCASTER_COMMANDS=1`) | when someone sends it in the chat: the current stats, as a reply in the same topic, at most once per `BROADCASTER_COMMAND_COOLDOWN_SEC` (default 30) |
| **Periodic update** (`BROADCASTER_SUMMARY_HOURS=N`) | the same stats as a new message every N hours, on UTC boundaries (6 → 00, 06, 12, 18 UTC) |

- **Edited is not posted.** Telegram does not notify anyone of an edit or
Expand Down Expand Up @@ -58,6 +59,23 @@ appear under the channel's name and avatar.
(`23563`). In a group, the bot posts as itself and needs no admin rights
when members may send (and, for the live message, pin) messages.

### `/pool_status`

With `BROADCASTER_COMMANDS=1` the service also reads the chat, with
`getUpdates`, and answers `/pool_status` with the live message's numbers
plus the last block. It answers only in `TELEGRAM_CHAT_ID`, so adding the
bot to another group does not make it answer there.

- It registers the command in the bot's `/` menu. A bot in privacy mode
(BotFather's default) always sees `/pool_status@yourbot`, which is what a
tap in that menu sends, but not always a bare `/pool_status`. Turn privacy
off (`/setprivacy` → Disable) if people type it by hand.
- Telegram allows one `getUpdates` reader per bot and none while a webhook
is set. Run one broadcaster per bot, and don't poll the bot from anywhere
else while it runs.
- The first start skips commands sent before it, and the read position is
kept in `broadcaster.json`, so a restart answers nothing twice.

Try it without a token first:

```bash
Expand All @@ -82,7 +100,9 @@ Environment only. The full list with defaults is in
| `BROADCASTER_POLL_MS` | default 60000 |
| `DIGEST_UTC_HOUR` | 0–23, or `off` (default 12) |
| `BROADCASTER_LIVE`, `BROADCASTER_LIVE_MS` | pinned live message (default off, 300000) |
| `BROADCASTER_LIVE_PIN` | `0` = post the live message without pinning it (default 1) |
| `BROADCASTER_SUMMARY_HOURS` | 1–168, or `off` (default off) |
| `BROADCASTER_COMMANDS`, `BROADCASTER_COMMAND_COOLDOWN_SEC` | answer `/pool_status` (default off, 30) |
| `STALE_SHARES_SEC` | default 900 |
| `HEALTH_DEBOUNCE` | default 3 |

Expand Down
32 changes: 23 additions & 9 deletions broadcaster/index.js
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@
* Reads the dashboard's public JSON API and posts to one channel: a daily
* digest, found / orphaned blocks, the pool going quiet and coming back,
* ledger health changes, and optionally a pinned message kept current.
* With BROADCASTER_COMMANDS=1 it also answers /pool_status in that chat.
*
* Config is environment-only — see lib/config.js for the full list.
*
Expand All @@ -17,10 +18,11 @@
import { readFileSync } from 'node:fs';

import { loadConfig } from './lib/config.js';
import { loadState } from './lib/state.js';
import { loadState, saveState } from './lib/state.js';
import { DashboardClient } from './lib/dashboard.js';
import { TelegramClient, DryRunClient } from './lib/telegram.js';
import { Broadcaster } from './lib/broadcaster.js';
import { Commands } from './lib/commands.js';

const cfg = loadConfig();
const { version } = JSON.parse(readFileSync(new URL('./package.json', import.meta.url), 'utf8'));
Expand All @@ -36,17 +38,29 @@ log.info(`simplepool-broadcaster ${version} starting ` +
`${cfg.threadId != null ? ` topic=${cfg.threadId}` : ''} state=${cfg.statePath}` +
`${cfg.dryRun ? ' DRY RUN' : ''})`);
log.info(` poll ${cfg.pollMs}ms, digest ${cfg.digestHour == null ? 'off' : `${cfg.digestHour}:00 UTC`}, ` +
`live ${cfg.live ? `every ${cfg.liveMs}ms` : 'off'}, ` +
`live ${cfg.live ? `every ${cfg.liveMs}ms${cfg.livePin ? '' : ' (unpinned)'}` : 'off'}, ` +
`summary ${cfg.summaryHours == null ? 'off' : `every ${cfg.summaryHours}h`}, ` +
`stale after ${cfg.staleSharesSec}s`);

const broadcaster = new Broadcaster({
dashboard: new DashboardClient({ url: cfg.dashboardUrl }),
telegram: cfg.dryRun ? new DryRunClient() : new TelegramClient({ token: cfg.token, chatId: cfg.chatId, threadId: cfg.threadId }),
state: loadState(cfg.statePath),
cfg,
log,
});
// One state object and one writer for both loops, so neither can save over
// the other's changes with a stale copy.
const state = loadState(cfg.statePath);
const persist = () => saveState(cfg.statePath, state);
const dashboard = new DashboardClient({ url: cfg.dashboardUrl });
const telegram = cfg.dryRun
? new DryRunClient()
: new TelegramClient({ token: cfg.token, chatId: cfg.chatId, threadId: cfg.threadId });

const broadcaster = new Broadcaster({ dashboard, telegram, state, cfg, log, persist });

if (cfg.commands && cfg.dryRun) {
log.info(' commands: off under BROADCASTER_DRY_RUN (they need a real bot to read messages)');
} else if (cfg.commands) {
const commands = new Commands({ dashboard, telegram, state, cfg, log, persist });
// Its own loop: a long poll must never delay a post, and a failure here
// must never stop the posts.
commands.run().catch((e) => log.error(`commands stopped: ${e.message}`));
}

let lastTickError = null;
async function loop() {
Expand Down
26 changes: 17 additions & 9 deletions broadcaster/lib/broadcaster.js
Original file line number Diff line number Diff line change
Expand Up @@ -7,7 +7,8 @@
* recovering, once the new state has held for `debounce` polls
* - posts the daily digest once per UTC day, at or after DIGEST_UTC_HOUR
* - posts a stats summary once per BROADCASTER_SUMMARY_HOURS period
* - refreshes the pinned live message (BROADCASTER_LIVE=1)
* - refreshes the live message (BROADCASTER_LIVE=1), pinned unless
* BROADCASTER_LIVE_PIN=0
*
* State is saved after every post, so a failure halfway through a tick
* retries only what was not sent.
Expand Down Expand Up @@ -45,9 +46,9 @@ export class Broadcaster {
await this.#liveness(status, nowSec);
await this.#health(status);
await this.#digest(status, blocks, nowMs);
await this.#summary(status, nowMs);
await this.#summary(status, blocks, nowMs);
}
await this.#live(status, nowMs);
await this.#live(status, blocks, nowMs);
}

/* First run: remember everything as already said. */
Expand Down Expand Up @@ -115,13 +116,14 @@ export class Broadcaster {
/* A new message per period, unlike the live one, which is edited. The
* first period seen (first run, or the option just turned on) is
* recorded without a post, so enabling it does not post on the spot. */
async #summary(status, nowMs) {
async #summary(status, blocks, nowMs) {
if (this.cfg.summaryHours == null) return;
const slot = this.#summarySlot(nowMs);
if (this.state.lastSummarySlot === slot) return;
if (this.state.lastSummarySlot != null) {
await this.#post(msg.summary({
status,
blocks,
nowSec: Math.floor(nowMs / 1000),
flowing: this.state.sharesFlowing !== false,
...this.#ctx(),
Expand All @@ -135,10 +137,11 @@ export class Broadcaster {
return Math.floor(nowMs / (this.cfg.summaryHours * 3_600_000));
}

async #live(status, nowMs) {
async #live(status, blocks, nowMs) {
if (!this.cfg.live || nowMs - this.state.liveUpdatedAt < this.cfg.liveMs) return;
const html = msg.live({
status,
blocks,
nowSec: Math.floor(nowMs / 1000),
flowing: this.state.sharesFlowing !== false,
...this.#ctx(),
Expand All @@ -156,10 +159,15 @@ export class Broadcaster {
if (this.state.liveMessageId == null) {
this.state.liveMessageId = await this.telegram.send(html);
this.persist();
try {
await this.telegram.pin(this.state.liveMessageId);
} catch (e) {
this.log.warn(`could not pin the live message (does the bot have "Pin messages"?): ${e.message}`);
// BROADCASTER_LIVE_PIN=0: a chat where the bot may post but not
// pin keeps the message, without a warning on every new one.
if (this.cfg.livePin !== false) {
try {
await this.telegram.pin(this.state.liveMessageId);
} catch (e) {
this.log.warn(`could not pin the live message (does the bot have "Pin messages"? ` +
`BROADCASTER_LIVE_PIN=0 stops trying): ${e.message}`);
}
}
}
this.state.liveUpdatedAt = nowMs;
Expand Down
143 changes: 143 additions & 0 deletions broadcaster/lib/commands.js
Original file line number Diff line number Diff line change
@@ -0,0 +1,143 @@
/* /pool_status: answer a command in the chat with the current stats.
*
* Runs beside the poster, on its own long-poll loop (BROADCASTER_COMMANDS=1).
* Only TELEGRAM_CHAT_ID is answered: the bot can be added anywhere, and a
* stranger's group must not get a reply. The answer goes to the topic the
* command was typed in, as a reply to it, and at most once per
* BROADCASTER_COMMAND_COOLDOWN_SEC per chat, so the command cannot be used
* to flood the chat.
*
* A bot in privacy mode (the default) is shown "/pool_status@<bot>" but not
* always a bare "/pool_status". Registering the command puts it in the "/"
* menu, and a tap there sends the addressed form.
*
* The update offset is saved in broadcaster.json, so a restart neither
* answers a command twice nor answers commands typed while it was down: the
* very first start skips the backlog.
*/

import * as msg from './messages.js';

export const COMMANDS = [
{ command: 'pool_status', description: 'Current pool stats' },
];

const RETRY_MS = 5000;

export class Commands {
constructor({ dashboard, telegram, state, cfg, log, persist, sleep = defaultSleep }) {
this.dashboard = dashboard;
this.telegram = telegram;
this.state = state;
this.cfg = cfg;
this.log = log;
this.persist = persist;
this.sleep = sleep;
this.username = null;
this.lastReply = new Map(); // chat id -> ms of the last reply
this.stopped = false;
}

/* getMe, register the menu, skip the backlog on a first start. */
async init() {
this.username = (await this.telegram.getMe()).username ?? null;
try {
await this.telegram.setMyCommands(COMMANDS);
} catch (e) {
this.log.warn(`could not register the command menu: ${e.message}`);
}
if (this.state.updateOffset == null) {
const pending = await this.telegram.getUpdates({ offset: -1, timeoutSec: 0 });
this.state.updateOffset = pending.length ? pending[pending.length - 1].update_id + 1 : 0;
this.persist();
}
this.log.info(`commands: answering /pool_status in ${this.cfg.chatId} as @${this.username}`);
}

/* One long poll, and every message it brought. */
async pollOnce(clock = Date.now) {
const updates = await this.telegram.getUpdates({ offset: this.state.updateOffset, timeoutSec: 25 });
for (const u of updates) {
// Advanced before handling: a reply that fails is not retried,
// because the person can ask again and a stale answer is worse.
this.state.updateOffset = u.update_id + 1;
this.persist();
if (u.message) await this.handle(u.message, clock());
}
}

async run() {
let ready = false;
let lastError = null;
while (!this.stopped) {
try {
// Retried like a poll: Telegram being unreachable at start
// must not switch the command off until the next restart.
if (!ready) { await this.init(); ready = true; }
await this.pollOnce();
lastError = null;
} catch (e) {
if (e.message !== lastError) this.log.warn(`commands: ${e.message}`);
lastError = e.message;
await this.sleep(RETRY_MS);
}
}
}

async handle(m, nowMs) {
if (!this.#isPoolStatus(m.text) || !this.#allowedChat(m.chat)) return;

const chatKey = String(m.chat.id);
const last = this.lastReply.get(chatKey);
if (last != null && nowMs - last < this.cfg.commandCooldownSec * 1000) return;
this.lastReply.set(chatKey, nowMs);

const where = {
chatId: m.chat.id,
// In a forum, a command typed in a topic is answered there; one
// typed in General carries no topic.
threadId: m.is_topic_message ? m.message_thread_id : null,
replyTo: m.message_id,
};
let html;
try {
const [status, blocks] = await Promise.all([this.dashboard.status(), this.dashboard.blocks()]);
const nowSec = Math.floor(nowMs / 1000);
html = msg.poolStatus({
status, blocks, nowSec,
flowing: this.#flowing(status, nowSec),
poolName: this.cfg.poolName, publicUrl: this.cfg.publicUrl,
});
} catch (e) {
this.log.warn(`commands: /pool_status could not read the dashboard: ${e.message}`);
html = msg.statusUnavailable({ poolName: this.cfg.poolName });
}
await this.telegram.send(html, where);
}

/* "/pool_status", "/pool_status@this_bot", with or without arguments;
* never "/pool_status@another_bot". */
#isPoolStatus(text) {
const match = /^\/pool_status(?:@(\w+))?(?:\s|$)/i.exec(text ?? '');
if (!match) return false;
return match[1] == null || this.username == null
|| match[1].toLowerCase() === this.username.toLowerCase();
}

#allowedChat(chat) {
const want = String(this.cfg.chatId);
if (String(chat.id) === want) return true;
return chat.username != null && want.toLowerCase() === `@${chat.username}`.toLowerCase();
}

/* Same rule as the poster's liveness check, without its debounce: this
* answers about right now. */
#flowing(status, nowSec) {
const last = status.pool?.last_share_ts;
return last != null && nowSec - last < this.cfg.staleSharesSec;
}
}

function defaultSleep(ms) {
return new Promise((r) => setTimeout(r, ms));
}
11 changes: 11 additions & 0 deletions broadcaster/lib/config.js
Original file line number Diff line number Diff line change
Expand Up @@ -31,6 +31,14 @@
* with current stats (default off)
* BROADCASTER_LIVE_MS how often the pinned message is refreshed
* (default 300000)
* BROADCASTER_LIVE_PIN '0' = post the live message without pinning it,
* for a chat where the bot may not pin (default 1)
* BROADCASTER_COMMANDS '1' = answer /pool_status in TELEGRAM_CHAT_ID with
* the current stats (default off). Reads messages
* with getUpdates, so nothing else may poll this
* bot and it must have no webhook.
* BROADCASTER_COMMAND_COOLDOWN_SEC at most one reply per chat this often
* (default 30), so the command cannot flood a chat
* BROADCASTER_SUMMARY_HOURS also post the current stats as a NEW message
* every N hours (1-168), on UTC boundaries: 6 means
* 00, 06, 12 and 18 UTC (exact for any N that
Expand Down Expand Up @@ -90,6 +98,9 @@ export function loadConfig() {
pollMs: num('BROADCASTER_POLL_MS', 60000, { min: 1000 }),
digestHour: digestRaw === 'off' ? null : num('DIGEST_UTC_HOUR', 12, { min: 0, max: 23 }),
live: process.env.BROADCASTER_LIVE === '1',
livePin: process.env.BROADCASTER_LIVE_PIN !== '0',
commands: process.env.BROADCASTER_COMMANDS === '1',
commandCooldownSec: num('BROADCASTER_COMMAND_COOLDOWN_SEC', 30, { min: 0 }),
liveMs: num('BROADCASTER_LIVE_MS', 300000, { min: 60000 }),
summaryHours: summaryRaw == null || summaryRaw === 'off' || summaryRaw === '0'
? null : num('BROADCASTER_SUMMARY_HOURS', null, { min: 1, max: 168 }),
Expand Down
Loading
Loading