diff --git a/src/ORM/_migrations/UserAgentReportNonzeroHashrate1781313000000.ts b/src/ORM/_migrations/UserAgentReportNonzeroHashrate1781313000000.ts new file mode 100644 index 0000000..b9fb335 --- /dev/null +++ b/src/ORM/_migrations/UserAgentReportNonzeroHashrate1781313000000.ts @@ -0,0 +1,54 @@ +import { MigrationInterface, QueryRunner } from 'typeorm'; + +export class UserAgentReportNonzeroHashrate1781313000000 implements MigrationInterface { + public name = 'UserAgentReportNonzeroHashrate1781313000000'; + public transaction = false; + + public async up(queryRunner: QueryRunner): Promise { + await queryRunner.query(`DROP MATERIALIZED VIEW IF EXISTS "user_agent_report_view"`); + await queryRunner.query(` + CREATE MATERIALIZED VIEW "user_agent_report_view" AS + SELECT + "userAgent", + COUNT("userAgent") AS "count", + MAX("bestDifficulty") AS "bestDifficulty", + SUM("hashRate") AS "totalHashRate" + FROM "client_entity" + WHERE "deletedAt" IS NULL + AND "hashRate" > 0 + GROUP BY "userAgent" + ORDER BY "totalHashRate" DESC + `); + await queryRunner.query(` + CREATE INDEX CONCURRENTLY IF NOT EXISTS "IDX_client_working_report_updated_user_agent" + ON "client_entity" ("updatedAt" DESC, "userAgent") + INCLUDE ("hashRate", "bestDifficulty") + WHERE "deletedAt" IS NULL + AND "hashRate" > 0 + `); + await queryRunner.query(`DROP INDEX CONCURRENTLY IF EXISTS "IDX_client_active_report_updated_user_agent"`); + } + + public async down(queryRunner: QueryRunner): Promise { + await queryRunner.query(`DROP INDEX CONCURRENTLY IF EXISTS "IDX_client_working_report_updated_user_agent"`); + await queryRunner.query(`DROP MATERIALIZED VIEW IF EXISTS "user_agent_report_view"`); + await queryRunner.query(` + CREATE MATERIALIZED VIEW "user_agent_report_view" AS + SELECT + "userAgent", + COUNT("userAgent") AS "count", + MAX("bestDifficulty") AS "bestDifficulty", + SUM("hashRate") AS "totalHashRate" + FROM "client_entity" + WHERE "deletedAt" IS NULL + GROUP BY "userAgent" + ORDER BY "totalHashRate" DESC + `); + await queryRunner.query(` + CREATE INDEX CONCURRENTLY IF NOT EXISTS "IDX_client_active_report_updated_user_agent" + ON "client_entity" ("updatedAt" DESC, "userAgent") + INCLUDE ("hashRate", "bestDifficulty") + WHERE "deletedAt" IS NULL + `); + } +} diff --git a/src/ORM/_views/user-agent-report/user-agent-report.service.spec.ts b/src/ORM/_views/user-agent-report/user-agent-report.service.spec.ts index cd9d664..2c2225d 100644 --- a/src/ORM/_views/user-agent-report/user-agent-report.service.spec.ts +++ b/src/ORM/_views/user-agent-report/user-agent-report.service.spec.ts @@ -55,6 +55,7 @@ describe('UserAgentReportService', () => { 'client.updatedAt > :activeSince', expect.objectContaining({ activeSince: expect.any(Date) }), ); + expect(clientRepository.queryBuilder.andWhere).toHaveBeenCalledWith('client.hashRate > 0'); }); it('falls back to the materialized view when no active database clients exist', async () => { diff --git a/src/ORM/_views/user-agent-report/user-agent-report.service.ts b/src/ORM/_views/user-agent-report/user-agent-report.service.ts index c494a3d..28d1e02 100644 --- a/src/ORM/_views/user-agent-report/user-agent-report.service.ts +++ b/src/ORM/_views/user-agent-report/user-agent-report.service.ts @@ -72,6 +72,7 @@ export class UserAgentReportService { .addSelect('COALESCE(SUM(client.hashRate), 0)', 'totalHashRate') .where('client.deletedAt IS NULL') .andWhere('client.updatedAt > :activeSince', { activeSince }) + .andWhere('client.hashRate > 0') .groupBy('COALESCE(NULLIF(client.userAgent, \'\'), \'Other\')') .orderBy('"totalHashRate"', 'DESC') .getRawMany(); diff --git a/src/ORM/_views/user-agent-report/user-agent-report.view.ts b/src/ORM/_views/user-agent-report/user-agent-report.view.ts index b3c6f49..09070cf 100644 --- a/src/ORM/_views/user-agent-report/user-agent-report.view.ts +++ b/src/ORM/_views/user-agent-report/user-agent-report.view.ts @@ -13,6 +13,7 @@ import { ClientEntity } from '../../client/client.entity'; .addSelect('SUM(client.hashRate)', 'totalHashRate') .from(ClientEntity, 'client') .where('client.deletedAt IS NULL') + .andWhere('client.hashRate > 0') .groupBy('client.userAgent') .orderBy('"totalHashRate"', 'DESC') }) diff --git a/src/database.config.ts b/src/database.config.ts index 4b82e47..8a1de75 100644 --- a/src/database.config.ts +++ b/src/database.config.ts @@ -17,6 +17,7 @@ import { BlocksSubmissionMetadata1781220000000 } from './ORM/_migrations/BlocksS import { PayoutModes1781300000000 } from './ORM/_migrations/PayoutModes1781300000000'; import { ShareRollupStoragePolicy1781305000000 } from './ORM/_migrations/ShareRollupStoragePolicy1781305000000'; import { AcceptedShareHighScores1781309000000 } from './ORM/_migrations/AcceptedShareHighScores1781309000000'; +import { UserAgentReportNonzeroHashrate1781313000000 } from './ORM/_migrations/UserAgentReportNonzeroHashrate1781313000000'; import { UserAgentReportView } from './ORM/_views/user-agent-report/user-agent-report.view'; import { AcceptedShareEntity } from './ORM/accepted-share/accepted-share.entity'; import { AddressSettingsEntity } from './ORM/address-settings/address-settings.entity'; @@ -60,6 +61,7 @@ export const databaseMigrations = [ PayoutModes1781300000000, ShareRollupStoragePolicy1781305000000, AcceptedShareHighScores1781309000000, + UserAgentReportNonzeroHashrate1781313000000, ]; export function createDatabaseOptions(env: NodeJS.ProcessEnv): TypeOrmModuleOptions & DataSourceOptions { diff --git a/src/services/datum.service.spec.ts b/src/services/datum.service.spec.ts index c473ac8..08bab77 100644 --- a/src/services/datum.service.spec.ts +++ b/src/services/datum.service.spec.ts @@ -389,6 +389,109 @@ describe('DatumService job validation', () => { ).valid).toBe(true); }); + 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); + const staleOutputs = (service as any).getDatumPayoutOutputs({ + blockData: { + payoutOutputs: [{ address: 'tb1q42vtlphyjjcun9wcv9f0d9pkhup9dcf5z9k4gh', amountSats: 596 }], + }, + }, 596, 'pplns'); + const submittedOutputs = (service as any).getDatumPayoutOutputs({ + blockData: { + payoutOutputs: [{ address: 'tb1q9r8gvnx3j4d6jvl0fqjrmy3dar4k4l3052af7q', amountSats: 596 }], + }, + }, 596, 'pplns'); + const coinbase = createDatumCoinbaseSplit([ + { address: 'tb1q9r8gvnx3j4d6jvl0fqjrmy3dar4k4l3052af7q', amountSats: 596 }, + ], extranonce); + const latestTemplate = { + blockData: { + coinbasevalue: 596, + payoutOutputs: [{ address: 'tb1q42vtlphyjjcun9wcv9f0d9pkhup9dcf5z9k4gh', amountSats: 596 }], + }, + }; + const datumJob = { + height: 4991366, + coinbaseValue: 596n, + expectedPayoutOutputs: staleOutputs, + coinbasePairs: new Map(), + }; + const state = { + payoutMode: 'pplns', + coinbaserPayoutContexts: new Map([ + [1, { + payoutOutputs: staleOutputs, + payoutSnapshotId: 'stale', + blockHeight: 4991366, + payoutMode: 'pplns', + }], + [2, { + payoutOutputs: submittedOutputs, + payoutSnapshotId: 'submitted', + blockHeight: 4991366, + payoutMode: 'pplns', + }], + ]), + }; + + const result = service.validateDatumCoinbasePayoutContext( + coinbase, + { extranonce, targetByte: 0 }, + latestTemplate, + datumJob, + state, + ); + + expect(result.validation.valid).toBe(true); + expect(result.context?.payoutSnapshotId).toBe('submitted'); + }); + + it('still rejects DATUM coinbases that do not match any pool-issued coinbaser context', () => { + const service = createService() as any; + const extranonce = Buffer.alloc(12, 1); + const poolOutputs = (service as any).getDatumPayoutOutputs({ + blockData: { + payoutOutputs: [{ address: 'tb1q42vtlphyjjcun9wcv9f0d9pkhup9dcf5z9k4gh', amountSats: 596 }], + }, + }, 596, 'pplns'); + const soloCoinbase = createDatumCoinbaseSplit([ + { address: 'tb1qdyjakeepue4trak9d3hvyelrd0aw7mwju2d0c2', amountSats: 596 }, + ], extranonce); + const latestTemplate = { + blockData: { + coinbasevalue: 596, + payoutOutputs: [{ address: 'tb1q42vtlphyjjcun9wcv9f0d9pkhup9dcf5z9k4gh', amountSats: 596 }], + }, + }; + const datumJob = { + height: 4991366, + coinbaseValue: 596n, + expectedPayoutOutputs: poolOutputs, + coinbasePairs: new Map(), + }; + const state = { + payoutMode: 'pplns', + coinbaserPayoutContexts: new Map([[1, { + payoutOutputs: poolOutputs, + payoutSnapshotId: 'pool', + blockHeight: 4991366, + payoutMode: 'pplns', + }]]), + }; + + const result = service.validateDatumCoinbasePayoutContext( + soloCoinbase, + { extranonce, targetByte: 0 }, + latestTemplate, + datumJob, + state, + ); + + expect(result.validation.valid).toBe(false); + expect(result.validation.error).toBe('output-mismatch'); + }); + it('carries the coinbaser payout snapshot id into the DATUM job cache', () => { const service = createService() as any; const payoutOutputs = [{ diff --git a/src/services/datum.service.ts b/src/services/datum.service.ts index c5d6ccc..e9f27ec 100644 --- a/src/services/datum.service.ts +++ b/src/services/datum.service.ts @@ -296,12 +296,16 @@ export class DatumService implements OnModuleInit { await this.sendShareResponse(socket, state, DatumShareResponseStatus.REJECTED, DatumRejectReason.OTHER, pow.nonce, pow.targetByte, pow.jobId); return; } - const payoutValidation = this.validateDatumCoinbasePayouts(coinbase, pow, latestTemplate, datumJob.coinbaseValue, datumJob.expectedPayoutOutputs, state.payoutMode); - if (!payoutValidation.valid) { - this.logDatumCoinbaseMismatchOnce(state, payoutValidation); + const payoutValidation = this.validateDatumCoinbasePayoutContext(coinbase, pow, latestTemplate, datumJob, state); + if (!payoutValidation.validation.valid) { + this.logDatumCoinbaseMismatchOnce(state, payoutValidation.validation); await this.sendShareResponse(socket, state, DatumShareResponseStatus.REJECTED, DatumRejectReason.BAD_COINBASER_ID, pow.nonce, pow.targetByte, pow.jobId); return; } + if (payoutValidation.context != null) { + datumJob.expectedPayoutOutputs = payoutValidation.context.payoutOutputs; + datumJob.payoutSnapshotId = payoutValidation.context.payoutSnapshotId; + } const shareDifficulty = this.getDatumSubmittedShareDifficulty(pow); const nBits = datumJob.nBits.readUInt32LE(0); @@ -623,6 +627,52 @@ export class DatumService implements OnModuleInit { }; } + private validateDatumCoinbasePayoutContext( + coinbase: { coinb1: Buffer; coinb2: Buffer }, + pow: DatumPowSubmit, + latestTemplate: IJobTemplate, + datumJob: DatumJobCache, + state: DatumClientState, + ): { validation: DatumCoinbasePayoutValidation; context?: DatumCoinbaserPayoutContext } { + const primaryValidation = this.validateDatumCoinbasePayouts( + coinbase, + pow, + latestTemplate, + datumJob.coinbaseValue, + datumJob.expectedPayoutOutputs, + state.payoutMode, + ); + if (primaryValidation.valid) { + return { validation: primaryValidation }; + } + + for (const context of state.coinbaserPayoutContexts.values()) { + if (context.payoutMode != null && context.payoutMode !== state.payoutMode) { + continue; + } + if (context.blockHeight != null && datumJob.height != null && context.blockHeight !== datumJob.height) { + continue; + } + if (context.payoutOutputs === datumJob.expectedPayoutOutputs) { + continue; + } + + const validation = this.validateDatumCoinbasePayouts( + coinbase, + pow, + latestTemplate, + datumJob.coinbaseValue, + context.payoutOutputs, + state.payoutMode, + ); + if (validation.valid) { + return { validation, context }; + } + } + + return { validation: primaryValidation }; + } + private getDatumPayoutOutputs(latestTemplate: { blockData?: { payoutOutputs?: AddressObject[] } }, rewardValue: number, payoutMode: PayoutMode = 'solo'): DatumPayoutOutput[] { const configuredOutputs = latestTemplate?.blockData?.payoutOutputs; const payoutAddresses = payoutMode === 'pplns' && configuredOutputs?.length > 0 diff --git a/test/timescale-redis.integration-spec.ts b/test/timescale-redis.integration-spec.ts index 0db0d6d..d0004ec 100644 --- a/test/timescale-redis.integration-spec.ts +++ b/test/timescale-redis.integration-spec.ts @@ -603,15 +603,24 @@ describe('TimescaleDB and Redis integration', () => { redisMessagingService, ); - const activeClient = await repository.save({ + await repository.save({ address: 'tb1qumezefzdeqqwn5zfvgdrhxjzc5ylr39uhuxcz4', - clientName: 'active-live-worker', + clientName: 'idle-live-worker', sessionId: '91b2c3d4', userAgent, startTime: new Date(), bestDifficulty: 11, hashRate: 0, }); + const activeClient = await repository.save({ + address: 'tb1qumezefzdeqqwn5zfvgdrhxjzc5ylr39uhuxcz4', + clientName: 'active-live-worker', + sessionId: '93b2c3d4', + userAgent, + startTime: new Date(), + bestDifficulty: 33, + hashRate: 300, + }); const staleClient = await repository.save({ address: 'tb1qumezefzdeqqwn5zfvgdrhxjzc5ylr39uhuxcz4', clientName: 'stale-live-worker', @@ -628,7 +637,7 @@ describe('TimescaleDB and Redis integration', () => { userAgent, count: '1', bestDifficulty: Number(activeClient.bestDifficulty), - totalHashRate: '0', + totalHashRate: String(activeClient.hashRate), }])); });