mirror of
https://github.com/benjamin-wilson/public-pool.git
synced 2026-09-29 09:05:06 -07:00
Block found notifications
This commit is contained in:
@@ -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}`);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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;
|
||||
}
|
||||
@@ -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<string, number>();
|
||||
|
||||
constructor(
|
||||
private readonly telegramService: TelegramService,
|
||||
private readonly discordService: DiscordService
|
||||
private readonly discordService: DiscordService,
|
||||
private readonly redisMessagingService: RedisMessagingService
|
||||
) { }
|
||||
|
||||
async onModuleInit(): Promise<void> {
|
||||
@@ -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<void> {
|
||||
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;
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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);
|
||||
|
||||
@@ -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<boolean> {
|
||||
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<void>,
|
||||
): Promise<void> {
|
||||
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<BlockFoundNotification>;
|
||||
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<boolean> {
|
||||
if (!this.connected) {
|
||||
try {
|
||||
|
||||
@@ -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;
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user