Compare commits

3 Commits
Author SHA1 Message Date
Ben 0b030fecca Split maintenance duties from notifier 2026-07-13 22:58:54 -04:00
Ben 448f7d4c52 Fix pool accounting summary fallback 2026-07-13 22:47:53 -04:00
Ben d830227178 Speed up prestaged block activation 2026-07-13 22:15:18 -04:00
22 changed files with 1135 additions and 35 deletions
+6
View File
@@ -68,6 +68,12 @@ STRATUM_MAX_SOCKET_BUFFER_BYTES=262144
# transaction template. SV2 uses the same activation to select its pre-staged
# future job; the full job then follows on the active tip.
SV1_SUBSIDY_BRIDGE_ENABLED=true
# Publish only authoritative fixed-width GBT fields on the first Redis command.
# Workers promote their already-built next-height coinbase/notify buffers.
SV1_COMPACT_PRESTAGE_ACTIVATION_ENABLED=true
# Rolling-deploy aid for old workers. Leave false after workers understand the
# compact activation protocol; true publishes the larger empty bridge afterward.
SV1_COMPACT_ACTIVATION_COMPATIBILITY_BRIDGE=false
# 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
+2
View File
@@ -1,6 +1,8 @@
# compiled output
/dist
/node_modules
__pycache__/
*.pyc
# Logs
logs
+21 -8
View File
@@ -93,13 +93,20 @@ 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.
before any urgent or canonical miner notification. PM2 runs the master through a
minimal, non-HTTP notifier entrypoint so API, Stratum, reporting, and integration
timers cannot delay its Core longpoll callback.
On a new tip, the notifier's first Redis command is a compact SV1 prestage
activation containing only Core-authoritative fixed-width header fields, height,
subsidy, and payout identity. It does not serialize the transaction template or
a second block-template object. Workers consume this on dedicated urgent Redis
publisher/subscriber connections and promote their retained next-height jobs
directly, bypassing the canonical template/RxJS preparation pipeline. It waits
only `SV1_BRIDGE_PUBLISH_BUDGET_MS` for acknowledgement before proceeding. A
timeout immediately retries the activation on the independent normal Redis
command socket. If compact delivery fails or the feature is disabled, the
self-contained subsidy-only bridge remains the fallback.
After canonical publication, the master publishes a durable placeholder-prevhash
empty template for the following height. Stratum workers replay it after restart
@@ -139,7 +146,10 @@ 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
key as JSON only after a PPLNS-safe compatibility template is ready. Deploying
workers before the master is still the preferred rollout order. When any PPLNS
workers before the notifier is required for compact activation. During a mixed
deployment, `SV1_COMPACT_ACTIVATION_COMPATIBILITY_BRIDGE=true` publishes the
larger empty bridge immediately after compact activation; disable it after every
worker supports the compact protocol. When any PPLNS
listener is configured and snapshot preparation fails, legacy workers are held
on their prior job instead of being woken with a miner-address fallback job.
@@ -151,6 +161,9 @@ Two structured log events expose the end-to-end timing:
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.
- `sv1_prestage_activation_miss` identifies a worker that could not match an
authoritative activation to retained prestage state and therefore waits for
the fallback bridge or canonical full job.
For upstream latency, place the Core nodes in different well-connected networks,
enable normal compact-block relay, and keep Redis and Stratum workers close
+17 -2
View File
@@ -45,11 +45,26 @@ module.exports = {
},
time: true,
},
// Master instance
// Minimal authoritative template/notifier instance. This entrypoint avoids
// initializing API, Stratum, reporting, and notification integrations.
{
...dockerLogConfig,
name: 'master',
script: './dist/main.js',
script: './dist/notifier-main.js',
instances: 1,
exec_mode: 'fork',
env: {
MASTER: 'true',
API_ENABLED: 'false',
NODE_CLUSTER_SCHED_POLICY: 'none',
},
time: true,
},
// Non-hot-path master duties: notifications, reporting, and cleanup.
{
...dockerLogConfig,
name: 'maintenance',
script: './dist/maintenance-main.js',
instances: 1,
exec_mode: 'fork',
env: {
@@ -272,6 +272,62 @@ describe('ShareAccountingService', () => {
);
});
it('should serve API-only pool accounting from rollups when the precomputed pool cache is missing', async () => {
process.env.API_ONLY = 'true';
const redis = {
getJsonCache: jest.fn()
.mockResolvedValueOnce(null)
.mockResolvedValueOnce(null),
setJsonCache: jest.fn().mockResolvedValue(undefined),
};
const repository = {
query: jest.fn()
.mockResolvedValueOnce([{
totalAcceptedShares: '12',
totalCreditedDifficulty: '384',
acceptedSharesLast10Minutes: '4',
creditedDifficultyLast10Minutes: '128',
acceptedSharesLastHour: '10',
creditedDifficultyLastHour: '320',
acceptedSharesLastDay: '12',
creditedDifficultyLastDay: '384',
hashRateLast10Minutes: '916259689.8',
hashRateLastHour: '381774870.2',
latestShareAt: new Date('2026-06-07T12:30:00Z'),
}])
.mockResolvedValueOnce([{
bestSubmissionDifficulty: '4096',
bestSubmissionDifficultyAt: new Date('2026-06-07T12:20:00Z'),
currentRoundAcceptedShares: '11',
workSinceLastBlock: '352',
currentRoundNetworkDifficulty: '1000',
}])
.mockResolvedValueOnce([]),
};
const service = new ShareAccountingService(repository as any, redis as any);
await expect(service.getPoolSummary('solo')).resolves.toEqual(expect.objectContaining({
totalAcceptedShares: 12,
totalCreditedDifficulty: 384,
acceptedSharesLast10Minutes: 4,
creditedDifficultyLast10Minutes: 128,
bestSubmissionDifficulty: 4096,
workSinceLastBlock: 352,
networkDifficultyPercent: 35.2,
}));
expect(repository.query).toHaveBeenNthCalledWith(
1,
expect.stringContaining('FROM "accepted_share_10m"'),
['solo'],
);
expect(repository.query).toHaveBeenNthCalledWith(
2,
expect.stringContaining('FROM "accepted_share_block_10m"'),
['solo'],
);
});
it('should use Redis cache for share accounting summaries across API workers', async () => {
process.env.SHARE_ACCOUNTING_SUMMARY_CACHE_MS = '0';
const repository = {
@@ -429,8 +485,8 @@ describe('ShareAccountingService', () => {
};
const service = new ShareAccountingService(repository as any);
await service.getPoolSummary();
await service.getPoolSummary();
await service.getAddressSummary('bc1qcached');
await service.getAddressSummary('bc1qcached');
expect(repository.query).toHaveBeenCalledTimes(1);
});
@@ -360,11 +360,7 @@ export class ShareAccountingService implements OnModuleInit, OnModuleDestroy {
return cached;
}
if (process.env.API_ONLY === 'true') {
return this.emptySummary();
}
return this.getSummary({ payoutMode: mode });
return this.withPoolRollupOverlay(await this.getSummary({ payoutMode: mode }), mode);
}
public async refreshPoolSummary(payoutMode?: PayoutMode): Promise<ShareAccountingSummary> {
@@ -380,10 +376,6 @@ export class ShareAccountingService implements OnModuleInit, OnModuleDestroy {
}
private async withPoolRollupOverlay(summary: ShareAccountingSummary, payoutMode?: PayoutMode): Promise<ShareAccountingSummary> {
if (process.env.API_ONLY === 'true') {
return summary;
}
const [currentRoundRow] = await this.acceptedShareRepository.query(`
WITH latest_found_block AS (
SELECT COALESCE(MAX("height"), 0) AS "height"
+26
View File
@@ -33,6 +33,32 @@ describe('PM2 worker sizing', () => {
expect(config.apps.find((app) => app.name === 'workers').instances).toBe(7);
});
it('runs the master through the isolated notifier entrypoint', () => {
const config = require('../ecosystem.config.js');
expect(config.apps.find((app) => app.name === 'master')).toEqual(
expect.objectContaining({
script: './dist/notifier-main.js',
instances: 1,
exec_mode: 'fork',
env: expect.objectContaining({ MASTER: 'true', API_ENABLED: 'false' }),
}),
);
});
it('runs maintenance duties outside the isolated notifier process', () => {
const config = require('../ecosystem.config.js');
expect(config.apps.find((app) => app.name === 'maintenance')).toEqual(
expect.objectContaining({
script: './dist/maintenance-main.js',
instances: 1,
exec_mode: 'fork',
env: expect.objectContaining({ MASTER: 'true', API_ENABLED: 'false' }),
}),
);
});
it('rejects invalid fixed worker counts instead of silently starting no workers', () => {
process.env.STRATUM_WORKERS = 'many';
+16
View File
@@ -0,0 +1,16 @@
import { NestFactory } from '@nestjs/core';
import { MaintenanceModule } from './maintenance.module';
async function bootstrap(): Promise<void> {
process.env.MASTER = 'true';
process.env.API_ENABLED = 'false';
const application = await NestFactory.createApplicationContext(
MaintenanceModule,
);
application.enableShutdownHooks();
console.log('Maintenance services started');
}
void bootstrap();
+53
View File
@@ -0,0 +1,53 @@
import { ConfigModule, ConfigService } from '@nestjs/config';
import { Module } from '@nestjs/common';
import { TypeOrmModule } from '@nestjs/typeorm';
import { UserAgentReportModule } from './ORM/_views/user-agent-report/user-agent-report.module';
import { ClientModule } from './ORM/client/client.module';
import { RpcBlocksModule } from './ORM/rpc-block/rpc-block.module';
import { ShareAccountingModule } from './ORM/share-accounting/share-accounting.module';
import { TelegramSubscriptionsModule } from './ORM/telegram-subscriptions/telegram-subscriptions.module';
import { createDatabaseOptions } from './database.config';
import { AppService } from './services/app.service';
import { DiscordService } from './services/discord.service';
import { NotificationService } from './services/notification.service';
import { RedisMessagingModule } from './services/redis-messaging.module';
import { TelegramService } from './services/telegram.service';
/**
* Background process for non-hot-path master duties. Keep this separate from
* NotifierModule so chat integrations, reporting, and cleanup timers cannot
* delay block-template notification.
*/
@Module({
imports: [
ConfigModule.forRoot(),
TypeOrmModule.forRootAsync({
imports: [ConfigModule],
inject: [ConfigService],
useFactory: (configService: ConfigService) => createDatabaseOptions({
...process.env,
DB_HOST: configService.get('DB_HOST'),
DB_PORT: configService.get('DB_PORT'),
DB_USERNAME: configService.get('DB_USERNAME'),
DB_PASSWORD: configService.get('DB_PASSWORD'),
DB_DATABASE: configService.get('DB_DATABASE'),
DB_LOGGING: configService.get('DB_LOGGING'),
DB_POOL_SIZE: configService.get('DB_POOL_SIZE'),
}),
}),
RedisMessagingModule,
ClientModule,
RpcBlocksModule,
UserAgentReportModule,
ShareAccountingModule,
TelegramSubscriptionsModule,
],
providers: [
AppService,
DiscordService,
NotificationService,
TelegramService,
],
})
export class MaintenanceModule { }
+10
View File
@@ -85,6 +85,8 @@ export class StratumV1Client {
private connectionClosed = false;
private lastSentMiningJobTimestamp: number = null;
private lastSentMiningJobSignature: string = null;
private lastSentMiningTipKey: string = null;
private lastSentMiningJobType: 'full' | 'empty' | null = null;
private lastHashRatePersistedAt = 0;
private readonly network: bitcoinjs.Network;
private readonly maxSocketBufferBytes: number;
@@ -506,6 +508,12 @@ export class StratumV1Client {
if (!force && signature === this.lastSentMiningJobSignature) {
return { status: 'skipped', bytes: 0, bufferedBytes: this.socket.writableLength ?? 0 };
}
if (!force
&& jobTemplate.blockData.jobType === 'empty'
&& this.lastSentMiningJobType === 'empty'
&& this.lastSentMiningTipKey === jobTemplate.blockData.tipKey) {
return { status: 'skipped', bytes: 0, bufferedBytes: this.socket.writableLength ?? 0 };
}
const maximumBufferedBytes = this.maxSocketBufferBytes;
const bufferedBeforeBuild = this.socket.writableLength ?? 0;
@@ -576,6 +584,8 @@ export class StratumV1Client {
const accepted = this.socket.write(payload);
this.lastSentMiningJobTimestamp = jobTemplate.block.timestamp;
this.lastSentMiningJobSignature = signature;
this.lastSentMiningTipKey = jobTemplate.blockData.tipKey;
this.lastSentMiningJobType = jobTemplate.blockData.jobType;
const bufferedAfterWrite = this.socket.writableLength ?? 0;
if (bufferedAfterWrite >= maximumBufferedBytes) {
this.closeSocket();
+15
View File
@@ -0,0 +1,15 @@
import { NestFactory } from '@nestjs/core';
import { NotifierModule } from './notifier.module';
async function bootstrap(): Promise<void> {
process.env.MASTER = 'true';
process.env.API_ENABLED = 'false';
const application = await NestFactory.createApplicationContext(
NotifierModule,
);
application.enableShutdownHooks();
console.log('Authoritative block notifier started');
}
void bootstrap();
+41
View File
@@ -0,0 +1,41 @@
import { Module } from '@nestjs/common';
import { ConfigModule, ConfigService } from '@nestjs/config';
import { TypeOrmModule } from '@nestjs/typeorm';
import { createDatabaseOptions } from './database.config';
import { PayoutSnapshotModule } from './ORM/payout-snapshot/payout-snapshot.module';
import { RpcBlocksModule } from './ORM/rpc-block/rpc-block.module';
import { BitcoinRpcService } from './services/bitcoin-rpc.service';
import { RedisMessagingModule } from './services/redis-messaging.module';
/**
* Minimal process graph for new-tip notification. Keeping API, Stratum,
* reporting, chat integrations, and accounting providers out of this Nest
* context prevents their timers and callbacks from delaying the Core longpoll.
*/
@Module({
imports: [
ConfigModule.forRoot(),
TypeOrmModule.forRootAsync({
imports: [ConfigModule],
inject: [ConfigService],
useFactory: (configService: ConfigService) => createDatabaseOptions({
...process.env,
DB_HOST: configService.get('DB_HOST'),
DB_PORT: configService.get('DB_PORT'),
DB_USERNAME: configService.get('DB_USERNAME'),
DB_PASSWORD: configService.get('DB_PASSWORD'),
DB_DATABASE: configService.get('DB_DATABASE'),
DB_LOGGING: configService.get('DB_LOGGING'),
DB_POOL_SIZE: configService.get('DB_POOL_SIZE'),
}),
}),
RedisMessagingModule,
RpcBlocksModule,
PayoutSnapshotModule,
],
providers: [
BitcoinRpcService,
],
})
export class NotifierModule { }
+154
View File
@@ -236,6 +236,104 @@ describe('BitcoinRpcService template publication', () => {
expect(order.indexOf('redis:publish:bridge:solo')).toBeLessThan(order.indexOf('redis:publish:solo'));
});
it('publishes compact authoritative activation before canonical template work', async () => {
const order: string[] = [];
const redis = createRedisMock(order);
redis.publishSv1PrestageActivation = jest.fn(async (activation: any) => {
order.push(`redis:publish:activation:${activation.payoutMode}`);
return true;
});
const template = createTemplateAtHeight(840_000, '68');
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');
expect(redis.publishSv1PrestageActivation).toHaveBeenCalledWith(expect.objectContaining({
schemaVersion: 1,
type: 'prestage-activation',
height: template.height,
previousBlockHash: template.previousblockhash,
version: template.version,
bits: template.bits,
minTime: template.mintime,
currentTime: template.curtime,
subsidySats: calculateBlockSubsidySats(template.height, 'mainnet'),
payoutMode: 'solo',
requiredVersionBits: template.vbrequired >>> 0,
}));
expect(redis.publishSv1BridgeUpdate).not.toHaveBeenCalled();
expect(order.indexOf('redis:publish:activation:solo'))
.toBeLessThan(order.indexOf('redis:set:solo'));
});
it('falls back to the full empty bridge when compact activation delivery fails', async () => {
const order: string[] = [];
const redis = createRedisMock(order);
redis.publishSv1PrestageActivation = jest.fn().mockResolvedValue(false);
const template = createTemplateAtHeight(840_000, '69');
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.publishSv1PrestageActivation).toHaveBeenCalledTimes(1);
expect(redis.publishSv1BridgeUpdate).toHaveBeenCalledWith(expect.objectContaining({
type: 'subsidy-bridge',
template: expect.objectContaining({ payoutMode: 'solo', jobType: 'empty' }),
}));
errorSpy.mockRestore();
});
it('retries a timed-out compact activation through the normal Redis lane', async () => {
const redis = createRedisMock([]);
let resolveUrgent: () => void;
const stalled = new Promise<void>(resolve => { resolveUrgent = resolve; });
redis.publishSv1PrestageActivation = jest.fn(async (
_activation: any,
lane?: 'fallback',
) => {
if (lane == null) {
await stalled;
}
return true;
});
const template = createTemplateAtHeight(840_000, '67');
const service = new BitcoinRpcService(
createConfig({ NETWORK: 'mainnet', SV1_BRIDGE_PUBLISH_BUDGET_MS: '1' }),
{ 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 warnSpy = jest.spyOn(console, 'warn').mockImplementation(() => undefined);
await service.getAndBroadcastLatestTemplate('new_block');
await flushPromises();
expect(redis.publishSv1PrestageActivation).toHaveBeenCalledWith(
expect.objectContaining({ type: 'prestage-activation' }),
'fallback',
);
expect(redis.setBlockTemplate).toHaveBeenCalled();
resolveUrgent!();
warnSpy.mockRestore();
});
it('rejects an urgent bridge when Core fee data cannot validate its subsidy', async () => {
const redis = createRedisMock([]);
const template = createTemplateAtHeight(840_000, '70');
@@ -1119,6 +1217,61 @@ describe('BitcoinRpcService template publication', () => {
expect(canonicals).toEqual([expect.objectContaining({ previousblockhash: canonical.previousblockhash })]);
});
it('delivers fresh compact activations once and ignores stale or duplicate events', () => {
const service = new BitcoinRpcService(
createConfig({}),
{} as any,
{} as any,
);
service.miningInfo = { blocks: 900_000 } as any;
const activations: any[] = [];
service.newSv1PrestageActivation$.subscribe(activation => activations.push(activation));
const activation = {
schemaVersion: 1,
type: 'prestage-activation',
eventId: 'fresh-activation',
height: 900_001,
previousBlockHash: '61'.repeat(32),
version: 0x20000000,
bits: '17034219',
minTime: 1_700_000_000,
currentTime: 1_700_000_001,
subsidySats: 312_500_000,
payoutMode: 'solo',
requiredVersionBits: 0,
publishedAtMs: Date.now(),
};
(service as any).handleSv1PrestageActivation({
...activation,
eventId: 'stale-activation',
height: 899_999,
});
(service as any).handleSv1PrestageActivation(activation);
(service as any).handleSv1PrestageActivation(activation);
(service as any).emitCanonicalTemplate({
...createTemplate(),
height: activation.height,
previousblockhash: activation.previousBlockHash,
payoutMode: 'solo',
jobType: 'full',
notificationEventId: 'canonical-before-late-activation',
notificationPublishedAtMs: activation.publishedAtMs + 1,
});
(service as any).handleSv1PrestageActivation({
...activation,
eventId: 'late-after-canonical',
publishedAtMs: activation.publishedAtMs + 2,
});
expect(activations).toEqual([
expect.objectContaining({
eventId: activation.eventId,
workerReceivedAtMs: expect.any(Number),
}),
]);
});
it('rejects live bridge updates older than or redundant with active canonical work', async () => {
const service = new BitcoinRpcService(
createConfig({}),
@@ -1489,6 +1642,7 @@ function createRedisMock(order: string[]) {
order.push(`redis:publish:bridge:${update.template.payoutMode}`);
return true;
}),
publishSv1PrestageActivation: undefined as jest.Mock | undefined,
publishSv1PrestageUpdate: jest.fn(async (update: { template: IBlockTemplate }) => {
order.push(`redis:publish:prestage:${update.template.payoutMode}`);
}),
+193 -5
View File
@@ -13,6 +13,7 @@ import {
BlockTemplateUpdate,
RedisMessagingService,
Sv1BridgeUpdate,
Sv1PrestageActivation,
Sv1PrestageUpdate,
} from './redis-messaging.service';
import {
@@ -104,6 +105,7 @@ export class BitcoinRpcService implements OnModuleInit {
private readonly auxiliaryTemplateSources: TemplateRpcSource[] = [];
private _newBlockTemplate$: BehaviorSubject<IBlockTemplate> = new BehaviorSubject(undefined);
private _newSv1BridgeTemplate$ = new ReplaySubject<IBlockTemplate>(1);
private _newSv1PrestageActivation$ = new ReplaySubject<Sv1PrestageActivation>(2);
private _newSv1PrestageTemplate$ = new ReplaySubject<IBlockTemplate>(2);
private resetTemplateInterval$ = new Subject<void>();
private rpcRequestId = 0;
@@ -119,6 +121,7 @@ export class BitcoinRpcService implements OnModuleInit {
private readonly latestBridgeTemplates = new Map<'solo' | 'pplns', IBlockTemplate>();
private readonly canonicalEmissionStates = new Map<'solo' | 'pplns' | 'all', CanonicalEmissionState>();
private readonly lastPublishedBridgeTipKeys = new Map<'solo' | 'pplns', BridgePublishReservation>();
private readonly lastPublishedActivationTipKeys = new Map<'solo' | 'pplns', BridgePublishReservation>();
private bridgePublishAttemptId = 0;
private readonly subsidyValidatedTemplates = new WeakSet<IBlockTemplate>();
private readonly lastPublishedPrestageKeys = new Map<'solo' | 'pplns', string>();
@@ -141,6 +144,9 @@ 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 newSv1PrestageActivation$ = this._newSv1PrestageActivation$.pipe(
shareReplay({ refCount: true, bufferSize: 2 }),
);
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$;
@@ -178,6 +184,11 @@ export class BitcoinRpcService implements OnModuleInit {
if (process.env.MASTER != 'true') {
await this.loadLatestMiningInfoForReplayProcess();
if (process.env.API_ONLY != 'true') {
if (typeof this.redisMessagingService.subscribeSv1PrestageActivations === 'function') {
await this.redisMessagingService.subscribeSv1PrestageActivations(async activation => {
this.handleSv1PrestageActivation(activation);
});
}
await this.redisMessagingService.subscribeSv1BridgeUpdates(async update => {
// A new-tip bridge is an interrupt, not canonical replay work.
// It has its own Redis socket and must never sit behind a full
@@ -964,9 +975,9 @@ export class BitcoinRpcService implements OnModuleInit {
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'),
// The compact activation is enqueued before any compatibility bridge.
this.publishSoloUrgentWork(authoritativeTemplate, trace, 'urgent'),
this.publishPplnsUrgentWork(authoritativeTemplate, trace, 'urgent'),
]);
const budgetMs = this.getPositiveIntegerEnv(
'SV1_BRIDGE_PUBLISH_BUDGET_MS',
@@ -1019,8 +1030,8 @@ export class BitcoinRpcService implements OnModuleInit {
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'),
this.publishSoloUrgentWork(authoritativeTemplate, trace, 'fallback'),
this.publishPplnsUrgentWork(authoritativeTemplate, trace, 'fallback'),
]);
if (results.includes('failed')) {
this.markTrace(trace, 'sv1_bridge_fallback_failed');
@@ -1042,6 +1053,131 @@ export class BitcoinRpcService implements OnModuleInit {
if (this.lastPublishedBridgeTipKeys.get(payoutMode)?.tipKey === tipKey) {
this.lastPublishedBridgeTipKeys.delete(payoutMode);
}
if (this.lastPublishedActivationTipKeys.get(payoutMode)?.tipKey === tipKey) {
this.lastPublishedActivationTipKeys.delete(payoutMode);
}
}
}
private async publishSoloUrgentWork(
authoritativeTemplate: IBlockTemplate,
trace: BlockNotificationTrace,
lane: 'urgent' | 'fallback',
): Promise<BridgePublishResult> {
const activation = await this.publishPrestageActivation(
authoritativeTemplate,
trace,
'solo',
lane,
);
if (activation === 'published') {
if (this.isSv1CompatibilityBridgeEnabled()) {
void this.publishSoloSubsidyBridge(authoritativeTemplate, trace, lane);
}
return 'published';
}
return this.publishSoloSubsidyBridge(authoritativeTemplate, trace, lane);
}
private async publishPplnsUrgentWork(
authoritativeTemplate: IBlockTemplate,
trace: BlockNotificationTrace,
lane: 'urgent' | 'fallback',
): Promise<BridgePublishResult> {
const activation = await this.publishPrestageActivation(
authoritativeTemplate,
trace,
'pplns',
lane,
);
if (activation === 'published') {
if (this.isSv1CompatibilityBridgeEnabled()) {
void this.publishPplnsSubsidyBridge(authoritativeTemplate, trace, lane);
}
return 'published';
}
return this.publishPplnsSubsidyBridge(authoritativeTemplate, trace, lane);
}
private async publishPrestageActivation(
authoritativeTemplate: IBlockTemplate,
trace: BlockNotificationTrace,
payoutMode: 'solo' | 'pplns',
lane: 'urgent' | 'fallback',
): Promise<BridgePublishResult> {
if (!this.isSv1CompactActivationEnabled()
|| !this.isSv1SubsidyBridgeEnabled()
|| !this.getSv1SubsidyBridgePayoutModes().has(payoutMode)
|| typeof this.redisMessagingService.publishSv1PrestageActivation !== 'function') {
return 'skipped';
}
let subsidySats: number;
try {
subsidySats = this.validateSubsidyAgainstAuthoritativeTemplate(
authoritativeTemplate,
);
} catch (error) {
console.error(`Skipping ${payoutMode} SV1 prestage activation: ${error.message}`);
return 'failed';
}
let payoutSnapshotId: string | undefined;
if (payoutMode === 'pplns') {
const seed = this.getFreshPplnsSubsidyBridgeSeed(authoritativeTemplate.height);
if (seed == null
|| seed.subsidySats !== subsidySats
|| seed.basisBits.toLowerCase() !== authoritativeTemplate.bits.toLowerCase()) {
return 'skipped';
}
payoutSnapshotId = seed.payoutSnapshotId;
}
const tipKey = `${authoritativeTemplate.height}:${authoritativeTemplate.previousblockhash}`;
if (this.lastPublishedActivationTipKeys.get(payoutMode)?.tipKey === tipKey) {
return 'skipped';
}
const reservation: BridgePublishReservation = {
tipKey,
attemptId: ++this.bridgePublishAttemptId,
};
this.lastPublishedActivationTipKeys.set(payoutMode, reservation);
try {
const publishedAtMs = Date.now();
const activation: Sv1PrestageActivation = {
schemaVersion: 1,
type: 'prestage-activation',
eventId: `${trace.eventId}:activate:${payoutMode}`,
height: authoritativeTemplate.height,
previousBlockHash: authoritativeTemplate.previousblockhash.toLowerCase(),
version: authoritativeTemplate.version,
bits: authoritativeTemplate.bits.toLowerCase(),
minTime: authoritativeTemplate.mintime,
currentTime: authoritativeTemplate.curtime,
subsidySats,
payoutMode,
...(payoutSnapshotId == null ? {} : { payoutSnapshotId }),
requiredVersionBits: authoritativeTemplate.vbrequired >>> 0,
sourceNotificationReceivedAtMs: trace.sourceNotificationReceivedAtMs,
publishedAtMs,
};
const delivered = lane === 'urgent'
? await this.redisMessagingService.publishSv1PrestageActivation(activation)
: await this.redisMessagingService.publishSv1PrestageActivation(
activation,
'fallback',
);
if (!delivered) {
throw new Error(`Redis ${lane} prestage activation publisher is unavailable`);
}
this.markTrace(trace, `sv1_${payoutMode}_activation_workers_notified`);
return 'published';
} catch (error) {
if (this.lastPublishedActivationTipKeys.get(payoutMode) === reservation) {
this.lastPublishedActivationTipKeys.delete(payoutMode);
}
console.error(`Skipping ${payoutMode} SV1 prestage activation: ${error.message}`);
return 'failed';
}
}
@@ -1495,6 +1631,46 @@ export class BitcoinRpcService implements OnModuleInit {
this._newSv1BridgeTemplate$.next(template);
}
private handleSv1PrestageActivation(activation: Sv1PrestageActivation): void {
if (this.processedTemplateEvents.has(activation.eventId)) {
return;
}
this.rememberProcessedTemplateEvent(activation.eventId);
if (this.miningInfo?.blocks != null
&& activation.height < this.miningInfo.blocks + 1) {
return;
}
if (this.isActivationSupersededByCanonical(activation)) {
return;
}
this._newSv1PrestageActivation$.next({
...activation,
workerReceivedAtMs: activation.workerReceivedAtMs ?? Date.now(),
});
}
private isActivationSupersededByCanonical(
activation: Sv1PrestageActivation,
): boolean {
const canonicalStates = [
this.canonicalEmissionStates.get(activation.payoutMode),
this.canonicalEmissionStates.get('all'),
].filter((state): state is CanonicalEmissionState => state != null);
return canonicalStates.some(canonical => {
if (canonical.height > activation.height) {
return true;
}
if (canonical.height < activation.height) {
return false;
}
if (canonical.previousBlockHash === activation.previousBlockHash) {
return true;
}
return canonical.publishedAtMs == null
|| canonical.publishedAtMs >= activation.publishedAtMs;
});
}
private handleSv1PrestageUpdate(update: Sv1PrestageUpdate): void {
if (this.processedTemplateEvents.has(update.eventId)) {
return;
@@ -1980,6 +2156,18 @@ export class BitcoinRpcService implements OnModuleInit {
return configured?.toLowerCase() !== 'false';
}
private isSv1CompactActivationEnabled(): boolean {
const configured = this.configService.get<string>('SV1_COMPACT_PRESTAGE_ACTIVATION_ENABLED')
?? process.env.SV1_COMPACT_PRESTAGE_ACTIVATION_ENABLED;
return configured?.toLowerCase() !== 'false';
}
private isSv1CompatibilityBridgeEnabled(): boolean {
const configured = this.configService.get<string>('SV1_COMPACT_ACTIVATION_COMPATIBILITY_BRIDGE')
?? process.env.SV1_COMPACT_ACTIVATION_COMPATIBILITY_BRIDGE;
return configured?.toLowerCase() === 'true';
}
private getSv1SubsidyBridgePayoutModes(): ReadonlySet<'solo' | 'pplns'> {
const configured = this.configService.get<string>('SV1_SUBSIDY_BRIDGE_PAYOUT_MODES')
?? process.env.SV1_SUBSIDY_BRIDGE_PAYOUT_MODES
@@ -258,6 +258,53 @@ describe('RedisMessagingService', () => {
expect(clientsByRole.urgentPublisher.publish).not.toHaveBeenCalled();
});
it('publishes compact prestage activation on the urgent socket', async () => {
await service.connect();
const handler = jest.fn().mockResolvedValue(undefined);
const activation = createPrestageActivation();
await service.subscribeSv1PrestageActivations(handler);
await expect(service.publishSv1PrestageActivation(activation)).resolves.toBe(true);
expect(clientsByRole.urgentPublisher.publish).toHaveBeenCalledWith(
'sv1-prestage.activate',
JSON.stringify(activation),
);
expect(handler).toHaveBeenCalledWith(expect.objectContaining({
...activation,
workerReceivedAtMs: expect.any(Number),
}));
});
it('uses the normal Redis socket for compact activation fallback', async () => {
await service.connect();
const activation = createPrestageActivation();
await expect(service.publishSv1PrestageActivation(
activation,
'fallback',
)).resolves.toBe(true);
expect(clientsByRole.publisher.publish).toHaveBeenCalledWith(
'sv1-prestage.activate',
JSON.stringify(activation),
);
});
it('rejects malformed compact prestage activation fields', async () => {
await service.connect();
const activation = createPrestageActivation();
await expect(service.publishSv1PrestageActivation({
...activation,
previousBlockHash: 'not-a-hash',
})).rejects.toThrow('unsupported SV1 prestage activation');
await expect(service.publishSv1PrestageActivation({
...activation,
payoutMode: 'pplns',
})).rejects.toThrow('unsupported SV1 prestage activation');
});
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'));
@@ -412,6 +459,25 @@ function createBridgeUpdate(payoutMode: 'solo' | 'pplns', eventId: string) {
};
}
function createPrestageActivation() {
return {
schemaVersion: 1 as const,
type: 'prestage-activation' as const,
eventId: 'activate:solo:900001',
height: 900001,
previousBlockHash: '55'.repeat(32),
version: 0x20000000,
bits: '17034219',
minTime: 1_700_000_000,
currentTime: 1_700_000_001,
subsidySats: 312_500_000,
payoutMode: 'solo' as const,
requiredVersionBits: 0,
sourceNotificationReceivedAtMs: 123,
publishedAtMs: 124,
};
}
const store = new Map<string, string>();
const sets = new Map<string, Set<string>>();
const clientsByRole: {
+95
View File
@@ -9,6 +9,7 @@ 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_ACTIVATION_CHANNEL = 'sv1-prestage.activate';
const SV1_PRESTAGE_CHANNEL = 'sv1-prestage.updated';
const BLOCK_FOUND_NOTIFICATION_CHANNEL = 'block-found.notification';
const MINING_INFO_KEY = 'mining-info:latest';
@@ -59,6 +60,31 @@ export interface Sv1PrestageUpdate {
preparedAtMs: number;
}
/**
* Header-only activation for work whose coinbase and notify buffers were
* prepared during the prior height. Every field comes from an authoritative
* GBT; no next-block consensus field is inferred from ZMQ.
*/
export interface Sv1PrestageActivation {
schemaVersion: 1;
type: 'prestage-activation';
eventId: string;
height: number;
previousBlockHash: string;
version: number;
bits: string;
minTime: number;
currentTime: number;
subsidySats: number;
payoutMode: PayoutMode;
payoutSnapshotId?: string;
requiredVersionBits: number;
sourceNotificationReceivedAtMs?: number;
publishedAtMs: number;
/** Local worker timestamp; populated after Redis delivery. */
workerReceivedAtMs?: number;
}
export interface BlockFoundNotification {
schemaVersion: 1;
eventId: string;
@@ -236,6 +262,40 @@ export class RedisMessagingService implements OnModuleInit, OnModuleDestroy {
return true;
}
public async publishSv1PrestageActivation(
activation: Sv1PrestageActivation,
lane: 'urgent' | 'fallback' = 'urgent',
): Promise<boolean> {
if (!await this.ensureConnected()) {
return false;
}
const serialized = JSON.stringify(activation);
this.parseSv1PrestageActivation(serialized);
const publisher = lane === 'urgent' ? this.urgentPublisher : this.publisher;
await publisher.publish(SV1_PRESTAGE_ACTIVATION_CHANNEL, serialized);
return true;
}
public async subscribeSv1PrestageActivations(
handler: (activation: Sv1PrestageActivation) => Promise<void>,
): Promise<void> {
if (!await this.ensureConnected()) {
return;
}
await this.urgentSubscriber.subscribe(
SV1_PRESTAGE_ACTIVATION_CHANNEL,
async message => {
try {
const activation = this.parseSv1PrestageActivation(message);
activation.workerReceivedAtMs = Date.now();
await handler(activation);
} catch (error) {
console.error(`Invalid Redis SV1 prestage activation: ${error.message}`);
}
},
);
}
public async subscribeSv1BridgeUpdates(handler: (update: Sv1BridgeUpdate) => Promise<void>): Promise<void> {
if (!await this.ensureConnected()) {
return;
@@ -568,6 +628,41 @@ export class RedisMessagingService implements OnModuleInit, OnModuleDestroy {
return update as Sv1PrestageUpdate;
}
private parseSv1PrestageActivation(message: string): Sv1PrestageActivation {
const activation = JSON.parse(message) as Partial<Sv1PrestageActivation>;
const payoutMode = activation.payoutMode;
const hasValidPayoutIdentity = payoutMode === 'solo'
? activation.payoutSnapshotId == null
: typeof activation.payoutSnapshotId === 'string'
&& activation.payoutSnapshotId.trim().length > 0;
if (activation.schemaVersion !== 1
|| activation.type !== 'prestage-activation'
|| typeof activation.eventId !== 'string'
|| activation.eventId.trim().length === 0
|| !Number.isSafeInteger(activation.height)
|| activation.height < 0
|| typeof activation.previousBlockHash !== 'string'
|| !/^[0-9a-f]{64}$/.test(activation.previousBlockHash)
|| !Number.isInteger(activation.version)
|| activation.version < -0x80000000
|| activation.version > 0x7fffffff
|| typeof activation.bits !== 'string'
|| !/^[0-9a-f]{8}$/.test(activation.bits)
|| !Number.isSafeInteger(activation.minTime)
|| !Number.isSafeInteger(activation.currentTime)
|| !Number.isSafeInteger(activation.subsidySats)
|| activation.subsidySats < 0
|| (payoutMode !== 'solo' && payoutMode !== 'pplns')
|| !hasValidPayoutIdentity
|| !Number.isInteger(activation.requiredVersionBits)
|| activation.requiredVersionBits < 0
|| activation.requiredVersionBits > 0xffffffff
|| !Number.isFinite(activation.publishedAtMs)) {
throw new Error('unsupported SV1 prestage activation');
}
return activation as Sv1PrestageActivation;
}
private parseBlockFoundNotification(message: string): BlockFoundNotification {
const notification = JSON.parse(message) as Partial<BlockFoundNotification>;
if (notification.schemaVersion !== 1
@@ -211,6 +211,96 @@ describe('StratumV1JobsService', () => {
expect(service.getJobById(staged.jobId)).toBe(staged);
});
it('promotes the latest prestage directly from a compact 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 = 'compact-activation-miner';
const staged = service.preStageJob(
bitcoinjs.networks.testnet,
payout,
detached,
payoutIdentity,
'solo',
);
const prebuiltNotify = staged.responseBuffer(detached);
const activation = {
schemaVersion: 1 as const,
type: 'prestage-activation' as const,
eventId: 'compact-authoritative-activation',
height: future.height,
previousBlockHash: 'ab'.repeat(32),
version: future.version,
bits: future.bits,
minTime: future.mintime,
currentTime: future.curtime,
subsidySats: future.coinbasevalue,
payoutMode: 'solo' as const,
requiredVersionBits: future.vbrequired >>> 0,
sourceNotificationReceivedAtMs: Date.now() - 2,
publishedAtMs: Date.now() - 1,
workerReceivedAtMs: Date.now(),
};
const activatedTemplate = service.activateLatestPrestage(activation);
const activatedJob = service.activatePreStagedJob(
activatedTemplate,
payoutIdentity,
'solo',
);
expect(activatedTemplate).toEqual(expect.objectContaining({
blockData: expect.objectContaining({
height: future.height,
tipKey: `${future.height}:${activation.previousBlockHash}`,
jobType: 'empty',
clearJobs: true,
notificationEventId: activation.eventId,
}),
}));
expect(activatedTemplate.block.prevHash.toString('hex')).toBe(
Buffer.from(activation.previousBlockHash, 'hex').reverse().toString('hex'),
);
expect(activatedTemplate.block.bits).toBe(parseInt(activation.bits, 16));
expect(activatedJob).toBe(staged);
expect(activatedJob.responseBuffer(activatedTemplate)).toBe(prebuiltNotify);
expect(service.getSubmissionContext(staged.jobId)?.status).toBe('current');
});
it('refuses compact activation when the worker has no matching prestage', async () => {
await firstValueFrom(service.newMiningJob$);
expect(service.activateLatestPrestage({
schemaVersion: 1,
type: 'prestage-activation',
eventId: 'missing-prestage',
height: MockRecording1.BLOCK_TEMPLATE.height + 1,
previousBlockHash: 'ab'.repeat(32),
version: MockRecording1.BLOCK_TEMPLATE.version,
bits: MockRecording1.BLOCK_TEMPLATE.bits,
minTime: MockRecording1.BLOCK_TEMPLATE.mintime,
currentTime: MockRecording1.BLOCK_TEMPLATE.curtime,
subsidySats: 312_500_000,
payoutMode: 'solo',
requiredVersionBits: 0,
publishedAtMs: Date.now(),
})).toBeNull();
});
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 }];
+83
View File
@@ -7,6 +7,7 @@ import { IBlockTemplate, IBlockTemplateTx } from '../models/bitcoin-rpc/IBlockTe
import { AddressObject, MiningJob, MiningNotifyHeaderFields } from '../models/MiningJob';
import { PayoutMode } from '../types/payout-mode';
import { BitcoinRpcService } from './bitcoin-rpc.service';
import { Sv1PrestageActivation } from './redis-messaging.service';
import {
createPreparedMiningJob,
PreparedMiningJobBodyReference,
@@ -305,6 +306,88 @@ export class StratumV1JobsService {
return this.latestPrestageJobTemplates.get(payoutMode) ?? null;
}
/**
* Promote a detached next-height job using only authoritative fixed-width
* header fields. This intentionally bypasses the full template/RxJS
* preparation pipeline on the new-tip event-loop turn.
*/
public activateLatestPrestage(
activation: Sv1PrestageActivation,
): IJobTemplate | null {
const payoutMode = activation.payoutMode;
const prestage = this.latestPrestageJobTemplates.get(payoutMode);
if (prestage == null
|| prestage.blockData.height !== activation.height
|| prestage.blockData.jobType !== 'empty'
|| prestage.blockData.payoutMode !== payoutMode
|| prestage.blockData.coinbasevalue !== activation.subsidySats
|| (prestage.blockData.payoutSnapshotId ?? undefined)
!== (activation.payoutSnapshotId ?? undefined)
|| !prestage.block.prevHash.equals(PLACEHOLDER_PREV_HASH)) {
return null;
}
const previousBlockHash = Buffer.from(activation.previousBlockHash, 'hex');
const bits = Buffer.from(activation.bits, 'hex');
if (previousBlockHash.length !== 32 || bits.length !== 4) {
return null;
}
const timestamp = Math.max(
activation.minTime,
activation.currentTime,
Math.floor(Date.now() / 1000),
);
const tipKey = `${activation.height}:${activation.previousBlockHash}`;
const latest = this.latestJobTemplates.get(payoutMode);
if (latest?.blockData.tipKey === tipKey) {
return latest.blockData.jobType === 'empty' ? latest : null;
}
const block = Object.assign(new bitcoinjs.Block(), prestage.block, {
prevHash: Buffer.from(previousBlockHash).reverse(),
version: activation.version,
bits: bits.readUInt32BE(0),
timestamp,
});
const id = this.getNextTemplateId();
this.latestJobTemplateId++;
const isNewBlock = this.lastPreviousBlockHashes.get(payoutMode)
!== activation.previousBlockHash;
this.lastPreviousBlockHashes.set(payoutMode, activation.previousBlockHash);
const activated: IJobTemplate = {
block,
merkle_branch: prestage.merkle_branch,
blockData: {
...prestage.blockData,
id,
creation: Date.now(),
networkDifficulty: this.calculateNetworkDifficulty(block.bits),
tipKey,
clearJobs: true,
isNewBlock,
requiredVersionBits: activation.requiredVersionBits >>> 0,
notificationEventId: activation.eventId,
sourceNotificationReceivedAtMs: activation.sourceNotificationReceivedAtMs,
notificationPreparedAtMs: activation.publishedAtMs,
notificationPublishedAtMs: activation.publishedAtMs,
notificationWorkerReceivedAtMs: activation.workerReceivedAtMs,
notificationWorkerHandledAtMs: Date.now(),
},
};
this.blocks[id] = activated;
const pinnedTemplateIds = this.pinnedCurrentTipTemplateIds.get(payoutMode);
if (this.currentTipKeys.get(payoutMode) !== tipKey) {
pinnedTemplateIds.clear();
}
this.currentTipKeys.set(payoutMode, tipKey);
pinnedTemplateIds.add(id);
this.latestJobTemplates.set(payoutMode, activated);
this.cleanupExpiredJobsAndTemplates();
return activated;
}
public preStageJob(
network: bitcoinjs.networks.Network,
payoutInformation: AddressObject[],
+45 -1
View File
@@ -21,6 +21,7 @@ describe('StratumV1Service', () => {
let redisMessagingService;
let miningJobs: Subject<any>;
let prestageJobs: Subject<any>;
let prestageActivations: Subject<any>;
let consoleLogSpy: jest.SpyInstance;
let consoleWarnSpy: jest.SpyInstance;
@@ -39,8 +40,9 @@ describe('StratumV1Service', () => {
redisMessagingService = {};
miningJobs = new Subject();
prestageJobs = new Subject();
prestageActivations = new Subject();
service = new StratumV1Service(
{} as any,
{ newSv1PrestageActivation$: prestageActivations.asObservable() } as any,
clientService,
{} as any,
{} as any,
@@ -216,6 +218,48 @@ describe('StratumV1Service', () => {
service.onModuleDestroy();
});
it('promotes and broadcasts compact prestage activations immediately', async () => {
process.env.MASTER = 'false';
process.env.STRATUM_PORTS = '';
process.env.STRATUM_SECURE = 'false';
const activatedJob = {
blockData: {
id: 'activated-2',
height: 900002,
tipKey: `900002:${'55'.repeat(32)}`,
payoutMode: 'solo',
jobType: 'empty',
isNewBlock: true,
clearJobs: true,
notificationEventId: 'activate-2',
},
};
const activateLatestPrestage = jest.fn().mockReturnValue(activatedJob);
(service as any).stratumV1JobsService.activateLatestPrestage = activateLatestPrestage;
const client = {
broadcastMiningJob: jest.fn().mockReturnValue({
status: 'written',
bytes: 256,
bufferedBytes: 0,
preStaged: true,
}),
};
(service as any).clients.add(client);
await service.onModuleInit();
const activation = {
eventId: 'activate-2',
height: 900002,
payoutMode: 'solo',
previousBlockHash: '55'.repeat(32),
};
prestageActivations.next(activation);
expect(activateLatestPrestage).toHaveBeenCalledWith(activation);
expect(client.broadcastMiningJob).toHaveBeenCalledWith(activatedJob);
service.onModuleDestroy();
});
it('does not log routine non-new-block fanout unless explicitly enabled', () => {
(service as any).clients.add({
broadcastMiningJob: jest.fn().mockReturnValue({
+25
View File
@@ -15,6 +15,7 @@ import { ShareAccountingService } from '../ORM/share-accounting/share-accounting
import { BitcoinRpcService } from './bitcoin-rpc.service';
import { NotificationService } from './notification.service';
import { RedisMessagingService } from './redis-messaging.service';
import { Sv1PrestageActivation } from './redis-messaging.service';
import { StratumV1JobsService } from './stratum-v1-jobs.service';
import { StratumV2Service } from './stratum-v2.service';
import { parsePayoutModePorts, PayoutMode } from '../types/payout-mode';
@@ -59,6 +60,7 @@ export class StratumV1Service implements OnModuleInit, OnModuleDestroy {
private healthyBackpressureChecks = 0;
private readonly clients = new Set<StratumV1Client>();
private jobBroadcastSubscription: Subscription | null = null;
private jobActivationSubscription: Subscription | null = null;
private jobPrestageSubscription: Subscription | null = null;
private readonly pendingPrestageJobs = new Map<PayoutMode, import('./stratum-v1-jobs.service').IJobTemplate>();
private prestageDrainRunning = false;
@@ -95,6 +97,11 @@ export class StratumV1Service implements OnModuleInit, OnModuleDestroy {
return;
}
this.jobActivationSubscription = this.bitcoinRpcService
.newSv1PrestageActivation$?.subscribe({
next: activation => this.activateAndBroadcastPrestage(activation),
error: error => console.error(`SV1 prestage activation subscription failed: ${error.message}`),
}) ?? null;
this.jobBroadcastSubscription = (
this.stratumV1JobsService.sv1MiningJob$
?? this.stratumV1JobsService.newMiningJob$
@@ -140,6 +147,8 @@ export class StratumV1Service implements OnModuleInit, OnModuleDestroy {
public onModuleDestroy(): void {
this.jobBroadcastSubscription?.unsubscribe();
this.jobBroadcastSubscription = null;
this.jobActivationSubscription?.unsubscribe();
this.jobActivationSubscription = null;
this.jobPrestageSubscription?.unsubscribe();
this.jobPrestageSubscription = null;
this.prestageGeneration++;
@@ -586,6 +595,22 @@ export class StratumV1Service implements OnModuleInit, OnModuleDestroy {
}));
}
private activateAndBroadcastPrestage(activation: Sv1PrestageActivation): void {
const jobTemplate = this.stratumV1JobsService.activateLatestPrestage(activation);
if (jobTemplate == null) {
console.warn(JSON.stringify({
event: 'sv1_prestage_activation_miss',
eventId: activation.eventId,
height: activation.height,
payoutMode: activation.payoutMode,
previousBlockHash: activation.previousBlockHash,
workerReceivedAtMs: activation.workerReceivedAtMs,
}));
return;
}
this.broadcastMiningJob(jobTemplate);
}
private queuePrestageMiningJob(
jobTemplate: import('./stratum-v1-jobs.service').IJobTemplate,
): void {
+65 -1
View File
@@ -4,6 +4,7 @@ import { Socket } from 'net';
import { Observable, Subject } from 'rxjs';
import { StratumV2Client } from '../models/StratumV2Client';
import { Sv1PrestageActivation } from './redis-messaging.service';
import { IJobTemplate } from './stratum-v1-jobs.service';
import { StratumV2Service } from './stratum-v2.service';
@@ -156,6 +157,45 @@ describe('StratumV2Service canonical job broadcaster', () => {
await service.onModuleDestroy();
});
it('converts compact activations into header-only SV2 work activation fanout', async () => {
const { service, compactActivations } = createService();
await service.onModuleInit();
const client = createClientMock();
service.registerClient(client);
compactActivations.next(createCompactActivation(900_002, '22'.repeat(32)));
expect(client.enqueueWorkActivation).toHaveBeenCalledTimes(1);
expect(client.enqueueWorkActivation).toHaveBeenCalledWith(expect.objectContaining({
height: 900_002,
previousblockhash: '22'.repeat(32),
transactions: [],
jobType: 'empty',
version: 0x20000000,
bits: '1d00ffff',
mintime: 1_700_000_002,
}));
expect(service.getLatestWorkActivationTemplate()).toEqual(expect.objectContaining({
height: 900_002,
notificationEventId: 'compact-900002',
}));
await service.onModuleDestroy();
});
it('deduplicates a legacy bridge that follows the compact activation for the same tip', async () => {
const { service, activations, compactActivations } = createService();
await service.onModuleInit();
const client = createClientMock();
service.registerClient(client);
const previousBlockHash = '33'.repeat(32);
compactActivations.next(createCompactActivation(900_003, previousBlockHash));
activations.next(createActivationTemplate(900_003, previousBlockHash, 'solo'));
expect(client.enqueueWorkActivation).toHaveBeenCalledTimes(1);
await service.onModuleDestroy();
});
it('registers created clients and unregisters them during client destruction', async () => {
const { service } = createService();
(service as any).noiseConfig = createNoiseConfig();
@@ -230,11 +270,13 @@ function createService(): {
templates: Subject<IJobTemplate>;
subscribeSpy: jest.Mock;
activations: Subject<any>;
compactActivations: Subject<Sv1PrestageActivation>;
activationSubscribeSpy: jest.Mock;
} {
const templates = new Subject<IJobTemplate>();
const subscribeSpy = jest.fn((observer) => templates.subscribe(observer));
const activations = new Subject<any>();
const compactActivations = new Subject<Sv1PrestageActivation>();
const activationSubscribeSpy = jest.fn((observer) => activations.subscribe(observer));
const jobsService = {
newMiningJob$: new Observable((subscriber) => subscribeSpy(subscriber)),
@@ -248,6 +290,7 @@ function createService(): {
};
const bitcoinRpcService = {
workActivationTemplate$: new Observable((subscriber) => activationSubscribeSpy(subscriber)),
newSv1PrestageActivation$: compactActivations.asObservable(),
};
const service = new StratumV2Service(
bitcoinRpcService as any,
@@ -263,7 +306,7 @@ function createService(): {
{} as any,
{} as any,
);
return { service, templates, subscribeSpy, activations, activationSubscribeSpy };
return { service, templates, subscribeSpy, activations, compactActivations, activationSubscribeSpy };
}
function createClientMock(
@@ -298,6 +341,27 @@ function createActivationTemplate(
} as any;
}
function createCompactActivation(
height: number,
previousBlockHash: string,
): Sv1PrestageActivation {
return {
schemaVersion: 1,
type: 'prestage-activation',
eventId: `compact-${height}`,
height,
previousBlockHash,
version: 0x20000000,
bits: '1d00ffff',
minTime: 1_700_000_001,
currentTime: 1_700_000_002,
subsidySats: 312_500_000,
payoutMode: 'solo',
requiredVersionBits: 0,
publishedAtMs: 1_700_000_000_000,
};
}
function createJobTemplate(
id: string,
payoutMode: 'solo' | 'pplns' | 'all',
+53 -7
View File
@@ -30,7 +30,7 @@ import {
import { BitcoinRpcService } from './bitcoin-rpc.service';
import { CustomWorkService } from './custom-work.service';
import { NotificationService } from './notification.service';
import { RedisMessagingService } from './redis-messaging.service';
import { RedisMessagingService, Sv1PrestageActivation } from './redis-messaging.service';
import { IJobTemplate, StratumV1JobsService } from './stratum-v1-jobs.service';
import { Sv2JobDeclarationRegistryService } from './sv2-job-declaration-registry.service';
import { parsePayoutModePorts, PayoutMode } from '../types/payout-mode';
@@ -45,6 +45,7 @@ export class StratumV2Service implements OnModuleInit, OnModuleDestroy {
private readonly latestCanonicalJobs = new Map<PayoutMode | 'all', IJobTemplate>();
private canonicalJobSubscription: Subscription = null;
private workActivationSubscription: Subscription = null;
private compactActivationSubscription: Subscription = null;
private latestWorkActivationTemplate: IBlockTemplate = null;
private latestWorkActivationKey: string = null;
private authorityPrivKey: Buffer;
@@ -98,6 +99,8 @@ export class StratumV2Service implements OnModuleInit, OnModuleDestroy {
this.canonicalJobSubscription = null;
this.workActivationSubscription?.unsubscribe();
this.workActivationSubscription = null;
this.compactActivationSubscription?.unsubscribe();
this.compactActivationSubscription = null;
const clients = Array.from(this.clients);
this.clients.clear();
@@ -261,14 +264,57 @@ export class StratumV2Service implements OnModuleInit, OnModuleDestroy {
}
private startWorkActivationBroadcaster(): void {
if (this.workActivationSubscription != null || this.bitcoinRpcService.workActivationTemplate$ == null) {
return;
if (
this.compactActivationSubscription == null
&& this.bitcoinRpcService.newSv1PrestageActivation$ != null
) {
this.compactActivationSubscription = this.bitcoinRpcService.newSv1PrestageActivation$.subscribe({
next: activation => this.broadcastWorkActivation(
this.createCompactWorkActivationTemplate(activation),
),
error: error => console.error(`SV2 compact work activation subscription failed: ${error.message}`),
});
}
this.workActivationSubscription = this.bitcoinRpcService.workActivationTemplate$.subscribe({
next: template => this.broadcastWorkActivation(template),
error: error => console.error(`SV2 work activation subscription failed: ${error.message}`),
});
if (this.workActivationSubscription == null && this.bitcoinRpcService.workActivationTemplate$ != null) {
this.workActivationSubscription = this.bitcoinRpcService.workActivationTemplate$.subscribe({
next: template => this.broadcastWorkActivation(template),
error: error => console.error(`SV2 work activation subscription failed: ${error.message}`),
});
}
}
private createCompactWorkActivationTemplate(activation: Sv1PrestageActivation): IBlockTemplate {
return {
version: activation.version,
rules: [],
vbavailable: {},
vbrequired: activation.requiredVersionBits,
previousblockhash: activation.previousBlockHash,
transactions: [],
coinbaseaux: {},
coinbasevalue: activation.subsidySats,
longpollid: '',
target: '',
mintime: Math.max(activation.minTime, activation.currentTime),
mutable: [],
noncerange: '',
sigoplimit: 0,
sizelimit: 0,
weightlimit: 0,
curtime: activation.currentTime,
bits: activation.bits,
height: activation.height,
default_witness_commitment: '',
capabilities: [],
payoutSnapshotId: activation.payoutSnapshotId,
jobType: 'empty',
payoutMode: activation.payoutMode,
notificationEventId: activation.eventId,
sourceNotificationReceivedAtMs: activation.sourceNotificationReceivedAtMs,
notificationPublishedAtMs: activation.publishedAtMs,
notificationWorkerReceivedAtMs: activation.workerReceivedAtMs,
};
}
private broadcastWorkActivation(template: IBlockTemplate): void {