diff --git a/.env.example b/.env.example index f69b7a4..45f617f 100644 --- a/.env.example +++ b/.env.example @@ -58,6 +58,8 @@ STRATUM_TCP_KEEPALIVE_INITIAL_DELAY_MS=60000 #PPLNS_DATUM_PORTS= #DATUM_POOL_PAYOUT_ADDRESS= #DATUM_SHARE_DIFFICULTY=1 +#DATUM_PAYOUT_MAX_COINBASE_OUTPUTS=10 +#DATUM_PAYOUT_COINBASE_WEIGHT_BUDGET=2000 # 32-byte hex seed for the DATUM server identity key. DATUM clients must be # configured with the matching logged public key. #DATUM_IDENTITY_SEED= diff --git a/docker-compose.external-db.yml b/docker-compose.external-db.yml index 48840cb..aee4e43 100644 --- a/docker-compose.external-db.yml +++ b/docker-compose.external-db.yml @@ -95,6 +95,8 @@ services: PPLNS_SV2_TDP_PORTS: ${PPLNS_SV2_TDP_PORTS:-} DATUM_PORTS: ${DATUM_PORTS:-} PPLNS_DATUM_PORTS: ${PPLNS_DATUM_PORTS:-} + DATUM_PAYOUT_MAX_COINBASE_OUTPUTS: ${DATUM_PAYOUT_MAX_COINBASE_OUTPUTS:-10} + DATUM_PAYOUT_COINBASE_WEIGHT_BUDGET: ${DATUM_PAYOUT_COINBASE_WEIGHT_BUDGET:-2000} SHARE_ACCOUNTING_BATCH_SIZE: ${SHARE_ACCOUNTING_BATCH_SIZE:-500} SHARE_ACCOUNTING_FLUSH_INTERVAL_MS: ${SHARE_ACCOUNTING_FLUSH_INTERVAL_MS:-25} SHARE_ACCOUNTING_MAX_QUEUE_SIZE: ${SHARE_ACCOUNTING_MAX_QUEUE_SIZE:-50000} diff --git a/docker-compose.yml b/docker-compose.yml index 072efb1..27732c6 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -124,6 +124,8 @@ services: PPLNS_SV2_TDP_PORTS: ${PPLNS_SV2_TDP_PORTS:-} DATUM_PORTS: ${DATUM_PORTS:-} PPLNS_DATUM_PORTS: ${PPLNS_DATUM_PORTS:-} + DATUM_PAYOUT_MAX_COINBASE_OUTPUTS: ${DATUM_PAYOUT_MAX_COINBASE_OUTPUTS:-10} + DATUM_PAYOUT_COINBASE_WEIGHT_BUDGET: ${DATUM_PAYOUT_COINBASE_WEIGHT_BUDGET:-2000} healthcheck: test: ["CMD-SHELL", "node -e \"const http=require('http'); const https=require('https'); const secure=process.env.API_SECURE==='true'; const client=secure?https:http; const req=client.get({hostname:'127.0.0.1',port:process.env.API_PORT||3334,path:'/api/network',rejectUnauthorized:false},res=>process.exit(res.statusCode<500?0:1)); req.on('error',()=>process.exit(1)); req.setTimeout(5000,()=>{req.destroy(); process.exit(1);});\""] interval: 30s diff --git a/full-setup/docker-compose-mainnet.yml b/full-setup/docker-compose-mainnet.yml index cf0726d..4f5e1d9 100644 --- a/full-setup/docker-compose-mainnet.yml +++ b/full-setup/docker-compose-mainnet.yml @@ -127,6 +127,8 @@ services: PAYOUT_FEE_ADDRESS: ${PAYOUT_FEE_ADDRESS:-} PAYOUT_FEE_PERCENT: ${PAYOUT_FEE_PERCENT:-0} PAYOUT_COINBASE_WEIGHT_BUDGET: ${PAYOUT_COINBASE_WEIGHT_BUDGET:-50000} + DATUM_PAYOUT_MAX_COINBASE_OUTPUTS: ${DATUM_PAYOUT_MAX_COINBASE_OUTPUTS:-10} + DATUM_PAYOUT_COINBASE_WEIGHT_BUDGET: ${DATUM_PAYOUT_COINBASE_WEIGHT_BUDGET:-2000} networks: bitcoin: diff --git a/src/ORM/payout-snapshot/payout-snapshot.service.spec.ts b/src/ORM/payout-snapshot/payout-snapshot.service.spec.ts index 1bbd99a..db56b91 100644 --- a/src/ORM/payout-snapshot/payout-snapshot.service.spec.ts +++ b/src/ORM/payout-snapshot/payout-snapshot.service.spec.ts @@ -95,6 +95,63 @@ describe('PayoutSnapshotService', () => { expect(manager.query.mock.calls[0][1]).toEqual([100, 'pplns']); }); + it('should allow protocol-specific snapshot methods and coinbase limits', async () => { + manager.query + .mockResolvedValueOnce([{ + startBatchId: '10', + endBatchId: '12', + windowStartShareIndex: '1000', + windowEndShareIndex: '2000', + totalCreditedDifficulty: '100', + totalAcceptedShareCount: '5', + }]) + .mockResolvedValueOnce([]) + .mockResolvedValueOnce([ + { address: ADDRESS_A, creditedDifficulty: 60, acceptedShareCount: 3 }, + { address: ADDRESS_B, creditedDifficulty: 40, acceptedShareCount: 2 }, + ]) + .mockResolvedValueOnce([]) + .mockResolvedValueOnce([{ id: '57' }]) + .mockResolvedValueOnce([]) + .mockResolvedValueOnce([]) + .mockResolvedValueOnce([{ + id: '57', + method: 'pplns-datum', + blockHeight: 900002, + coinbaseValueSats: '1000', + windowStartShareIndex: '1000', + windowEndShareIndex: '2000', + totalCreditedDifficulty: 100, + totalAcceptedShareCount: '5', + eligibleAddressCount: 2, + includedOutputCount: 1, + distributedSats: '1000', + unallocatedRemainderSats: '0', + }]) + .mockResolvedValueOnce([ + { address: ADDRESS_A, payoutSats: '1000' }, + ]); + + const snapshot = await service.createSnapshotForTemplate({ + blockHeight: 900002, + coinbaseValueSats: 1000, + networkDifficulty: 25, + method: 'pplns-datum', + maxCoinbaseOutputs: 1, + coinbaseWeightBudget: 2000, + }); + + expect(snapshot).toEqual(expect.objectContaining({ + id: '57', + method: 'pplns-datum', + includedOutputCount: 1, + })); + expect(manager.query.mock.calls[1][1][0]).toBe('pplns-datum'); + expect(manager.query.mock.calls[4][1][0]).toBe('pplns-datum'); + expect(manager.query.mock.calls[4][1][10]).toBe(2000); + expect(manager.query.mock.calls[4][1][18]).toBe(1); + }); + it('should bootstrap the PPLNS window from paid block count up to the configured factor', async () => { process.env.PAYOUT_BOOTSTRAP_WINDOW = 'true'; service = new PayoutSnapshotService(dataSource as unknown as DataSource); @@ -303,7 +360,7 @@ describe('PayoutSnapshotService', () => { createdAt: new Date('2026-06-16T12:00:00.000Z'), percent: 60, }); - expect(dataSource.query.mock.calls[0][1]).toEqual([ADDRESS_A, 'pplns']); + expect(dataSource.query.mock.calls[0][1]).toEqual([ADDRESS_A, 'pplns', 'pplns']); expect(dataSource.query.mock.calls[0][0]).toContain('WITH latest_snapshot AS'); }); diff --git a/src/ORM/payout-snapshot/payout-snapshot.service.ts b/src/ORM/payout-snapshot/payout-snapshot.service.ts index d3638d4..380b0c0 100644 --- a/src/ORM/payout-snapshot/payout-snapshot.service.ts +++ b/src/ORM/payout-snapshot/payout-snapshot.service.ts @@ -51,6 +51,15 @@ export interface ExpectedPayout { percent: number; } +export interface CreatePayoutSnapshotInput { + blockHeight: number; + coinbaseValueSats: number; + networkDifficulty: number; + method?: string; + maxCoinbaseOutputs?: number; + coinbaseWeightBudget?: number; +} + const DEFAULT_MAX_COINBASE_OUTPUTS = 10; const DEFAULT_MIN_OUTPUT_SATS = 546; const DEFAULT_PAYOUT_METHOD = 'pplns'; @@ -75,15 +84,15 @@ export class PayoutSnapshotService { private readonly dataSource: DataSource, ) { } - public async createSnapshotForTemplate(input: { - blockHeight: number; - coinbaseValueSats: number; - networkDifficulty: number; - }): Promise { + public async createSnapshotForTemplate(input: CreatePayoutSnapshotInput): Promise { if (!this.snapshotsEnabled || input.coinbaseValueSats <= 0 || input.networkDifficulty <= 0) { return null; } + const method = this.normalizeMethod(input.method); + const maxCoinbaseOutputs = this.resolvePositiveInt(input.maxCoinbaseOutputs, this.maxCoinbaseOutputs); + const coinbaseWeightBudget = this.resolvePositiveInt(input.coinbaseWeightBudget, this.coinbaseWeightBudget); + return this.dataSource.transaction(async manager => { const effectiveWindowFactor = await this.getEffectiveWindowFactor(manager); const windowTargetDifficulty = input.networkDifficulty * effectiveWindowFactor; @@ -93,6 +102,7 @@ export class PayoutSnapshotService { } const existing = await this.getExistingSnapshot(manager, { + method, blockHeight: input.blockHeight, coinbaseValueSats: input.coinbaseValueSats, windowEndShareIndex: window.windowEndShareIndex, @@ -114,9 +124,9 @@ export class PayoutSnapshotService { feeAddress: this.feeAddress, feePercent: this.feePercent, minOutputSats: this.minOutputSats, - coinbaseWeightBudget: this.coinbaseWeightBudget, + coinbaseWeightBudget, }); - const entries = this.limitCoinbaseOutputs(distribution.entries); + const entries = this.limitCoinbaseOutputs(distribution.entries, maxCoinbaseOutputs); if (entries.every(entry => !entry.includedInCoinbase)) { return null; } @@ -176,7 +186,7 @@ export class PayoutSnapshotService { ) RETURNING "id"::text AS "id" `, [ - this.method, + method, PPLNS_PAYOUT_MODE, input.blockHeight, input.coinbaseValueSats.toString(), @@ -186,7 +196,7 @@ export class PayoutSnapshotService { this.feeAddress, distribution.feeSats.toString(), this.minOutputSats, - this.coinbaseWeightBudget, + coinbaseWeightBudget, window.startBatchId, window.endBatchId, window.windowStartShareIndex, @@ -390,9 +400,10 @@ export class PayoutSnapshotService { FROM "payout_snapshot" WHERE "status" = 'finalized' AND "payoutMode" = $1 + AND "method" = $2 ORDER BY "createdAt" DESC, "id" DESC LIMIT 1 - `, [PPLNS_PAYOUT_MODE]); + `, [PPLNS_PAYOUT_MODE, this.method]); if (snapshot?.id == null) { return null; } @@ -414,6 +425,7 @@ export class PayoutSnapshotService { FROM "payout_snapshot" WHERE "status" = 'finalized' AND "payoutMode" = $2 + AND "method" = $3 ORDER BY "createdAt" DESC, "id" DESC LIMIT 1 ) @@ -437,7 +449,7 @@ export class PayoutSnapshotService { AND e."includedInCoinbase" = true AND e."payoutSats" > 0 LIMIT 1 - `, [address, PPLNS_PAYOUT_MODE]); + `, [address, PPLNS_PAYOUT_MODE, this.method]); if (row == null) { return null; @@ -554,7 +566,7 @@ export class PayoutSnapshotService { private async getExistingSnapshot( manager: EntityManager, - input: { blockHeight: number; coinbaseValueSats: number; windowEndShareIndex: string }, + input: { method: string; blockHeight: number; coinbaseValueSats: number; windowEndShareIndex: string }, ): Promise { const [snapshot] = await manager.query(` SELECT "id"::text AS "id" @@ -568,7 +580,7 @@ export class PayoutSnapshotService { ORDER BY "id" DESC LIMIT 1 `, [ - this.method, + input.method, input.blockHeight, input.coinbaseValueSats.toString(), input.windowEndShareIndex, @@ -689,7 +701,7 @@ export class PayoutSnapshotService { ]); } - private limitCoinbaseOutputs(entries: PayoutDistributionEntry[]): PayoutDistributionEntry[] { + private limitCoinbaseOutputs(entries: PayoutDistributionEntry[], maxCoinbaseOutputs: number): PayoutDistributionEntry[] { let included = 0; let removedPayoutSats = 0; const limitedEntries = entries.map(entry => { @@ -697,7 +709,7 @@ export class PayoutSnapshotService { return entry; } included++; - if (included <= this.maxCoinbaseOutputs) { + if (included <= maxCoinbaseOutputs) { return entry; } removedPayoutSats += entry.payoutSats; @@ -761,6 +773,17 @@ export class PayoutSnapshotService { return Number.isInteger(value) && value > 0 ? value : defaultValue; } + private resolvePositiveInt(value: number | undefined, defaultValue: number): number { + return Number.isInteger(value) && value > 0 ? value : defaultValue; + } + + private normalizeMethod(method?: string): string { + const normalized = method?.trim(); + return normalized == null || normalized.length === 0 + ? this.method + : normalized.slice(0, 32); + } + private readNonNegativeInt(name: string, defaultValue: number): number { const value = Number(process.env[name]); return Number.isInteger(value) && value >= 0 ? value : defaultValue; diff --git a/src/services/datum.service.spec.ts b/src/services/datum.service.spec.ts index 08bab77..7fc6747 100644 --- a/src/services/datum.service.spec.ts +++ b/src/services/datum.service.spec.ts @@ -9,6 +9,7 @@ function createService(overrides: { configService?: any; clientService?: any; redisMessagingService?: any; + payoutSnapshotService?: any; } = {}): DatumService { const configService = overrides.configService ?? { get: jest.fn((key: string) => { @@ -43,6 +44,7 @@ function createService(overrides: { {} as any, templateProvider as unknown as TemplateProviderService, overrides.redisMessagingService, + overrides.payoutSnapshotService, ); } @@ -389,6 +391,60 @@ describe('DatumService job validation', () => { ).valid).toBe(true); }); + it('creates DATUM-sized payout snapshots for PPLNS coinbaser fetches', async () => { + const payoutSnapshotService = { + createSnapshotForTemplate: jest.fn().mockResolvedValue({ + id: 'datum-snapshot-1', + payoutOutputs: [ + { address: 'tb1q42vtlphyjjcun9wcv9f0d9pkhup9dcf5z9k4gh', amountSats: 596 }, + ], + }), + }; + const configService = { + get: jest.fn((key: string) => { + if (key === 'NETWORK') { + return 'testnet'; + } + if (key === 'DATUM_PAYOUT_MAX_COINBASE_OUTPUTS') { + return '6'; + } + if (key === 'DATUM_PAYOUT_COINBASE_WEIGHT_BUDGET') { + return '1500'; + } + return undefined; + }), + }; + const service = createService({ configService, payoutSnapshotService }) as any; + + const result = await service.getDatumCoinbaserPayoutContext({ + blockData: { + height: 5010000, + networkDifficulty: 42, + payoutSnapshotId: 'template-snapshot', + payoutOutputs: [ + { address: 'tb1q9r8gvnx3j4d6jvl0fqjrmy3dar4k4l3052af7q', amountSats: 596 }, + ], + }, + }, 596, 'pplns'); + + expect(payoutSnapshotService.createSnapshotForTemplate).toHaveBeenCalledWith({ + blockHeight: 5010000, + coinbaseValueSats: 596, + networkDifficulty: 42, + method: 'pplns-datum', + maxCoinbaseOutputs: 6, + coinbaseWeightBudget: 1500, + }); + expect(result.payoutSnapshotId).toBe('datum-snapshot-1'); + expect(result.payoutOutputs).toEqual([{ + value: 596n, + scriptPubKey: bitcoinjs.address.toOutputScript( + 'tb1q42vtlphyjjcun9wcv9f0d9pkhup9dcf5z9k4gh', + bitcoinjs.networks.testnet, + ), + }]); + }); + it('accepts DATUM coinbases matching another recent pool-issued coinbaser context for the same height', () => { const service = createService() as any; const extranonce = Buffer.alloc(12, 1); diff --git a/src/services/datum.service.ts b/src/services/datum.service.ts index e9f27ec..f6c5e8c 100644 --- a/src/services/datum.service.ts +++ b/src/services/datum.service.ts @@ -45,6 +45,9 @@ import { parsePayoutModePorts, PayoutMode } from '../types/payout-mode'; const DEFAULT_DATUM_SHARE_DIFFICULTY = 1; const DEFAULT_DATUM_PING_INTERVAL_MS = 30_000; const DEFAULT_CLIENT_HASHRATE_PERSIST_INTERVAL_MS = 60_000; +const DEFAULT_DATUM_PAYOUT_MAX_COINBASE_OUTPUTS = 10; +const DEFAULT_DATUM_PAYOUT_COINBASE_WEIGHT_BUDGET = 2_000; +const DATUM_PAYOUT_METHOD = 'pplns-datum'; @Injectable() export class DatumService implements OnModuleInit { @@ -221,7 +224,8 @@ export class DatumService implements OnModuleInit { private async handleCoinbaserFetch(socket: Socket, state: DatumClientState, payload: Buffer): Promise { const fetch = deserializeDatumCoinbaserFetch(payload); const latestTemplate = await firstValueFrom(this.jobsService.newMiningJob$); - const payoutOutputs = this.getDatumPayoutOutputs(latestTemplate, Number(fetch.rewardValue), state.payoutMode); + const payoutContext = await this.getDatumCoinbaserPayoutContext(latestTemplate, Number(fetch.rewardValue), state.payoutMode); + const payoutOutputs = payoutContext.payoutOutputs; if (payoutOutputs.length === 0) { throw new Error('DATUM_POOL_PAYOUT_ADDRESS or DEV_FEE_ADDRESS must be set before DATUM coinbaser fetches can be served'); } @@ -229,7 +233,7 @@ export class DatumService implements OnModuleInit { const coinbaserId = this.nextDatumCoinbaserId(state); state.coinbaserPayoutContexts.set(coinbaserId, { payoutOutputs, - payoutSnapshotId: latestTemplate.blockData.payoutSnapshotId ?? null, + payoutSnapshotId: payoutContext.payoutSnapshotId, blockHeight: latestTemplate.blockData.height, payoutMode: state.payoutMode, }); @@ -244,6 +248,40 @@ export class DatumService implements OnModuleInit { await this.writeRaw(socket, state.session.encryptChannelFrame(DatumProtocolCommand.MINING, response)); } + private async getDatumCoinbaserPayoutContext( + latestTemplate: IJobTemplate, + rewardValue: number, + payoutMode: PayoutMode, + ): Promise<{ payoutOutputs: DatumPayoutOutput[]; payoutSnapshotId?: string | null }> { + if (payoutMode === 'pplns' && this.payoutSnapshotService != null) { + try { + const snapshot = await this.payoutSnapshotService.createSnapshotForTemplate({ + blockHeight: latestTemplate.blockData.height, + coinbaseValueSats: rewardValue, + networkDifficulty: latestTemplate.blockData.networkDifficulty, + method: DATUM_PAYOUT_METHOD, + maxCoinbaseOutputs: this.getDatumPayoutMaxCoinbaseOutputs(), + coinbaseWeightBudget: this.getDatumPayoutCoinbaseWeightBudget(), + }); + if (snapshot?.payoutOutputs?.length > 0) { + return { + payoutOutputs: this.getDatumPayoutOutputs({ blockData: { payoutOutputs: snapshot.payoutOutputs } }, rewardValue, payoutMode), + payoutSnapshotId: snapshot.id, + }; + } + } catch (error) { + console.error(`[DATUM] Error creating DATUM payout snapshot: ${error.message}`); + } + } + + return { + payoutOutputs: this.getDatumPayoutOutputs(latestTemplate, rewardValue, payoutMode), + payoutSnapshotId: payoutMode === 'pplns' + ? latestTemplate.blockData.payoutSnapshotId ?? null + : null, + }; + } + private async handlePowSubmit(socket: Socket, state: DatumClientState, payload: Buffer): Promise { const pow = deserializeDatumPowSubmit(payload); const { address, workerName } = this.parseUserIdentity(pow.username); @@ -909,6 +947,20 @@ export class DatumService implements OnModuleInit { return Math.pow(2, pow.targetByte); } + private getDatumPayoutMaxCoinbaseOutputs(): number { + const configured = Number(this.configService.get('DATUM_PAYOUT_MAX_COINBASE_OUTPUTS')); + return Number.isInteger(configured) && configured > 0 + ? configured + : DEFAULT_DATUM_PAYOUT_MAX_COINBASE_OUTPUTS; + } + + private getDatumPayoutCoinbaseWeightBudget(): number { + const configured = Number(this.configService.get('DATUM_PAYOUT_COINBASE_WEIGHT_BUDGET')); + return Number.isInteger(configured) && configured > 0 + ? configured + : DEFAULT_DATUM_PAYOUT_COINBASE_WEIGHT_BUDGET; + } + private mapDatumTemplateRejectReason(errorCode?: string): DatumRejectReason { switch (errorCode) { case 'prevhash-mismatch':