diff --git a/.env.example b/.env.example index 9b24575..9a3c264 100644 --- a/.env.example +++ b/.env.example @@ -15,6 +15,11 @@ BITCOIN_RPC_TIMEOUT=10000 BLOCK_TEMPLATE_LONGPOLL_TIMEOUT_MS=600000 BLOCK_TEMPLATE_RETRY_INITIAL_MS=25 BLOCK_TEMPLATE_RETRY_MAX_MS=1000 +# Optional comma-separated Bitcoin Core RPC endpoints. Each keeps an independent +# GBT longpoll open, but its template is published only when the primary Core's +# best block hash exactly authorizes the same tip. Endpoints currently share the +# primary RPC credentials/port. +#BITCOIN_RPC_AUX_URLS=http://bitcoin-core-2,http://bitcoin-core-3 # Enable in bitcoin.conf with: # zmqpubrawblock=tcp://0.0.0.0:3000 @@ -43,7 +48,9 @@ SV2_AUTH_FAILURE_LOG_ENABLED=false STRATUM_PORTS=3333,3332,3331,3330 # Optional PPLNS-mode SV1/SV2 autodetect ports. #PPLNS_STRATUM_PORTS= -STRATUM_WORKERS=2 +# `auto` assigns available CPUs after reserving API_WORKERS plus one master CPU. +# For fixed topology, set an explicit positive integer. +STRATUM_WORKERS=auto STRATUM_WORKER_MAX_MEMORY_RESTART=4096M # High-byte base for pool-assigned SV2 extranonce prefixes. PM2 reserves two # values per NODE_APP_INSTANCE so an old/new reload pair cannot overlap. Keep @@ -63,6 +70,16 @@ SV1_SUBSIDY_BRIDGE_ENABLED=true # Comma-separated: solo,pplns. PPLNS requires a fresh precomputed snapshot; # solo remains the safe default and is always published first when both are set. SV1_SUBSIDY_BRIDGE_PAYOUT_MODES=solo +# Give the dedicated urgent Redis socket a short opportunity to acknowledge the +# empty bridge before full-template hashing/serialization begins. On timeout the +# master retries immediately through the independent normal command socket. +SV1_BRIDGE_PUBLISH_BUDGET_MS=10 +# Build next-height per-miner coinbases in background batches and yield between +# batches so normal share processing remains responsive. +SV1_PRESTAGE_BATCH_SIZE=500 +# Operational target only: fanout logs flag workers above this client count so +# STRATUM_WORKERS can be raised before one event loop becomes the bottleneck. +STRATUM_FANOUT_TARGET_CLIENTS_PER_WORKER=10000 SV1_SUBSIDY_BRIDGE_PPLNS_SEED_MAX_AGE_MS=300000 STRATUM_JOB_RETENTION_MS=300000 STRATUM_SUBMISSION_DEDUP_TTL_MS=300000 diff --git a/README.md b/README.md index 3bf4b5b..061a6fe 100644 --- a/README.md +++ b/README.md @@ -59,6 +59,16 @@ $ NODE_CLUSTER_SCHED_POLICY=none pm2 start ecosystem.config.js Cluster-mode connection dropping requires Node.js `22.12.0` or newer. +`STRATUM_WORKERS=auto` is the default. It uses the container's available CPU +count after reserving `API_WORKERS` CPUs plus one for the master. A numeric value +still pins an exact topology. Size fixed deployments from the busiest worker, +not just total connection capacity: keep roughly 10,000 or fewer miners per +worker when low new-tip fanout latency matters. The `stratum_job_fanout` +`overTargetClients` field uses `STRATUM_FANOUT_TARGET_CLIENTS_PER_WORKER` to +make an undersized worker tier visible. Pools beyond one host's practical CPU or +socket capacity should shard listeners across multiple worker hosts connected to +the same urgent Redis channel. + `STRATUM_MAX_CONNECTIONS_PER_LISTENER` is enforced per worker and Stratum port. Size it using the busiest port: `worker count * limit`. For example, 28 workers with the default limit of `10000` allow up to `280000` connections on one port. @@ -78,11 +88,28 @@ reserve two namespace values per configured Stratum worker. ### New-block notification path The master keeps an authoritative Bitcoin Core `getblocktemplate` longpoll open; -rawblock ZMQ remains a watchdog and duplicate results are discarded. On a new tip, -the master publishes a compact subsidy-only SV1 job before serializing the full -transaction template. Workers fan that job out through one process-level socket -broadcaster, then issue the full fee-paying job as a second clean switch. Payout -snapshot creation and Postgres persistence run after the immediate solo publish. +rawblock ZMQ remains a watchdog and duplicate results are discarded. Optional +endpoints in `BITCOIN_RPC_AUX_URLS` keep independent longpolls open and race the +primary source, reducing dependence on one node's block-relay peers. An auxiliary +template is only eligible after the primary Core's `getbestblockhash` exactly +matches its previous block hash; a mismatch or unavailable primary fails closed +before any urgent or canonical miner notification. On a new tip, the master +publishes a compact subsidy-only SV1 job over dedicated urgent Redis +publisher/subscriber connections. It waits only +`SV1_BRIDGE_PUBLISH_BUDGET_MS` for acknowledgement before proceeding. A timeout +immediately starts the same bridge on the independent normal Redis command +socket, so a stuck urgent socket cannot suppress delivery or block canonical +full-template publication. + +After canonical publication, the master publishes a durable placeholder-prevhash +empty template for the following height. Stratum workers replay it after restart +and prebuild every connected miner's height- and payout-bound coinbase in batches +of `SV1_PRESTAGE_BATCH_SIZE`, yielding between batches. A newly authorized miner +also stages against the latest template. The next authoritative bridge promotes +the same cached job, patches the authoritative header fields into its already +serialized notify buffer, and writes it without rebuilding the coinbase or JSON. +Canonical full work then follows as a second clean switch. Payout +snapshot creation and Postgres persistence remain outside the urgent path. SV2 solo channels pre-stage a native subsidy-only future job for the next height. When the authoritative header arrives, the pool activates that job with only @@ -104,7 +131,10 @@ support requires a precomputed subsidy-valued payout snapshot. Retained jobs are kept for `STRATUM_JOB_RETENTION_MS` so a late network-target candidate can still be reconstructed and submitted, while ordinary old-tip shares are rejected. PPLNS seeds use a non-active snapshot status and are skipped if the next -authoritative `nBits` differs from their preparation basis. +authoritative `nBits` differs from their preparation basis. Before any bridge is +published, the master verifies that Core's `coinbasevalue` equals the locally +calculated consensus subsidy plus every GBT transaction fee. Master startup also +fails if `NETWORK` does not match Core's reported chain. The Redis protocol remains rolling-deploy compatible: new workers retain the legacy mining-info reload path, while the master writes the historical latest @@ -116,7 +146,17 @@ on their prior job instead of being woken with a miner-address fallback job. Two structured log events expose the end-to-end timing: - `block_notification_trace` reports Core, bridge, Redis, PPLNS, and persistence stages, separated by payout mode and job type. -- `stratum_job_fanout` reports client count, bytes, backpressure, and p50/p95/p99/last enqueue time, correlated by `eventId`. +- `stratum_job_fanout` reports true source-to-fanout-start time, master publish, + worker receipt/handling, prestage hits, client count, bytes, backpressure, and + p50/p95/p99/last enqueue time, correlated by `eventId`. +- `sv1_job_prestage` reports the number of miners prepared for the next height and + the background preparation duration. + +For upstream latency, place the Core nodes in different well-connected networks, +enable normal compact-block relay, and keep Redis and Stratum workers close +together. Auxiliary nodes remain untrusted candidates for tip detection: the +primary Core authorizes their exact tip before publication. Do not point +`BITCOIN_RPC_AUX_URLS` at third-party RPC services. ## Docker diff --git a/docker-compose.external-db.yml b/docker-compose.external-db.yml index c55f85b..55eb93f 100644 --- a/docker-compose.external-db.yml +++ b/docker-compose.external-db.yml @@ -75,7 +75,7 @@ services: API_SECURE: ${API_SECURE:-false} API_WORKERS: ${API_WORKERS:-4} PM2_ENABLED: ${PM2_ENABLED:-true} - STRATUM_WORKERS: ${STRATUM_WORKERS:-2} + STRATUM_WORKERS: ${STRATUM_WORKERS:-auto} STRATUM_PORTS: ${STRATUM_PORTS:-3333,3332,3331,3330} PPLNS_STRATUM_PORTS: ${PPLNS_STRATUM_PORTS:-} STRATUM_MIN_DIFFICULTY: ${STRATUM_MIN_DIFFICULTY:-1} diff --git a/docker-compose.yml b/docker-compose.yml index 072efb1..ad1b758 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -105,7 +105,7 @@ services: API_SECURE: ${API_SECURE:-false} API_WORKERS: ${API_WORKERS:-4} PM2_ENABLED: ${PM2_ENABLED:-true} - STRATUM_WORKERS: ${STRATUM_WORKERS:-2} + STRATUM_WORKERS: ${STRATUM_WORKERS:-auto} STRATUM_PORTS: ${STRATUM_PORTS:-3333,3332,3331,3330} PPLNS_STRATUM_PORTS: ${PPLNS_STRATUM_PORTS:-} STRATUM_MIN_DIFFICULTY: ${STRATUM_MIN_DIFFICULTY:-1} diff --git a/ecosystem.config.js b/ecosystem.config.js index c4fd413..706ee3e 100644 --- a/ecosystem.config.js +++ b/ecosystem.config.js @@ -1,54 +1,79 @@ +const { availableParallelism, cpus } = require('os'); + const dockerLogConfig = { out_file: '/dev/stdout', error_file: '/dev/stderr', merge_logs: true, }; +const cpuCount = + typeof availableParallelism === 'function' + ? availableParallelism() + : cpus().length; + +const positiveInteger = (name, value, fallback) => { + const candidate = value == null || value === '' ? fallback : Number(value); + if (!Number.isInteger(candidate) || candidate <= 0) { + throw new Error(`${name} must be a positive integer`); + } + return candidate; +}; + +const apiWorkers = positiveInteger('API_WORKERS', process.env.API_WORKERS, 4); +const automaticStratumWorkers = Math.max(1, cpuCount - apiWorkers - 1); +const stratumWorkers = + process.env.STRATUM_WORKERS == null || + process.env.STRATUM_WORKERS === '' || + process.env.STRATUM_WORKERS.toLowerCase() === 'auto' + ? automaticStratumWorkers + : positiveInteger('STRATUM_WORKERS', process.env.STRATUM_WORKERS); + module.exports = { - apps: [ - // API instance - { - ...dockerLogConfig, - name: 'api', - script: './dist/main.js', - instances: parseInt(process.env.API_WORKERS || '4', 10), - exec_mode: 'cluster', - env: { - MASTER: 'false', - API_ONLY: 'true', - API_ENABLED: 'true', - NODE_CLUSTER_SCHED_POLICY: 'none', - }, - time: true + apps: [ + // API instance + { + ...dockerLogConfig, + name: 'api', + script: './dist/main.js', + instances: apiWorkers, + exec_mode: 'cluster', + env: { + MASTER: 'false', + API_ONLY: 'true', + API_ENABLED: 'true', + NODE_CLUSTER_SCHED_POLICY: 'none', }, - // Master instance - { - ...dockerLogConfig, - name: 'master', - script: './dist/main.js', - instances: 1, - exec_mode: 'fork', - env: { - MASTER: 'true', - API_ENABLED: 'false', - NODE_CLUSTER_SCHED_POLICY: 'none', - }, - time: true + time: true, + }, + // Master instance + { + ...dockerLogConfig, + name: 'master', + script: './dist/main.js', + instances: 1, + exec_mode: 'fork', + env: { + MASTER: 'true', + API_ENABLED: 'false', + NODE_CLUSTER_SCHED_POLICY: 'none', }, - // Worker instances - { - ...dockerLogConfig, - name: 'workers', - script: './dist/main.js', - instances: parseInt(process.env.STRATUM_WORKERS || '2', 10), - exec_mode: "cluster", - max_memory_restart: process.env.STRATUM_WORKER_MAX_MEMORY_RESTART || '4096M', - env: { - MASTER: 'false', - API_ENABLED: 'false', - NODE_CLUSTER_SCHED_POLICY: 'none', - }, - time: true + time: true, + }, + // Worker instances + { + ...dockerLogConfig, + name: 'workers', + script: './dist/main.js', + instances: stratumWorkers, + exec_mode: 'cluster', + max_memory_restart: + process.env.STRATUM_WORKER_MAX_MEMORY_RESTART || '4096M', + env: { + MASTER: 'false', + API_ENABLED: 'false', + NODE_CLUSTER_SCHED_POLICY: 'none', }, - ], - }; + time: true, + }, + ], +}; diff --git a/src/ecosystem.config.spec.ts b/src/ecosystem.config.spec.ts new file mode 100644 index 0000000..54406d7 --- /dev/null +++ b/src/ecosystem.config.spec.ts @@ -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; +} diff --git a/src/models/MiningJob.ts b/src/models/MiningJob.ts index 56432ee..61928fe 100644 --- a/src/models/MiningJob.ts +++ b/src/models/MiningJob.ts @@ -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(); 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) { diff --git a/src/models/StratumV1Client.ts b/src/models/StratumV1Client.ts index 56b2e0e..74987d0 100644 --- a/src/models/StratumV1Client.ts +++ b/src/models/StratumV1Client.ts @@ -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 }), }; } } diff --git a/src/models/bitcoin-rpc/IBlockTemplate.ts b/src/models/bitcoin-rpc/IBlockTemplate.ts index d774943..21f5fab 100644 --- a/src/models/bitcoin-rpc/IBlockTemplate.ts +++ b/src/models/bitcoin-rpc/IBlockTemplate.ts @@ -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; diff --git a/src/services/bitcoin-rpc.service.spec.ts b/src/services/bitcoin-rpc.service.spec.ts index be6ee43..5f47172 100644 --- a/src/services/bitcoin-rpc.service.spec.ts +++ b/src/services/bitcoin-rpc.service.spec.ts @@ -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(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}`); diff --git a/src/services/bitcoin-rpc.service.ts b/src/services/bitcoin-rpc.service.ts index 3c9d794..b72f4eb 100644 --- a/src/services/bitcoin-rpc.service.ts +++ b/src/services/bitcoin-rpc.service.ts @@ -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; } +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 = new BehaviorSubject(undefined); private _newSv1BridgeTemplate$ = new ReplaySubject(1); + private _newSv1PrestageTemplate$ = new ReplaySubject(2); private resetTemplateInterval$ = new Subject(); private rpcRequestId = 0; private readonly processedTemplateEvents = new Set(); @@ -94,7 +118,10 @@ export class BitcoinRpcService implements OnModuleInit { private pplnsTemplatePublicationPromise: Promise | 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(); + private readonly lastPublishedPrestageKeys = new Map<'solo' | 'pplns', string>(); private readonly pplnsSubsidyBridgeSeeds = new Map(); private pplnsSeedPrecomputeTail: Promise = 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 { 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 { 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('getblocktemplate', [ + blockTemplate = await this.callRpcWithClient( + 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 { - if (this.longpollLoopStarted || process.env.MASTER !== 'true') { + private async listenForLongpollTemplates(templateSource?: TemplateRpcSource): Promise { + 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 { + 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, + ): Promise { + // 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 { + lane: 'urgent' | 'fallback' = 'urgent', + ): Promise { 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 { + lane: 'urgent' | 'fallback' = 'urgent', + ): Promise { 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 { + 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 { @@ -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 { - const response = await this.client.post('', { + return this.callRpcWithClient(this.client, method, params, timeoutMs); + } + + private async callRpcWithClient( + client: AxiosInstance, + method: string, + params: unknown[] = [], + timeoutMs?: number, + ): Promise { + 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('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 { + try { + const primaryBestBlockHash = await this.callRpc('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'}`, + ); + } + } } diff --git a/src/services/redis-messaging.service.spec.ts b/src/services/redis-messaging.service.spec.ts index 1d35d8d..1be2fea 100644 --- a/src/services/redis-messaging.service.spec.ts +++ b/src/services/redis-messaging.service.spec.ts @@ -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(); const sets = new Map>(); -const clientsByRole: { publisher?: any; subscriber?: any } = {}; +const clientsByRole: { + publisher?: any; + subscriber?: any; + urgentPublisher?: any; + urgentSubscriber?: any; +} = {}; const subscriptions = new Map Promise>(); 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; } diff --git a/src/services/redis-messaging.service.ts b/src/services/redis-messaging.service.ts index c40df29..b1be3b5 100644 --- a/src/services/redis-messaging.service.ts +++ b/src/services/redis-messaging.service.ts @@ -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 { + public async publishSv1BridgeUpdate( + update: Sv1BridgeUpdate, + lane: 'urgent' | 'fallback' = 'urgent', + ): Promise { 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): Promise { 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 { + 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, + ): Promise { + 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 { + if (!await this.ensureConnected()) { + return []; + } + const values = await this.publisher.mGet([ + sv1PrestageLatestKey('solo'), + sv1PrestageLatestKey('pplns'), + ]) as Array; + 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; + const template = update.template as Partial | 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 { if (!this.connected) { try { diff --git a/src/services/stratum-v1-jobs.service.spec.ts b/src/services/stratum-v1-jobs.service.spec.ts index 8713e01..ee38ad7 100644 --- a/src/services/stratum-v1-jobs.service.spec.ts +++ b/src/services/stratum-v1-jobs.service.spec.ts @@ -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; let bridgeTemplate$: Subject; - let bitcoinRpcService: { newBlockTemplate$: any, newSv1BridgeTemplate$: any, miningInfo: { blocks: number } }; + let prestageTemplate$: Subject; + 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(); + prestageTemplate$ = new Subject(); 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 }]; diff --git a/src/services/stratum-v1-jobs.service.ts b/src/services/stratum-v1-jobs.service.ts index 5816ead..c23b08d 100644 --- a/src/services/stratum-v1-jobs.service.ts +++ b/src/services/stratum-v1-jobs.service.ts @@ -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; /** Ordered SV1-only stream: subsidy bridge first, canonical full job second. */ public sv1MiningJob$: Observable; + /** Detached next-height empty jobs prepared without changing active work. */ + public sv1PrestageJob$: Observable; 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(); private readonly currentTipKeys = new Map(); private readonly cachedJobs = new Map(); + private readonly notifyHeaderFields = new WeakMap(); private readonly latestJobTemplates = new Map(); + private readonly latestPrestageJobTemplates = new Map(); + private readonly stagedJobs = new Map(); private readonly pinnedCurrentTipTemplateIds = new Map>([ ['solo', new Set()], ['pplns', new Set()], @@ -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]; diff --git a/src/services/stratum-v1.service.spec.ts b/src/services/stratum-v1.service.spec.ts index 10039c8..5b760e9 100644 --- a/src/services/stratum-v1.service.spec.ts +++ b/src/services/stratum-v1.service.spec.ts @@ -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; + let prestageJobs: Subject; 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({ diff --git a/src/services/stratum-v1.service.ts b/src/services/stratum-v1.service.ts index 6fdc7a5..949f916 100644 --- a/src/services/stratum-v1.service.ts +++ b/src/services/stratum-v1.service.ts @@ -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(); private jobBroadcastSubscription: Subscription | null = null; + private jobPrestageSubscription: Subscription | null = null; + private readonly pendingPrestageJobs = new Map(); + 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 { + 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(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;