mirror of
https://github.com/benjamin-wilson/public-pool.git
synced 2026-09-29 09:05:06 -07:00
Compare commits
4
Commits
d847e2666a
...
a693501648
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
a693501648 | ||
|
|
f259af947e | ||
|
|
f29081de36 | ||
|
|
0acc1910db |
@@ -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 = {
|
||||
getJsonCache: jest.fn().mockResolvedValue(null),
|
||||
setJsonCache: jest.fn().mockResolvedValue(undefined),
|
||||
@@ -259,27 +259,19 @@ describe('ShareAccountingService', () => {
|
||||
.mockResolvedValueOnce([{
|
||||
totalAcceptedShares: '3',
|
||||
totalCreditedDifficulty: '96',
|
||||
acceptedSharesLast10Minutes: '0',
|
||||
creditedDifficultyLast10Minutes: '0',
|
||||
acceptedSharesLast10Minutes: '2',
|
||||
creditedDifficultyLast10Minutes: '64',
|
||||
acceptedSharesLastHour: '3',
|
||||
creditedDifficultyLastHour: '96',
|
||||
acceptedSharesLastDay: '3',
|
||||
creditedDifficultyLastDay: '96',
|
||||
hashRateLast10Minutes: '0',
|
||||
hashRateLast10Minutes: '458129844.9',
|
||||
hashRateLastHour: '114532461.2',
|
||||
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([{
|
||||
bestSubmissionDifficulty: '4096',
|
||||
bestSubmissionDifficultyAt: new Date('2026-06-07T12:19:00Z'),
|
||||
}])
|
||||
.mockResolvedValueOnce([{
|
||||
bestSubmissionDifficultyAt: new Date('2026-06-07T12:10:00Z'),
|
||||
currentRoundAcceptedShares: '11',
|
||||
workSinceLastBlock: '352',
|
||||
currentRoundNetworkDifficulty: '1000',
|
||||
@@ -288,25 +280,24 @@ describe('ShareAccountingService', () => {
|
||||
const service = new ShareAccountingService(repository as any, redis as any);
|
||||
|
||||
await expect(service.refreshPoolSummary()).resolves.toEqual(expect.objectContaining({
|
||||
acceptedSharesLast10Minutes: 7,
|
||||
creditedDifficultyLast10Minutes: 224,
|
||||
hashRateLast10Minutes: 1603451170.77,
|
||||
acceptedSharesLast10Minutes: 2,
|
||||
creditedDifficultyLast10Minutes: 64,
|
||||
hashRateLast10Minutes: 458129844.9,
|
||||
bestSubmissionDifficulty: 4096,
|
||||
bestSubmissionDifficultyAt: '2026-06-07T12:19:00.000Z',
|
||||
bestSubmissionDifficultyAt: '2026-06-07T12:10:00.000Z',
|
||||
workSinceLastBlock: 352,
|
||||
currentRoundAcceptedShares: 11,
|
||||
currentRoundNetworkDifficulty: 1000,
|
||||
networkDifficultyPercent: 35.2,
|
||||
latestShareAt: '2026-06-07T12:20:00.000Z',
|
||||
latestShareAt: '2026-06-07T12:10:00.000Z',
|
||||
}));
|
||||
expect(repository.query).toHaveBeenNthCalledWith(
|
||||
3,
|
||||
expect.stringContaining('WHERE "blockHeight" > latest_found_block."height"'),
|
||||
);
|
||||
expect(repository.query).toHaveBeenNthCalledWith(
|
||||
4,
|
||||
2,
|
||||
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 () => {
|
||||
|
||||
@@ -156,7 +156,7 @@ export class ShareAccountingService implements OnModuleDestroy {
|
||||
}
|
||||
|
||||
public async refreshPoolSummary(): Promise<ShareAccountingSummary> {
|
||||
const summary = await this.withPoolLiveOverlay(await this.getSummary({}));
|
||||
const summary = await this.withPoolRollupOverlay(await this.getSummary({}));
|
||||
await this.redisMessagingService
|
||||
?.setJsonCache(this.poolSummaryCacheKey, summary, 10 * 60 * 1000)
|
||||
.catch(error => {
|
||||
@@ -165,68 +165,52 @@ export class ShareAccountingService implements OnModuleDestroy {
|
||||
return summary;
|
||||
}
|
||||
|
||||
private async withPoolLiveOverlay(summary: ShareAccountingSummary): Promise<ShareAccountingSummary> {
|
||||
private async withPoolRollupOverlay(summary: ShareAccountingSummary): Promise<ShareAccountingSummary> {
|
||||
if (process.env.API_ONLY === 'true') {
|
||||
return summary;
|
||||
}
|
||||
|
||||
const [[liveWindow], [bestDifficultyRow], [currentRoundRow]] = await Promise.all([
|
||||
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(`
|
||||
const [currentRoundRow] = await this.acceptedShareRepository.query(`
|
||||
WITH latest_found_block AS (
|
||||
SELECT COALESCE(MAX("height"), 0) AS "height"
|
||||
FROM "blocks_entity"
|
||||
)
|
||||
SELECT
|
||||
COALESCE("submissionDifficulty", 0)::float AS "bestSubmissionDifficulty",
|
||||
"acceptedAt" AS "bestSubmissionDifficultyAt"
|
||||
FROM "accepted_share_entity", latest_found_block
|
||||
),
|
||||
filtered_rows AS (
|
||||
SELECT "accepted_share_block_10m".*
|
||||
FROM "accepted_share_block_10m", latest_found_block
|
||||
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
|
||||
`),
|
||||
this.acceptedShareRepository.query(`
|
||||
WITH latest_found_block AS (
|
||||
SELECT COALESCE(MAX("height"), 0) AS "height"
|
||||
FROM "blocks_entity"
|
||||
)
|
||||
SELECT
|
||||
COALESCE(SUM("acceptedCount"), 0)::int AS "currentRoundAcceptedShares",
|
||||
COALESCE(SUM("shares"), 0)::float AS "workSinceLastBlock",
|
||||
COALESCE(MAX("networkDifficulty"), 0)::float AS "currentRoundNetworkDifficulty"
|
||||
FROM "accepted_share_block_10m", latest_found_block
|
||||
WHERE "blockHeight" > latest_found_block."height"
|
||||
`),
|
||||
]);
|
||||
COALESCE(MAX("networkDifficulty"), 0)::float AS "currentRoundNetworkDifficulty",
|
||||
COALESCE((SELECT "bestSubmissionDifficulty" FROM best_share), 0)::float AS "bestSubmissionDifficulty",
|
||||
(SELECT "bucket" FROM best_share) AS "bestSubmissionDifficultyAt"
|
||||
FROM filtered_rows
|
||||
`);
|
||||
const currentRoundNetworkDifficulty = this.toNumber(currentRoundRow?.currentRoundNetworkDifficulty);
|
||||
const workSinceLastBlock = this.toNumber(currentRoundRow?.workSinceLastBlock);
|
||||
|
||||
return {
|
||||
...summary,
|
||||
acceptedSharesLast10Minutes: this.toNumber(liveWindow?.acceptedSharesLast10Minutes),
|
||||
creditedDifficultyLast10Minutes: this.toNumber(liveWindow?.creditedDifficultyLast10Minutes),
|
||||
hashRateLast10Minutes: this.toNumber(liveWindow?.hashRateLast10Minutes),
|
||||
bestSubmissionDifficulty: this.toNumber(bestDifficultyRow?.bestSubmissionDifficulty),
|
||||
bestSubmissionDifficultyAt: bestDifficultyRow?.bestSubmissionDifficultyAt == null
|
||||
bestSubmissionDifficulty: this.toNumber(currentRoundRow?.bestSubmissionDifficulty),
|
||||
bestSubmissionDifficultyAt: currentRoundRow?.bestSubmissionDifficultyAt == null
|
||||
? null
|
||||
: new Date(bestDifficultyRow.bestSubmissionDifficultyAt).toISOString(),
|
||||
: new Date(currentRoundRow.bestSubmissionDifficultyAt).toISOString(),
|
||||
workSinceLastBlock,
|
||||
currentRoundAcceptedShares: this.toNumber(currentRoundRow?.currentRoundAcceptedShares),
|
||||
currentRoundNetworkDifficulty,
|
||||
networkDifficultyPercent: currentRoundNetworkDifficulty > 0
|
||||
? this.roundPercent((workSinceLastBlock / currentRoundNetworkDifficulty) * 100)
|
||||
: 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(`
|
||||
WITH bounds AS (
|
||||
WITH clock AS MATERIALIZED (
|
||||
SELECT
|
||||
time_bucket(INTERVAL '10 minutes', NOW()) AS "currentBucket",
|
||||
time_bucket(INTERVAL '10 minutes', NOW()) - INTERVAL '10 minutes' AS "lastCompletedBucket"
|
||||
time_bucket(INTERVAL '10 minutes', NOW()) AS "currentBucket"
|
||||
),
|
||||
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 (
|
||||
SELECT "accepted_share_10m".*, bounds."lastCompletedBucket"
|
||||
SELECT "accepted_share_10m".*, bounds."latestCompletedBucket"
|
||||
FROM "accepted_share_10m", bounds
|
||||
WHERE "clientId" = ANY($1::uuid[])
|
||||
)
|
||||
SELECT
|
||||
"clientId",
|
||||
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
|
||||
GROUP BY "clientId"
|
||||
`, [uniqueClientIds]);
|
||||
@@ -357,27 +356,44 @@ export class ShareAccountingService implements OnModuleDestroy {
|
||||
private async loadSummary(filter: AccountingFilter): Promise<ShareAccountingSummary> {
|
||||
const { whereSql, params } = this.buildWhereClause(filter);
|
||||
const [summary] = await this.acceptedShareRepository.query(`
|
||||
WITH bounds AS (
|
||||
WITH clock AS MATERIALIZED (
|
||||
SELECT
|
||||
time_bucket(INTERVAL '10 minutes', NOW()) AS "currentBucket",
|
||||
time_bucket(INTERVAL '10 minutes', NOW()) - INTERVAL '10 minutes' AS "lastCompletedBucket"
|
||||
time_bucket(INTERVAL '10 minutes', NOW()) AS "currentBucket"
|
||||
),
|
||||
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 (
|
||||
SELECT "accepted_share_10m".*, bounds."currentBucket", bounds."lastCompletedBucket"
|
||||
SELECT "accepted_share_10m".*, bounds."currentBucket", bounds."latestCompletedBucket"
|
||||
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
|
||||
COALESCE(SUM("acceptedCount"), 0)::int AS "totalAcceptedShares",
|
||||
COALESCE(SUM("shares"), 0)::float AS "totalCreditedDifficulty",
|
||||
COALESCE(SUM("acceptedCount") FILTER (WHERE "bucket" = "lastCompletedBucket"), 0)::int AS "acceptedSharesLast10Minutes",
|
||||
COALESCE(SUM("shares") FILTER (WHERE "bucket" = "lastCompletedBucket"), 0)::float AS "creditedDifficultyLast10Minutes",
|
||||
COALESCE(SUM("acceptedCount") FILTER (WHERE "bucket" >= "currentBucket" - INTERVAL '1 hour' AND "bucket" < "currentBucket"), 0)::int AS "acceptedSharesLastHour",
|
||||
COALESCE(SUM("shares") FILTER (WHERE "bucket" >= "currentBucket" - INTERVAL '1 hour' AND "bucket" < "currentBucket"), 0)::float AS "creditedDifficultyLastHour",
|
||||
COALESCE(SUM("acceptedCount") FILTER (WHERE "bucket" >= "currentBucket" - INTERVAL '1 day' AND "bucket" < "currentBucket"), 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" = "lastCompletedBucket") * ${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("acceptedCount") FILTER (WHERE "bucket" = "latestCompletedBucket"), 0)::int AS "acceptedSharesLast10Minutes",
|
||||
COALESCE(SUM("shares") FILTER (WHERE "bucket" = "latestCompletedBucket"), 0)::float AS "creditedDifficultyLast10Minutes",
|
||||
COALESCE(SUM("acceptedCount") FILTER (WHERE "bucket" > "latestCompletedBucket" - INTERVAL '1 hour' AND "bucket" <= "latestCompletedBucket"), 0)::int AS "acceptedSharesLastHour",
|
||||
COALESCE(SUM("shares") FILTER (WHERE "bucket" > "latestCompletedBucket" - INTERVAL '1 hour' AND "bucket" <= "latestCompletedBucket"), 0)::float AS "creditedDifficultyLastHour",
|
||||
COALESCE(SUM("acceptedCount") FILTER (WHERE "bucket" > "latestCompletedBucket" - INTERVAL '1 day' AND "bucket" <= "latestCompletedBucket"), 0)::int AS "acceptedSharesLastDay",
|
||||
COALESCE(SUM("shares") FILTER (WHERE "bucket" > "latestCompletedBucket" - INTERVAL '1 day' AND "bucket" <= "latestCompletedBucket"), 0)::float AS "creditedDifficultyLastDay",
|
||||
COALESCE((SUM("shares") FILTER (WHERE "bucket" = "latestCompletedBucket") * ${HASHES_PER_DIFFICULTY}) / ${ROLLUP_BUCKET_SECONDS}, 0)::float AS "hashRateLast10Minutes",
|
||||
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"
|
||||
FROM filtered_rows
|
||||
`, params);
|
||||
|
||||
Reference in New Issue
Block a user