Compare commits

4 Commits
3 changed files with 102 additions and 81 deletions
@@ -0,0 +1,14 @@
import { MigrationInterface, QueryRunner } from 'typeorm';
export class RemoveAcceptedShareRetention1780959000000 implements MigrationInterface {
public name = 'RemoveAcceptedShareRetention1780959000000';
public transaction = false;
public async up(queryRunner: QueryRunner): Promise<void> {
await queryRunner.query(`SELECT remove_retention_policy('accepted_share_entity', if_exists => TRUE)`);
}
public async down(queryRunner: QueryRunner): Promise<void> {
await queryRunner.query(`SELECT add_retention_policy('accepted_share_entity', INTERVAL '30 days', if_not_exists => TRUE)`);
}
}
@@ -249,7 +249,7 @@ describe('ShareAccountingService', () => {
); );
}); });
it('should overlay live pool data and best share from the current round', async () => { it('should refresh pool summaries from completed rollup buckets and current round rollups', async () => {
const redis = { const redis = {
getJsonCache: jest.fn().mockResolvedValue(null), getJsonCache: jest.fn().mockResolvedValue(null),
setJsonCache: jest.fn().mockResolvedValue(undefined), setJsonCache: jest.fn().mockResolvedValue(undefined),
@@ -259,27 +259,19 @@ describe('ShareAccountingService', () => {
.mockResolvedValueOnce([{ .mockResolvedValueOnce([{
totalAcceptedShares: '3', totalAcceptedShares: '3',
totalCreditedDifficulty: '96', totalCreditedDifficulty: '96',
acceptedSharesLast10Minutes: '0', acceptedSharesLast10Minutes: '2',
creditedDifficultyLast10Minutes: '0', creditedDifficultyLast10Minutes: '64',
acceptedSharesLastHour: '3', acceptedSharesLastHour: '3',
creditedDifficultyLastHour: '96', creditedDifficultyLastHour: '96',
acceptedSharesLastDay: '3', acceptedSharesLastDay: '3',
creditedDifficultyLastDay: '96', creditedDifficultyLastDay: '96',
hashRateLast10Minutes: '0', hashRateLast10Minutes: '458129844.9',
hashRateLastHour: '114532461.2', hashRateLastHour: '114532461.2',
latestShareAt: new Date('2026-06-07T12:10:00Z'), latestShareAt: new Date('2026-06-07T12:10:00Z'),
}]) }])
.mockResolvedValueOnce([{
acceptedSharesLast10Minutes: '7',
creditedDifficultyLast10Minutes: '224',
hashRateLast10Minutes: '1603451170.77',
latestShareAt: new Date('2026-06-07T12:20:00Z'),
}])
.mockResolvedValueOnce([{ .mockResolvedValueOnce([{
bestSubmissionDifficulty: '4096', bestSubmissionDifficulty: '4096',
bestSubmissionDifficultyAt: new Date('2026-06-07T12:19:00Z'), bestSubmissionDifficultyAt: new Date('2026-06-07T12:10:00Z'),
}])
.mockResolvedValueOnce([{
currentRoundAcceptedShares: '11', currentRoundAcceptedShares: '11',
workSinceLastBlock: '352', workSinceLastBlock: '352',
currentRoundNetworkDifficulty: '1000', currentRoundNetworkDifficulty: '1000',
@@ -288,25 +280,24 @@ describe('ShareAccountingService', () => {
const service = new ShareAccountingService(repository as any, redis as any); const service = new ShareAccountingService(repository as any, redis as any);
await expect(service.refreshPoolSummary()).resolves.toEqual(expect.objectContaining({ await expect(service.refreshPoolSummary()).resolves.toEqual(expect.objectContaining({
acceptedSharesLast10Minutes: 7, acceptedSharesLast10Minutes: 2,
creditedDifficultyLast10Minutes: 224, creditedDifficultyLast10Minutes: 64,
hashRateLast10Minutes: 1603451170.77, hashRateLast10Minutes: 458129844.9,
bestSubmissionDifficulty: 4096, bestSubmissionDifficulty: 4096,
bestSubmissionDifficultyAt: '2026-06-07T12:19:00.000Z', bestSubmissionDifficultyAt: '2026-06-07T12:10:00.000Z',
workSinceLastBlock: 352, workSinceLastBlock: 352,
currentRoundAcceptedShares: 11, currentRoundAcceptedShares: 11,
currentRoundNetworkDifficulty: 1000, currentRoundNetworkDifficulty: 1000,
networkDifficultyPercent: 35.2, networkDifficultyPercent: 35.2,
latestShareAt: '2026-06-07T12:20:00.000Z', latestShareAt: '2026-06-07T12:10:00.000Z',
})); }));
expect(repository.query).toHaveBeenNthCalledWith( expect(repository.query).toHaveBeenNthCalledWith(
3, 2,
expect.stringContaining('WHERE "blockHeight" > latest_found_block."height"'),
);
expect(repository.query).toHaveBeenNthCalledWith(
4,
expect.stringContaining('FROM "accepted_share_block_10m"'), expect.stringContaining('FROM "accepted_share_block_10m"'),
); );
expect(repository.query).not.toHaveBeenCalledWith(
expect.stringContaining('FROM "accepted_share_entity"'),
);
}); });
it('should cache accounting summaries briefly to protect hot dashboard endpoints', async () => { it('should cache accounting summaries briefly to protect hot dashboard endpoints', async () => {
@@ -156,7 +156,7 @@ export class ShareAccountingService implements OnModuleDestroy {
} }
public async refreshPoolSummary(): Promise<ShareAccountingSummary> { public async refreshPoolSummary(): Promise<ShareAccountingSummary> {
const summary = await this.withPoolLiveOverlay(await this.getSummary({})); const summary = await this.withPoolRollupOverlay(await this.getSummary({}));
await this.redisMessagingService await this.redisMessagingService
?.setJsonCache(this.poolSummaryCacheKey, summary, 10 * 60 * 1000) ?.setJsonCache(this.poolSummaryCacheKey, summary, 10 * 60 * 1000)
.catch(error => { .catch(error => {
@@ -165,68 +165,52 @@ export class ShareAccountingService implements OnModuleDestroy {
return summary; return summary;
} }
private async withPoolLiveOverlay(summary: ShareAccountingSummary): Promise<ShareAccountingSummary> { private async withPoolRollupOverlay(summary: ShareAccountingSummary): Promise<ShareAccountingSummary> {
if (process.env.API_ONLY === 'true') { if (process.env.API_ONLY === 'true') {
return summary; return summary;
} }
const [[liveWindow], [bestDifficultyRow], [currentRoundRow]] = await Promise.all([ const [currentRoundRow] = await this.acceptedShareRepository.query(`
this.acceptedShareRepository.query(`
SELECT
COUNT(*)::int AS "acceptedSharesLast10Minutes",
COALESCE(SUM("creditedDifficulty"), 0)::float AS "creditedDifficultyLast10Minutes",
COALESCE((SUM("creditedDifficulty") * ${HASHES_PER_DIFFICULTY}) / 600, 0)::float AS "hashRateLast10Minutes",
MAX("acceptedAt") AS "latestShareAt"
FROM "accepted_share_entity"
WHERE "acceptedAt" > NOW() - INTERVAL '10 minutes'
`),
this.acceptedShareRepository.query(`
WITH latest_found_block AS ( WITH latest_found_block AS (
SELECT COALESCE(MAX("height"), 0) AS "height" SELECT COALESCE(MAX("height"), 0) AS "height"
FROM "blocks_entity" FROM "blocks_entity"
) ),
SELECT filtered_rows AS (
COALESCE("submissionDifficulty", 0)::float AS "bestSubmissionDifficulty", SELECT "accepted_share_block_10m".*
"acceptedAt" AS "bestSubmissionDifficultyAt" FROM "accepted_share_block_10m", latest_found_block
FROM "accepted_share_entity", latest_found_block
WHERE "blockHeight" > latest_found_block."height" WHERE "blockHeight" > latest_found_block."height"
ORDER BY "submissionDifficulty" DESC, "acceptedAt" DESC ),
best_share AS (
SELECT
"bestSubmissionDifficulty",
"bucket"
FROM filtered_rows
ORDER BY "bestSubmissionDifficulty" DESC, "bucket" DESC
LIMIT 1 LIMIT 1
`),
this.acceptedShareRepository.query(`
WITH latest_found_block AS (
SELECT COALESCE(MAX("height"), 0) AS "height"
FROM "blocks_entity"
) )
SELECT SELECT
COALESCE(SUM("acceptedCount"), 0)::int AS "currentRoundAcceptedShares", COALESCE(SUM("acceptedCount"), 0)::int AS "currentRoundAcceptedShares",
COALESCE(SUM("shares"), 0)::float AS "workSinceLastBlock", COALESCE(SUM("shares"), 0)::float AS "workSinceLastBlock",
COALESCE(MAX("networkDifficulty"), 0)::float AS "currentRoundNetworkDifficulty" COALESCE(MAX("networkDifficulty"), 0)::float AS "currentRoundNetworkDifficulty",
FROM "accepted_share_block_10m", latest_found_block COALESCE((SELECT "bestSubmissionDifficulty" FROM best_share), 0)::float AS "bestSubmissionDifficulty",
WHERE "blockHeight" > latest_found_block."height" (SELECT "bucket" FROM best_share) AS "bestSubmissionDifficultyAt"
`), FROM filtered_rows
]); `);
const currentRoundNetworkDifficulty = this.toNumber(currentRoundRow?.currentRoundNetworkDifficulty); const currentRoundNetworkDifficulty = this.toNumber(currentRoundRow?.currentRoundNetworkDifficulty);
const workSinceLastBlock = this.toNumber(currentRoundRow?.workSinceLastBlock); const workSinceLastBlock = this.toNumber(currentRoundRow?.workSinceLastBlock);
return { return {
...summary, ...summary,
acceptedSharesLast10Minutes: this.toNumber(liveWindow?.acceptedSharesLast10Minutes), bestSubmissionDifficulty: this.toNumber(currentRoundRow?.bestSubmissionDifficulty),
creditedDifficultyLast10Minutes: this.toNumber(liveWindow?.creditedDifficultyLast10Minutes), bestSubmissionDifficultyAt: currentRoundRow?.bestSubmissionDifficultyAt == null
hashRateLast10Minutes: this.toNumber(liveWindow?.hashRateLast10Minutes),
bestSubmissionDifficulty: this.toNumber(bestDifficultyRow?.bestSubmissionDifficulty),
bestSubmissionDifficultyAt: bestDifficultyRow?.bestSubmissionDifficultyAt == null
? null ? null
: new Date(bestDifficultyRow.bestSubmissionDifficultyAt).toISOString(), : new Date(currentRoundRow.bestSubmissionDifficultyAt).toISOString(),
workSinceLastBlock, workSinceLastBlock,
currentRoundAcceptedShares: this.toNumber(currentRoundRow?.currentRoundAcceptedShares), currentRoundAcceptedShares: this.toNumber(currentRoundRow?.currentRoundAcceptedShares),
currentRoundNetworkDifficulty, currentRoundNetworkDifficulty,
networkDifficultyPercent: currentRoundNetworkDifficulty > 0 networkDifficultyPercent: currentRoundNetworkDifficulty > 0
? this.roundPercent((workSinceLastBlock / currentRoundNetworkDifficulty) * 100) ? this.roundPercent((workSinceLastBlock / currentRoundNetworkDifficulty) * 100)
: 0, : 0,
latestShareAt: liveWindow?.latestShareAt == null
? summary.latestShareAt
: new Date(liveWindow.latestShareAt).toISOString(),
}; };
} }
@@ -274,20 +258,35 @@ export class ShareAccountingService implements OnModuleDestroy {
} }
const rows = await this.acceptedShareRepository.query(` const rows = await this.acceptedShareRepository.query(`
WITH bounds AS ( WITH clock AS MATERIALIZED (
SELECT SELECT
time_bucket(INTERVAL '10 minutes', NOW()) AS "currentBucket", time_bucket(INTERVAL '10 minutes', NOW()) AS "currentBucket"
time_bucket(INTERVAL '10 minutes', NOW()) - INTERVAL '10 minutes' AS "lastCompletedBucket" ),
latest_bucket AS MATERIALIZED (
SELECT "bucket"
FROM "accepted_share_10m", clock
WHERE "bucket" < clock."currentBucket"
ORDER BY "bucket" DESC
LIMIT 1
),
bounds AS MATERIALIZED (
SELECT
"currentBucket",
COALESCE(
(SELECT "bucket" FROM latest_bucket),
"currentBucket" - INTERVAL '10 minutes'
) AS "latestCompletedBucket"
FROM clock
), ),
filtered_rows AS ( filtered_rows AS (
SELECT "accepted_share_10m".*, bounds."lastCompletedBucket" SELECT "accepted_share_10m".*, bounds."latestCompletedBucket"
FROM "accepted_share_10m", bounds FROM "accepted_share_10m", bounds
WHERE "clientId" = ANY($1::uuid[]) WHERE "clientId" = ANY($1::uuid[])
) )
SELECT SELECT
"clientId", "clientId",
MAX("bucket") AS "latestShareAt", MAX("bucket") AS "latestShareAt",
COALESCE((SUM("shares") FILTER (WHERE "bucket" = "lastCompletedBucket") * ${HASHES_PER_DIFFICULTY}) / ${ROLLUP_BUCKET_SECONDS}, 0)::float AS "hashRateLast10Minutes" COALESCE((SUM("shares") FILTER (WHERE "bucket" = "latestCompletedBucket") * ${HASHES_PER_DIFFICULTY}) / ${ROLLUP_BUCKET_SECONDS}, 0)::float AS "hashRateLast10Minutes"
FROM filtered_rows FROM filtered_rows
GROUP BY "clientId" GROUP BY "clientId"
`, [uniqueClientIds]); `, [uniqueClientIds]);
@@ -357,27 +356,44 @@ export class ShareAccountingService implements OnModuleDestroy {
private async loadSummary(filter: AccountingFilter): Promise<ShareAccountingSummary> { private async loadSummary(filter: AccountingFilter): Promise<ShareAccountingSummary> {
const { whereSql, params } = this.buildWhereClause(filter); const { whereSql, params } = this.buildWhereClause(filter);
const [summary] = await this.acceptedShareRepository.query(` const [summary] = await this.acceptedShareRepository.query(`
WITH bounds AS ( WITH clock AS MATERIALIZED (
SELECT SELECT
time_bucket(INTERVAL '10 minutes', NOW()) AS "currentBucket", time_bucket(INTERVAL '10 minutes', NOW()) AS "currentBucket"
time_bucket(INTERVAL '10 minutes', NOW()) - INTERVAL '10 minutes' AS "lastCompletedBucket" ),
latest_bucket AS MATERIALIZED (
SELECT "bucket"
FROM "accepted_share_10m", clock
WHERE "bucket" < clock."currentBucket"
ORDER BY "bucket" DESC
LIMIT 1
),
bounds AS MATERIALIZED (
SELECT
"currentBucket",
COALESCE(
(SELECT "bucket" FROM latest_bucket),
"currentBucket" - INTERVAL '10 minutes'
) AS "latestCompletedBucket"
FROM clock
), ),
filtered_rows AS ( filtered_rows AS (
SELECT "accepted_share_10m".*, bounds."currentBucket", bounds."lastCompletedBucket" SELECT "accepted_share_10m".*, bounds."currentBucket", bounds."latestCompletedBucket"
FROM "accepted_share_10m", bounds FROM "accepted_share_10m", bounds
${whereSql} ${whereSql.length > 0
? `${whereSql} AND "accepted_share_10m"."bucket" <= bounds."latestCompletedBucket"`
: `WHERE "accepted_share_10m"."bucket" <= bounds."latestCompletedBucket"`}
) )
SELECT SELECT
COALESCE(SUM("acceptedCount"), 0)::int AS "totalAcceptedShares", COALESCE(SUM("acceptedCount"), 0)::int AS "totalAcceptedShares",
COALESCE(SUM("shares"), 0)::float AS "totalCreditedDifficulty", COALESCE(SUM("shares"), 0)::float AS "totalCreditedDifficulty",
COALESCE(SUM("acceptedCount") FILTER (WHERE "bucket" = "lastCompletedBucket"), 0)::int AS "acceptedSharesLast10Minutes", COALESCE(SUM("acceptedCount") FILTER (WHERE "bucket" = "latestCompletedBucket"), 0)::int AS "acceptedSharesLast10Minutes",
COALESCE(SUM("shares") FILTER (WHERE "bucket" = "lastCompletedBucket"), 0)::float AS "creditedDifficultyLast10Minutes", COALESCE(SUM("shares") FILTER (WHERE "bucket" = "latestCompletedBucket"), 0)::float AS "creditedDifficultyLast10Minutes",
COALESCE(SUM("acceptedCount") FILTER (WHERE "bucket" >= "currentBucket" - INTERVAL '1 hour' AND "bucket" < "currentBucket"), 0)::int AS "acceptedSharesLastHour", COALESCE(SUM("acceptedCount") FILTER (WHERE "bucket" > "latestCompletedBucket" - INTERVAL '1 hour' AND "bucket" <= "latestCompletedBucket"), 0)::int AS "acceptedSharesLastHour",
COALESCE(SUM("shares") FILTER (WHERE "bucket" >= "currentBucket" - INTERVAL '1 hour' AND "bucket" < "currentBucket"), 0)::float AS "creditedDifficultyLastHour", COALESCE(SUM("shares") FILTER (WHERE "bucket" > "latestCompletedBucket" - INTERVAL '1 hour' AND "bucket" <= "latestCompletedBucket"), 0)::float AS "creditedDifficultyLastHour",
COALESCE(SUM("acceptedCount") FILTER (WHERE "bucket" >= "currentBucket" - INTERVAL '1 day' AND "bucket" < "currentBucket"), 0)::int AS "acceptedSharesLastDay", COALESCE(SUM("acceptedCount") FILTER (WHERE "bucket" > "latestCompletedBucket" - INTERVAL '1 day' AND "bucket" <= "latestCompletedBucket"), 0)::int AS "acceptedSharesLastDay",
COALESCE(SUM("shares") FILTER (WHERE "bucket" >= "currentBucket" - INTERVAL '1 day' AND "bucket" < "currentBucket"), 0)::float AS "creditedDifficultyLastDay", COALESCE(SUM("shares") FILTER (WHERE "bucket" > "latestCompletedBucket" - INTERVAL '1 day' AND "bucket" <= "latestCompletedBucket"), 0)::float AS "creditedDifficultyLastDay",
COALESCE((SUM("shares") FILTER (WHERE "bucket" = "lastCompletedBucket") * ${HASHES_PER_DIFFICULTY}) / ${ROLLUP_BUCKET_SECONDS}, 0)::float AS "hashRateLast10Minutes", COALESCE((SUM("shares") FILTER (WHERE "bucket" = "latestCompletedBucket") * ${HASHES_PER_DIFFICULTY}) / ${ROLLUP_BUCKET_SECONDS}, 0)::float AS "hashRateLast10Minutes",
COALESCE((SUM("shares") FILTER (WHERE "bucket" >= "currentBucket" - INTERVAL '1 hour' AND "bucket" < "currentBucket") * ${HASHES_PER_DIFFICULTY}) / 3600, 0)::float AS "hashRateLastHour", COALESCE((SUM("shares") FILTER (WHERE "bucket" > "latestCompletedBucket" - INTERVAL '1 hour' AND "bucket" <= "latestCompletedBucket") * ${HASHES_PER_DIFFICULTY}) / 3600, 0)::float AS "hashRateLastHour",
MAX("bucket") AS "latestShareAt" MAX("bucket") AS "latestShareAt"
FROM filtered_rows FROM filtered_rows
`, params); `, params);