Optimize SV1 block notify fanout

This commit is contained in:
Ben
2026-07-13 11:55:26 -04:00
parent 3d58802feb
commit c60264bbe5
17 changed files with 1799 additions and 169 deletions
+51
View File
@@ -0,0 +1,51 @@
/* eslint-disable @typescript-eslint/no-var-requires */
describe('PM2 worker sizing', () => {
const originalApiWorkers = process.env.API_WORKERS;
const originalStratumWorkers = process.env.STRATUM_WORKERS;
afterEach(() => {
restoreEnv('API_WORKERS', originalApiWorkers);
restoreEnv('STRATUM_WORKERS', originalStratumWorkers);
jest.resetModules();
jest.unmock('os');
});
it('assigns remaining CPUs to Stratum workers in auto mode', () => {
process.env.API_WORKERS = '4';
process.env.STRATUM_WORKERS = 'auto';
jest.doMock('os', () => ({
availableParallelism: () => 14,
cpus: () => Array.from({ length: 14 }),
}));
const config = require('../ecosystem.config.js');
expect(config.apps.find((app) => app.name === 'api').instances).toBe(4);
expect(config.apps.find((app) => app.name === 'workers').instances).toBe(9);
});
it('preserves an explicit Stratum worker count', () => {
process.env.API_WORKERS = '2';
process.env.STRATUM_WORKERS = '7';
const config = require('../ecosystem.config.js');
expect(config.apps.find((app) => app.name === 'workers').instances).toBe(7);
});
it('rejects invalid fixed worker counts instead of silently starting no workers', () => {
process.env.STRATUM_WORKERS = 'many';
expect(() => require('../ecosystem.config.js')).toThrow(
'STRATUM_WORKERS must be a positive integer',
);
});
});
function restoreEnv(key: string, value: string | undefined) {
if (value == null) {
delete process.env[key];
return;
}
process.env[key] = value;
}
+128 -9
View File
@@ -20,10 +20,29 @@ export interface MiningJobOwnership {
payoutIdentity: string;
}
export interface MiningNotifyHeaderFields {
previousBlockHash: string;
version: string;
bits: string;
timestamp: string;
cleanJobs: boolean;
}
interface PreStagedNotifyLayout {
fields: MiningNotifyHeaderFields;
offsets: {
previousBlockHash: number;
version: number;
bits: number;
timestamp: number;
};
}
export class MiningJob {
private static readonly paymentScriptCache = new Map<string, Buffer>();
private static readonly paymentScriptCacheMaxEntries = 250_000;
private static readonly placeholderPrevHash = Buffer.alloc(32, 0);
private coinbaseTransaction: bitcoinjs.Transaction;
private coinbasePart1: string;
@@ -33,6 +52,7 @@ export class MiningJob {
private merkleBranchBuffers: Buffer[];
private miningNotifyResponse: string;
private miningNotifyResponseBuffer: Buffer;
private preStagedNotifyLayout: PreStagedNotifyLayout;
public jobTemplateId: string;
public tipKey: string;
@@ -50,6 +70,7 @@ export class MiningJob {
this.creation = new Date().getTime();
this.jobTemplateId = jobTemplate.blockData.id;
this.tipKey = jobTemplate.blockData.tipKey;
this.networkDifficulty = jobTemplate.blockData.networkDifficulty;
this.merkleBranchBuffers = jobTemplate.merkle_branch.map(branch => Buffer.from(branch, 'hex'));
this.coinbaseTransaction = this.createCoinbaseTransaction(payoutInformation, jobTemplate.blockData.coinbasevalue);
@@ -104,6 +125,60 @@ export class MiningJob {
return Buffer.from(this.coinbasePart2Buffer);
}
/**
* Promote a coinbase that was built during the prior height to the
* authoritative empty template. The coinbase bytes are height/payout bound;
* header fields and submission ownership are replaced only at activation.
*/
public activatePreStagedTemplate(
jobTemplate: IJobTemplate,
headerFields = MiningJob.createNotifyHeaderFields(jobTemplate),
): void {
this.jobTemplateId = jobTemplate.blockData.id;
this.tipKey = jobTemplate.blockData.tipKey;
this.networkDifficulty = jobTemplate.blockData.networkDifficulty;
this.creation = Date.now();
if (!this.patchPreStagedNotify(headerFields)) {
this.miningNotifyResponseBuffer = null;
}
this.miningNotifyResponse = null;
this.preStagedNotifyLayout = null;
}
/**
* Serialize the miner-specific notify while the prior height is active.
* Activation can then overwrite the fixed-width header fields in place,
* avoiding JSON serialization and Buffer allocation in the new-tip loop.
*/
public preparePreStagedNotify(
jobTemplate: IJobTemplate,
headerFields = MiningJob.createNotifyHeaderFields(jobTemplate),
): void {
if (!jobTemplate.block.prevHash.equals(MiningJob.placeholderPrevHash)) {
throw new Error('Pre-staged notify requires a placeholder prevhash');
}
this.responseBuffer(jobTemplate, headerFields);
const response = this.miningNotifyResponse;
let cursor = response.indexOf(`"${this.jobId}"`);
const locate = (value: string): number => {
const offset = response.indexOf(`"${value}"`, cursor);
if (offset < 0) {
throw new Error('Unable to locate pre-staged notify field');
}
cursor = offset + value.length + 2;
return offset + 1;
};
this.preStagedNotifyLayout = {
fields: { ...headerFields },
offsets: {
previousBlockHash: locate(headerFields.previousBlockHash),
version: locate(headerFields.version),
bits: locate(headerFields.bits),
timestamp: locate(headerFields.timestamp),
},
};
}
public buildCoinbaseMerkleRoot(extraNonce: string, extraNonce2: string): Buffer {
const coinbaseBuffer = Buffer.concat([
this.coinbasePart1Buffer,
@@ -317,25 +392,32 @@ export class MiningJob {
return paymentScript;
}
public response(jobTemplate: IJobTemplate): string {
public response(
jobTemplate: IJobTemplate,
headerFields = MiningJob.createNotifyHeaderFields(jobTemplate),
): string {
if (this.miningNotifyResponse != null) {
return this.miningNotifyResponse;
}
if (this.miningNotifyResponseBuffer != null) {
this.miningNotifyResponse = this.miningNotifyResponseBuffer.toString();
return this.miningNotifyResponse;
}
const job: IMiningNotify = {
id: null,
method: eResponseMethod.MINING_NOTIFY,
params: [
this.jobId,
this.swapEndianWords(jobTemplate.block.prevHash).toString('hex'),
headerFields.previousBlockHash,
this.coinbasePart1,
this.coinbasePart2,
jobTemplate.merkle_branch,
jobTemplate.block.version.toString(16),
jobTemplate.block.bits.toString(16),
jobTemplate.block.timestamp.toString(16),
jobTemplate.blockData.clearJobs
headerFields.version,
headerFields.bits,
headerFields.timestamp,
headerFields.cleanJobs
]
};
@@ -344,15 +426,52 @@ export class MiningJob {
return this.miningNotifyResponse;
}
public responseBuffer(jobTemplate: IJobTemplate): Buffer {
public responseBuffer(
jobTemplate: IJobTemplate,
headerFields = MiningJob.createNotifyHeaderFields(jobTemplate),
): Buffer {
if (this.miningNotifyResponseBuffer == null) {
this.response(jobTemplate);
this.response(jobTemplate, headerFields);
}
return this.miningNotifyResponseBuffer;
}
public static createNotifyHeaderFields(jobTemplate: IJobTemplate): MiningNotifyHeaderFields {
return {
previousBlockHash: MiningJob.swapEndianWords(jobTemplate.block.prevHash).toString('hex'),
version: jobTemplate.block.version.toString(16),
bits: jobTemplate.block.bits.toString(16),
timestamp: jobTemplate.block.timestamp.toString(16),
cleanJobs: jobTemplate.blockData.clearJobs,
};
}
private swapEndianWords(buffer: Buffer): Buffer {
private patchPreStagedNotify(headerFields: MiningNotifyHeaderFields): boolean {
const layout = this.preStagedNotifyLayout;
const response = this.miningNotifyResponseBuffer;
if (layout == null
|| response == null
|| layout.fields.cleanJobs !== headerFields.cleanJobs) {
return false;
}
const fields = [
['previousBlockHash', headerFields.previousBlockHash],
['version', headerFields.version],
['bits', headerFields.bits],
['timestamp', headerFields.timestamp],
] as const;
if (fields.some(([name, value]) =>
value.length !== layout.fields[name].length)) {
return false;
}
for (const [name, value] of fields) {
response.write(value, layout.offsets[name], value.length, 'ascii');
}
return true;
}
private static swapEndianWords(buffer: Buffer): Buffer {
const swappedBuffer = Buffer.alloc(buffer.length);
for (let i = 0; i < buffer.length; i += 4) {
+53 -4
View File
@@ -52,6 +52,7 @@ export interface MiningJobBroadcastResult {
status: 'written' | 'backpressured' | 'skipped' | 'closed' | 'error';
bytes: number;
bufferedBytes: number;
preStaged?: boolean;
}
export class StratumV1Client {
@@ -440,6 +441,10 @@ export class StratumV1Client {
this.stratumInitialized = true;
const latestJobTemplate = await this.getLatestPayoutJobTemplate();
this.broadcastMiningJob(latestJobTemplate);
const latestPrestage = this.stratumV1JobsService.getLatestPrestageJobTemplate(this.payoutMode);
if (latestPrestage != null) {
this.preStageMiningJob(latestPrestage);
}
this.backgroundWork.push(
setInterval(async () => {
@@ -461,6 +466,29 @@ export class StratumV1Client {
&& !this.socket.writableEnded;
}
public preStageMiningJob(jobTemplate: IJobTemplate): boolean {
if (!this.isReadyForMiningJobs()
|| jobTemplate.blockData.jobType !== 'empty'
|| jobTemplate.blockData.payoutMode !== this.payoutMode) {
return false;
}
const payoutInformation = this.getPayoutInformation(
jobTemplate,
this.clientAuthorization.address,
);
if (payoutInformation == null) {
return false;
}
const payoutIdentity = this.getPayoutIdentity(jobTemplate);
return this.stratumV1JobsService.preStageJob(
this.network,
payoutInformation,
jobTemplate,
payoutIdentity,
this.payoutMode,
) != null;
}
public broadcastMiningJob(jobTemplate: IJobTemplate, force = false): MiningJobBroadcastResult {
if (!this.isReadyForMiningJobs()) {
return { status: 'skipped', bytes: 0, bufferedBytes: this.socket.writableLength ?? 0 };
@@ -512,19 +540,37 @@ export class StratumV1Client {
// ];
// }
const job = this.stratumV1JobsService.getOrCreateJob(
const payoutIdentity = this.getPayoutIdentity(jobTemplate);
const preStagedJob = (
jobTemplate.blockData.jobType === 'empty'
? this.stratumV1JobsService.activatePreStagedJob(
jobTemplate,
payoutIdentity,
this.payoutMode,
)
: null
);
const job = preStagedJob ?? this.stratumV1JobsService.getOrCreateJob(
this.network,
payoutInformation,
jobTemplate,
this.getPayoutIdentity(jobTemplate),
payoutIdentity,
this.payoutMode,
);
const payload = job.responseBuffer(jobTemplate);
const payload = job.responseBuffer(
jobTemplate,
this.stratumV1JobsService.getNotifyHeaderFields?.(jobTemplate),
);
const bufferedBeforeWrite = this.socket.writableLength ?? 0;
if (bufferedBeforeWrite >= maximumBufferedBytes
|| payload.length >= maximumBufferedBytes - bufferedBeforeWrite) {
this.closeSocket();
return { status: 'closed', bytes: 0, bufferedBytes: bufferedBeforeWrite };
return {
status: 'closed',
bytes: 0,
bufferedBytes: bufferedBeforeWrite,
...(preStagedJob == null ? {} : { preStaged: true }),
};
}
try {
const accepted = this.socket.write(payload);
@@ -537,12 +583,14 @@ export class StratumV1Client {
status: 'closed',
bytes: payload.length,
bufferedBytes: bufferedAfterWrite,
...(preStagedJob == null ? {} : { preStaged: true }),
};
}
return {
status: accepted ? 'written' : 'backpressured',
bytes: payload.length,
bufferedBytes: bufferedAfterWrite,
...(preStagedJob == null ? {} : { preStaged: true }),
};
} catch (error) {
this.closeSocket();
@@ -550,6 +598,7 @@ export class StratumV1Client {
status: 'error',
bytes: 0,
bufferedBytes: this.socket.writableLength ?? 0,
...(preStagedJob == null ? {} : { preStaged: true }),
};
}
}
+7
View File
@@ -56,7 +56,14 @@ export interface IBlockTemplate {
notificationEventId?: string;
/** Wall-clock time when the master first observed the source block notification. */
sourceNotificationReceivedAtMs?: number;
/** Wall-clock time when the compact job finished construction on the master. */
notificationPreparedAtMs?: number;
/** Wall-clock time immediately before the Redis publish command was issued. */
notificationPublishedAtMs?: number;
/** Worker-local wall-clock time when the urgent Redis subscriber received the bridge. */
notificationWorkerReceivedAtMs?: number;
/** Worker-local wall-clock time immediately before the bridge entered job preparation. */
notificationWorkerHandledAtMs?: number;
/** Timestamp of the rolling PPLNS snapshot seed used by an empty bridge. */
payoutBridgeSeedCreatedAtMs?: number;
+239 -11
View File
@@ -199,7 +199,7 @@ describe('BitcoinRpcService template publication', () => {
.toEqual([true, true]);
});
it('publishes a validated compact subsidy bridge before serializing the full template', async () => {
it('publishes a consensus-subsidy bridge before traversing or serializing the full body', async () => {
const order: string[] = [];
const redis = createRedisMock(order);
const template = createTemplate();
@@ -236,6 +236,49 @@ describe('BitcoinRpcService template publication', () => {
expect(order.indexOf('redis:publish:bridge:solo')).toBeLessThan(order.indexOf('redis:publish:solo'));
});
it('rejects an urgent bridge when Core fee data cannot validate its subsidy', async () => {
const redis = createRedisMock([]);
const template = createTemplateAtHeight(840_000, '70');
Object.defineProperty(template.transactions[0], 'fee', {
get: () => {
throw new Error('full transaction body was traversed');
},
});
const service = new BitcoinRpcService(
createConfig({ NETWORK: 'mainnet' }),
{} as any,
redis as any,
);
const trace = (service as any).startTrace('new_block', Date.now());
const consoleSpy = jest.spyOn(console, 'error').mockImplementation(() => undefined);
await expect((service as any).publishSoloSubsidyBridge(template, trace))
.resolves.toBe('failed');
expect(redis.publishSv1BridgeUpdate).not.toHaveBeenCalled();
consoleSpy.mockRestore();
});
it('rejects an urgent bridge when Core coinbasevalue is not subsidy plus fees', async () => {
const redis = createRedisMock([]);
const template = createTemplateAtHeight(840_000, '79');
template.coinbasevalue += 1;
const service = new BitcoinRpcService(
createConfig({ NETWORK: 'mainnet' }),
{} as any,
redis as any,
);
const trace = (service as any).startTrace('new_block', Date.now());
const consoleSpy = jest.spyOn(console, 'error').mockImplementation(() => undefined);
await expect((service as any).publishSoloSubsidyBridge(template, trace))
.resolves.toBe('failed');
expect(redis.publishSv1BridgeUpdate).not.toHaveBeenCalled();
expect(consoleSpy).toHaveBeenCalledWith(expect.stringContaining(
'does not equal subsidy',
));
consoleSpy.mockRestore();
});
it('precomputes an exact next-height PPLNS subsidy seed off the canonical path', async () => {
const order: string[] = [];
const canonical = createTemplateAtHeight(839_999, '66');
@@ -500,7 +543,7 @@ describe('BitcoinRpcService template publication', () => {
await expect((service as any).publishPplnsSubsidyBridge(
nextTemplate,
trace,
)).resolves.toBe(false);
)).resolves.toBe('skipped');
expect(redis.publishSv1BridgeUpdate).not.toHaveBeenCalled();
consoleSpy.mockRestore();
});
@@ -536,7 +579,7 @@ describe('BitcoinRpcService template publication', () => {
await expect((service as any).publishPplnsSubsidyBridge(
template,
trace,
)).resolves.toBe(false);
)).resolves.toBe('skipped');
(service as any).storePplnsSubsidyBridgeSeed({
...baseSeed,
@@ -546,7 +589,7 @@ describe('BitcoinRpcService template publication', () => {
await expect((service as any).publishPplnsSubsidyBridge(
template,
trace,
)).resolves.toBe(false);
)).resolves.toBe('failed');
(service as any).storePplnsSubsidyBridgeSeed({
...baseSeed,
@@ -556,7 +599,7 @@ describe('BitcoinRpcService template publication', () => {
await expect((service as any).publishPplnsSubsidyBridge(
template,
trace,
)).resolves.toBe(false);
)).resolves.toBe('skipped');
expect(redis.publishSv1BridgeUpdate).not.toHaveBeenCalled();
});
@@ -569,6 +612,7 @@ describe('BitcoinRpcService template publication', () => {
if (update.template.payoutMode === 'pplns') {
throw new Error('PPLNS Redis unavailable');
}
return true;
});
const service = new BitcoinRpcService(
createConfig({
@@ -633,6 +677,102 @@ describe('BitcoinRpcService template publication', () => {
expect((service as any).latestLongpollId).toBe('next-longpoll-id');
});
it('accepts the first authoritative template from an independent auxiliary longpoll source', async () => {
const redis = createRedisMock([]);
const template = createTemplateAtHeight(840_000, '73');
template.longpollid = 'aux-next-longpoll';
const auxiliaryClient = {
post: jest.fn().mockResolvedValue({
data: { result: template, error: null },
}),
};
const source = {
name: 'aux-1',
client: auxiliaryClient,
latestLongpollId: 'aux-prior-longpoll',
longpollLoopStarted: true,
};
const service = new BitcoinRpcService(
createConfig({ NETWORK: 'mainnet', BLOCK_TEMPLATE_LONGPOLL_TIMEOUT_MS: '1000' }),
{ saveBlock: jest.fn().mockResolvedValue(undefined), getSavedBlockTemplate: jest.fn() } as any,
redis as any,
);
service.miningInfo = { blocks: template.height - 1 } as any;
const primaryPost = jest.fn().mockResolvedValue({
data: { result: template.previousblockhash, error: null },
});
(service as any).client = { post: primaryPost };
const logSpy = jest.spyOn(console, 'log').mockImplementation(() => undefined);
await (service as any).getAndBroadcastLatestTemplateOnce(
'longpoll',
source.latestLongpollId,
undefined,
source,
);
expect(auxiliaryClient.post).toHaveBeenCalledWith('', expect.objectContaining({
method: 'getblocktemplate',
params: [expect.objectContaining({ longpollid: 'aux-prior-longpoll' })],
}), { timeout: 1000 });
expect(source.latestLongpollId).toBe('aux-next-longpoll');
expect(primaryPost).toHaveBeenCalledWith('', expect.objectContaining({
method: 'getbestblockhash',
}), undefined);
expect(redis.publishSv1BridgeUpdate).toHaveBeenCalled();
expect(redis.setBlockTemplate).toHaveBeenCalledWith(
template.height - 1,
expect.objectContaining({ previousblockhash: template.previousblockhash }),
);
const sourceLog = logSpy.mock.calls
.map(call => call[0])
.find(value => typeof value === 'string' && value.includes('block_source_notification'));
expect(JSON.parse(sourceLog)).toEqual(expect.objectContaining({
source: 'longpoll',
templateSource: 'aux-1',
}));
logSpy.mockRestore();
});
it('rejects auxiliary work before urgent and canonical publication when primary Core has another tip', async () => {
const redis = createRedisMock([]);
const template = createTemplateAtHeight(840_000, '74');
const source = {
name: 'aux-1',
client: {
post: jest.fn().mockResolvedValue({
data: { result: template, error: null },
}),
},
latestLongpollId: 'aux-prior-longpoll',
longpollLoopStarted: true,
};
const service = new BitcoinRpcService(
createConfig({ NETWORK: 'mainnet' }),
{ saveBlock: jest.fn(), getSavedBlockTemplate: jest.fn() } as any,
redis as any,
);
service.miningInfo = { blocks: template.height - 1 } as any;
(service as any).client = {
post: jest.fn().mockResolvedValue({
data: { result: '75'.repeat(32), error: null },
}),
};
const warnSpy = jest.spyOn(console, 'warn').mockImplementation(() => undefined);
await (service as any).getAndBroadcastLatestTemplateOnce(
'longpoll',
source.latestLongpollId,
undefined,
source,
);
expect(redis.publishSv1BridgeUpdate).not.toHaveBeenCalled();
expect(redis.setBlockTemplate).not.toHaveBeenCalled();
expect(warnSpy).toHaveBeenCalledWith(expect.stringContaining('aux_template_rejected'));
warnSpy.mockRestore();
});
it('deduplicates the same authoritative template returned by ZMQ and longpoll', async () => {
const redis = createRedisMock([]);
const template = createTemplate();
@@ -694,10 +834,16 @@ describe('BitcoinRpcService template publication', () => {
const redis = createRedisMock(order);
let resolveBridge: () => void;
const stalledBridge = new Promise<void>(resolve => { resolveBridge = resolve; });
redis.publishSv1BridgeUpdate.mockImplementation(async (update: { template: IBlockTemplate }) => {
order.push(`redis:publish:bridge:${update.template.payoutMode}`);
await stalledBridge;
});
redis.publishSv1BridgeUpdate.mockImplementation((async (
update: { template: IBlockTemplate },
lane: 'fallback' | undefined,
) => {
order.push(`redis:publish:bridge:${update.template.payoutMode}:${lane ?? 'urgent'}`);
if (lane == null) {
await stalledBridge;
}
return true;
}) as any);
const template = createTemplateAtHeight(840_000, '71');
const service = new BitcoinRpcService(
createConfig({ NETWORK: 'mainnet' }),
@@ -710,7 +856,12 @@ describe('BitcoinRpcService template publication', () => {
await service.getAndBroadcastLatestTemplate('new_block');
expect(redis.publishSv1BridgeUpdate).toHaveBeenCalledTimes(1);
expect(redis.publishSv1BridgeUpdate).toHaveBeenCalledTimes(2);
expect(redis.publishSv1BridgeUpdate).toHaveBeenNthCalledWith(
2,
expect.objectContaining({ type: 'subsidy-bridge' }),
'fallback',
);
expect(redis.setBlockTemplate).toHaveBeenCalledWith(
template.height - 1,
expect.objectContaining({ payoutMode: 'solo', jobType: 'full' }),
@@ -718,12 +869,85 @@ describe('BitcoinRpcService template publication', () => {
expect(redis.publishBlockTemplateUpdate).toHaveBeenCalledWith(expect.objectContaining({
payoutMode: 'solo',
}));
expect(order.indexOf('redis:publish:bridge:solo')).toBeLessThan(order.indexOf('redis:set:solo'));
expect(order.indexOf('redis:publish:bridge:solo:fallback')).toBeLessThan(order.indexOf('redis:set:solo'));
resolveBridge!();
await flushPromises();
});
it('retries the bridge from the canonical path after an immediate Redis delivery failure', async () => {
const redis = createRedisMock([]);
redis.publishSv1BridgeUpdate
.mockResolvedValueOnce(false)
.mockResolvedValueOnce(true);
const template = createTemplateAtHeight(840_000, '78');
const service = new BitcoinRpcService(
createConfig({ NETWORK: 'mainnet' }),
{ saveBlock: jest.fn().mockResolvedValue(undefined), getSavedBlockTemplate: jest.fn() } as any,
redis as any,
{ createSnapshotForTemplate: jest.fn().mockResolvedValue(null) } as any,
);
service.miningInfo = { blocks: template.height - 1 } as any;
jest.spyOn(service as any, 'fetchBlockTemplate').mockResolvedValue(template);
const errorSpy = jest.spyOn(console, 'error').mockImplementation(() => undefined);
await service.getAndBroadcastLatestTemplate('new_block');
expect(redis.publishSv1BridgeUpdate).toHaveBeenCalledTimes(2);
expect(redis.setBlockTemplate).toHaveBeenCalledWith(
template.height - 1,
expect.objectContaining({ payoutMode: 'solo', jobType: 'full' }),
);
errorSpy.mockRestore();
});
it.each([
['mainnet', 'main'],
['testnet', 'test'],
['regtest', 'regtest'],
] as const)('accepts NETWORK=%s only for the matching Core chain', (network, chain) => {
const service = new BitcoinRpcService(
createConfig({ NETWORK: network }),
{} as any,
{} as any,
);
expect(() => (service as any).validateConfiguredNetworkAgainstCore(chain))
.not.toThrow();
expect(() => (service as any).validateConfiguredNetworkAgainstCore(
chain === 'main' ? 'test' : 'main',
)).toThrow(`NETWORK=${network} does not match Bitcoin Core chain=`);
});
it('publishes a durable next-height empty prestage after canonical work', async () => {
const redis = createRedisMock([]);
const template = createTemplateAtHeight(840_000, '72');
const service = new BitcoinRpcService(
createConfig({ NETWORK: 'mainnet' }),
{ saveBlock: jest.fn().mockResolvedValue(undefined), getSavedBlockTemplate: jest.fn() } as any,
redis as any,
{ createSnapshotForTemplate: jest.fn().mockResolvedValue(null) } as any,
);
service.miningInfo = { blocks: template.height - 1 } as any;
jest.spyOn(service as any, 'fetchBlockTemplate').mockResolvedValue(template);
await service.getAndBroadcastLatestTemplate('new_block');
await flushPromises();
expect(redis.publishSv1PrestageUpdate).toHaveBeenCalledWith(expect.objectContaining({
schemaVersion: 1,
type: 'subsidy-prestage',
template: expect.objectContaining({
height: template.height + 1,
previousblockhash: '0'.repeat(64),
payoutMode: 'solo',
jobType: 'empty',
transactions: [],
coinbasevalue: calculateBlockSubsidySats(template.height + 1, 'mainnet'),
}),
}));
});
it('verifies same-height reorgs against Core and never switches back to an orphan', async () => {
const redis = createRedisMock([]);
const first = createTemplateAtHeight(840_000, '31');
@@ -1263,6 +1487,10 @@ function createRedisMock(order: string[]) {
}),
publishSv1BridgeUpdate: jest.fn(async (update: { template: IBlockTemplate }) => {
order.push(`redis:publish:bridge:${update.template.payoutMode}`);
return true;
}),
publishSv1PrestageUpdate: jest.fn(async (update: { template: IBlockTemplate }) => {
order.push(`redis:publish:prestage:${update.template.payoutMode}`);
}),
setLegacyBlockTemplate: jest.fn(async (_height: number, template: IBlockTemplate) => {
order.push(`redis:set:legacy:${template.payoutMode}`);
+516 -66
View File
@@ -9,7 +9,12 @@ import * as zmq from 'zeromq';
import { IBlockTemplate } from '../models/bitcoin-rpc/IBlockTemplate';
import { IMiningInfo } from '../models/bitcoin-rpc/IMiningInfo';
import { BlockTemplateUpdate, RedisMessagingService, Sv1BridgeUpdate } from './redis-messaging.service';
import {
BlockTemplateUpdate,
RedisMessagingService,
Sv1BridgeUpdate,
Sv1PrestageUpdate,
} from './redis-messaging.service';
import {
BitcoinNetworkName,
calculateBlockSubsidySats,
@@ -24,9 +29,25 @@ interface BlockNotificationTrace {
startedWallMs: number;
startedMonotonic: bigint;
sourceNotificationReceivedAtMs?: number;
templateSource?: string;
stages: Record<string, number>;
}
interface TemplateRpcSource {
name: string;
client: AxiosInstance;
latestLongpollId: string | null;
longpollLoopStarted: boolean;
}
type BridgePublishResult = 'published' | 'skipped' | 'failed';
type UrgentBridgePublishResult = 'acknowledged' | 'fallback-started' | 'failed';
interface BridgePublishReservation {
tipKey: string;
attemptId: number;
}
interface StoredBlockTemplateEnvelope {
schemaVersion: 1;
templates: IBlockTemplate[];
@@ -65,6 +86,7 @@ interface PendingPplnsTemplatePublication {
}
const DEFAULT_PPLNS_BRIDGE_SEED_MAX_AGE_MS = 5 * 60 * 1000;
const DEFAULT_SV1_BRIDGE_PUBLISH_BUDGET_MS = 10;
const PPLNS_LISTENER_CONFIG_KEYS = [
'PPLNS_STRATUM_PORTS',
'PPLNS_SECURE_STRATUM_PORTS',
@@ -79,8 +101,10 @@ export class BitcoinRpcService implements OnModuleInit {
private client: AxiosInstance;
private readonly auxiliaryTemplateSources: TemplateRpcSource[] = [];
private _newBlockTemplate$: BehaviorSubject<IBlockTemplate> = new BehaviorSubject(undefined);
private _newSv1BridgeTemplate$ = new ReplaySubject<IBlockTemplate>(1);
private _newSv1PrestageTemplate$ = new ReplaySubject<IBlockTemplate>(2);
private resetTemplateInterval$ = new Subject<void>();
private rpcRequestId = 0;
private readonly processedTemplateEvents = new Set<string>();
@@ -94,7 +118,10 @@ export class BitcoinRpcService implements OnModuleInit {
private pplnsTemplatePublicationPromise: Promise<void> | null = null;
private readonly latestBridgeTemplates = new Map<'solo' | 'pplns', IBlockTemplate>();
private readonly canonicalEmissionStates = new Map<'solo' | 'pplns' | 'all', CanonicalEmissionState>();
private readonly lastPublishedBridgeTipKeys = new Map<'solo' | 'pplns', string>();
private readonly lastPublishedBridgeTipKeys = new Map<'solo' | 'pplns', BridgePublishReservation>();
private bridgePublishAttemptId = 0;
private readonly subsidyValidatedTemplates = new WeakSet<IBlockTemplate>();
private readonly lastPublishedPrestageKeys = new Map<'solo' | 'pplns', string>();
private readonly pplnsSubsidyBridgeSeeds = new Map<number, PplnsSubsidyBridgeSeed>();
private pplnsSeedPrecomputeTail: Promise<void> = Promise.resolve();
private pplnsSeedPrecomputeGeneration = 0;
@@ -103,6 +130,8 @@ export class BitcoinRpcService implements OnModuleInit {
private lastPublishedTemplateSignature: string | null = null;
private lastPublishedTipKey: string | null = null;
private lastPublishedCandidateHeight = -1;
/** Atomically reserves the fastest forward height before any Redis await. */
private highestUrgentCandidateHeight = -1;
private highestPublishedCandidateHeight = -1;
private reorgVerificationFenceHeight = -1;
private latestLongpollId: string | null = null;
@@ -112,6 +141,7 @@ export class BitcoinRpcService implements OnModuleInit {
public miningInfo: IMiningInfo;
public newBlockTemplate$ = this._newBlockTemplate$.pipe(filter(block => block != null), shareReplay({ refCount: true, bufferSize: 1 }));
public newSv1BridgeTemplate$ = this._newSv1BridgeTemplate$.pipe(shareReplay({ refCount: true, bufferSize: 1 }));
public newSv1PrestageTemplate$ = this._newSv1PrestageTemplate$.pipe(shareReplay({ refCount: true, bufferSize: 2 }));
/** Core-authoritative early header activation; currently carries the solo subsidy bridge body. */
public workActivationTemplate$ = this.newSv1BridgeTemplate$;
@@ -142,17 +172,25 @@ export class BitcoinRpcService implements OnModuleInit {
password: pass
}
});
this.configureAuxiliaryTemplateSources({ user, pass, port, timeout });
console.log(`MASTER? ${process.env.MASTER}`)
if (process.env.MASTER != 'true') {
await this.loadLatestMiningInfoForReplayProcess();
if (process.env.API_ONLY != 'true') {
await this.redisMessagingService.subscribeSv1BridgeUpdates(async update => {
this.workerTemplateTail = this.workerTemplateTail
.then(() => this.handleSv1BridgeUpdate(update))
.catch(error => console.error(`Unable to load SV1 bridge update: ${error.message}`));
await this.workerTemplateTail;
// A new-tip bridge is an interrupt, not canonical replay work.
// It has its own Redis socket and must never sit behind a full
// template GET/JSON parse on workerTemplateTail.
await this.handleSv1BridgeUpdate(update).catch(error => {
console.error(`Unable to load SV1 bridge update: ${error.message}`);
});
});
if (typeof this.redisMessagingService.subscribeSv1PrestageUpdates === 'function') {
await this.redisMessagingService.subscribeSv1PrestageUpdates(async update => {
this.handleSv1PrestageUpdate(update);
});
}
await this.redisMessagingService.subscribeBlockTemplateUpdates(async update => {
this.workerTemplateTail = this.workerTemplateTail
.then(() => this.loadTemplateUpdateForWorker(update))
@@ -179,6 +217,11 @@ export class BitcoinRpcService implements OnModuleInit {
// state. Replaying one before canonical work can briefly put a
// restarting worker on an orphan; load only authoritative jobs.
await this.loadLatestTemplateForWorker();
if (typeof this.redisMessagingService.getLatestSv1PrestageUpdates === 'function') {
for (const update of await this.redisMessagingService.getLatestSv1PrestageUpdates()) {
this.handleSv1PrestageUpdate(update);
}
}
});
await this.workerTemplateTail;
}
@@ -194,6 +237,7 @@ export class BitcoinRpcService implements OnModuleInit {
});
this.miningInfo = await this.getMiningInfo();
this.validateConfiguredNetworkAgainstCore(this.miningInfo.chain);
console.log('Using ZMQ');
const sock = new zmq.Subscriber;
@@ -212,6 +256,9 @@ export class BitcoinRpcService implements OnModuleInit {
await this.getAndBroadcastLatestTemplate('startup');
void this.listenForLongpollTemplates();
this.auxiliaryTemplateSources.forEach(source => {
void this.listenForLongpollTemplates(source);
});
// Between new blocks we want refresh jobs with the latest transactions
this.resetTemplateInterval$.pipe(
@@ -287,21 +334,53 @@ export class BitcoinRpcService implements OnModuleInit {
reason: TemplateRefreshReason,
longpollId?: string,
sourceNotificationReceivedAtMs?: number,
templateSource?: TemplateRpcSource,
): Promise<void> {
if (this.miningInfo?.blocks == null) {
console.warn('Skipping block template broadcast because mining info is not available');
return;
}
const trace = this.startTrace(reason, sourceNotificationReceivedAtMs);
const blockTemplate = await this.fetchBlockTemplate(trace, longpollId);
const trace = this.startTrace(
reason,
sourceNotificationReceivedAtMs,
templateSource?.name ?? 'primary',
);
const blockTemplate = await this.fetchBlockTemplate(trace, longpollId, templateSource);
if (blockTemplate == null) {
console.warn(`Skipping block template broadcast for height ${this.miningInfo.blocks}; block template is not available`);
return;
}
if (templateSource != null
&& !await this.isAuxiliaryTemplateAuthorizedByPrimary(blockTemplate, templateSource)) {
return;
}
const tipKey = `${blockTemplate.height}:${blockTemplate.previousblockhash}`;
const canInterruptPublicationTail = tipKey !== this.lastPublishedTipKey
&& blockTemplate.height > Math.max(
this.lastPublishedCandidateHeight,
this.highestUrgentCandidateHeight,
);
let urgentBridgeHandled = false;
if (canInterruptPublicationTail) {
// Reserve synchronously. Two racing Core sources at the same height
// must not announce conflicting prevhashes before canonical reorg
// verification has selected one of them.
this.highestUrgentCandidateHeight = blockTemplate.height;
this.markTrace(trace, 'template_ready');
this.logSourceNotification(trace, blockTemplate);
urgentBridgeHandled = await this.publishUrgentSubsidyBridges(blockTemplate, trace)
!== 'failed';
}
const publication = this.templatePublicationTail.then(() =>
this.publishFetchedBlockTemplate(blockTemplate, reason, trace),
this.publishFetchedBlockTemplate(
blockTemplate,
reason,
trace,
urgentBridgeHandled,
),
);
this.templatePublicationTail = publication.catch(error => {
console.error(`Unable to publish fetched block template: ${error.message}`);
@@ -313,6 +392,7 @@ export class BitcoinRpcService implements OnModuleInit {
blockTemplate: IBlockTemplate,
reason: TemplateRefreshReason,
trace: BlockNotificationTrace,
urgentBridgeHandled = false,
): Promise<void> {
const tipHeight = blockTemplate.height - 1;
const tipKey = `${blockTemplate.height}:${blockTemplate.previousblockhash}`;
@@ -342,18 +422,20 @@ export class BitcoinRpcService implements OnModuleInit {
}
const isNewTip = tipKey !== this.lastPublishedTipKey;
this.miningInfo = { ...this.miningInfo, blocks: tipHeight };
this.markTrace(trace, 'template_ready');
if (isNewTip) {
if (!urgentBridgeHandled) {
this.markTrace(trace, 'template_ready');
}
if (isNewTip && !urgentBridgeHandled) {
this.logSourceNotification(trace, blockTemplate);
}
if (isNewTip) {
// Start bridge publication first, but never await it on the
// authoritative path. With a healthy Redis connection its command
// is enqueued first; if Redis stalls, canonical storage/publication
// can still make progress independently.
void this.publishSoloSubsidyBridge(blockTemplate, trace);
void this.publishPplnsSubsidyBridge(blockTemplate, trace);
if (isNewTip && !urgentBridgeHandled) {
// Yield until the urgent Redis socket acknowledges the compact jobs,
// but enforce a tiny budget so degraded Redis cannot hold canonical
// storage indefinitely. This guarantees that full-body hashing and
// JSON serialization do not occupy the same event-loop tick before
// the bridge command has had an opportunity to leave the process.
await this.publishUrgentSubsidyBridges(blockTemplate, trace);
}
const templateSignature = this.createTemplateSignature(blockTemplate);
@@ -408,6 +490,9 @@ export class BitcoinRpcService implements OnModuleInit {
}
this.queueTemplatePersistence(tipHeight, [soloTemplate], trace);
this.queuePplnsSubsidyBridgeSeedPrecompute(blockTemplate);
void this.publishNextHeightPrestage(blockTemplate, 'solo').catch(error => {
console.error(`Unable to publish next-height solo prestage: ${error.message}`);
});
this.queuePplnsTemplatePublication({
blockTemplate,
soloTemplate,
@@ -611,7 +696,11 @@ export class BitcoinRpcService implements OnModuleInit {
return [];
}
private async fetchBlockTemplate(trace: BlockNotificationTrace, longpollId?: string) {
private async fetchBlockTemplate(
trace: BlockNotificationTrace,
longpollId?: string,
templateSource?: TemplateRpcSource,
) {
console.log(`Master fetching block template after tip ${this.miningInfo?.blocks}`);
@@ -620,7 +709,9 @@ export class BitcoinRpcService implements OnModuleInit {
const maxRetryDelayMs = this.getPositiveIntegerEnv('BLOCK_TEMPLATE_RETRY_MAX_MS', 1000);
while (blockTemplate == null) {
try {
blockTemplate = await this.callRpc<IBlockTemplate>('getblocktemplate', [
blockTemplate = await this.callRpcWithClient<IBlockTemplate>(
templateSource?.client ?? this.client,
'getblocktemplate', [
{
rules: ['segwit'],
mode: 'template',
@@ -646,27 +737,66 @@ export class BitcoinRpcService implements OnModuleInit {
} else {
this.markTrace(trace, 'getblocktemplate_complete');
}
this.latestLongpollId = blockTemplate.longpollid;
if (templateSource == null) {
this.latestLongpollId = blockTemplate.longpollid;
} else {
templateSource.latestLongpollId = blockTemplate.longpollid;
}
return blockTemplate;
}
private async listenForLongpollTemplates(): Promise<void> {
if (this.longpollLoopStarted || process.env.MASTER !== 'true') {
private async listenForLongpollTemplates(templateSource?: TemplateRpcSource): Promise<void> {
if (process.env.MASTER !== 'true') {
return;
}
this.longpollLoopStarted = true;
if (templateSource == null) {
if (this.longpollLoopStarted) {
return;
}
this.longpollLoopStarted = true;
} else {
if (templateSource.longpollLoopStarted) {
return;
}
templateSource.longpollLoopStarted = true;
}
while (process.env.MASTER === 'true') {
const longpollId = this.latestLongpollId;
const longpollId = templateSource == null
? this.latestLongpollId
: templateSource.latestLongpollId;
if (longpollId == null) {
await new Promise(resolve => setTimeout(resolve, 1000));
try {
await this.getAndBroadcastLatestTemplateOnce(
'startup',
undefined,
undefined,
templateSource,
);
} catch (error) {
console.error(
`Block template source ${templateSource?.name ?? 'primary'} startup failed: ${error.message}`,
);
await new Promise(resolve => setTimeout(resolve, 1000));
}
continue;
}
try {
await this.getAndBroadcastLatestTemplate('longpoll', longpollId);
if (templateSource == null) {
await this.getAndBroadcastLatestTemplate('longpoll', longpollId);
} else {
await this.getAndBroadcastLatestTemplateOnce(
'longpoll',
longpollId,
undefined,
templateSource,
);
}
} catch (error) {
console.error(`Block template longpoll failed: ${error.message}`);
console.error(
`Block template longpoll failed for ${templateSource?.name ?? 'primary'}: ${error.message}`,
);
await new Promise(resolve => setTimeout(resolve, 1000));
}
}
@@ -817,36 +947,125 @@ export class BitcoinRpcService implements OnModuleInit {
});
}
private async publishUrgentSubsidyBridges(
authoritativeTemplate: IBlockTemplate,
trace: BlockNotificationTrace,
): Promise<UrgentBridgePublishResult> {
const startedAt = Date.now();
this.markTrace(trace, 'sv1_bridge_publish_started');
const publication = Promise.all([
// Preserve solo first on the dedicated Redis command socket.
this.publishSoloSubsidyBridge(authoritativeTemplate, trace, 'urgent'),
this.publishPplnsSubsidyBridge(authoritativeTemplate, trace, 'urgent'),
]);
const budgetMs = this.getPositiveIntegerEnv(
'SV1_BRIDGE_PUBLISH_BUDGET_MS',
DEFAULT_SV1_BRIDGE_PUBLISH_BUDGET_MS,
);
let timeout: NodeJS.Timeout | null = null;
const result = await Promise.race([
publication.then(results => ({ timedOut: false as const, results })),
new Promise<{ timedOut: true; results?: never }>(resolve => {
timeout = setTimeout(() => resolve({ timedOut: true }), budgetMs);
}),
]);
if (timeout != null) {
clearTimeout(timeout);
}
if (result.timedOut) {
this.markTrace(trace, 'sv1_bridge_publish_budget_exhausted');
console.warn(JSON.stringify({
event: 'sv1_bridge_publish_budget_exhausted',
eventId: trace.eventId,
budgetMs,
elapsedMs: Date.now() - startedAt,
templateHeight: authoritativeTemplate.height,
previousBlockHash: authoritativeTemplate.previousblockhash,
}));
// The urgent socket may be stuck rather than merely slow. Release
// only this attempt's reservations and immediately retry through the
// independent normal command socket. Duplicate delivery is safe: the
// envelope event id is identical and workers de-duplicate it.
this.releaseBridgeReservations(authoritativeTemplate);
void this.publishFallbackSubsidyBridges(authoritativeTemplate, trace, publication);
return 'fallback-started';
}
if (result.results.includes('failed')) {
this.markTrace(trace, 'sv1_bridge_publish_failed');
return 'failed';
}
this.markTrace(trace, 'sv1_bridge_publish_acknowledged');
return 'acknowledged';
}
private async publishFallbackSubsidyBridges(
authoritativeTemplate: IBlockTemplate,
trace: BlockNotificationTrace,
originalPublication: Promise<BridgePublishResult[]>,
): Promise<void> {
// Observe the original attempt even if it rejects unexpectedly.
void originalPublication.catch(error => {
console.error(`Urgent SV1 bridge publication failed after timeout: ${error.message}`);
});
const results = await Promise.all([
this.publishSoloSubsidyBridge(authoritativeTemplate, trace, 'fallback'),
this.publishPplnsSubsidyBridge(authoritativeTemplate, trace, 'fallback'),
]);
if (results.includes('failed')) {
this.markTrace(trace, 'sv1_bridge_fallback_failed');
console.error(JSON.stringify({
event: 'sv1_bridge_fallback_failed',
eventId: trace.eventId,
templateHeight: authoritativeTemplate.height,
previousBlockHash: authoritativeTemplate.previousblockhash,
results,
}));
return;
}
this.markTrace(trace, 'sv1_bridge_fallback_acknowledged');
}
private releaseBridgeReservations(authoritativeTemplate: IBlockTemplate): void {
const tipKey = `${authoritativeTemplate.height}:${authoritativeTemplate.previousblockhash}`;
for (const payoutMode of ['solo', 'pplns'] as const) {
if (this.lastPublishedBridgeTipKeys.get(payoutMode)?.tipKey === tipKey) {
this.lastPublishedBridgeTipKeys.delete(payoutMode);
}
}
}
private async publishSoloSubsidyBridge(
authoritativeTemplate: IBlockTemplate,
trace: BlockNotificationTrace,
): Promise<boolean> {
lane: 'urgent' | 'fallback' = 'urgent',
): Promise<BridgePublishResult> {
if (!this.isSv1SubsidyBridgeEnabled()
|| !this.getSv1SubsidyBridgePayoutModes().has('solo')) {
return false;
return 'skipped';
}
const tipKey = `${authoritativeTemplate.height}:${authoritativeTemplate.previousblockhash}`;
if (this.lastPublishedBridgeTipKeys.get('solo') === tipKey) {
return false;
if (this.lastPublishedBridgeTipKeys.get('solo')?.tipKey === tipKey) {
return 'skipped';
}
// Reserve the tip before the Redis await so duplicate ZMQ/longpoll
// callbacks cannot launch a second in-flight bridge publication.
this.lastPublishedBridgeTipKeys.set('solo', tipKey);
const reservation: BridgePublishReservation = {
tipKey,
attemptId: ++this.bridgePublishAttemptId,
};
this.lastPublishedBridgeTipKeys.set('solo', reservation);
try {
this.validateSubsidyAgainstAuthoritativeTemplate(authoritativeTemplate);
const bridgeTemplate = createSubsidyOnlyBlockTemplate({
authoritativeTemplate,
network: this.getNetworkName(),
payoutMode: 'solo',
halvingInterval: this.getOptionalPositiveIntegerEnv('SUBSIDY_HALVING_INTERVAL'),
});
this.validateSubsidyAgainstAuthoritativeTemplate(
authoritativeTemplate,
bridgeTemplate.coinbasevalue,
);
bridgeTemplate.mintime = Math.max(
authoritativeTemplate.mintime,
authoritativeTemplate.curtime,
@@ -854,51 +1073,62 @@ export class BitcoinRpcService implements OnModuleInit {
);
bridgeTemplate.notificationEventId = `${trace.eventId}:bridge`;
bridgeTemplate.sourceNotificationReceivedAtMs = trace.sourceNotificationReceivedAtMs;
bridgeTemplate.notificationPreparedAtMs = Date.now();
const publishedAtMs = Date.now();
bridgeTemplate.notificationPublishedAtMs = publishedAtMs;
const update: Sv1BridgeUpdate = {
schemaVersion: 1,
type: 'subsidy-bridge',
eventId: `${trace.eventId}:bridge`,
template: bridgeTemplate,
publishedAtMs: Date.now(),
publishedAtMs,
};
await this.redisMessagingService.publishSv1BridgeUpdate(update);
const delivered = lane === 'urgent'
? await this.redisMessagingService.publishSv1BridgeUpdate(update)
: await this.redisMessagingService.publishSv1BridgeUpdate(update, 'fallback');
if (!delivered) {
throw new Error(`Redis ${lane} bridge publisher is unavailable`);
}
this.markTrace(trace, 'sv1_bridge_workers_notified');
return true;
return 'published';
} catch (error) {
if (this.lastPublishedBridgeTipKeys.get('solo') === tipKey) {
if (this.lastPublishedBridgeTipKeys.get('solo') === reservation) {
this.lastPublishedBridgeTipKeys.delete('solo');
}
console.error(`Skipping SV1 subsidy bridge: ${error.message}`);
return false;
return 'failed';
}
}
private async publishPplnsSubsidyBridge(
authoritativeTemplate: IBlockTemplate,
trace: BlockNotificationTrace,
): Promise<boolean> {
lane: 'urgent' | 'fallback' = 'urgent',
): Promise<BridgePublishResult> {
if (!this.isSv1SubsidyBridgeEnabled()
|| !this.getSv1SubsidyBridgePayoutModes().has('pplns')) {
return false;
return 'skipped';
}
const tipKey = `${authoritativeTemplate.height}:${authoritativeTemplate.previousblockhash}`;
if (this.lastPublishedBridgeTipKeys.get('pplns') === tipKey) {
return false;
if (this.lastPublishedBridgeTipKeys.get('pplns')?.tipKey === tipKey) {
return 'skipped';
}
const seed = this.getFreshPplnsSubsidyBridgeSeed(authoritativeTemplate.height);
if (seed == null) {
return false;
return 'skipped';
}
this.lastPublishedBridgeTipKeys.set('pplns', tipKey);
const reservation: BridgePublishReservation = {
tipKey,
attemptId: ++this.bridgePublishAttemptId,
};
this.lastPublishedBridgeTipKeys.set('pplns', reservation);
try {
const expectedSubsidy = calculateBlockSubsidySats(
authoritativeTemplate.height,
this.getNetworkName(),
this.getOptionalPositiveIntegerEnv('SUBSIDY_HALVING_INTERVAL'),
const expectedSubsidy = this.validateSubsidyAgainstAuthoritativeTemplate(
authoritativeTemplate,
);
if (seed.subsidySats !== expectedSubsidy) {
throw new Error(
@@ -910,11 +1140,6 @@ export class BitcoinRpcService implements OnModuleInit {
`PPLNS bridge seed bits ${seed.basisBits} do not match authoritative bits ${authoritativeTemplate.bits}`,
);
}
this.validateSubsidyAgainstAuthoritativeTemplate(
authoritativeTemplate,
expectedSubsidy,
);
const bridgeTemplate = createSubsidyOnlyBlockTemplate({
authoritativeTemplate,
network: this.getNetworkName(),
@@ -933,23 +1158,118 @@ export class BitcoinRpcService implements OnModuleInit {
bridgeTemplate.notificationEventId = `${trace.eventId}:bridge:pplns`;
bridgeTemplate.sourceNotificationReceivedAtMs = trace.sourceNotificationReceivedAtMs;
bridgeTemplate.payoutBridgeSeedCreatedAtMs = seed.preparedAtMs;
bridgeTemplate.notificationPreparedAtMs = Date.now();
const publishedAtMs = Date.now();
bridgeTemplate.notificationPublishedAtMs = publishedAtMs;
const update: Sv1BridgeUpdate = {
schemaVersion: 1,
type: 'subsidy-bridge',
eventId: `${trace.eventId}:bridge:pplns`,
template: bridgeTemplate,
publishedAtMs: Date.now(),
publishedAtMs,
};
await this.redisMessagingService.publishSv1BridgeUpdate(update);
const delivered = lane === 'urgent'
? await this.redisMessagingService.publishSv1BridgeUpdate(update)
: await this.redisMessagingService.publishSv1BridgeUpdate(update, 'fallback');
if (!delivered) {
throw new Error(`Redis ${lane} bridge publisher is unavailable`);
}
this.markTrace(trace, 'sv1_pplns_bridge_workers_notified');
return true;
return 'published';
} catch (error) {
if (this.lastPublishedBridgeTipKeys.get('pplns') === tipKey) {
if (this.lastPublishedBridgeTipKeys.get('pplns') === reservation) {
this.lastPublishedBridgeTipKeys.delete('pplns');
}
console.error(`Skipping PPLNS SV1 subsidy bridge: ${error.message}`);
return 'failed';
}
}
private async publishNextHeightPrestage(
authoritativeTemplate: IBlockTemplate,
payoutMode: 'solo' | 'pplns',
seed?: PplnsSubsidyBridgeSeed,
): Promise<boolean> {
if (!this.isSv1SubsidyBridgeEnabled()
|| !this.getSv1SubsidyBridgePayoutModes().has(payoutMode)) {
return false;
}
const candidateHeight = authoritativeTemplate.height + 1;
const subsidySats = calculateBlockSubsidySats(
candidateHeight,
this.getNetworkName(),
this.getOptionalPositiveIntegerEnv('SUBSIDY_HALVING_INTERVAL'),
);
if (payoutMode === 'pplns') {
if (seed == null
|| seed.candidateHeight !== candidateHeight
|| seed.subsidySats !== subsidySats
|| seed.basisBits.toLowerCase() !== authoritativeTemplate.bits.toLowerCase()) {
return false;
}
}
const prestageKey = [
candidateHeight,
payoutMode,
subsidySats,
authoritativeTemplate.bits.toLowerCase(),
seed?.payoutSnapshotId ?? '',
].join(':');
if (this.lastPublishedPrestageKeys.get(payoutMode) === prestageKey) {
return false;
}
this.lastPublishedPrestageKeys.set(payoutMode, prestageKey);
try {
const futureShape: IBlockTemplate = {
...authoritativeTemplate,
height: candidateHeight,
previousblockhash: '0'.repeat(64),
transactions: [],
coinbasevalue: subsidySats,
longpollid: `prestage:${candidateHeight}:${authoritativeTemplate.bits}`,
curtime: Math.max(
authoritativeTemplate.curtime,
Math.floor(Date.now() / 1000),
),
mintime: Math.max(
authoritativeTemplate.mintime,
authoritativeTemplate.curtime,
),
};
const template = createSubsidyOnlyBlockTemplate({
authoritativeTemplate: futureShape,
network: this.getNetworkName(),
payoutMode,
halvingInterval: this.getOptionalPositiveIntegerEnv('SUBSIDY_HALVING_INTERVAL'),
...(payoutMode === 'pplns'
? {
payoutSnapshot: {
id: seed!.payoutSnapshotId,
payoutOutputs: seed!.payoutOutputs,
},
}
: {}),
});
template.notificationEventId = `prestage:${payoutMode}:${candidateHeight}:${Date.now()}`;
template.notificationPreparedAtMs = Date.now();
const update: Sv1PrestageUpdate = {
schemaVersion: 1,
type: 'subsidy-prestage',
eventId: template.notificationEventId,
template,
preparedAtMs: template.notificationPreparedAtMs,
};
await this.redisMessagingService.publishSv1PrestageUpdate(update);
return true;
} catch (error) {
if (this.lastPublishedPrestageKeys.get(payoutMode) === prestageKey) {
this.lastPublishedPrestageKeys.delete(payoutMode);
}
throw error;
}
}
private queuePplnsSubsidyBridgeSeedPrecompute(
@@ -1026,6 +1346,12 @@ export class BitcoinRpcService implements OnModuleInit {
payoutOutputs,
preparedAtMs: Date.now(),
});
const seed = this.pplnsSubsidyBridgeSeeds.get(candidateHeight);
if (seed != null) {
void this.publishNextHeightPrestage(template, 'pplns', seed).catch(error => {
console.error(`Unable to publish next-height PPLNS prestage: ${error.message}`);
});
}
} catch (error) {
if (generation === this.pplnsSeedPrecomputeGeneration) {
this.pplnsSubsidyBridgeSeeds.delete(candidateHeight);
@@ -1094,19 +1420,32 @@ export class BitcoinRpcService implements OnModuleInit {
private validateSubsidyAgainstAuthoritativeTemplate(
authoritativeTemplate: IBlockTemplate,
subsidySats: number,
): void {
): number {
const subsidySats = calculateBlockSubsidySats(
authoritativeTemplate.height,
this.getNetworkName(),
this.getOptionalPositiveIntegerEnv('SUBSIDY_HALVING_INTERVAL'),
);
if (this.subsidyValidatedTemplates.has(authoritativeTemplate)) {
return subsidySats;
}
const totalFees = authoritativeTemplate.transactions.reduce((sum, transaction) => {
if (!Number.isSafeInteger(transaction.fee) || transaction.fee < 0) {
throw new Error('GBT transaction fee is missing or invalid');
}
return sum + transaction.fee;
const next = sum + transaction.fee;
if (!Number.isSafeInteger(next)) {
throw new Error('GBT transaction fee total exceeds safe integer range');
}
return next;
}, 0);
if (subsidySats + totalFees !== authoritativeTemplate.coinbasevalue) {
throw new Error(
`GBT coinbase value ${authoritativeTemplate.coinbasevalue} does not equal subsidy ${subsidySats} plus fees ${totalFees}`,
);
}
this.subsidyValidatedTemplates.add(authoritativeTemplate);
return subsidySats;
}
private async handleSv1BridgeUpdate(update: Sv1BridgeUpdate): Promise<void> {
@@ -1117,6 +1456,8 @@ export class BitcoinRpcService implements OnModuleInit {
const template = {
...update.template,
notificationPublishedAtMs: update.publishedAtMs,
notificationWorkerReceivedAtMs: update.workerReceivedAtMs,
notificationWorkerHandledAtMs: Date.now(),
};
const payoutMode = template.payoutMode === 'pplns' ? 'pplns' : 'solo';
if (this.miningInfo?.blocks != null
@@ -1144,6 +1485,25 @@ export class BitcoinRpcService implements OnModuleInit {
this._newSv1BridgeTemplate$.next(template);
}
private handleSv1PrestageUpdate(update: Sv1PrestageUpdate): void {
if (this.processedTemplateEvents.has(update.eventId)) {
return;
}
this.rememberProcessedTemplateEvent(update.eventId);
const template = {
...update.template,
notificationPreparedAtMs: update.preparedAtMs,
};
if (template.previousblockhash !== '0'.repeat(64)
|| template.jobType !== 'empty'
|| template.transactions.length !== 0
|| (this.miningInfo?.blocks != null
&& template.height <= this.miningInfo.blocks + 1)) {
return;
}
this._newSv1PrestageTemplate$.next(template);
}
private isBridgeSupersededByCanonical(
bridge: IBlockTemplate,
payoutMode: 'solo' | 'pplns',
@@ -1409,7 +1769,16 @@ export class BitcoinRpcService implements OnModuleInit {
params: unknown[] = [],
timeoutMs?: number,
): Promise<T> {
const response = await this.client.post('', {
return this.callRpcWithClient(this.client, method, params, timeoutMs);
}
private async callRpcWithClient<T>(
client: AxiosInstance,
method: string,
params: unknown[] = [],
timeoutMs?: number,
): Promise<T> {
const response = await client.post('', {
jsonrpc: '1.0',
id: ++this.rpcRequestId,
method,
@@ -1431,6 +1800,67 @@ export class BitcoinRpcService implements OnModuleInit {
return maxTarget / target;
}
private configureAuxiliaryTemplateSources(options: {
user: string;
pass: string;
port: number;
timeout: number;
}): void {
const configured = this.configService.get<string>('BITCOIN_RPC_AUX_URLS')
?? process.env.BITCOIN_RPC_AUX_URLS
?? '';
const urls = [...new Set(configured
.split(',')
.map(value => value.trim())
.filter(Boolean))];
urls.forEach((url, index) => {
this.auxiliaryTemplateSources.push({
name: `aux-${index + 1}`,
client: axios.create({
baseURL: this.buildRpcUrl(url, options.port),
timeout: options.timeout,
auth: {
username: options.user,
password: options.pass,
},
}),
latestLongpollId: null,
longpollLoopStarted: false,
});
});
}
private async isAuxiliaryTemplateAuthorizedByPrimary(
template: IBlockTemplate,
source: TemplateRpcSource,
): Promise<boolean> {
try {
const primaryBestBlockHash = await this.callRpc<string>('getbestblockhash');
if (primaryBestBlockHash === template.previousblockhash) {
return true;
}
console.warn(JSON.stringify({
event: 'aux_template_rejected',
templateSource: source.name,
templateHeight: template.height,
previousBlockHash: template.previousblockhash,
primaryBestBlockHash,
reason: 'primary-tip-mismatch',
}));
return false;
} catch (error) {
console.warn(JSON.stringify({
event: 'aux_template_rejected',
templateSource: source.name,
templateHeight: template.height,
previousBlockHash: template.previousblockhash,
reason: 'primary-verification-failed',
error: error.message ?? String(error),
}));
return false;
}
}
private buildRpcUrl(url: string, port: number): string {
const normalizedUrl = /^https?:\/\//i.test(url) ? url : `http://${url}`;
const rpcUrl = new URL(normalizedUrl);
@@ -1443,6 +1873,7 @@ export class BitcoinRpcService implements OnModuleInit {
private startTrace(
reason: TemplateRefreshReason,
sourceNotificationReceivedAtMs?: number,
templateSource = 'primary',
): BlockNotificationTrace {
return {
eventId: `${reason}:${Date.now()}:${this.rpcRequestId + 1}`,
@@ -1450,6 +1881,7 @@ export class BitcoinRpcService implements OnModuleInit {
startedWallMs: Date.now(),
startedMonotonic: process.hrtime.bigint(),
sourceNotificationReceivedAtMs,
templateSource,
stages: { start: 0 },
};
}
@@ -1468,6 +1900,7 @@ export class BitcoinRpcService implements OnModuleInit {
event: 'block_source_notification',
eventId: trace.eventId,
source: trace.reason === 'new_block' ? 'zmq' : 'longpoll',
templateSource: trace.templateSource,
receivedAt: new Date(trace.sourceNotificationReceivedAtMs).toISOString(),
receivedAtMs: trace.sourceNotificationReceivedAtMs,
tipHeight: blockTemplate.height - 1,
@@ -1487,6 +1920,7 @@ export class BitcoinRpcService implements OnModuleInit {
event: 'block_notification_trace',
eventId: trace.eventId,
reason: trace.reason,
templateSource: trace.templateSource,
tipHeight: blockTemplate.height - 1,
templateHeight: blockTemplate.height,
previousBlockHash: blockTemplate.previousblockhash,
@@ -1560,4 +1994,20 @@ export class BitcoinRpcService implements OnModuleInit {
}
throw new Error(`Invalid NETWORK configuration: ${network ?? 'unset'}`);
}
private validateConfiguredNetworkAgainstCore(
coreChain: IMiningInfo['chain'],
): void {
const configuredNetwork = this.getNetworkName();
const expectedChain: IMiningInfo['chain'] = configuredNetwork === 'mainnet'
? 'main'
: configuredNetwork === 'testnet'
? 'test'
: 'regtest';
if (coreChain !== expectedChain) {
throw new Error(
`NETWORK=${configuredNetwork} does not match Bitcoin Core chain=${coreChain ?? 'unset'}`,
);
}
}
}
+93 -8
View File
@@ -17,7 +17,14 @@ describe('RedisMessagingService', () => {
subscriptions.clear();
clientsByRole.publisher = null;
clientsByRole.subscriber = null;
clients = [createRedisClient(), createRedisClient()];
clientsByRole.urgentPublisher = null;
clientsByRole.urgentSubscriber = null;
clients = [
createRedisClient(),
createRedisClient(),
createRedisClient(),
createRedisClient(),
];
(createClient as jest.Mock).mockImplementation(() => clients.shift());
service = new RedisMessagingService({
get: jest.fn((key: string) => key === 'REDIS_URL' ? 'redis://test-redis:6379' : null),
@@ -193,12 +200,39 @@ describe('RedisMessagingService', () => {
};
await service.subscribeSv1BridgeUpdates(handler);
await service.publishSv1BridgeUpdate(update);
await expect(service.publishSv1BridgeUpdate(update)).resolves.toBe(true);
expect(handler).toHaveBeenCalledWith(update);
expect(handler).toHaveBeenCalledWith(expect.objectContaining({
...update,
workerReceivedAtMs: expect.any(Number),
}));
expect(await service.getLatestSv1BridgeUpdate()).toEqual(update);
});
it('uses the normal command socket as an explicit bridge fallback lane', async () => {
await service.connect();
const update = createBridgeUpdate('solo', 'fallback-bridge');
await expect(service.publishSv1BridgeUpdate(update, 'fallback')).resolves.toBe(true);
expect(clientsByRole.publisher.publish).toHaveBeenCalledWith(
'sv1-bridge.updated',
JSON.stringify(update),
);
expect(clientsByRole.urgentPublisher.publish).not.toHaveBeenCalled();
});
it('reports bridge delivery failure when Redis cannot connect', async () => {
const errorSpy = jest.spyOn(console, 'error').mockImplementation(() => undefined);
clients[0].connect.mockRejectedValueOnce(new Error('redis unavailable'));
await expect(service.publishSv1BridgeUpdate(
createBridgeUpdate('solo', 'unavailable-bridge'),
)).resolves.toBe(false);
errorSpy.mockRestore();
});
it('stores and replays the latest SV1 bridge independently by payout mode', async () => {
await service.connect();
const solo = createBridgeUpdate('solo', 'solo-bridge');
@@ -213,6 +247,51 @@ describe('RedisMessagingService', () => {
expect(store.get('sv1-bridge:latest:pplns')).toBe(JSON.stringify(pplns));
});
it('publishes and durably replays next-height SV1 prestage templates', async () => {
await service.connect();
const handler = jest.fn().mockResolvedValue(undefined);
const update = {
schemaVersion: 1 as const,
type: 'subsidy-prestage' as const,
eventId: 'prestage:solo:900002',
preparedAtMs: 456,
template: {
height: 900002,
previousblockhash: '0'.repeat(64),
payoutMode: 'solo' as const,
jobType: 'empty' as const,
transactions: [],
} as any,
};
await service.subscribeSv1PrestageUpdates(handler);
await service.publishSv1PrestageUpdate(update);
expect(handler).toHaveBeenCalledWith(update);
expect(await service.getLatestSv1PrestageUpdates()).toEqual([update]);
});
it('rejects a prestage template that claims an authoritative prevhash', async () => {
await service.connect();
const update = {
schemaVersion: 1 as const,
type: 'subsidy-prestage' as const,
eventId: 'unsafe-prestage',
preparedAtMs: 456,
template: {
height: 900002,
previousblockhash: '11'.repeat(32),
payoutMode: 'solo' as const,
jobType: 'empty' as const,
transactions: [],
} as any,
};
await expect(service.publishSv1PrestageUpdate(update)).rejects.toThrow(
'unsupported SV1 prestage update',
);
});
it('rejects PPLNS bridges without an explicit snapshot and fixed-value outputs', async () => {
await service.connect();
const valid = createBridgeUpdate('pplns', 'pplns-valid');
@@ -299,7 +378,12 @@ function createBridgeUpdate(payoutMode: 'solo' | 'pplns', eventId: string) {
const store = new Map<string, string>();
const sets = new Map<string, Set<string>>();
const clientsByRole: { publisher?: any; subscriber?: any } = {};
const clientsByRole: {
publisher?: any;
subscriber?: any;
urgentPublisher?: any;
urgentSubscriber?: any;
} = {};
const subscriptions = new Map<string, (message: string) => Promise<void>>();
function createRedisClient() {
@@ -360,11 +444,12 @@ function createRedisClient() {
}),
};
if (clientsByRole.publisher == null) {
clientsByRole.publisher = client;
} else {
clientsByRole.subscriber = client;
const role = (['publisher', 'subscriber', 'urgentPublisher', 'urgentSubscriber'] as const)
.find(candidate => clientsByRole[candidate] == null);
if (role == null) {
throw new Error('Unexpected extra Redis test client');
}
clientsByRole[role] = client;
return client;
}
+127 -6
View File
@@ -9,12 +9,16 @@ import { PayoutMode } from '../types/payout-mode';
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 MINING_INFO_KEY = 'mining-info:latest';
const BLOCK_TEMPLATE_LATEST_KEY = 'block-template:latest';
const SV1_BRIDGE_LATEST_KEY = 'sv1-bridge:latest';
const SV1_PRESTAGE_LATEST_KEY = 'sv1-prestage:latest';
const BLOCK_TEMPLATE_CACHE_TTL_SECONDS = 60 * 60;
const sv1BridgeLatestKey = (payoutMode: PayoutMode) =>
`${SV1_BRIDGE_LATEST_KEY}:${payoutMode}`;
const sv1PrestageLatestKey = (payoutMode: PayoutMode) =>
`${SV1_PRESTAGE_LATEST_KEY}:${payoutMode}`;
const blockTemplateLegacyKey = (height: number) => `block-template:${height}`;
const blockTemplateKey = (
height: number,
@@ -42,12 +46,26 @@ export interface Sv1BridgeUpdate {
eventId: string;
template: IBlockTemplate;
publishedAtMs: number;
/** Local worker timestamp; populated after Redis delivery, never serialized by the master. */
workerReceivedAtMs?: number;
}
export interface Sv1PrestageUpdate {
schemaVersion: 1;
type: 'subsidy-prestage';
eventId: string;
template: IBlockTemplate;
preparedAtMs: number;
}
@Injectable()
export class RedisMessagingService implements OnModuleInit, OnModuleDestroy {
private publisher: RedisClientType;
private subscriber: RedisClientType;
/** Dedicated command socket so full-template writes cannot queue ahead of a new-tip bridge. */
private urgentPublisher: RedisClientType;
/** Dedicated Pub/Sub socket so canonical/replay callbacks cannot delay bridge receipt. */
private urgentSubscriber: RedisClientType;
private connected = false;
constructor(
@@ -64,6 +82,8 @@ export class RedisMessagingService implements OnModuleInit, OnModuleDestroy {
await Promise.all([
this.publisher?.quit().catch(() => undefined),
this.subscriber?.quit().catch(() => undefined),
this.urgentPublisher?.quit().catch(() => undefined),
this.urgentSubscriber?.quit().catch(() => undefined),
]);
}
@@ -78,13 +98,24 @@ export class RedisMessagingService implements OnModuleInit, OnModuleDestroy {
};
this.publisher = createClient({ url, socket });
this.subscriber = createClient({ url, socket });
this.urgentPublisher = createClient({ url, socket });
this.urgentSubscriber = createClient({ url, socket });
this.publisher.on('error', error => console.error(`Redis publisher error: ${error.message}`));
this.subscriber.on('error', error => console.error(`Redis subscriber error: ${error.message}`));
this.urgentPublisher.on('error', error => console.error(`Redis urgent publisher error: ${error.message}`));
this.urgentSubscriber.on('error', error => console.error(`Redis urgent subscriber error: ${error.message}`));
this.publisher.on('end', () => { this.connected = false; });
this.subscriber.on('end', () => { this.connected = false; });
this.urgentPublisher.on('end', () => { this.connected = false; });
this.urgentSubscriber.on('end', () => { this.connected = false; });
await Promise.all([this.publisher.connect(), this.subscriber.connect()]);
await Promise.all([
this.publisher.connect(),
this.subscriber.connect(),
this.urgentPublisher.connect(),
this.urgentSubscriber.connect(),
]);
this.connected = true;
}
@@ -136,16 +167,20 @@ export class RedisMessagingService implements OnModuleInit, OnModuleDestroy {
});
}
public async publishSv1BridgeUpdate(update: Sv1BridgeUpdate): Promise<void> {
public async publishSv1BridgeUpdate(
update: Sv1BridgeUpdate,
lane: 'urgent' | 'fallback' = 'urgent',
): Promise<boolean> {
if (!await this.ensureConnected()) {
return;
return false;
}
const serialized = JSON.stringify(update);
const validated = this.parseSv1BridgeUpdate(serialized);
const payoutMode = validated.template.payoutMode === 'pplns' ? 'pplns' : 'solo';
// The bridge envelope is self-contained. Deliver it before spending a
// second Redis round trip on best-effort replay metadata.
await this.publisher.publish(SV1_BRIDGE_CHANNEL, serialized);
const publisher = lane === 'urgent' ? this.urgentPublisher : this.publisher;
await publisher.publish(SV1_BRIDGE_CHANNEL, serialized);
void Promise.all([
this.publisher.setEx(
sv1BridgeLatestKey(payoutMode),
@@ -162,15 +197,18 @@ export class RedisMessagingService implements OnModuleInit, OnModuleDestroy {
]).catch(error => {
console.error(`Unable to cache latest SV1 bridge: ${error.message}`);
});
return true;
}
public async subscribeSv1BridgeUpdates(handler: (update: Sv1BridgeUpdate) => Promise<void>): Promise<void> {
if (!await this.ensureConnected()) {
return;
}
await this.subscriber.subscribe(SV1_BRIDGE_CHANNEL, async message => {
await this.urgentSubscriber.subscribe(SV1_BRIDGE_CHANNEL, async message => {
try {
await handler(this.parseSv1BridgeUpdate(message));
const update = this.parseSv1BridgeUpdate(message);
update.workerReceivedAtMs = Date.now();
await handler(update);
} catch (error) {
console.error(`Invalid Redis SV1 bridge update: ${error.message}`);
}
@@ -201,6 +239,53 @@ export class RedisMessagingService implements OnModuleInit, OnModuleDestroy {
}
}
public async publishSv1PrestageUpdate(update: Sv1PrestageUpdate): Promise<void> {
if (!await this.ensureConnected()) {
return;
}
const serialized = JSON.stringify(update);
const validated = this.parseSv1PrestageUpdate(serialized);
const payoutMode = validated.template.payoutMode === 'pplns' ? 'pplns' : 'solo';
await Promise.all([
// Keep the current next-height seed until it is replaced. A Bitcoin
// block interval can exceed the canonical-template cache TTL, and a
// worker restart late in that interval must still be able to stage.
this.publisher.set(
sv1PrestageLatestKey(payoutMode),
serialized,
),
this.publisher.publish(SV1_PRESTAGE_CHANNEL, serialized),
]);
}
public async subscribeSv1PrestageUpdates(
handler: (update: Sv1PrestageUpdate) => Promise<void>,
): Promise<void> {
if (!await this.ensureConnected()) {
return;
}
await this.subscriber.subscribe(SV1_PRESTAGE_CHANNEL, async message => {
try {
await handler(this.parseSv1PrestageUpdate(message));
} catch (error) {
console.error(`Invalid Redis SV1 prestage update: ${error.message}`);
}
});
}
public async getLatestSv1PrestageUpdates(): Promise<Sv1PrestageUpdate[]> {
if (!await this.ensureConnected()) {
return [];
}
const values = await this.publisher.mGet([
sv1PrestageLatestKey('solo'),
sv1PrestageLatestKey('pplns'),
]) as Array<string | null>;
return values
.filter((value): value is string => value != null)
.map(value => this.parseSv1PrestageUpdate(value));
}
public async setLatestMiningInfo(miningInfo: IMiningInfo) {
if (!await this.ensureConnected()) {
return;
@@ -411,6 +496,42 @@ export class RedisMessagingService implements OnModuleInit, OnModuleDestroy {
return update as Sv1BridgeUpdate;
}
private parseSv1PrestageUpdate(message: string): Sv1PrestageUpdate {
const update = JSON.parse(message) as Partial<Sv1PrestageUpdate>;
const template = update.template as Partial<IBlockTemplate> | undefined;
const payoutMode = template?.payoutMode;
const payoutOutputs = template?.payoutOutputs;
const hasExplicitPplnsPayout = payoutMode !== 'pplns' || (
typeof template?.payoutSnapshotId === 'string'
&& template.payoutSnapshotId.trim().length > 0
&& Array.isArray(payoutOutputs)
&& payoutOutputs.length > 0
&& Number.isSafeInteger(template.coinbasevalue)
&& payoutOutputs.every(output => (
typeof output.address === 'string'
&& output.address.trim().length > 0
&& Number.isSafeInteger(output.amountSats)
&& output.amountSats >= 0
))
&& payoutOutputs.reduce((sum, output) => sum + output.amountSats, 0)
=== template.coinbasevalue
);
if (update.schemaVersion !== 1
|| update.type !== 'subsidy-prestage'
|| typeof update.eventId !== 'string'
|| !Number.isFinite(update.preparedAtMs)
|| template?.jobType !== 'empty'
|| (payoutMode !== 'solo' && payoutMode !== 'pplns')
|| !Array.isArray(template?.transactions)
|| template.transactions.length !== 0
|| template.previousblockhash !== '0'.repeat(64)
|| !Number.isInteger(template.height)
|| !hasExplicitPplnsPayout) {
throw new Error('unsupported SV1 prestage update');
}
return update as Sv1PrestageUpdate;
}
private async ensureConnected(): Promise<boolean> {
if (!this.connected) {
try {
+71 -1
View File
@@ -5,11 +5,13 @@ import { MockRecording1 } from '../../test/models/MockRecording1';
import { IBlockTemplate } from '../models/bitcoin-rpc/IBlockTemplate';
import { createPreparedMiningJob } from './prepared-mining-job.factory';
import { StratumV1JobsService } from './stratum-v1-jobs.service';
import { EMPTY_DEFAULT_WITNESS_COMMITMENT } from './subsidy-only-template.factory';
describe('StratumV1JobsService', () => {
let blockTemplate$: BehaviorSubject<IBlockTemplate>;
let bridgeTemplate$: Subject<IBlockTemplate>;
let bitcoinRpcService: { newBlockTemplate$: any, newSv1BridgeTemplate$: any, miningInfo: { blocks: number } };
let prestageTemplate$: Subject<IBlockTemplate>;
let bitcoinRpcService: { newBlockTemplate$: any, newSv1BridgeTemplate$: any, newSv1PrestageTemplate$: any, miningInfo: { blocks: number } };
let service: StratumV1JobsService;
let consoleLogSpy: jest.SpyInstance;
@@ -26,9 +28,11 @@ describe('StratumV1JobsService', () => {
blockTemplate$ = new BehaviorSubject(createTemplate());
bridgeTemplate$ = new Subject<IBlockTemplate>();
prestageTemplate$ = new Subject<IBlockTemplate>();
bitcoinRpcService = {
newBlockTemplate$: blockTemplate$.asObservable(),
newSv1BridgeTemplate$: bridgeTemplate$.asObservable(),
newSv1PrestageTemplate$: prestageTemplate$.asObservable(),
miningInfo: { blocks: MockRecording1.BLOCK_TEMPLATE.height }
};
service = new StratumV1JobsService(bitcoinRpcService as any);
@@ -141,6 +145,72 @@ describe('StratumV1JobsService', () => {
expect(full.merkle_branch.length).toBeGreaterThan(0);
});
it('prebuilds a next-height coinbase and promotes the same job on authoritative activation', async () => {
await firstValueFrom(service.newMiningJob$);
const future = createTemplate(MockRecording1.BLOCK_TEMPLATE.height + 1);
future.previousblockhash = '0'.repeat(64);
future.transactions = [];
future.coinbasevalue = 312_500_000;
future.default_witness_commitment = EMPTY_DEFAULT_WITNESS_COMMITMENT;
future.jobType = 'empty';
future.payoutMode = 'solo';
future.forceCleanJobs = true;
const detachedResult = firstValueFrom(service.sv1PrestageJob$);
prestageTemplate$.next(future);
const detached = await detachedResult;
const payout = [{
address: 'tb1qumezefzdeqqwn5zfvgdrhxjzc5ylr39uhuxcz4',
percent: 100,
}];
const payoutIdentity = 'solo-prestaged-miner';
const staged = service.preStageJob(
bitcoinjs.networks.testnet,
payout,
detached,
payoutIdentity,
'solo',
);
expect(staged).not.toBeNull();
expect(service.getJobById(staged.jobId)).toBeUndefined();
const prebuiltNotifyBuffer = staged.responseBuffer(detached);
expect(JSON.parse(prebuiltNotifyBuffer.toString()).params[1]).toBe('0'.repeat(64));
const activatedSource = {
...future,
previousblockhash: 'ab'.repeat(32),
notificationEventId: 'authoritative-activation',
};
const activatedResult = firstValueFrom(service.sv1MiningJob$.pipe(skip(1)));
bridgeTemplate$.next(activatedSource);
const activatedTemplate = await activatedResult;
const activated = service.activatePreStagedJob(
activatedTemplate,
payoutIdentity,
'solo',
);
expect(activated).toBe(staged);
expect(staged.responseBuffer(activatedTemplate)).toBe(prebuiltNotifyBuffer);
expect(service.getJobById(staged.jobId)).toBe(staged);
expect(service.getSubmissionContext(staged.jobId)?.status).toBe('current');
const notify = JSON.parse(staged.response(activatedTemplate));
expect(notify.params[1]).toBe(
Buffer.from(activatedTemplate.block.prevHash).swap32().toString('hex'),
);
expect(notify.params.slice(5, 9)).toEqual([
activatedTemplate.block.version.toString(16),
activatedTemplate.block.bits.toString(16),
activatedTemplate.block.timestamp.toString(16),
true,
]);
expect((service as any).stagedJobs.size).toBe(1);
(service as any).cleanupPrestageJobs(activatedTemplate.blockData.height + 1);
expect((service as any).stagedJobs.size).toBe(0);
expect(service.getJobById(staged.jobId)).toBe(staged);
});
it('should keep same-tip bridge and full jobs current, then mark both stale on the next tip', async () => {
const bridgeTemplate = await firstValueFrom(service.newMiningJob$);
const payout = [{ address: 'tb1qumezefzdeqqwn5zfvgdrhxjzc5ylr39uhuxcz4', percent: 100 }];
+212 -7
View File
@@ -4,7 +4,7 @@ import * as crypto from 'crypto';
import { EMPTY, filter, map, merge, Observable, shareReplay, tap } from 'rxjs';
import { IBlockTemplate, IBlockTemplateTx } from '../models/bitcoin-rpc/IBlockTemplate';
import { AddressObject, MiningJob } from '../models/MiningJob';
import { AddressObject, MiningJob, MiningNotifyHeaderFields } from '../models/MiningJob';
import { PayoutMode } from '../types/payout-mode';
import { BitcoinRpcService } from './bitcoin-rpc.service';
import {
@@ -29,7 +29,10 @@ export interface IJobTemplate {
payoutMode: PayoutMode | 'all';
notificationEventId?: string;
sourceNotificationReceivedAtMs?: number;
notificationPreparedAtMs?: number;
notificationPublishedAtMs?: number;
notificationWorkerReceivedAtMs?: number;
notificationWorkerHandledAtMs?: number;
payoutSnapshotId?: string;
payoutOutputs?: AddressObject[];
transactions?: IBlockTemplateTx[];
@@ -49,6 +52,7 @@ export interface IJobSubmissionContext {
const DEFAULT_JOB_RETENTION_MS = 5 * 60 * 1000;
const PAYOUT_MODES: readonly PayoutMode[] = ['solo', 'pplns'];
const PLACEHOLDER_PREV_HASH = Buffer.alloc(32, 0);
export function createPayoutOutputIdentity(
payoutMode: PayoutMode,
@@ -70,6 +74,8 @@ export class StratumV1JobsService {
public newMiningJob$: Observable<IJobTemplate>;
/** Ordered SV1-only stream: subsidy bridge first, canonical full job second. */
public sv1MiningJob$: Observable<IJobTemplate>;
/** Detached next-height empty jobs prepared without changing active work. */
public sv1PrestageJob$: Observable<IJobTemplate>;
public latestJobId: number = 1;
public latestJobTemplateId: number = 1;
public jobs: { [jobId: string]: MiningJob } = {};
@@ -79,7 +85,13 @@ export class StratumV1JobsService {
private readonly lastWorkSignatures = new Map<PayoutMode, string>();
private readonly currentTipKeys = new Map<PayoutMode, string>();
private readonly cachedJobs = new Map<string, MiningJob>();
private readonly notifyHeaderFields = new WeakMap<IJobTemplate, MiningNotifyHeaderFields>();
private readonly latestJobTemplates = new Map<PayoutMode | 'all', IJobTemplate>();
private readonly latestPrestageJobTemplates = new Map<PayoutMode, IJobTemplate>();
private readonly stagedJobs = new Map<string, {
job: MiningJob;
activatedTemplateId?: string;
}>();
private readonly pinnedCurrentTipTemplateIds = new Map<PayoutMode, Set<string>>([
['solo', new Set<string>()],
['pplns', new Set<string>()],
@@ -152,7 +164,10 @@ export class StratumV1JobsService {
isNewBlock,
notificationEventId: blockTemplate.notificationEventId,
sourceNotificationReceivedAtMs: blockTemplate.sourceNotificationReceivedAtMs,
notificationPreparedAtMs: blockTemplate.notificationPreparedAtMs,
notificationPublishedAtMs: blockTemplate.notificationPublishedAtMs,
notificationWorkerReceivedAtMs: blockTemplate.notificationWorkerReceivedAtMs,
notificationWorkerHandledAtMs: blockTemplate.notificationWorkerHandledAtMs,
rawTransactions: blockTemplate.transactions,
sigoplimit: blockTemplate.sigoplimit,
sizelimit: blockTemplate.sizelimit,
@@ -161,7 +176,7 @@ export class StratumV1JobsService {
};
}),
filter(next => next != null),
map(({ prepared, timestamp, networkDifficulty, clearJobs, isNewBlock, notificationEventId, sourceNotificationReceivedAtMs, notificationPublishedAtMs, rawTransactions, sigoplimit, sizelimit, weightlimit, requiredVersionBits }) => {
map(({ prepared, timestamp, networkDifficulty, clearJobs, isNewBlock, notificationEventId, sourceNotificationReceivedAtMs, notificationPreparedAtMs, notificationPublishedAtMs, notificationWorkerReceivedAtMs, notificationWorkerHandledAtMs, rawTransactions, sigoplimit, sizelimit, weightlimit, requiredVersionBits }) => {
const block = new bitcoinjs.Block();
// Keep only a placeholder coinbase on the hot path. The full raw body
@@ -201,7 +216,10 @@ export class StratumV1JobsService {
payoutMode: prepared.payoutMode,
notificationEventId,
sourceNotificationReceivedAtMs,
notificationPreparedAtMs,
notificationPublishedAtMs,
notificationWorkerReceivedAtMs,
notificationWorkerHandledAtMs,
payoutSnapshotId: prepared.coinbase.payoutSnapshotId,
payoutOutputs: prepared.coinbase.payoutOutputs?.map(output => ({ ...output })),
transactions: rawTransactions,
@@ -244,9 +262,20 @@ export class StratumV1JobsService {
this.sv1MiningJob$ = merge(bridgeMiningJob$, this.newMiningJob$).pipe(
shareReplay({ refCount: true, bufferSize: 1 }),
);
this.sv1PrestageJob$ = (
this.bitcoinRpcService.newSv1PrestageTemplate$ ?? EMPTY
).pipe(
map(blockTemplate => this.createDetachedPrestageJobTemplate(blockTemplate)),
tap(jobTemplate => {
const payoutMode = jobTemplate.blockData.payoutMode === 'pplns' ? 'pplns' : 'solo';
this.latestPrestageJobTemplates.set(payoutMode, jobTemplate);
}),
shareReplay({ refCount: true, bufferSize: 2 }),
);
if (process.env.API_ONLY !== 'true' && process.env.MASTER !== 'true') {
this.sv1MiningJob$.subscribe();
this.sv1PrestageJob$.subscribe();
}
}
@@ -272,6 +301,80 @@ export class StratumV1JobsService {
?? null;
}
public getLatestPrestageJobTemplate(payoutMode: PayoutMode): IJobTemplate | null {
return this.latestPrestageJobTemplates.get(payoutMode) ?? null;
}
public preStageJob(
network: bitcoinjs.networks.Network,
payoutInformation: AddressObject[],
jobTemplate: IJobTemplate,
payoutIdentity: string,
payoutMode: PayoutMode,
): MiningJob | null {
if (jobTemplate.blockData.jobType !== 'empty'
|| jobTemplate.blockData.payoutMode !== payoutMode
|| !jobTemplate.block.prevHash.equals(PLACEHOLDER_PREV_HASH)) {
return null;
}
const key = this.getPrestageJobKey(jobTemplate, payoutIdentity, payoutMode);
const existing = this.stagedJobs.get(key);
if (existing != null) {
return existing.job;
}
const job = new MiningJob(
network,
this.getNextId(),
payoutInformation,
jobTemplate,
{ payoutMode, payoutIdentity },
);
job.preparePreStagedNotify(
jobTemplate,
this.getNotifyHeaderFields(jobTemplate),
);
// Reserve the id now; staged jobs are intentionally absent from `jobs`
// until their authoritative prevhash arrives and miners can submit them.
this.latestJobId++;
this.stagedJobs.set(key, { job });
this.cleanupPrestageJobs(jobTemplate.blockData.height);
return job;
}
public activatePreStagedJob(
jobTemplate: IJobTemplate,
payoutIdentity: string,
payoutMode: PayoutMode,
): MiningJob | null {
if (jobTemplate.blockData.jobType !== 'empty') {
return null;
}
const key = this.getPrestageJobKey(jobTemplate, payoutIdentity, payoutMode);
const staged = this.stagedJobs.get(key);
if (staged == null) {
return null;
}
if (staged.activatedTemplateId != null
&& staged.activatedTemplateId !== jobTemplate.blockData.id) {
// A same-height reorg needs a new job id/payload so a cached response
// can never retain the orphaned prevhash.
return null;
}
if (staged.activatedTemplateId == null) {
staged.job.activatePreStagedTemplate(
jobTemplate,
this.getNotifyHeaderFields(jobTemplate),
);
staged.activatedTemplateId = jobTemplate.blockData.id;
this.jobs[staged.job.jobId] = staged.job;
this.cachedJobs.set(
this.getJobCacheKey(jobTemplate, payoutMode, payoutIdentity),
staged.job,
);
}
return staged.job;
}
public addJob(job: MiningJob) {
this.jobs[job.jobId] = job;
this.latestJobId++;
@@ -288,13 +391,11 @@ export class StratumV1JobsService {
?? (jobTemplate.blockData.payoutMode === 'pplns' ? 'pplns' : 'solo');
const effectivePayoutIdentity = payoutIdentity
?? createPayoutOutputIdentity(effectivePayoutMode, payoutInformation);
const cacheKey = [
jobTemplate.blockData.id,
jobTemplate.block.timestamp,
jobTemplate.blockData.clearJobs,
const cacheKey = this.getJobCacheKey(
jobTemplate,
effectivePayoutMode,
effectivePayoutIdentity,
].join(':');
);
const cached = this.cachedJobs.get(cacheKey);
if (cached != null) {
return cached;
@@ -388,6 +489,110 @@ export class StratumV1JobsService {
return payoutMode === 'all' ? PAYOUT_MODES : [payoutMode];
}
private createDetachedPrestageJobTemplate(blockTemplate: IBlockTemplate): IJobTemplate {
const prepared = createPreparedMiningJob(blockTemplate);
if (prepared.jobType !== 'empty'
|| prepared.header.previousBlockHash !== '0'.repeat(64)
|| (prepared.payoutMode !== 'solo' && prepared.payoutMode !== 'pplns')) {
throw new Error('SV1 prestage must be a payout-specific empty template with a placeholder prevhash');
}
const block = new bitcoinjs.Block();
const tempCoinbaseTx = new bitcoinjs.Transaction();
tempCoinbaseTx.version = 2;
tempCoinbaseTx.addInput(Buffer.alloc(32, 0), 0xffffffff, 0xffffffff);
tempCoinbaseTx.ins[0].witness = [Buffer.alloc(32, 0)];
block.prevHash = Buffer.alloc(32, 0);
block.version = prepared.header.version;
block.bits = prepared.header.bits;
block.timestamp = Math.max(
prepared.header.minTime,
Math.floor(Date.now() / 1000),
);
block.transactions = [tempCoinbaseTx];
block.merkleRoot = tempCoinbaseTx.getHash(false);
block.witnessCommit = Buffer.from(prepared.coinbase.witnessCommitmentHash, 'hex');
const id = `prestage-${this.getNextTemplateId()}`;
this.latestJobTemplateId++;
return {
block,
merkle_branch: [],
blockData: {
id,
creation: Date.now(),
coinbasevalue: prepared.coinbase.valueSats,
networkDifficulty: this.calculateNetworkDifficulty(prepared.header.bits),
height: prepared.height,
tipKey: `prestage:${prepared.height}`,
clearJobs: true,
isNewBlock: false,
jobType: 'empty',
payoutMode: prepared.payoutMode,
notificationEventId: blockTemplate.notificationEventId,
notificationPreparedAtMs: blockTemplate.notificationPreparedAtMs,
payoutSnapshotId: prepared.coinbase.payoutSnapshotId,
payoutOutputs: prepared.coinbase.payoutOutputs?.map(output => ({ ...output })),
transactions: [],
bodyReference: prepared.body,
sigoplimit: blockTemplate.sigoplimit,
sizelimit: blockTemplate.sizelimit,
weightlimit: blockTemplate.weightlimit,
requiredVersionBits: blockTemplate.vbrequired >>> 0,
},
};
}
private getJobCacheKey(
jobTemplate: IJobTemplate,
payoutMode: PayoutMode,
payoutIdentity: string,
): string {
return [
jobTemplate.blockData.id,
jobTemplate.block.timestamp,
jobTemplate.blockData.clearJobs,
payoutMode,
payoutIdentity,
].join(':');
}
public getNotifyHeaderFields(jobTemplate: IJobTemplate): MiningNotifyHeaderFields {
const cached = this.notifyHeaderFields.get(jobTemplate);
if (cached != null) {
return cached;
}
const fields = MiningJob.createNotifyHeaderFields(jobTemplate);
this.notifyHeaderFields.set(jobTemplate, fields);
return fields;
}
private getPrestageJobKey(
jobTemplate: IJobTemplate,
payoutIdentity: string,
payoutMode: PayoutMode,
): string {
return [
jobTemplate.blockData.height,
jobTemplate.blockData.coinbasevalue,
jobTemplate.block.witnessCommit.toString('hex'),
payoutMode,
payoutIdentity,
jobTemplate.blockData.payoutSnapshotId ?? '',
].join(':');
}
private cleanupPrestageJobs(latestCandidateHeight: number): void {
for (const key of this.stagedJobs.keys()) {
const height = Number(key.slice(0, key.indexOf(':')));
if (height < latestCandidateHeight) {
// Activated jobs remain available through `jobs`/`cachedJobs` for
// candidate reconstruction and normal retention. The staging
// index is only needed through that height's fanout and would
// otherwise retain every miner payout identity forever.
this.stagedJobs.delete(key);
}
}
}
private getJobPayoutModes(job: MiningJob, jobTemplate: IJobTemplate): readonly PayoutMode[] {
if (job.ownership?.payoutMode != null) {
return [job.ownership.payoutMode];
+41 -1
View File
@@ -12,6 +12,7 @@ describe('StratumV1Service', () => {
const originalTlsHandshakeTimeoutMs = process.env.STRATUM_TLS_HANDSHAKE_TIMEOUT_MS;
const originalSocketTimeoutMs = process.env.STRATUM_SOCKET_TIMEOUT_MS;
const originalTcpKeepAliveInitialDelayMs = process.env.STRATUM_TCP_KEEPALIVE_INITIAL_DELAY_MS;
const originalFanoutTargetClientsPerWorker = process.env.STRATUM_FANOUT_TARGET_CLIENTS_PER_WORKER;
let service: StratumV1Service;
let clientService;
@@ -19,6 +20,7 @@ describe('StratumV1Service', () => {
let stratumV2Service;
let redisMessagingService;
let miningJobs: Subject<any>;
let prestageJobs: Subject<any>;
let consoleLogSpy: jest.SpyInstance;
let consoleWarnSpy: jest.SpyInstance;
@@ -36,13 +38,17 @@ describe('StratumV1Service', () => {
};
redisMessagingService = {};
miningJobs = new Subject();
prestageJobs = new Subject();
service = new StratumV1Service(
{} as any,
clientService,
{} as any,
{} as any,
{} as any,
{ newMiningJob$: miningJobs.asObservable() } as any,
{
newMiningJob$: miningJobs.asObservable(),
sv1PrestageJob$: prestageJobs.asObservable(),
} as any,
{} as any,
stratumV2Service as any,
userAgentReportService as any,
@@ -63,6 +69,7 @@ describe('StratumV1Service', () => {
restoreEnv('STRATUM_TLS_HANDSHAKE_TIMEOUT_MS', originalTlsHandshakeTimeoutMs);
restoreEnv('STRATUM_SOCKET_TIMEOUT_MS', originalSocketTimeoutMs);
restoreEnv('STRATUM_TCP_KEEPALIVE_INITIAL_DELAY_MS', originalTcpKeepAliveInitialDelayMs);
restoreEnv('STRATUM_FANOUT_TARGET_CLIENTS_PER_WORKER', originalFanoutTargetClientsPerWorker);
consoleLogSpy.mockRestore();
consoleWarnSpy.mockRestore();
jest.useRealTimers();
@@ -137,6 +144,7 @@ describe('StratumV1Service', () => {
});
it('enqueues a 100,000-client fanout without per-client async serialization', () => {
process.env.STRATUM_FANOUT_TARGET_CLIENTS_PER_WORKER = '10000';
const clientCount = 100_000;
let writes = 0;
const broadcastMiningJob = () => {
@@ -166,6 +174,8 @@ describe('StratumV1Service', () => {
event: 'stratum_job_fanout',
eventId: 'load-test',
clients: clientCount,
targetClientsPerWorker: 10_000,
overTargetClients: 90_000,
written: clientCount,
}));
expect(trace.milestoneMs).toEqual(expect.objectContaining({
@@ -176,6 +186,36 @@ describe('StratumV1Service', () => {
expect(trace.totalMs).toBeLessThan(1_000);
});
it('pre-stages next-height jobs for connected miners without broadcasting them', async () => {
process.env.MASTER = 'false';
process.env.STRATUM_PORTS = '';
process.env.STRATUM_SECURE = 'false';
const clients = Array.from({ length: 3 }, () => ({
preStageMiningJob: jest.fn().mockReturnValue(true),
broadcastMiningJob: jest.fn(),
}));
clients.forEach(client => (service as any).clients.add(client));
await service.onModuleInit();
const prestage = {
blockData: {
id: 'prestage-2',
height: 900002,
payoutMode: 'solo',
notificationEventId: 'prestage-event',
},
};
prestageJobs.next(prestage);
await Promise.resolve();
clients.forEach(client => {
expect(client.preStageMiningJob).toHaveBeenCalledWith(prestage);
expect(client.broadcastMiningJob).not.toHaveBeenCalled();
});
expect(consoleLogSpy).toHaveBeenCalledWith(expect.stringContaining('sv1_job_prestage'));
service.onModuleDestroy();
});
it('does not log routine non-new-block fanout unless explicitly enabled', () => {
(service as any).clients.add({
broadcastMiningJob: jest.fn().mockReturnValue({
+125 -2
View File
@@ -41,6 +41,8 @@ const DEFAULT_MAX_CONNECTIONS_PER_LISTENER = 10000;
const DEFAULT_TLS_HANDSHAKE_TIMEOUT_MS = 10000;
const DEFAULT_SOCKET_TIMEOUT_MS = 1000 * 60 * 60;
const DEFAULT_TCP_KEEPALIVE_INITIAL_DELAY_MS = 1000 * 60;
const DEFAULT_PRESTAGE_BATCH_SIZE = 500;
const DEFAULT_FANOUT_TARGET_CLIENTS_PER_WORKER = 10000;
@@ -57,6 +59,10 @@ export class StratumV1Service implements OnModuleInit, OnModuleDestroy {
private healthyBackpressureChecks = 0;
private readonly clients = new Set<StratumV1Client>();
private jobBroadcastSubscription: Subscription | null = null;
private jobPrestageSubscription: Subscription | null = null;
private readonly pendingPrestageJobs = new Map<PayoutMode, import('./stratum-v1-jobs.service').IJobTemplate>();
private prestageDrainRunning = false;
private prestageGeneration = 0;
constructor(
private readonly bitcoinRpcService: BitcoinRpcService,
@@ -96,6 +102,10 @@ export class StratumV1Service implements OnModuleInit, OnModuleDestroy {
next: jobTemplate => this.broadcastMiningJob(jobTemplate),
error: error => console.error(`SV1 job broadcast subscription failed: ${error.message}`),
});
this.jobPrestageSubscription = this.stratumV1JobsService.sv1PrestageJob$?.subscribe({
next: jobTemplate => this.queuePrestageMiningJob(jobTemplate),
error: error => console.error(`SV1 job prestage subscription failed: ${error.message}`),
}) ?? null;
// wait for all the other processes to init for an even connection distribution
setTimeout(() => {
@@ -130,6 +140,10 @@ export class StratumV1Service implements OnModuleInit, OnModuleDestroy {
public onModuleDestroy(): void {
this.jobBroadcastSubscription?.unsubscribe();
this.jobBroadcastSubscription = null;
this.jobPrestageSubscription?.unsubscribe();
this.jobPrestageSubscription = null;
this.prestageGeneration++;
this.pendingPrestageJobs.clear();
if (this.backpressureMonitor != null) {
clearInterval(this.backpressureMonitor);
this.backpressureMonitor = null;
@@ -459,8 +473,13 @@ export class StratumV1Service implements OnModuleInit, OnModuleDestroy {
}
private broadcastMiningJob(jobTemplate: import('./stratum-v1-jobs.service').IJobTemplate): void {
const fanoutStartedAtMs = Date.now();
const startedAt = process.hrtime.bigint();
const totalClients = this.clients.size;
const targetClientsPerWorker = this.getPositiveIntegerEnv(
'STRATUM_FANOUT_TARGET_CLIENTS_PER_WORKER',
DEFAULT_FANOUT_TARGET_CLIENTS_PER_WORKER,
);
const milestoneIndexes = {
p50: Math.max(1, Math.ceil(totalClients * 0.5)),
p95: Math.max(1, Math.ceil(totalClients * 0.95)),
@@ -473,6 +492,7 @@ export class StratumV1Service implements OnModuleInit, OnModuleDestroy {
let backpressured = 0;
let closed = 0;
let errors = 0;
let preStaged = 0;
let bytesQueued = 0;
let maxBufferedBytes = 0;
@@ -481,6 +501,9 @@ export class StratumV1Service implements OnModuleInit, OnModuleDestroy {
visited++;
try {
const result = client.broadcastMiningJob(jobTemplate);
if (result.preStaged) {
preStaged++;
}
bytesQueued += result.bytes;
maxBufferedBytes = Math.max(maxBufferedBytes, result.bufferedBytes);
switch (result.status) {
@@ -515,21 +538,47 @@ export class StratumV1Service implements OnModuleInit, OnModuleDestroy {
eventId: jobTemplate.blockData.notificationEventId,
sourceToFanoutStartMs: jobTemplate.blockData.sourceNotificationReceivedAtMs == null
? undefined
: Date.now() - jobTemplate.blockData.sourceNotificationReceivedAtMs,
: fanoutStartedAtMs - jobTemplate.blockData.sourceNotificationReceivedAtMs,
masterPrepareMs: jobTemplate.blockData.notificationPreparedAtMs == null
|| jobTemplate.blockData.sourceNotificationReceivedAtMs == null
? undefined
: jobTemplate.blockData.notificationPreparedAtMs
- jobTemplate.blockData.sourceNotificationReceivedAtMs,
masterPublishRequestMs: jobTemplate.blockData.notificationPublishedAtMs == null
|| jobTemplate.blockData.sourceNotificationReceivedAtMs == null
? undefined
: jobTemplate.blockData.notificationPublishedAtMs
- jobTemplate.blockData.sourceNotificationReceivedAtMs,
masterToWorkerReceiveMs: jobTemplate.blockData.notificationPublishedAtMs == null
|| jobTemplate.blockData.notificationWorkerReceivedAtMs == null
? undefined
: jobTemplate.blockData.notificationWorkerReceivedAtMs
- jobTemplate.blockData.notificationPublishedAtMs,
workerReceiveToHandleMs: jobTemplate.blockData.notificationWorkerReceivedAtMs == null
|| jobTemplate.blockData.notificationWorkerHandledAtMs == null
? undefined
: jobTemplate.blockData.notificationWorkerHandledAtMs
- jobTemplate.blockData.notificationWorkerReceivedAtMs,
workerHandleToFanoutStartMs: jobTemplate.blockData.notificationWorkerHandledAtMs == null
? undefined
: fanoutStartedAtMs - jobTemplate.blockData.notificationWorkerHandledAtMs,
redisToFanoutStartMs: jobTemplate.blockData.notificationPublishedAtMs == null
? undefined
: Date.now() - jobTemplate.blockData.notificationPublishedAtMs,
: fanoutStartedAtMs - jobTemplate.blockData.notificationPublishedAtMs,
templateId: jobTemplate.blockData.id,
height: jobTemplate.blockData.height,
jobType: jobTemplate.blockData.jobType,
isNewBlock: jobTemplate.blockData.isNewBlock,
cleanJobs: jobTemplate.blockData.clearJobs,
clients: totalClients,
targetClientsPerWorker,
overTargetClients: Math.max(0, totalClients - targetClientsPerWorker),
written,
skipped,
backpressured,
closed,
errors,
preStaged,
bytesQueued,
maxBufferedBytes,
milestoneMs,
@@ -537,6 +586,80 @@ export class StratumV1Service implements OnModuleInit, OnModuleDestroy {
}));
}
private queuePrestageMiningJob(
jobTemplate: import('./stratum-v1-jobs.service').IJobTemplate,
): void {
const payoutMode = jobTemplate.blockData.payoutMode === 'pplns' ? 'pplns' : 'solo';
this.pendingPrestageJobs.set(payoutMode, jobTemplate);
if (this.prestageDrainRunning) {
return;
}
this.prestageDrainRunning = true;
const generation = this.prestageGeneration;
void this.drainPrestageMiningJobs(generation).finally(() => {
this.prestageDrainRunning = false;
if (generation === this.prestageGeneration
&& this.pendingPrestageJobs.size > 0) {
const latest = this.pendingPrestageJobs.values().next().value;
if (latest != null) {
this.queuePrestageMiningJob(latest);
}
}
});
}
private async drainPrestageMiningJobs(generation: number): Promise<void> {
while (generation === this.prestageGeneration
&& this.pendingPrestageJobs.size > 0) {
const [payoutMode, jobTemplate] = this.pendingPrestageJobs.entries().next().value as [
PayoutMode,
import('./stratum-v1-jobs.service').IJobTemplate,
];
this.pendingPrestageJobs.delete(payoutMode);
const clients = [...this.clients];
const batchSize = this.getPositiveIntegerEnv(
'SV1_PRESTAGE_BATCH_SIZE',
DEFAULT_PRESTAGE_BATCH_SIZE,
);
const startedAt = process.hrtime.bigint();
let staged = 0;
let skipped = 0;
let errors = 0;
for (let offset = 0; offset < clients.length; offset += batchSize) {
if (generation !== this.prestageGeneration) {
return;
}
for (const client of clients.slice(offset, offset + batchSize)) {
try {
if (client.preStageMiningJob(jobTemplate)) {
staged++;
} else {
skipped++;
}
} catch {
errors++;
}
}
// Staging is deliberately background work. Yield between bounded
// batches so share parsing and urgent bridge callbacks stay live.
if (offset + batchSize < clients.length) {
await new Promise<void>(resolve => setImmediate(resolve));
}
}
console.log(JSON.stringify({
event: 'sv1_job_prestage',
eventId: jobTemplate.blockData.notificationEventId,
height: jobTemplate.blockData.height,
payoutMode,
clients: clients.length,
staged,
skipped,
errors,
totalMs: Number(process.hrtime.bigint() - startedAt) / 1e6,
}));
}
}
private shouldLogJobFanout(isNewBlock: boolean, errors: number): boolean {
if (errors > 0 || isNewBlock) {
return true;