mirror of
https://github.com/benjamin-wilson/public-pool.git
synced 2026-09-29 17:15:03 -07:00
update postgresql branch
This commit is contained in:
@@ -1,6 +1,6 @@
|
||||
import { Injectable, OnModuleInit } from '@nestjs/common';
|
||||
import { ConfigService } from '@nestjs/config';
|
||||
import { RPCClient } from 'rpc-bitcoin';
|
||||
import axios, { AxiosInstance } from 'axios';
|
||||
import { asyncScheduler, BehaviorSubject, delay, filter, from, interval, scheduled, shareReplay, startWith, Subject, switchMap } from 'rxjs';
|
||||
import { RpcBlockService } from '../ORM/rpc-block/rpc-block.service';
|
||||
import * as zmq from 'zeromq';
|
||||
@@ -13,10 +13,11 @@ import * as PGPubsub from 'pg-pubsub';
|
||||
export class BitcoinRpcService implements OnModuleInit {
|
||||
|
||||
|
||||
private client: RPCClient;
|
||||
private client: AxiosInstance;
|
||||
private _newBlockTemplate$: BehaviorSubject<IBlockTemplate> = new BehaviorSubject(undefined);
|
||||
private pubsubInstance: PGPubsub;
|
||||
private resetTemplateInterval$ = new Subject<void>();
|
||||
private rpcRequestId = 0;
|
||||
|
||||
public miningInfo: IMiningInfo;
|
||||
public newBlockTemplate$ = this._newBlockTemplate$.pipe(filter(block => block != null), shareReplay({ refCount: true, bufferSize: 1 }));
|
||||
@@ -39,9 +40,17 @@ export class BitcoinRpcService implements OnModuleInit {
|
||||
const port = parseInt(this.configService.get('BITCOIN_RPC_PORT'));
|
||||
const timeout = parseInt(this.configService.get('BITCOIN_RPC_TIMEOUT'));
|
||||
|
||||
this.client = new RPCClient({ url, port, timeout, user, pass });
|
||||
const baseURL = this.buildRpcUrl(url, port);
|
||||
this.client = axios.create({
|
||||
baseURL,
|
||||
timeout,
|
||||
auth: {
|
||||
username: user,
|
||||
password: pass
|
||||
}
|
||||
});
|
||||
|
||||
this.client.getrpcinfo().then((res) => {
|
||||
this.callRpc('getrpcinfo').then((res) => {
|
||||
console.log('Bitcoin RPC connected');
|
||||
}, () => {
|
||||
console.error('Could not reach RPC host');
|
||||
@@ -109,13 +118,13 @@ export class BitcoinRpcService implements OnModuleInit {
|
||||
|
||||
let blockTemplate: IBlockTemplate;
|
||||
while (blockTemplate == null) {
|
||||
blockTemplate = await this.client.getblocktemplate({
|
||||
template_request: {
|
||||
blockTemplate = await this.callRpc<IBlockTemplate>('getblocktemplate', [
|
||||
{
|
||||
rules: ['segwit'],
|
||||
mode: 'template',
|
||||
capabilities: ['serverlist', 'proposal']
|
||||
}
|
||||
});
|
||||
]);
|
||||
}
|
||||
|
||||
try {
|
||||
@@ -131,7 +140,7 @@ export class BitcoinRpcService implements OnModuleInit {
|
||||
|
||||
public async getMiningInfo(): Promise<IMiningInfo> {
|
||||
try {
|
||||
return await this.client.getmininginfo();
|
||||
return await this.callRpc<IMiningInfo>('getmininginfo');
|
||||
} catch (e) {
|
||||
console.error('Error getmininginfo', e.message);
|
||||
return null;
|
||||
@@ -142,9 +151,7 @@ export class BitcoinRpcService implements OnModuleInit {
|
||||
public async SUBMIT_BLOCK(hexdata: string): Promise<string> {
|
||||
let response: string = 'unknown';
|
||||
try {
|
||||
response = await this.client.submitblock({
|
||||
hexdata
|
||||
});
|
||||
response = await this.callRpc<string>('submitblock', [hexdata]);
|
||||
if (response == null) {
|
||||
response = 'SUCCESS!';
|
||||
}
|
||||
@@ -158,4 +165,28 @@ export class BitcoinRpcService implements OnModuleInit {
|
||||
return response;
|
||||
|
||||
}
|
||||
|
||||
private async callRpc<T>(method: string, params: unknown[] = []): Promise<T> {
|
||||
const response = await this.client.post('', {
|
||||
jsonrpc: '1.0',
|
||||
id: ++this.rpcRequestId,
|
||||
method,
|
||||
params
|
||||
});
|
||||
|
||||
if (response.data.error != null) {
|
||||
throw response.data.error;
|
||||
}
|
||||
|
||||
return response.data.result;
|
||||
}
|
||||
|
||||
private buildRpcUrl(url: string, port: number): string {
|
||||
const normalizedUrl = /^https?:\/\//i.test(url) ? url : `http://${url}`;
|
||||
const rpcUrl = new URL(normalizedUrl);
|
||||
if (Number.isFinite(port) && port > 0) {
|
||||
rpcUrl.port = port.toString();
|
||||
}
|
||||
return rpcUrl.toString();
|
||||
}
|
||||
}
|
||||
|
||||
@@ -0,0 +1,92 @@
|
||||
import { BehaviorSubject, firstValueFrom, skip } from 'rxjs';
|
||||
|
||||
import { MockRecording1 } from '../../test/models/MockRecording1';
|
||||
import { IBlockTemplate } from '../models/bitcoin-rpc/IBlockTemplate';
|
||||
import { StratumV1JobsService } from './stratum-v1-jobs.service';
|
||||
|
||||
describe('StratumV1JobsService', () => {
|
||||
let blockTemplate$: BehaviorSubject<IBlockTemplate>;
|
||||
let bitcoinRpcService: { newBlockTemplate$: any, miningInfo: { blocks: number } };
|
||||
let service: StratumV1JobsService;
|
||||
let consoleLogSpy: jest.SpyInstance;
|
||||
|
||||
const createTemplate = (height = MockRecording1.BLOCK_TEMPLATE.height): IBlockTemplate => ({
|
||||
...MockRecording1.BLOCK_TEMPLATE,
|
||||
height
|
||||
});
|
||||
|
||||
beforeEach(() => {
|
||||
jest.useFakeTimers();
|
||||
jest.setSystemTime(new Date(parseInt(MockRecording1.TIME, 16) * 1000));
|
||||
consoleLogSpy = jest.spyOn(console, 'log').mockImplementation(() => undefined);
|
||||
|
||||
blockTemplate$ = new BehaviorSubject(createTemplate());
|
||||
bitcoinRpcService = {
|
||||
newBlockTemplate$: blockTemplate$.asObservable(),
|
||||
miningInfo: { blocks: MockRecording1.BLOCK_TEMPLATE.height }
|
||||
};
|
||||
service = new StratumV1JobsService(bitcoinRpcService as any);
|
||||
});
|
||||
|
||||
afterEach(() => {
|
||||
consoleLogSpy.mockRestore();
|
||||
jest.useRealTimers();
|
||||
});
|
||||
|
||||
it('should create job templates from block templates', async () => {
|
||||
const jobTemplate = await firstValueFrom(service.newMiningJob$);
|
||||
|
||||
expect(jobTemplate.blockData).toEqual(expect.objectContaining({
|
||||
id: '1',
|
||||
height: MockRecording1.BLOCK_TEMPLATE.height,
|
||||
clearJobs: true,
|
||||
coinbasevalue: MockRecording1.BLOCK_TEMPLATE.coinbasevalue
|
||||
}));
|
||||
expect(jobTemplate.merkle_branch.length).toBeGreaterThan(0);
|
||||
expect(jobTemplate.block.transactions[0].ins[0].witness[0]).toHaveLength(32);
|
||||
expect(service.getJobTemplateById('1')).toBe(jobTemplate);
|
||||
});
|
||||
|
||||
it('should clear jobs when the block height changes', async () => {
|
||||
const firstTemplate = await firstValueFrom(service.newMiningJob$);
|
||||
service.addJob({ jobId: 'old-job', creation: Date.now() } as any);
|
||||
|
||||
bitcoinRpcService.miningInfo.blocks = MockRecording1.BLOCK_TEMPLATE.height + 1;
|
||||
const nextTemplate = firstValueFrom(service.newMiningJob$.pipe(skip(1)));
|
||||
blockTemplate$.next(createTemplate(MockRecording1.BLOCK_TEMPLATE.height + 1));
|
||||
const jobTemplate = await nextTemplate;
|
||||
|
||||
expect(jobTemplate.blockData.clearJobs).toBe(true);
|
||||
expect(service.getJobById('old-job')).toBeUndefined();
|
||||
expect(service.getJobTemplateById(firstTemplate.blockData.id)).toBeUndefined();
|
||||
expect(service.getJobTemplateById(jobTemplate.blockData.id)).toBe(jobTemplate);
|
||||
});
|
||||
|
||||
it('should age old jobs and templates after five minutes', async () => {
|
||||
await firstValueFrom(service.newMiningJob$);
|
||||
const oldCreation = Date.now() - (1000 * 60 * 11);
|
||||
service.jobs['old-job'] = { jobId: 'old-job', creation: oldCreation } as any;
|
||||
service.blocks['old-template'] = {
|
||||
blockData: { creation: oldCreation }
|
||||
} as any;
|
||||
|
||||
bitcoinRpcService.miningInfo.blocks = MockRecording1.BLOCK_TEMPLATE.height;
|
||||
const nextTemplate = firstValueFrom(service.newMiningJob$.pipe(skip(1)));
|
||||
blockTemplate$.next(createTemplate());
|
||||
const jobTemplate = await nextTemplate;
|
||||
|
||||
expect(jobTemplate.blockData.clearJobs).toBe(false);
|
||||
expect(service.getJobById('old-job')).toBeUndefined();
|
||||
expect(service.getJobTemplateById('old-template')).toBeUndefined();
|
||||
expect(service.getJobTemplateById(jobTemplate.blockData.id)).toBe(jobTemplate);
|
||||
});
|
||||
|
||||
it('should increment job ids when jobs are added', () => {
|
||||
expect(service.getNextId()).toBe('1');
|
||||
|
||||
service.addJob({ jobId: '1', creation: Date.now() } as any);
|
||||
|
||||
expect(service.getNextId()).toBe('2');
|
||||
expect(service.getJobById('1')).toEqual(expect.objectContaining({ jobId: '1' }));
|
||||
});
|
||||
});
|
||||
@@ -162,6 +162,32 @@ export class StratumV1JobsService {
|
||||
return this.blocks[jobTemplateId];
|
||||
}
|
||||
|
||||
public cleanup(clearJobs: boolean, now: number = Date.now()) {
|
||||
if (clearJobs) {
|
||||
this.blocks = {};
|
||||
this.jobs = {};
|
||||
return;
|
||||
}
|
||||
|
||||
let templatesDeleted = 0;
|
||||
let jobsDeleted = 0;
|
||||
|
||||
for (const templateId in this.blocks) {
|
||||
if (now - this.blocks[templateId].blockData.creation > (1000 * 60 * 5)) {
|
||||
delete this.blocks[templateId];
|
||||
templatesDeleted++;
|
||||
}
|
||||
}
|
||||
|
||||
for (const jobId in this.jobs) {
|
||||
if (now - this.jobs[jobId].creation > (1000 * 60 * 5)) {
|
||||
delete this.jobs[jobId];
|
||||
jobsDeleted++;
|
||||
}
|
||||
}
|
||||
//console.log(`Deleted ${templatesDeleted} templates and ${jobsDeleted} jobs.`)
|
||||
}
|
||||
|
||||
public addJob(job: MiningJob) {
|
||||
this.jobs[job.jobId] = job;
|
||||
this.latestJobId++;
|
||||
@@ -178,5 +204,4 @@ export class StratumV1JobsService {
|
||||
return this.latestJobId.toString(16);
|
||||
}
|
||||
|
||||
|
||||
}
|
||||
|
||||
@@ -0,0 +1,78 @@
|
||||
import { StratumV1Service } from './stratum-v1.service';
|
||||
|
||||
describe('StratumV1Service', () => {
|
||||
const originalMaster = process.env.MASTER;
|
||||
const originalStratumPorts = process.env.STRATUM_PORTS;
|
||||
const originalStratumSecure = process.env.STRATUM_SECURE;
|
||||
const originalSecureStratumPorts = process.env.SECURE_STRATUM_PORTS;
|
||||
|
||||
let service: StratumV1Service;
|
||||
let clientService;
|
||||
let consoleLogSpy: jest.SpyInstance;
|
||||
|
||||
beforeEach(() => {
|
||||
jest.useFakeTimers();
|
||||
clientService = {
|
||||
deleteAll: jest.fn().mockResolvedValue(undefined)
|
||||
};
|
||||
service = new StratumV1Service(
|
||||
{} as any,
|
||||
clientService,
|
||||
{} as any,
|
||||
{} as any,
|
||||
{} as any,
|
||||
{} as any,
|
||||
{} as any,
|
||||
{} as any
|
||||
);
|
||||
consoleLogSpy = jest.spyOn(console, 'log').mockImplementation(() => undefined);
|
||||
});
|
||||
|
||||
afterEach(() => {
|
||||
restoreEnv('MASTER', originalMaster);
|
||||
restoreEnv('STRATUM_PORTS', originalStratumPorts);
|
||||
restoreEnv('STRATUM_SECURE', originalStratumSecure);
|
||||
restoreEnv('SECURE_STRATUM_PORTS', originalSecureStratumPorts);
|
||||
consoleLogSpy.mockRestore();
|
||||
jest.useRealTimers();
|
||||
});
|
||||
|
||||
it('should skip Stratum listeners in the master process', async () => {
|
||||
process.env.MASTER = 'true';
|
||||
const startSocketServerSpy = jest.spyOn(service as any, 'startSocketServer');
|
||||
const startSecureSocketServerSpy = jest.spyOn(service as any, 'startSecureSocketServer');
|
||||
|
||||
await service.onModuleInit();
|
||||
jest.runOnlyPendingTimers();
|
||||
|
||||
expect(clientService.deleteAll).toHaveBeenCalled();
|
||||
expect(startSocketServerSpy).not.toHaveBeenCalled();
|
||||
expect(startSecureSocketServerSpy).not.toHaveBeenCalled();
|
||||
expect(consoleLogSpy).toHaveBeenCalledWith('Master process skipping Stratum socket listeners');
|
||||
});
|
||||
|
||||
it('should start Stratum listeners in worker processes', async () => {
|
||||
process.env.MASTER = 'false';
|
||||
process.env.STRATUM_PORTS = '3333,3334';
|
||||
process.env.STRATUM_SECURE = 'true';
|
||||
process.env.SECURE_STRATUM_PORTS = '4333';
|
||||
const startSocketServerSpy = jest.spyOn(service as any, 'startSocketServer').mockImplementation(() => undefined);
|
||||
const startSecureSocketServerSpy = jest.spyOn(service as any, 'startSecureSocketServer').mockImplementation(() => undefined);
|
||||
|
||||
await service.onModuleInit();
|
||||
jest.advanceTimersByTime(10000);
|
||||
|
||||
expect(clientService.deleteAll).not.toHaveBeenCalled();
|
||||
expect(startSocketServerSpy).toHaveBeenCalledWith(3333);
|
||||
expect(startSocketServerSpy).toHaveBeenCalledWith(3334);
|
||||
expect(startSecureSocketServerSpy).toHaveBeenCalledWith(4333);
|
||||
});
|
||||
|
||||
function restoreEnv(key: string, value: string | undefined) {
|
||||
if (value == null) {
|
||||
delete process.env[key];
|
||||
return;
|
||||
}
|
||||
process.env[key] = value;
|
||||
}
|
||||
});
|
||||
@@ -42,6 +42,8 @@ export class StratumV1Service implements OnModuleInit {
|
||||
|
||||
if (process.env.MASTER == 'true') {
|
||||
await this.clientService.deleteAll();
|
||||
console.log('Master process skipping Stratum socket listeners');
|
||||
return;
|
||||
}
|
||||
|
||||
// wait for all the other processes to init for an even connection distribution
|
||||
@@ -208,4 +210,4 @@ export class StratumV1Service implements OnModuleInit {
|
||||
}
|
||||
|
||||
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1,8 +1,8 @@
|
||||
import { Injectable, OnModuleInit } from '@nestjs/common';
|
||||
import { ConfigService } from '@nestjs/config';
|
||||
import axios, { AxiosInstance } from 'axios';
|
||||
import { validate } from 'bitcoin-address-validation';
|
||||
import { Block } from 'bitcoinjs-lib';
|
||||
import * as TelegramBot from 'node-telegram-bot-api';
|
||||
|
||||
import { TelegramSubscriptionsService } from '../ORM/telegram-subscriptions/telegram-subscriptions.service';
|
||||
|
||||
@@ -10,7 +10,9 @@ import { TelegramSubscriptionsService } from '../ORM/telegram-subscriptions/tele
|
||||
@Injectable()
|
||||
export class TelegramService implements OnModuleInit {
|
||||
|
||||
private bot: TelegramBot;
|
||||
private bot: AxiosInstance;
|
||||
private updateOffset = 0;
|
||||
private pollingTimer: NodeJS.Timeout;
|
||||
|
||||
constructor(
|
||||
private readonly configService: ConfigService,
|
||||
@@ -20,7 +22,10 @@ export class TelegramService implements OnModuleInit {
|
||||
if (token == null || token.length < 1) {
|
||||
return;
|
||||
}
|
||||
this.bot = new TelegramBot(token, { polling: true });
|
||||
this.bot = axios.create({
|
||||
baseURL: `https://api.telegram.org/bot${token}/`,
|
||||
timeout: 10000
|
||||
});
|
||||
console.log('Telegram bot init');
|
||||
|
||||
|
||||
@@ -32,23 +37,10 @@ export class TelegramService implements OnModuleInit {
|
||||
return;
|
||||
}
|
||||
|
||||
this.bot.onText(/\/subscribe/, async (msg) => {
|
||||
const address = msg.text.split('/subscribe ')[1];
|
||||
if (validate(address) == false) {
|
||||
this.bot.sendMessage(msg.chat.id, "Invalid address.");
|
||||
return;
|
||||
}
|
||||
await this.telegramSubscriptionsService.saveSubscription(msg.chat.id, address);
|
||||
this.bot.sendMessage(msg.chat.id, "Subscribed!");
|
||||
});
|
||||
|
||||
this.bot.onText(/\/start/, (msg) => {
|
||||
this.bot.sendMessage(msg.chat.id, "Welcome to the public-pool bot. /subscribe <address> to get notified.");
|
||||
});
|
||||
|
||||
this.bot.on('message', (msg) => {
|
||||
console.log(msg);
|
||||
});
|
||||
await this.pollUpdates();
|
||||
this.pollingTimer = setInterval(async () => {
|
||||
await this.pollUpdates();
|
||||
}, 2000);
|
||||
}
|
||||
|
||||
public async notifySubscribersBlockFound(address: string, height: number, block: Block, message: string) {
|
||||
@@ -57,8 +49,57 @@ export class TelegramService implements OnModuleInit {
|
||||
}
|
||||
|
||||
const subscribers = await this.telegramSubscriptionsService.getSubscriptions(address);
|
||||
subscribers.forEach(subscriber => {
|
||||
this.bot.sendMessage(subscriber.telegramChatId, `Block Found! Result: ${message}, Height: ${height}`);
|
||||
await Promise.all(subscribers.map(subscriber => {
|
||||
return this.sendMessage(subscriber.telegramChatId, `Block Found! Result: ${message}, Height: ${height}`);
|
||||
}));
|
||||
}
|
||||
|
||||
private async pollUpdates() {
|
||||
try {
|
||||
const response = await this.bot.get('getUpdates', {
|
||||
params: {
|
||||
offset: this.updateOffset,
|
||||
timeout: 0
|
||||
}
|
||||
});
|
||||
|
||||
for (const update of response.data.result ?? []) {
|
||||
this.updateOffset = update.update_id + 1;
|
||||
await this.handleMessage(update.message);
|
||||
}
|
||||
} catch (e) {
|
||||
console.error('Telegram polling failed', e.message);
|
||||
}
|
||||
}
|
||||
|
||||
private async handleMessage(msg: any) {
|
||||
if (msg?.text == null) {
|
||||
return;
|
||||
}
|
||||
|
||||
if (msg.text.startsWith('/subscribe')) {
|
||||
const address = msg.text.split('/subscribe ')[1];
|
||||
if (validate(address) == false) {
|
||||
await this.sendMessage(msg.chat.id, 'Invalid address.');
|
||||
return;
|
||||
}
|
||||
await this.telegramSubscriptionsService.saveSubscription(msg.chat.id, address);
|
||||
await this.sendMessage(msg.chat.id, 'Subscribed!');
|
||||
return;
|
||||
}
|
||||
|
||||
if (msg.text.startsWith('/start')) {
|
||||
await this.sendMessage(msg.chat.id, 'Welcome to the public-pool bot. /subscribe <address> to get notified.');
|
||||
return;
|
||||
}
|
||||
|
||||
console.log(msg);
|
||||
}
|
||||
|
||||
private async sendMessage(chatId: number | string, text: string) {
|
||||
await this.bot.post('sendMessage', {
|
||||
chat_id: chatId,
|
||||
text
|
||||
});
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user