From 65101a98a2ebce55a9ec72018af5d89742d27582 Mon Sep 17 00:00:00 2001 From: Ben Date: Mon, 13 Jul 2026 12:20:37 -0400 Subject: [PATCH] Block found notifications --- src/services/discord.service.ts | 4 +- src/services/notification.service.spec.ts | 147 +++++++++++++++++++ src/services/notification.service.ts | 94 +++++++++++- src/services/redis-messaging.service.spec.ts | 36 +++++ src/services/redis-messaging.service.ts | 53 +++++++ src/services/telegram.service.ts | 2 +- 6 files changed, 330 insertions(+), 6 deletions(-) create mode 100644 src/services/notification.service.spec.ts diff --git a/src/services/discord.service.ts b/src/services/discord.service.ts index 522f706..1ce4895 100644 --- a/src/services/discord.service.ts +++ b/src/services/discord.service.ts @@ -124,7 +124,7 @@ export class DiscordService implements OnModuleInit { } } - public async notifySubscribersBlockFound(height: number, block: Block, message: string) { + public async notifySubscribersBlockFound(height: number, block: Block | undefined, message: string) { if (process.env.MASTER == 'true') { if (this.bot == null) { return; @@ -135,4 +135,4 @@ export class DiscordService implements OnModuleInit { channel.send(`Block Found! Result: ${message}, Height: ${height}`); } } -} \ No newline at end of file +} diff --git a/src/services/notification.service.spec.ts b/src/services/notification.service.spec.ts new file mode 100644 index 0000000..4e69516 --- /dev/null +++ b/src/services/notification.service.spec.ts @@ -0,0 +1,147 @@ +import { Block } from 'bitcoinjs-lib'; + +import { NotificationService } from './notification.service'; + +describe('NotificationService', () => { + const originalMaster = process.env.MASTER; + let telegramService: { notifySubscribersBlockFound: jest.Mock }; + let discordService: { + notifyRestarted: jest.Mock; + notifySubscribersBlockFound: jest.Mock; + }; + let redisMessagingService: { + publishBlockFoundNotification: jest.Mock; + subscribeBlockFoundNotifications: jest.Mock; + }; + + beforeEach(() => { + telegramService = { + notifySubscribersBlockFound: jest.fn().mockResolvedValue(undefined), + }; + discordService = { + notifyRestarted: jest.fn().mockResolvedValue(undefined), + notifySubscribersBlockFound: jest.fn().mockResolvedValue(undefined), + }; + redisMessagingService = { + publishBlockFoundNotification: jest.fn().mockResolvedValue(true), + subscribeBlockFoundNotifications: jest.fn().mockResolvedValue(undefined), + }; + }); + + afterEach(() => { + process.env.MASTER = originalMaster; + jest.clearAllMocks(); + }); + + it('publishes block found notifications from worker processes', async () => { + process.env.MASTER = 'false'; + const service = createService(); + + await service.notifySubscribersBlockFound( + 'bc1qminer', + 900001, + createBlock('11'.repeat(32)), + 'accepted', + ); + + expect(redisMessagingService.publishBlockFoundNotification).toHaveBeenCalledWith({ + schemaVersion: 1, + eventId: 'block-found:900001:1111111111111111111111111111111111111111111111111111111111111111:bc1qminer', + address: 'bc1qminer', + height: 900001, + blockHash: '11'.repeat(32), + message: 'accepted', + publishedAtMs: expect.any(Number), + }); + expect(discordService.notifySubscribersBlockFound).not.toHaveBeenCalled(); + expect(telegramService.notifySubscribersBlockFound).not.toHaveBeenCalled(); + }); + + it('subscribes on the master process and dispatches received notifications', async () => { + process.env.MASTER = 'true'; + const service = createService(); + + await service.onModuleInit(); + const handler = redisMessagingService.subscribeBlockFoundNotifications.mock.calls[0][0]; + await handler({ + schemaVersion: 1, + eventId: 'block-found:900001:blockhash:bc1qminer', + address: 'bc1qminer', + height: 900001, + blockHash: 'blockhash', + message: 'accepted', + publishedAtMs: 123, + }); + + expect(redisMessagingService.subscribeBlockFoundNotifications).toHaveBeenCalledTimes(1); + expect(discordService.notifyRestarted).toHaveBeenCalledTimes(1); + expect(discordService.notifySubscribersBlockFound).toHaveBeenCalledWith( + 900001, + undefined, + 'accepted', + ); + expect(telegramService.notifySubscribersBlockFound).toHaveBeenCalledWith( + 'bc1qminer', + 900001, + undefined, + 'accepted', + ); + }); + + it('deduplicates repeated master notifications for the same block event', async () => { + process.env.MASTER = 'true'; + const service = createService(); + const notification = { + schemaVersion: 1 as const, + eventId: 'block-found:900001:blockhash:bc1qminer', + address: 'bc1qminer', + height: 900001, + blockHash: 'blockhash', + message: 'accepted', + publishedAtMs: 123, + }; + + await service.onModuleInit(); + const handler = redisMessagingService.subscribeBlockFoundNotifications.mock.calls[0][0]; + await handler(notification); + await handler(notification); + + expect(discordService.notifySubscribersBlockFound).toHaveBeenCalledTimes(1); + expect(telegramService.notifySubscribersBlockFound).toHaveBeenCalledTimes(1); + }); + + it('dispatches directly when called in the master process', async () => { + process.env.MASTER = 'true'; + const service = createService(); + const block = createBlock('22'.repeat(32)); + + await service.notifySubscribersBlockFound('bc1qminer', 900002, block, 'accepted'); + + expect(redisMessagingService.publishBlockFoundNotification).not.toHaveBeenCalled(); + expect(discordService.notifySubscribersBlockFound).toHaveBeenCalledWith( + 900002, + block, + 'accepted', + ); + expect(telegramService.notifySubscribersBlockFound).toHaveBeenCalledWith( + 'bc1qminer', + 900002, + block, + 'accepted', + ); + }); + + function createService(): NotificationService { + return new NotificationService( + telegramService as any, + discordService as any, + redisMessagingService as any, + ); + } +}); + +function createBlock(blockHash: string): Block { + return { + getId: () => blockHash, + } as unknown as Block; +} diff --git a/src/services/notification.service.ts b/src/services/notification.service.ts index 4b1355b..4c5a864 100644 --- a/src/services/notification.service.ts +++ b/src/services/notification.service.ts @@ -2,15 +2,20 @@ import { Injectable, OnModuleInit } from '@nestjs/common'; import { Block } from 'bitcoinjs-lib'; import { DiscordService } from './discord.service'; +import { BlockFoundNotification, RedisMessagingService } from './redis-messaging.service'; import { TelegramService } from './telegram.service'; +const BLOCK_NOTIFICATION_DEDUPE_TTL_MS = 10 * 60 * 1000; +const BLOCK_NOTIFICATION_DEDUPE_MAX_ENTRIES = 1000; @Injectable() export class NotificationService implements OnModuleInit { + private readonly processedBlockNotifications = new Map(); constructor( private readonly telegramService: TelegramService, - private readonly discordService: DiscordService + private readonly discordService: DiscordService, + private readonly redisMessagingService: RedisMessagingService ) { } async onModuleInit(): Promise { @@ -18,11 +23,94 @@ export class NotificationService implements OnModuleInit { return; } + await this.redisMessagingService.subscribeBlockFoundNotifications(async notification => { + await this.dispatchBlockFoundNotification(notification); + }); await this.discordService.notifyRestarted(); } public async notifySubscribersBlockFound(address: string, height: number, block: Block, message: string) { - await this.discordService.notifySubscribersBlockFound(height, block, message); - await this.telegramService.notifySubscribersBlockFound(address, height, block, message); + const blockHash = this.getBlockHash(block); + const notification: BlockFoundNotification = { + schemaVersion: 1, + eventId: this.createBlockFoundEventId(address, height, blockHash, message), + address, + height, + blockHash, + message, + publishedAtMs: Date.now(), + }; + + if (process.env.MASTER === 'true') { + await this.dispatchBlockFoundNotification(notification, block); + return; + } + + const published = await this.redisMessagingService.publishBlockFoundNotification(notification); + if (!published) { + console.error(`Unable to publish block found notification ${notification.eventId}`); + } + } + + private async dispatchBlockFoundNotification( + notification: BlockFoundNotification, + block?: Block, + ): Promise { + if (this.hasProcessedBlockNotification(notification.eventId)) { + return; + } + await this.discordService.notifySubscribersBlockFound(notification.height, block, notification.message); + await this.telegramService.notifySubscribersBlockFound( + notification.address, + notification.height, + block, + notification.message, + ); + } + + private createBlockFoundEventId( + address: string, + height: number, + blockHash: string, + message: string, + ): string { + if (blockHash !== 'unknown') { + return `block-found:${height}:${blockHash}:${address}`; + } + return `block-found:${height}:${address}:${message}`; + } + + private getBlockHash(block: Block): string { + try { + const maybeBlock = block as Block & { + getId?: () => string; + getHash?: () => Buffer; + }; + if (typeof maybeBlock.getId === 'function') { + const id = maybeBlock.getId(); + return typeof id === 'string' && id.length > 0 ? id : 'unknown'; + } + if (typeof maybeBlock.getHash === 'function') { + return Buffer.from(maybeBlock.getHash()).reverse().toString('hex'); + } + } catch { + return 'unknown'; + } + return 'unknown'; + } + + private hasProcessedBlockNotification(eventId: string): boolean { + const now = Date.now(); + for (const [key, processedAtMs] of this.processedBlockNotifications) { + if (now - processedAtMs > BLOCK_NOTIFICATION_DEDUPE_TTL_MS + || this.processedBlockNotifications.size > BLOCK_NOTIFICATION_DEDUPE_MAX_ENTRIES) { + this.processedBlockNotifications.delete(key); + } + } + if (this.processedBlockNotifications.has(eventId)) { + return true; + } + this.processedBlockNotifications.set(eventId, now); + return false; } } diff --git a/src/services/redis-messaging.service.spec.ts b/src/services/redis-messaging.service.spec.ts index 1be2fea..f2343a0 100644 --- a/src/services/redis-messaging.service.spec.ts +++ b/src/services/redis-messaging.service.spec.ts @@ -182,6 +182,42 @@ describe('RedisMessagingService', () => { consoleSpy.mockRestore(); }); + it('publishes validated block found notifications', async () => { + await service.connect(); + const handler = jest.fn().mockResolvedValue(undefined); + const notification = { + schemaVersion: 1 as const, + eventId: 'block-found:900001:blockhash:bc1qminer', + address: 'bc1qminer', + height: 900001, + blockHash: 'aa'.repeat(32), + message: 'accepted', + publishedAtMs: 123, + }; + + await service.subscribeBlockFoundNotifications(handler); + await expect(service.publishBlockFoundNotification(notification)).resolves.toBe(true); + + expect(handler).toHaveBeenCalledWith(notification); + }); + + it('ignores malformed block found notifications', async () => { + await service.connect(); + const handler = jest.fn(); + const consoleSpy = jest.spyOn(console, 'error').mockImplementation(() => undefined); + await service.subscribeBlockFoundNotifications(handler); + + await subscriptions.get('block-found.notification')!(JSON.stringify({ + schemaVersion: 1, + eventId: '', + height: 900001, + })); + + expect(handler).not.toHaveBeenCalled(); + expect(consoleSpy).toHaveBeenCalledWith(expect.stringContaining('Invalid Redis block found notification')); + consoleSpy.mockRestore(); + }); + it('stores, publishes, and replays compact SV1 bridge updates', async () => { await service.connect(); const handler = jest.fn().mockResolvedValue(undefined); diff --git a/src/services/redis-messaging.service.ts b/src/services/redis-messaging.service.ts index b1be3b5..60a3187 100644 --- a/src/services/redis-messaging.service.ts +++ b/src/services/redis-messaging.service.ts @@ -10,6 +10,7 @@ const MINING_INFO_CHANNEL = 'mining-info.updated'; const BLOCK_TEMPLATE_CHANNEL = 'block-template.updated'; const SV1_BRIDGE_CHANNEL = 'sv1-bridge.updated'; const SV1_PRESTAGE_CHANNEL = 'sv1-prestage.updated'; +const BLOCK_FOUND_NOTIFICATION_CHANNEL = 'block-found.notification'; const MINING_INFO_KEY = 'mining-info:latest'; const BLOCK_TEMPLATE_LATEST_KEY = 'block-template:latest'; const SV1_BRIDGE_LATEST_KEY = 'sv1-bridge:latest'; @@ -58,6 +59,16 @@ export interface Sv1PrestageUpdate { preparedAtMs: number; } +export interface BlockFoundNotification { + schemaVersion: 1; + eventId: string; + address: string; + height: number; + blockHash: string; + message: string; + publishedAtMs: number; +} + @Injectable() export class RedisMessagingService implements OnModuleInit, OnModuleDestroy { private publisher: RedisClientType; @@ -167,6 +178,31 @@ export class RedisMessagingService implements OnModuleInit, OnModuleDestroy { }); } + public async publishBlockFoundNotification(notification: BlockFoundNotification): Promise { + if (!await this.ensureConnected()) { + return false; + } + const serialized = JSON.stringify(notification); + const validated = this.parseBlockFoundNotification(serialized); + await this.publisher.publish(BLOCK_FOUND_NOTIFICATION_CHANNEL, JSON.stringify(validated)); + return true; + } + + public async subscribeBlockFoundNotifications( + handler: (notification: BlockFoundNotification) => Promise, + ): Promise { + if (!await this.ensureConnected()) { + return; + } + await this.subscriber.subscribe(BLOCK_FOUND_NOTIFICATION_CHANNEL, async message => { + try { + await handler(this.parseBlockFoundNotification(message)); + } catch (error) { + console.error(`Invalid Redis block found notification: ${error.message}`); + } + }); + } + public async publishSv1BridgeUpdate( update: Sv1BridgeUpdate, lane: 'urgent' | 'fallback' = 'urgent', @@ -532,6 +568,23 @@ export class RedisMessagingService implements OnModuleInit, OnModuleDestroy { return update as Sv1PrestageUpdate; } + private parseBlockFoundNotification(message: string): BlockFoundNotification { + const notification = JSON.parse(message) as Partial; + if (notification.schemaVersion !== 1 + || typeof notification.eventId !== 'string' + || notification.eventId.trim().length < 1 + || typeof notification.address !== 'string' + || notification.address.trim().length < 1 + || !Number.isInteger(notification.height) + || typeof notification.blockHash !== 'string' + || notification.blockHash.trim().length < 1 + || typeof notification.message !== 'string' + || !Number.isFinite(notification.publishedAtMs)) { + throw new Error('unsupported block found notification'); + } + return notification as BlockFoundNotification; + } + private async ensureConnected(): Promise { if (!this.connected) { try { diff --git a/src/services/telegram.service.ts b/src/services/telegram.service.ts index 2bdf78b..180a7c7 100644 --- a/src/services/telegram.service.ts +++ b/src/services/telegram.service.ts @@ -50,7 +50,7 @@ export class TelegramService implements OnModuleInit { }, 2000); } - public async notifySubscribersBlockFound(address: string, height: number, block: Block, message: string) { + public async notifySubscribersBlockFound(address: string, height: number, block: Block | undefined, message: string) { if (this.bot == null) { return; }