Compare commits

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