mirror of
https://github.com/benjamin-wilson/public-pool.git
synced 2026-09-29 17:15:03 -07:00
client unit tests
This commit is contained in:
@@ -28,7 +28,6 @@ export class MiningJob {
|
||||
public blockTemplate: IBlockTemplate,
|
||||
public clean_jobs: boolean) {
|
||||
|
||||
console.log(JSON.stringify(blockTemplate))
|
||||
|
||||
this.jobId = id;
|
||||
this.block.prevHash = this.convertToLittleEndian(blockTemplate.previousblockhash);
|
||||
|
||||
@@ -1,21 +1,44 @@
|
||||
import { ConfigService } from '@nestjs/config';
|
||||
import { Test, TestingModule } from '@nestjs/testing';
|
||||
import { TypeOrmModule } from '@nestjs/typeorm';
|
||||
import { PromiseSocket } from 'promise-socket';
|
||||
import { BehaviorSubject } from 'rxjs';
|
||||
import { DataSource } from 'typeorm';
|
||||
|
||||
import { MockRecording1 } from '../../test/models/MockRecording1';
|
||||
import { BlocksService } from '../ORM/blocks/blocks.service';
|
||||
import { ClientStatisticsEntity } from '../ORM/client-statistics/client-statistics.entity';
|
||||
import { ClientStatisticsModule } from '../ORM/client-statistics/client-statistics.module';
|
||||
import { ClientStatisticsService } from '../ORM/client-statistics/client-statistics.service';
|
||||
import { ClientEntity } from '../ORM/client/client.entity';
|
||||
import { ClientModule } from '../ORM/client/client.module';
|
||||
import { ClientService } from '../ORM/client/client.service';
|
||||
import { BitcoinRpcService } from '../services/bitcoin-rpc.service';
|
||||
import { BitcoinRpcService as MockBitcoinRpcService } from '../services/bitcoin-rpc.service';
|
||||
import { BlockTemplateService } from '../services/block-template.service';
|
||||
import { NotificationService } from '../services/notification.service';
|
||||
import { StratumV1JobsService } from '../services/stratum-v1-jobs.service';
|
||||
import { IMiningInfo } from './bitcoin-rpc/IMiningInfo';
|
||||
import { StratumV1Client } from './StratumV1Client';
|
||||
|
||||
|
||||
|
||||
|
||||
|
||||
jest.mock('../services/bitcoin-rpc.service')
|
||||
|
||||
jest.mock('./validators/bitcoin-address.validator', () => ({
|
||||
IsBitcoinAddress() {
|
||||
return jest.fn();
|
||||
},
|
||||
}));
|
||||
|
||||
|
||||
describe('StratumV1Client', () => {
|
||||
|
||||
let promiseSocket: PromiseSocket<any> = new PromiseSocket();
|
||||
let stratumV1JobsService: StratumV1JobsService;
|
||||
let bitcoinRpcService: BitcoinRpcService;
|
||||
|
||||
let promiseSocket: PromiseSocket<any>;
|
||||
let stratumV1JobsService: StratumV1JobsService = new StratumV1JobsService();
|
||||
let bitcoinRpcService: MockBitcoinRpcService;
|
||||
let blockTemplateService: BlockTemplateService;
|
||||
let clientService: ClientService;
|
||||
let clientStatisticsService: ClientStatisticsService;
|
||||
@@ -25,10 +48,79 @@ describe('StratumV1Client', () => {
|
||||
|
||||
let client: StratumV1Client;
|
||||
|
||||
let socketEmitter: (data: Buffer) => void;
|
||||
|
||||
let newBlockEmitter: BehaviorSubject<IMiningInfo> = new BehaviorSubject(null);
|
||||
|
||||
let moduleRef: TestingModule;
|
||||
|
||||
beforeAll(async () => {
|
||||
moduleRef = await Test.createTestingModule({
|
||||
imports: [
|
||||
TypeOrmModule.forRoot({
|
||||
type: 'sqlite',
|
||||
database: './DB/public-pool.test.sqlite',
|
||||
synchronize: true,
|
||||
autoLoadEntities: true,
|
||||
cache: true,
|
||||
logging: false
|
||||
}),
|
||||
ClientModule,
|
||||
ClientStatisticsModule
|
||||
],
|
||||
providers: [
|
||||
{
|
||||
provide: ConfigService,
|
||||
useValue: {
|
||||
get: jest.fn((key: string) => {
|
||||
switch (key) {
|
||||
case 'DEV_FEE_ADDRESS':
|
||||
return 'tb1qumezefzdeqqwn5zfvgdrhxjzc5ylr39uhuxcz4';
|
||||
case 'NETWORK':
|
||||
return 'testnet';
|
||||
}
|
||||
return null;
|
||||
})
|
||||
}
|
||||
}
|
||||
],
|
||||
}).compile();
|
||||
|
||||
|
||||
})
|
||||
|
||||
|
||||
beforeEach(async () => {
|
||||
|
||||
jest.spyOn(promiseSocket.socket, 'on').mockImplementation();
|
||||
|
||||
|
||||
clientService = moduleRef.get<ClientService>(ClientService);
|
||||
|
||||
const dataSource = moduleRef.get<DataSource>(DataSource);
|
||||
|
||||
dataSource.getRepository(ClientEntity).delete({});
|
||||
dataSource.getRepository(ClientStatisticsEntity).delete({});
|
||||
|
||||
|
||||
clientStatisticsService = moduleRef.get<ClientStatisticsService>(ClientStatisticsService);
|
||||
|
||||
configService = moduleRef.get<ConfigService>(ConfigService);
|
||||
|
||||
bitcoinRpcService = new MockBitcoinRpcService(null);
|
||||
|
||||
jest.spyOn(bitcoinRpcService, 'getBlockTemplate').mockReturnValue(Promise.resolve(MockRecording1.BLOCK_TEMPLATE));
|
||||
bitcoinRpcService.newBlock$ = newBlockEmitter.asObservable();
|
||||
|
||||
blockTemplateService = new BlockTemplateService(bitcoinRpcService);
|
||||
|
||||
|
||||
promiseSocket = new PromiseSocket();
|
||||
jest.spyOn(promiseSocket.socket, 'on').mockImplementation((event: string, fn: (data: Buffer) => void) => {
|
||||
socketEmitter = fn;
|
||||
});
|
||||
|
||||
promiseSocket.end = jest.fn();
|
||||
|
||||
|
||||
client = new StratumV1Client(
|
||||
promiseSocket,
|
||||
@@ -42,12 +134,126 @@ describe('StratumV1Client', () => {
|
||||
configService
|
||||
);
|
||||
|
||||
client.extraNonceAndSessionId = MockRecording1.EXTRA_NONCE;
|
||||
|
||||
jest.useFakeTimers({ advanceTimers: true })
|
||||
});
|
||||
|
||||
afterEach(async () => {
|
||||
client.destroy();
|
||||
jest.useRealTimers();
|
||||
})
|
||||
|
||||
|
||||
it('should subscribe to socket', () => {
|
||||
expect(promiseSocket.socket.on).toHaveBeenCalled();
|
||||
});
|
||||
|
||||
it('should close socket on invalid JSON', () => {
|
||||
socketEmitter(Buffer.from('INVALID'));
|
||||
expect(promiseSocket.end).toHaveBeenCalledTimes(1);
|
||||
});
|
||||
|
||||
it('should respond to mining.subscribe', async () => {
|
||||
jest.spyOn(promiseSocket, 'write').mockImplementation((data) => Promise.resolve(1));
|
||||
|
||||
expect(promiseSocket.socket.on).toHaveBeenCalled();
|
||||
socketEmitter(Buffer.from(MockRecording1.MINING_SUBSCRIBE));
|
||||
|
||||
await new Promise((r) => setTimeout(r, 1));
|
||||
|
||||
expect(promiseSocket.write).toHaveBeenCalledWith(`{"id":1,"error":null,"result":[[["mining.notify","${client.extraNonceAndSessionId}"]],"${client.extraNonceAndSessionId}",4]}\n`);
|
||||
|
||||
});
|
||||
|
||||
|
||||
it('should parse message', () => {
|
||||
it('should respond to mining.configure', async () => {
|
||||
|
||||
jest.spyOn(promiseSocket, 'write').mockImplementation((data) => Promise.resolve(1));
|
||||
|
||||
expect(promiseSocket.socket.on).toHaveBeenCalled();
|
||||
socketEmitter(Buffer.from(MockRecording1.MINING_CONFIGURE));
|
||||
await new Promise((r) => setTimeout(r, 1));
|
||||
expect(promiseSocket.write).toHaveBeenCalledWith(`{"id":2,"error":null,"result":{"version-rolling":true,"version-rolling.mask":"1fffe000"}}\n`);
|
||||
});
|
||||
|
||||
it('should respond to mining.authorize', async () => {
|
||||
|
||||
jest.spyOn(promiseSocket, 'write').mockImplementation((data) => Promise.resolve(1));
|
||||
|
||||
expect(promiseSocket.socket.on).toHaveBeenCalled();
|
||||
socketEmitter(Buffer.from(MockRecording1.MINING_AUTHORIZE));
|
||||
await new Promise((r) => setTimeout(r, 1));
|
||||
expect(promiseSocket.write).toHaveBeenCalledWith('{"id":3,"error":null,"result":true}\n');
|
||||
});
|
||||
|
||||
it('should respond to mining.suggest_difficulty', async () => {
|
||||
jest.spyOn(promiseSocket, 'write').mockImplementation((data) => Promise.resolve(1));
|
||||
|
||||
expect(promiseSocket.socket.on).toHaveBeenCalled();
|
||||
socketEmitter(Buffer.from(MockRecording1.MINING_SUGGEST_DIFFICULTY));
|
||||
await new Promise((r) => setTimeout(r, 1));
|
||||
expect(promiseSocket.write).toHaveBeenCalledWith(`{"id":4,"method":"mining.set_difficulty","params":[512]}\n`);
|
||||
});
|
||||
|
||||
it('should set difficulty', async () => {
|
||||
jest.spyOn(promiseSocket, 'write').mockImplementation((data) => Promise.resolve(1));
|
||||
|
||||
socketEmitter(Buffer.from(MockRecording1.MINING_SUBSCRIBE));
|
||||
socketEmitter(Buffer.from(MockRecording1.MINING_AUTHORIZE));
|
||||
await new Promise((r) => setTimeout(r, 100));
|
||||
|
||||
expect(promiseSocket.write).toHaveBeenCalledWith(`{"id":null,"method":"mining.set_difficulty","params":[32768]}\n`);
|
||||
|
||||
});
|
||||
|
||||
it('should save client', async () => {
|
||||
jest.spyOn(promiseSocket, 'write').mockImplementation((data) => Promise.resolve(1));
|
||||
|
||||
socketEmitter(Buffer.from(MockRecording1.MINING_SUBSCRIBE));
|
||||
socketEmitter(Buffer.from(MockRecording1.MINING_AUTHORIZE));
|
||||
await new Promise((r) => setTimeout(r, 100));
|
||||
|
||||
const clientCount = await clientService.connectedClientCount();
|
||||
expect(clientCount).toBe(1);
|
||||
|
||||
});
|
||||
|
||||
|
||||
|
||||
|
||||
it('should send job and accept submission', async () => {
|
||||
|
||||
|
||||
|
||||
const date = new Date(parseInt(MockRecording1.TIME, 16) * 1000);
|
||||
|
||||
|
||||
jest.setSystemTime(date);
|
||||
|
||||
jest.spyOn(promiseSocket, 'write').mockImplementation((data) => Promise.resolve(1));
|
||||
|
||||
|
||||
socketEmitter(Buffer.from(MockRecording1.MINING_SUBSCRIBE));
|
||||
socketEmitter(Buffer.from(MockRecording1.MINING_SUGGEST_DIFFICULTY));
|
||||
socketEmitter(Buffer.from(MockRecording1.MINING_AUTHORIZE));
|
||||
|
||||
|
||||
|
||||
await new Promise((r) => setTimeout(r, 100));
|
||||
|
||||
|
||||
|
||||
|
||||
expect(promiseSocket.write).lastCalledWith(`{"id":null,"method":"mining.notify","params":["3","171592f223740e92d223f6e68bff25279af7ac4f2246451e0000000200000000","02000000010000000000000000000000000000000000000000000000000000000000000000ffffffff1903c943255c7075626c69632d706f6f6c5c","ffffffff037a90000000000000160014e6f22ca44dc800e9d049621a3b9a42c509f1c4bc3b0f250000000000160014e6f22ca44dc800e9d049621a3b9a42c509f1c4bc0000000000000000266a24aa21a9edbd3d1d916aa0b57326a2d88ebe1b68a1d7c48585f26d8335fe6a94b62755f64c00000000",["175335649d5e8746982969ec88f52e85ac9917106fba5468e699c8879ab974a1","d5644ab3e708c54cd68dc5aedc92b8d3037449687f92ec41ed6e37673d969d4a","5c9ec187517edc0698556cca5ce27e54c96acb014770599ed9df4d4937fbf2b0"],"20000000","192495f8","${MockRecording1.TIME}",false]}\n`);
|
||||
|
||||
|
||||
socketEmitter(Buffer.from(MockRecording1.MINING_SUBMIT));
|
||||
|
||||
jest.useRealTimers();
|
||||
await new Promise((r) => setTimeout(r, 100));
|
||||
|
||||
|
||||
});
|
||||
|
||||
|
||||
|
||||
@@ -44,7 +44,7 @@ export class StratumV1Client extends EasyUnsubscribe {
|
||||
private sessionDifficulty: number = 32768;
|
||||
private entity: ClientEntity;
|
||||
|
||||
public extraNonce: string;
|
||||
public extraNonceAndSessionId: string;
|
||||
|
||||
constructor(
|
||||
public readonly promiseSocket: PromiseSocket<Socket>,
|
||||
@@ -60,9 +60,9 @@ export class StratumV1Client extends EasyUnsubscribe {
|
||||
super();
|
||||
|
||||
this.statistics = new StratumV1ClientStatistics(this.clientStatisticsService);
|
||||
this.extraNonce = this.getRandomHexString();
|
||||
this.extraNonceAndSessionId = this.getRandomHexString();
|
||||
|
||||
console.log(`New client ID: : ${this.extraNonce}`);
|
||||
console.log(`New client ID: : ${this.extraNonceAndSessionId}`);
|
||||
|
||||
this.promiseSocket.socket.on('data', (data: Buffer) => {
|
||||
data.toString()
|
||||
@@ -81,7 +81,7 @@ export class StratumV1Client extends EasyUnsubscribe {
|
||||
|
||||
|
||||
private async handleMessage(message: string) {
|
||||
console.log(`Received from ${this.extraNonce}`, message);
|
||||
console.log(`Received from ${this.extraNonceAndSessionId}`, message);
|
||||
|
||||
// Parse the message and check if it's the initial subscription message
|
||||
let parsedMessage = null;
|
||||
@@ -111,7 +111,7 @@ export class StratumV1Client extends EasyUnsubscribe {
|
||||
if (errors.length === 0) {
|
||||
this.clientSubscription = subscriptionMessage;
|
||||
|
||||
await this.promiseSocket.write(JSON.stringify(this.clientSubscription.response(this.extraNonce)) + '\n');
|
||||
await this.promiseSocket.write(JSON.stringify(this.clientSubscription.response(this.extraNonceAndSessionId)) + '\n');
|
||||
} else {
|
||||
const err = new StratumErrorMessage(
|
||||
subscriptionMessage.id,
|
||||
@@ -257,16 +257,18 @@ export class StratumV1Client extends EasyUnsubscribe {
|
||||
&& this.clientAuthorization != null
|
||||
&& this.stratumInitialized == false) {
|
||||
|
||||
this.stratumInitialized = true;
|
||||
|
||||
if (this.clientSuggestedDifficulty == null) {
|
||||
console.log(`Setting difficulty to ${this.sessionDifficulty}`)
|
||||
const setDifficulty = JSON.stringify(new SuggestDifficulty().response(this.sessionDifficulty));
|
||||
await this.promiseSocket.write(setDifficulty + '\n');
|
||||
}
|
||||
|
||||
this.stratumInitialized = true;
|
||||
|
||||
|
||||
this.entity = await this.clientService.save({
|
||||
sessionId: this.extraNonce,
|
||||
sessionId: this.extraNonceAndSessionId,
|
||||
address: this.clientAuthorization.address,
|
||||
clientName: this.clientAuthorization.worker,
|
||||
startTime: new Date(),
|
||||
@@ -275,8 +277,13 @@ export class StratumV1Client extends EasyUnsubscribe {
|
||||
let lastIntervalCount = undefined;
|
||||
let skipNext = false;
|
||||
combineLatest([this.blockTemplateService.currentBlockTemplate$, interval(60000).pipe(startWith(-1))])
|
||||
.pipe(takeUntil(this.easyUnsubscribe))
|
||||
.pipe(
|
||||
takeUntil(this.easyUnsubscribe)
|
||||
)
|
||||
.subscribe(async ([{ blockTemplate }, interValCount]) => {
|
||||
|
||||
|
||||
|
||||
let clearJobs = false;
|
||||
if (lastIntervalCount === interValCount) {
|
||||
clearJobs = true;
|
||||
@@ -303,7 +310,7 @@ export class StratumV1Client extends EasyUnsubscribe {
|
||||
|
||||
private async sendNewMiningJob(blockTemplate: IBlockTemplate, clearJobs: boolean) {
|
||||
|
||||
const hashRate = await this.clientStatisticsService.getHashRateForSession(this.clientAuthorization.address, this.clientAuthorization.worker, this.extraNonce);
|
||||
const hashRate = await this.clientStatisticsService.getHashRateForSession(this.clientAuthorization.address, this.clientAuthorization.worker, this.extraNonceAndSessionId);
|
||||
|
||||
let payoutInformation;
|
||||
//10Th/s
|
||||
@@ -320,6 +327,7 @@ export class StratumV1Client extends EasyUnsubscribe {
|
||||
];
|
||||
}
|
||||
|
||||
|
||||
const job = new MiningJob(
|
||||
this.configService.get('NETWORK') === 'mainnet' ? bitcoinjs.networks.bitcoin : bitcoinjs.networks.testnet,
|
||||
this.stratumV1JobsService.getNextId(),
|
||||
@@ -331,9 +339,14 @@ export class StratumV1Client extends EasyUnsubscribe {
|
||||
this.stratumV1JobsService.addJob(job, clearJobs);
|
||||
|
||||
|
||||
await this.promiseSocket.write(job.response());
|
||||
|
||||
console.log(`Sent new job to ${this.clientAuthorization.worker}.${this.extraNonce}. (clearJobs: ${clearJobs}, fee?: ${!noFee})`)
|
||||
try {
|
||||
await this.promiseSocket.write(job.response());
|
||||
} catch (e) {
|
||||
await this.promiseSocket.end();
|
||||
}
|
||||
|
||||
console.log(`Sent new job to ${this.clientAuthorization.worker}.${this.extraNonceAndSessionId}. (clearJobs: ${clearJobs}, fee?: ${!noFee})`)
|
||||
|
||||
}
|
||||
|
||||
@@ -354,14 +367,14 @@ export class StratumV1Client extends EasyUnsubscribe {
|
||||
const updatedJobBlock = job.copyAndUpdateBlock(
|
||||
parseInt(submission.versionMask, 16),
|
||||
parseInt(submission.nonce, 16),
|
||||
this.extraNonce,
|
||||
this.extraNonceAndSessionId,
|
||||
submission.extraNonce2,
|
||||
parseInt(submission.ntime, 16)
|
||||
);
|
||||
const header = updatedJobBlock.toBuffer(true);
|
||||
const { submissionDifficulty, submissionHash } = this.calculateDifficulty(header);
|
||||
|
||||
console.log(`DIFF: ${submissionDifficulty} of ${this.sessionDifficulty} from ${this.clientAuthorization.worker + '.' + this.extraNonce}`);
|
||||
console.log(`DIFF: ${submissionDifficulty} of ${this.sessionDifficulty} from ${this.clientAuthorization.worker + '.' + this.extraNonceAndSessionId}`);
|
||||
console.log(`Header: ${header.toString('hex')}`);
|
||||
|
||||
if (submissionDifficulty >= this.sessionDifficulty) {
|
||||
@@ -374,7 +387,7 @@ export class StratumV1Client extends EasyUnsubscribe {
|
||||
height: job.blockTemplate.height,
|
||||
minerAddress: this.clientAuthorization.address,
|
||||
worker: this.clientAuthorization.worker,
|
||||
sessionId: this.extraNonce,
|
||||
sessionId: this.extraNonceAndSessionId,
|
||||
blockData: blockHex
|
||||
});
|
||||
await this.notificationService.notifySubscribersBlockFound(this.clientAuthorization.address, job.blockTemplate.height, updatedJobBlock, result);
|
||||
@@ -382,6 +395,7 @@ export class StratumV1Client extends EasyUnsubscribe {
|
||||
try {
|
||||
await this.statistics.addSubmission(this.entity, submissionHash, this.sessionDifficulty);
|
||||
} catch (e) {
|
||||
console.log(e);
|
||||
const err = new StratumErrorMessage(
|
||||
submission.id,
|
||||
eStratumErrorCode.DuplicateShare,
|
||||
@@ -392,7 +406,7 @@ export class StratumV1Client extends EasyUnsubscribe {
|
||||
}
|
||||
|
||||
if (submissionDifficulty > this.entity.bestDifficulty) {
|
||||
await this.clientService.updateBestDifficulty(this.extraNonce, submissionDifficulty);
|
||||
await this.clientService.updateBestDifficulty(this.extraNonceAndSessionId, submissionDifficulty);
|
||||
this.entity.bestDifficulty = submissionDifficulty;
|
||||
}
|
||||
|
||||
|
||||
@@ -14,13 +14,13 @@ export interface IBlockTemplate {
|
||||
|
||||
vbavailable: { // (json object) set of pending, supported versionbit (BIP 9) softfork deployments
|
||||
rulename: number, // (numeric) identifies the bit number as indicating acceptance and readiness for the named softfork rule
|
||||
},
|
||||
} | {},
|
||||
vbrequired: number, // (numeric) bit mask of versionbits the server requires set in submissions
|
||||
previousblockhash: string, // (string) The hash of current highest block
|
||||
transactions: IBlockTemplateTx[], // (json array) contents of non-coinbase transactions that should be included in the next block
|
||||
coinbaseaux: { // (json object) data that should be included in the coinbase's scriptSig content
|
||||
key: string; //'hex', // (string) values must be in the coinbase (keys may be ignored)
|
||||
},
|
||||
} | {},
|
||||
coinbasevalue: number, // (numeric) maximum allowable input to coinbase transaction, including the generation award and transaction fees (in satoshis)
|
||||
longpollid: string, // (string) an id to include with a request to longpoll on an update to this template
|
||||
target: string, // (string) The hash target
|
||||
@@ -34,6 +34,6 @@ export interface IBlockTemplate {
|
||||
bits: string, // (string) compressed target of next block
|
||||
height: number, // (numeric) The height of the next block
|
||||
default_witness_commitment: string // (string, optional) a valid witness commitment for the unmodified block template
|
||||
|
||||
capabilities: string[]
|
||||
|
||||
}
|
||||
@@ -15,11 +15,11 @@ export class BitcoinRpcService {
|
||||
public newBlock$ = this._newBlock$.pipe(filter(block => block != null));
|
||||
|
||||
constructor(private readonly configService: ConfigService) {
|
||||
const url = configService.get('BITCOIN_RPC_URL');
|
||||
const user = configService.get('BITCOIN_RPC_USER');
|
||||
const pass = configService.get('BITCOIN_RPC_PASSWORD');
|
||||
const port = parseInt(configService.get('BITCOIN_RPC_PORT'));
|
||||
const timeout = parseInt(configService.get('BITCOIN_RPC_TIMEOUT'));
|
||||
const url = this.configService.get('BITCOIN_RPC_URL');
|
||||
const user = this.configService.get('BITCOIN_RPC_USER');
|
||||
const pass = this.configService.get('BITCOIN_RPC_PASSWORD');
|
||||
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 });
|
||||
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
import { Injectable } from '@nestjs/common';
|
||||
import { from, map, Observable, shareReplay, switchMap, tap } from 'rxjs';
|
||||
import { from, map, Observable, shareReplay, switchMap } from 'rxjs';
|
||||
|
||||
import { IBlockTemplate } from '../models/bitcoin-rpc/IBlockTemplate';
|
||||
import { BitcoinRpcService } from './bitcoin-rpc.service';
|
||||
@@ -7,14 +7,12 @@ import { BitcoinRpcService } from './bitcoin-rpc.service';
|
||||
@Injectable()
|
||||
export class BlockTemplateService {
|
||||
|
||||
public currentBlockTemplate: IBlockTemplate;
|
||||
|
||||
public currentBlockTemplate$: Observable<{ blockTemplate: IBlockTemplate }>;
|
||||
|
||||
constructor(private readonly bitcoinRpcService: BitcoinRpcService) {
|
||||
this.currentBlockTemplate$ = this.bitcoinRpcService.newBlock$.pipe(
|
||||
switchMap((miningInfo) => from(this.bitcoinRpcService.getBlockTemplate()).pipe(map(blockTemplate => { return { miningInfo, blockTemplate } }))),
|
||||
tap(({ blockTemplate }) => this.currentBlockTemplate = blockTemplate),
|
||||
shareReplay({ refCount: true, bufferSize: 1 })
|
||||
);
|
||||
}
|
||||
|
||||
@@ -62,16 +62,20 @@ export class StratumV1Service implements OnModuleInit {
|
||||
promiseSocket.socket.on('end', async (error: Error) => {
|
||||
// Handle socket disconnection
|
||||
client.destroy();
|
||||
await this.clientService.delete(client.extraNonce);
|
||||
promiseSocket.destroy();
|
||||
await this.clientService.delete(client.extraNonceAndSessionId);
|
||||
|
||||
const clientCount = await this.clientService.connectedClientCount();
|
||||
|
||||
console.log(`Client disconnected: ${promiseSocket.socket.remoteAddress}, ${clientCount} total clients`);
|
||||
});
|
||||
|
||||
promiseSocket.socket.on('error', async (error: Error) => {
|
||||
|
||||
client.destroy();
|
||||
await this.clientService.delete(client.extraNonce);
|
||||
promiseSocket.destroy();
|
||||
await this.clientService.delete(client.extraNonceAndSessionId);
|
||||
|
||||
const clientCount = await this.clientService.connectedClientCount();
|
||||
console.error(`Socket error:`, error);
|
||||
console.log(`Client disconnected: ${promiseSocket.socket.remoteAddress}, ${clientCount} total clients`);
|
||||
|
||||
Reference in New Issue
Block a user