mirror of
https://github.com/benjamin-wilson/public-pool.git
synced 2026-09-29 09:05:06 -07:00
Compare commits
3
Commits
fd226f206f
...
f31cdef7f0
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
f31cdef7f0 | ||
|
|
943dd60eef | ||
|
|
4150d1e425 |
@@ -0,0 +1,250 @@
|
||||
import { MigrationInterface, QueryRunner } from 'typeorm';
|
||||
|
||||
export class AcceptedShareHighScoreNumericRetention1781402000000 implements MigrationInterface {
|
||||
public name = 'AcceptedShareHighScoreNumericRetention1781402000000';
|
||||
public transaction = false;
|
||||
|
||||
public async up(queryRunner: QueryRunner): Promise<void> {
|
||||
await queryRunner.query(`
|
||||
CREATE INDEX IF NOT EXISTS "IDX_accepted_share_block_10m_best"
|
||||
ON "accepted_share_block_10m" ("bestSubmissionDifficulty" DESC, "bucket" DESC)
|
||||
WHERE "bestSubmissionDifficulty" > 0
|
||||
`);
|
||||
await queryRunner.query(`
|
||||
CREATE INDEX IF NOT EXISTS "IDX_accepted_share_block_10m_mode_best"
|
||||
ON "accepted_share_block_10m" ("payoutMode", "bestSubmissionDifficulty" DESC, "bucket" DESC)
|
||||
WHERE "bestSubmissionDifficulty" > 0
|
||||
`);
|
||||
await queryRunner.query(`
|
||||
CREATE INDEX IF NOT EXISTS "IDX_client_best_difficulty_mode"
|
||||
ON "client_entity" ("payoutMode", "bestDifficulty" DESC, "updatedAt" DESC)
|
||||
WHERE "bestDifficulty" > 0
|
||||
`);
|
||||
await queryRunner.query(`
|
||||
CREATE INDEX IF NOT EXISTS "IDX_address_settings_best_difficulty"
|
||||
ON "address_settings_entity" ("bestDifficulty" DESC, "updatedAt" DESC)
|
||||
WHERE "bestDifficulty" > 0
|
||||
`);
|
||||
|
||||
await this.backfillRollupAllTime(queryRunner);
|
||||
await this.backfillClientAllTime(queryRunner);
|
||||
await this.backfillAddressSettingsAllTime(queryRunner);
|
||||
}
|
||||
|
||||
public async down(queryRunner: QueryRunner): Promise<void> {
|
||||
await queryRunner.query(`DROP INDEX IF EXISTS "IDX_address_settings_best_difficulty"`);
|
||||
await queryRunner.query(`DROP INDEX IF EXISTS "IDX_client_best_difficulty_mode"`);
|
||||
await queryRunner.query(`DROP INDEX IF EXISTS "IDX_accepted_share_block_10m_mode_best"`);
|
||||
await queryRunner.query(`DROP INDEX IF EXISTS "IDX_accepted_share_block_10m_best"`);
|
||||
}
|
||||
|
||||
private async backfillRollupAllTime(queryRunner: QueryRunner): Promise<void> {
|
||||
await queryRunner.query(`
|
||||
WITH candidates AS (
|
||||
SELECT
|
||||
"payoutMode",
|
||||
"bucket",
|
||||
"blockHeight",
|
||||
"bestSubmissionDifficulty"
|
||||
FROM "accepted_share_block_10m"
|
||||
WHERE "bestSubmissionDifficulty" IS NOT NULL
|
||||
AND "bestSubmissionDifficulty" > 0
|
||||
|
||||
UNION ALL
|
||||
|
||||
SELECT
|
||||
'all' AS "payoutMode",
|
||||
"bucket",
|
||||
"blockHeight",
|
||||
"bestSubmissionDifficulty"
|
||||
FROM "accepted_share_block_10m"
|
||||
WHERE "bestSubmissionDifficulty" IS NOT NULL
|
||||
AND "bestSubmissionDifficulty" > 0
|
||||
),
|
||||
best_rows AS (
|
||||
SELECT DISTINCT ON ("payoutMode")
|
||||
"payoutMode",
|
||||
"bucket",
|
||||
"blockHeight",
|
||||
"bestSubmissionDifficulty"
|
||||
FROM candidates
|
||||
ORDER BY "payoutMode", "bestSubmissionDifficulty" DESC, "bucket" DESC
|
||||
)
|
||||
INSERT INTO "accepted_share_high_score" (
|
||||
"scope",
|
||||
"payoutMode",
|
||||
"bucketDate",
|
||||
"submissionDifficulty",
|
||||
"acceptedAt",
|
||||
"bucket",
|
||||
"blockHeight",
|
||||
"address",
|
||||
"clientName",
|
||||
"protocol"
|
||||
)
|
||||
SELECT
|
||||
'all_time',
|
||||
"payoutMode",
|
||||
DATE '1970-01-01',
|
||||
"bestSubmissionDifficulty",
|
||||
"bucket",
|
||||
"bucket",
|
||||
"blockHeight",
|
||||
NULL,
|
||||
NULL,
|
||||
NULL
|
||||
FROM best_rows
|
||||
ON CONFLICT ("scope", "payoutMode", "bucketDate")
|
||||
DO UPDATE SET
|
||||
"submissionDifficulty" = EXCLUDED."submissionDifficulty",
|
||||
"acceptedAt" = EXCLUDED."acceptedAt",
|
||||
"bucket" = EXCLUDED."bucket",
|
||||
"blockHeight" = EXCLUDED."blockHeight",
|
||||
"address" = EXCLUDED."address",
|
||||
"clientName" = EXCLUDED."clientName",
|
||||
"protocol" = EXCLUDED."protocol",
|
||||
"updatedAt" = NOW()
|
||||
WHERE EXCLUDED."submissionDifficulty" > "accepted_share_high_score"."submissionDifficulty"
|
||||
OR (
|
||||
EXCLUDED."submissionDifficulty" = "accepted_share_high_score"."submissionDifficulty"
|
||||
AND EXCLUDED."acceptedAt" > COALESCE("accepted_share_high_score"."acceptedAt", '-infinity'::timestamptz)
|
||||
)
|
||||
`);
|
||||
}
|
||||
|
||||
private async backfillClientAllTime(queryRunner: QueryRunner): Promise<void> {
|
||||
await queryRunner.query(`
|
||||
WITH candidates AS (
|
||||
SELECT
|
||||
"payoutMode",
|
||||
"bestDifficulty",
|
||||
"updatedAt",
|
||||
"address",
|
||||
"clientName",
|
||||
"userAgent"
|
||||
FROM "client_entity"
|
||||
WHERE "bestDifficulty" IS NOT NULL
|
||||
AND "bestDifficulty" > 0
|
||||
|
||||
UNION ALL
|
||||
|
||||
SELECT
|
||||
'all' AS "payoutMode",
|
||||
"bestDifficulty",
|
||||
"updatedAt",
|
||||
"address",
|
||||
"clientName",
|
||||
"userAgent"
|
||||
FROM "client_entity"
|
||||
WHERE "bestDifficulty" IS NOT NULL
|
||||
AND "bestDifficulty" > 0
|
||||
),
|
||||
best_rows AS (
|
||||
SELECT DISTINCT ON ("payoutMode")
|
||||
"payoutMode",
|
||||
"bestDifficulty",
|
||||
"updatedAt",
|
||||
"address",
|
||||
"clientName",
|
||||
"userAgent"
|
||||
FROM candidates
|
||||
ORDER BY "payoutMode", "bestDifficulty" DESC, "updatedAt" DESC
|
||||
)
|
||||
INSERT INTO "accepted_share_high_score" (
|
||||
"scope",
|
||||
"payoutMode",
|
||||
"bucketDate",
|
||||
"submissionDifficulty",
|
||||
"acceptedAt",
|
||||
"bucket",
|
||||
"blockHeight",
|
||||
"address",
|
||||
"clientName",
|
||||
"protocol"
|
||||
)
|
||||
SELECT
|
||||
'all_time',
|
||||
"payoutMode",
|
||||
DATE '1970-01-01',
|
||||
"bestDifficulty",
|
||||
"updatedAt",
|
||||
NULL,
|
||||
NULL,
|
||||
"address",
|
||||
"clientName",
|
||||
"userAgent"
|
||||
FROM best_rows
|
||||
ON CONFLICT ("scope", "payoutMode", "bucketDate")
|
||||
DO UPDATE SET
|
||||
"submissionDifficulty" = EXCLUDED."submissionDifficulty",
|
||||
"acceptedAt" = EXCLUDED."acceptedAt",
|
||||
"bucket" = EXCLUDED."bucket",
|
||||
"blockHeight" = EXCLUDED."blockHeight",
|
||||
"address" = EXCLUDED."address",
|
||||
"clientName" = EXCLUDED."clientName",
|
||||
"protocol" = EXCLUDED."protocol",
|
||||
"updatedAt" = NOW()
|
||||
WHERE EXCLUDED."submissionDifficulty" > "accepted_share_high_score"."submissionDifficulty"
|
||||
OR (
|
||||
EXCLUDED."submissionDifficulty" = "accepted_share_high_score"."submissionDifficulty"
|
||||
AND EXCLUDED."acceptedAt" > COALESCE("accepted_share_high_score"."acceptedAt", '-infinity'::timestamptz)
|
||||
)
|
||||
`);
|
||||
}
|
||||
|
||||
private async backfillAddressSettingsAllTime(queryRunner: QueryRunner): Promise<void> {
|
||||
await queryRunner.query(`
|
||||
WITH best_row AS (
|
||||
SELECT
|
||||
"bestDifficulty",
|
||||
"updatedAt",
|
||||
"address",
|
||||
"bestDifficultyUserAgent"
|
||||
FROM "address_settings_entity"
|
||||
WHERE "bestDifficulty" IS NOT NULL
|
||||
AND "bestDifficulty" > 0
|
||||
ORDER BY "bestDifficulty" DESC, "updatedAt" DESC
|
||||
LIMIT 1
|
||||
)
|
||||
INSERT INTO "accepted_share_high_score" (
|
||||
"scope",
|
||||
"payoutMode",
|
||||
"bucketDate",
|
||||
"submissionDifficulty",
|
||||
"acceptedAt",
|
||||
"bucket",
|
||||
"blockHeight",
|
||||
"address",
|
||||
"clientName",
|
||||
"protocol"
|
||||
)
|
||||
SELECT
|
||||
'all_time',
|
||||
'all',
|
||||
DATE '1970-01-01',
|
||||
"bestDifficulty",
|
||||
"updatedAt",
|
||||
NULL,
|
||||
NULL,
|
||||
"address",
|
||||
NULL,
|
||||
"bestDifficultyUserAgent"
|
||||
FROM best_row
|
||||
ON CONFLICT ("scope", "payoutMode", "bucketDate")
|
||||
DO UPDATE SET
|
||||
"submissionDifficulty" = EXCLUDED."submissionDifficulty",
|
||||
"acceptedAt" = EXCLUDED."acceptedAt",
|
||||
"bucket" = EXCLUDED."bucket",
|
||||
"blockHeight" = EXCLUDED."blockHeight",
|
||||
"address" = EXCLUDED."address",
|
||||
"clientName" = EXCLUDED."clientName",
|
||||
"protocol" = EXCLUDED."protocol",
|
||||
"updatedAt" = NOW()
|
||||
WHERE EXCLUDED."submissionDifficulty" > "accepted_share_high_score"."submissionDifficulty"
|
||||
OR (
|
||||
EXCLUDED."submissionDifficulty" = "accepted_share_high_score"."submissionDifficulty"
|
||||
AND EXCLUDED."acceptedAt" > COALESCE("accepted_share_high_score"."acceptedAt", '-infinity'::timestamptz)
|
||||
)
|
||||
`);
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,50 @@
|
||||
import { MigrationInterface, QueryRunner } from 'typeorm';
|
||||
|
||||
export class PoolSummaryContinuousAggregate1781400000000 implements MigrationInterface {
|
||||
public name = 'PoolSummaryContinuousAggregate1781400000000';
|
||||
public transaction = false;
|
||||
|
||||
public async up(queryRunner: QueryRunner): Promise<void> {
|
||||
await queryRunner.query(`
|
||||
CREATE MATERIALIZED VIEW IF NOT EXISTS "accepted_share_pool_10m"
|
||||
WITH (timescaledb.continuous) AS
|
||||
SELECT
|
||||
time_bucket(INTERVAL '10 minutes', "acceptedAt") AS "bucket",
|
||||
"payoutMode",
|
||||
SUM("creditedDifficulty") AS "shares",
|
||||
COUNT(*) AS "acceptedCount"
|
||||
FROM "accepted_share_entity"
|
||||
GROUP BY "bucket", "payoutMode"
|
||||
WITH NO DATA
|
||||
`);
|
||||
await queryRunner.query(`
|
||||
SELECT add_continuous_aggregate_policy(
|
||||
'accepted_share_pool_10m',
|
||||
start_offset => INTERVAL '25 hours',
|
||||
end_offset => INTERVAL '1 minute',
|
||||
schedule_interval => INTERVAL '1 minute',
|
||||
if_not_exists => TRUE
|
||||
)
|
||||
`);
|
||||
await queryRunner.query(`
|
||||
CREATE INDEX IF NOT EXISTS "IDX_accepted_share_pool_10m_bucket"
|
||||
ON "accepted_share_pool_10m" ("bucket" DESC)
|
||||
`);
|
||||
await queryRunner.query(`
|
||||
CREATE INDEX IF NOT EXISTS "IDX_accepted_share_pool_10m_mode_bucket"
|
||||
ON "accepted_share_pool_10m" ("payoutMode", "bucket" DESC)
|
||||
`);
|
||||
}
|
||||
|
||||
public async down(queryRunner: QueryRunner): Promise<void> {
|
||||
await queryRunner.query(`DROP INDEX IF EXISTS "IDX_accepted_share_pool_10m_mode_bucket"`);
|
||||
await queryRunner.query(`DROP INDEX IF EXISTS "IDX_accepted_share_pool_10m_bucket"`);
|
||||
await queryRunner.query(`
|
||||
SELECT remove_continuous_aggregate_policy(
|
||||
'accepted_share_pool_10m',
|
||||
if_exists => TRUE
|
||||
)
|
||||
`);
|
||||
await queryRunner.query(`DROP MATERIALIZED VIEW IF EXISTS "accepted_share_pool_10m"`);
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,42 @@
|
||||
import { MigrationInterface, QueryRunner } from 'typeorm';
|
||||
|
||||
export class PoolSummaryRefreshWindow1781401000000 implements MigrationInterface {
|
||||
public name = 'PoolSummaryRefreshWindow1781401000000';
|
||||
public transaction = false;
|
||||
|
||||
public async up(queryRunner: QueryRunner): Promise<void> {
|
||||
await queryRunner.query(`
|
||||
SELECT remove_continuous_aggregate_policy(
|
||||
'accepted_share_pool_10m',
|
||||
if_exists => TRUE
|
||||
)
|
||||
`);
|
||||
await queryRunner.query(`
|
||||
SELECT add_continuous_aggregate_policy(
|
||||
'accepted_share_pool_10m',
|
||||
start_offset => INTERVAL '25 hours',
|
||||
end_offset => INTERVAL '1 minute',
|
||||
schedule_interval => INTERVAL '1 minute',
|
||||
if_not_exists => TRUE
|
||||
)
|
||||
`);
|
||||
}
|
||||
|
||||
public async down(queryRunner: QueryRunner): Promise<void> {
|
||||
await queryRunner.query(`
|
||||
SELECT remove_continuous_aggregate_policy(
|
||||
'accepted_share_pool_10m',
|
||||
if_exists => TRUE
|
||||
)
|
||||
`);
|
||||
await queryRunner.query(`
|
||||
SELECT add_continuous_aggregate_policy(
|
||||
'accepted_share_pool_10m',
|
||||
start_offset => INTERVAL '2 days',
|
||||
end_offset => INTERVAL '1 minute',
|
||||
schedule_interval => INTERVAL '1 minute',
|
||||
if_not_exists => TRUE
|
||||
)
|
||||
`);
|
||||
}
|
||||
}
|
||||
@@ -318,7 +318,7 @@ describe('ShareAccountingService', () => {
|
||||
|
||||
expect(repository.query).toHaveBeenNthCalledWith(
|
||||
1,
|
||||
expect.stringContaining('FROM "accepted_share_10m"'),
|
||||
expect.stringContaining('FROM "accepted_share_pool_10m"'),
|
||||
['solo'],
|
||||
);
|
||||
expect(repository.query).toHaveBeenNthCalledWith(
|
||||
@@ -413,6 +413,11 @@ describe('ShareAccountingService', () => {
|
||||
networkDifficultyPercent: 35.2,
|
||||
latestShareAt: '2026-06-07T12:10:00.000Z',
|
||||
}));
|
||||
expect(repository.query).toHaveBeenNthCalledWith(
|
||||
1,
|
||||
expect.stringContaining('FROM "accepted_share_pool_10m"'),
|
||||
[],
|
||||
);
|
||||
expect(repository.query).toHaveBeenNthCalledWith(
|
||||
2,
|
||||
expect.stringContaining('FROM "accepted_share_block_10m"'),
|
||||
|
||||
@@ -605,6 +605,16 @@ export class ShareAccountingService implements OnModuleInit, OnModuleDestroy {
|
||||
}
|
||||
|
||||
private async loadSummary(filter: AccountingFilter): Promise<ShareAccountingSummary> {
|
||||
if (this.isPoolSummaryFilter(filter)) {
|
||||
try {
|
||||
return await this.loadPoolSummary(filter.payoutMode);
|
||||
} catch (error) {
|
||||
if (error?.code !== '42P01') {
|
||||
throw error;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
const { whereSql, params } = this.buildWhereClause(filter);
|
||||
const [summary] = await this.acceptedShareRepository.query(`
|
||||
WITH clock AS MATERIALIZED (
|
||||
@@ -649,6 +659,59 @@ export class ShareAccountingService implements OnModuleInit, OnModuleDestroy {
|
||||
FROM filtered_rows
|
||||
`, params);
|
||||
|
||||
return this.mapSummaryRow(summary);
|
||||
}
|
||||
|
||||
private async loadPoolSummary(payoutMode?: PayoutMode): Promise<ShareAccountingSummary> {
|
||||
const params = payoutMode == null ? [] : [payoutMode];
|
||||
const whereSql = payoutMode == null
|
||||
? 'WHERE "accepted_share_pool_10m"."bucket" <= bounds."latestCompletedBucket"'
|
||||
: 'WHERE "accepted_share_pool_10m"."payoutMode" = $1 AND "accepted_share_pool_10m"."bucket" <= bounds."latestCompletedBucket"';
|
||||
const [summary] = await this.acceptedShareRepository.query(`
|
||||
WITH clock AS MATERIALIZED (
|
||||
SELECT
|
||||
time_bucket(INTERVAL '10 minutes', NOW()) AS "currentBucket"
|
||||
),
|
||||
latest_bucket AS MATERIALIZED (
|
||||
SELECT "bucket"
|
||||
FROM "accepted_share_pool_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_pool_10m".*, bounds."currentBucket", bounds."latestCompletedBucket"
|
||||
FROM "accepted_share_pool_10m", bounds
|
||||
${whereSql}
|
||||
)
|
||||
SELECT
|
||||
COALESCE(SUM("acceptedCount"), 0)::int AS "totalAcceptedShares",
|
||||
COALESCE(SUM("shares"), 0)::float AS "totalCreditedDifficulty",
|
||||
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);
|
||||
|
||||
return this.mapSummaryRow(summary);
|
||||
}
|
||||
|
||||
private mapSummaryRow(summary: Record<string, unknown> | undefined): ShareAccountingSummary {
|
||||
return {
|
||||
totalAcceptedShares: this.toNumber(summary?.totalAcceptedShares),
|
||||
totalCreditedDifficulty: this.toNumber(summary?.totalCreditedDifficulty),
|
||||
@@ -669,11 +732,17 @@ export class ShareAccountingService implements OnModuleInit, OnModuleDestroy {
|
||||
blockCandidateCount: 0,
|
||||
latestShareAt: summary?.latestShareAt == null
|
||||
? null
|
||||
: new Date(summary.latestShareAt).toISOString(),
|
||||
: new Date(summary.latestShareAt as string | Date).toISOString(),
|
||||
protocolBreakdown: [],
|
||||
};
|
||||
}
|
||||
|
||||
private isPoolSummaryFilter(filter: AccountingFilter): boolean {
|
||||
return filter.address == null
|
||||
&& filter.clientName == null
|
||||
&& filter.clientId == null;
|
||||
}
|
||||
|
||||
private scheduleFlush(): void {
|
||||
if (this.flushTimer != null || this.activeFlush != null) {
|
||||
return;
|
||||
|
||||
@@ -0,0 +1,67 @@
|
||||
import { DataSource, EntityManager } from 'typeorm';
|
||||
|
||||
import { ShareHighScoreService } from './share-high-score.service';
|
||||
|
||||
describe('ShareHighScoreService', () => {
|
||||
it('orders aggregate candidates numerically and retains all-time client/address highs', async () => {
|
||||
const upserts: unknown[][] = [];
|
||||
const manager = {
|
||||
query: jest.fn(async (sql: string, params: unknown[] = []) => {
|
||||
if (sql.includes('pg_try_advisory_xact_lock')) {
|
||||
return [{ locked: true }];
|
||||
}
|
||||
if (sql.includes('WITH candidate_rows')) {
|
||||
expect(sql).toContain('"bestSubmissionDifficulty"::numeric DESC');
|
||||
return [];
|
||||
}
|
||||
if (sql.includes('FROM "accepted_share_block_10m"')) {
|
||||
return [];
|
||||
}
|
||||
if (sql.includes('FROM "client_entity"')) {
|
||||
const payoutMode = params[0] as string | undefined;
|
||||
const difficultyByMode = {
|
||||
solo: '10828324014691.838',
|
||||
pplns: '2688488439606.743',
|
||||
all: '10828324014691.838',
|
||||
};
|
||||
const mode = payoutMode ?? 'all';
|
||||
return [{
|
||||
submissionDifficulty: difficultyByMode[mode],
|
||||
acceptedAt: new Date('2026-07-16T06:33:41.270Z'),
|
||||
address: 'bc1qe4zjanpz5tg96a278ew3l2g0h9h0qddgcdvp43',
|
||||
clientName: 'bitaxe',
|
||||
protocol: 'bitaxe',
|
||||
}];
|
||||
}
|
||||
if (sql.includes('FROM "address_settings_entity"')) {
|
||||
return [{
|
||||
submissionDifficulty: '294141974944674.8',
|
||||
acceptedAt: new Date('2026-07-10T03:31:15.987Z'),
|
||||
address: 'bc1q0pp74ghs25vpn2ah6auz4vkehvzy8z8ddyzts7',
|
||||
protocol: 'bitaxe',
|
||||
}];
|
||||
}
|
||||
if (sql.includes('INSERT INTO "accepted_share_high_score"')) {
|
||||
upserts.push(params);
|
||||
return [{ inserted: 1 }];
|
||||
}
|
||||
throw new Error(`Unexpected query: ${sql}`);
|
||||
}),
|
||||
} as unknown as EntityManager;
|
||||
const dataSource = {
|
||||
transaction: jest.fn((callback: (manager: EntityManager) => Promise<unknown>) => callback(manager)),
|
||||
} as unknown as DataSource;
|
||||
const service = new ShareHighScoreService(dataSource);
|
||||
|
||||
await expect(service.refreshHighScores()).resolves.toEqual({
|
||||
processed: true,
|
||||
updatedRows: 4,
|
||||
});
|
||||
|
||||
expect(upserts).toEqual(expect.arrayContaining([
|
||||
expect.arrayContaining(['all_time', 'all', '1970-01-01', '294141974944674.8']),
|
||||
expect.arrayContaining(['all_time', 'solo', '1970-01-01', '10828324014691.838']),
|
||||
expect.arrayContaining(['all_time', 'pplns', '1970-01-01', '2688488439606.743']),
|
||||
]));
|
||||
});
|
||||
});
|
||||
@@ -139,6 +139,11 @@ export class ShareHighScoreService implements OnModuleInit, OnModuleDestroy {
|
||||
});
|
||||
}
|
||||
|
||||
const allTimeRecords = await this.loadAllTimeRecords(manager, payoutMode);
|
||||
for (const record of allTimeRecords) {
|
||||
updatedRows += await this.upsertHighScore(manager, record);
|
||||
}
|
||||
|
||||
return updatedRows;
|
||||
}
|
||||
|
||||
@@ -167,10 +172,130 @@ export class ShareHighScoreService implements OnModuleInit, OnModuleDestroy {
|
||||
"blockHeight"::text AS "blockHeight",
|
||||
"bestSubmissionDifficulty"::text AS "bestSubmissionDifficulty"
|
||||
FROM candidate_rows
|
||||
ORDER BY "bucketDate", "bestSubmissionDifficulty" DESC, "bucket" DESC
|
||||
ORDER BY "bucketDate", "bestSubmissionDifficulty"::numeric DESC, "bucket" DESC
|
||||
`, params);
|
||||
}
|
||||
|
||||
private async loadAllTimeRecords(manager: EntityManager, payoutMode: HighScorePayoutMode): Promise<HighScoreRecord[]> {
|
||||
const records: HighScoreRecord[] = [];
|
||||
const rollupCandidate = await this.loadRollupAllTimeCandidate(manager, payoutMode);
|
||||
if (rollupCandidate != null) {
|
||||
records.push(this.toRecord('all_time', payoutMode, ALL_TIME_BUCKET_DATE, rollupCandidate, null));
|
||||
}
|
||||
|
||||
const clientRecord = await this.loadClientAllTimeRecord(manager, payoutMode);
|
||||
if (clientRecord != null) {
|
||||
records.push(clientRecord);
|
||||
}
|
||||
|
||||
if (payoutMode === 'all') {
|
||||
const addressRecord = await this.loadAddressSettingsAllTimeRecord(manager);
|
||||
if (addressRecord != null) {
|
||||
records.push(addressRecord);
|
||||
}
|
||||
}
|
||||
|
||||
return records;
|
||||
}
|
||||
|
||||
private async loadRollupAllTimeCandidate(
|
||||
manager: EntityManager,
|
||||
payoutMode: HighScorePayoutMode,
|
||||
): Promise<HighScoreCandidate | null> {
|
||||
const params: unknown[] = [];
|
||||
const modeFilter = payoutMode === 'all'
|
||||
? ''
|
||||
: `AND "payoutMode" = $${params.push(payoutMode)}`;
|
||||
const rows = await manager.query(`
|
||||
SELECT
|
||||
date_trunc('day', "bucket")::date::text AS "bucketDate",
|
||||
"bucket",
|
||||
"blockHeight"::text AS "blockHeight",
|
||||
"bestSubmissionDifficulty"::text AS "bestSubmissionDifficulty"
|
||||
FROM "accepted_share_block_10m"
|
||||
WHERE "bestSubmissionDifficulty" IS NOT NULL
|
||||
AND "bestSubmissionDifficulty" > 0
|
||||
${modeFilter}
|
||||
ORDER BY "bestSubmissionDifficulty" DESC, "bucket" DESC
|
||||
LIMIT 1
|
||||
`, params);
|
||||
|
||||
return rows[0] ?? null;
|
||||
}
|
||||
|
||||
private async loadClientAllTimeRecord(
|
||||
manager: EntityManager,
|
||||
payoutMode: HighScorePayoutMode,
|
||||
): Promise<HighScoreRecord | null> {
|
||||
const params: unknown[] = [];
|
||||
const modeFilter = payoutMode === 'all'
|
||||
? ''
|
||||
: `AND "payoutMode" = $${params.push(payoutMode)}`;
|
||||
const rows = await manager.query(`
|
||||
SELECT
|
||||
"bestDifficulty"::text AS "submissionDifficulty",
|
||||
"updatedAt" AS "acceptedAt",
|
||||
"address",
|
||||
"clientName",
|
||||
"userAgent" AS "protocol"
|
||||
FROM "client_entity"
|
||||
WHERE "bestDifficulty" IS NOT NULL
|
||||
AND "bestDifficulty" > 0
|
||||
${modeFilter}
|
||||
ORDER BY "bestDifficulty" DESC, "updatedAt" DESC
|
||||
LIMIT 1
|
||||
`, params);
|
||||
const row = rows[0];
|
||||
if (row == null) {
|
||||
return null;
|
||||
}
|
||||
|
||||
return {
|
||||
scope: 'all_time',
|
||||
payoutMode,
|
||||
bucketDate: ALL_TIME_BUCKET_DATE,
|
||||
submissionDifficulty: row.submissionDifficulty,
|
||||
acceptedAt: row.acceptedAt ?? null,
|
||||
bucket: null,
|
||||
blockHeight: null,
|
||||
address: row.address ?? null,
|
||||
clientName: row.clientName ?? null,
|
||||
protocol: row.protocol ?? null,
|
||||
};
|
||||
}
|
||||
|
||||
private async loadAddressSettingsAllTimeRecord(manager: EntityManager): Promise<HighScoreRecord | null> {
|
||||
const rows = await manager.query(`
|
||||
SELECT
|
||||
"bestDifficulty"::text AS "submissionDifficulty",
|
||||
"updatedAt" AS "acceptedAt",
|
||||
"address",
|
||||
"bestDifficultyUserAgent" AS "protocol"
|
||||
FROM "address_settings_entity"
|
||||
WHERE "bestDifficulty" IS NOT NULL
|
||||
AND "bestDifficulty" > 0
|
||||
ORDER BY "bestDifficulty" DESC, "updatedAt" DESC
|
||||
LIMIT 1
|
||||
`);
|
||||
const row = rows[0];
|
||||
if (row == null) {
|
||||
return null;
|
||||
}
|
||||
|
||||
return {
|
||||
scope: 'all_time',
|
||||
payoutMode: 'all',
|
||||
bucketDate: ALL_TIME_BUCKET_DATE,
|
||||
submissionDifficulty: row.submissionDifficulty,
|
||||
acceptedAt: row.acceptedAt ?? null,
|
||||
bucket: null,
|
||||
blockHeight: null,
|
||||
address: row.address ?? null,
|
||||
clientName: null,
|
||||
protocol: row.protocol ?? null,
|
||||
};
|
||||
}
|
||||
|
||||
private async loadExactHighScore(
|
||||
manager: EntityManager,
|
||||
candidate: HighScoreCandidate,
|
||||
|
||||
@@ -18,6 +18,9 @@ import { PayoutModes1781300000000 } from './ORM/_migrations/PayoutModes178130000
|
||||
import { ShareRollupStoragePolicy1781305000000 } from './ORM/_migrations/ShareRollupStoragePolicy1781305000000';
|
||||
import { AcceptedShareHighScores1781309000000 } from './ORM/_migrations/AcceptedShareHighScores1781309000000';
|
||||
import { UserAgentReportNonzeroHashrate1781313000000 } from './ORM/_migrations/UserAgentReportNonzeroHashrate1781313000000';
|
||||
import { PoolSummaryContinuousAggregate1781400000000 } from './ORM/_migrations/PoolSummaryContinuousAggregate1781400000000';
|
||||
import { PoolSummaryRefreshWindow1781401000000 } from './ORM/_migrations/PoolSummaryRefreshWindow1781401000000';
|
||||
import { AcceptedShareHighScoreNumericRetention1781402000000 } from './ORM/_migrations/AcceptedShareHighScoreNumericRetention1781402000000';
|
||||
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';
|
||||
@@ -62,6 +65,9 @@ export const databaseMigrations = [
|
||||
ShareRollupStoragePolicy1781305000000,
|
||||
AcceptedShareHighScores1781309000000,
|
||||
UserAgentReportNonzeroHashrate1781313000000,
|
||||
PoolSummaryContinuousAggregate1781400000000,
|
||||
PoolSummaryRefreshWindow1781401000000,
|
||||
AcceptedShareHighScoreNumericRetention1781402000000,
|
||||
];
|
||||
|
||||
export function createDatabaseOptions(env: NodeJS.ProcessEnv): TypeOrmModuleOptions & DataSourceOptions {
|
||||
|
||||
@@ -76,7 +76,7 @@ describe('TimescaleDB and Redis integration', () => {
|
||||
const aggregates = await dataSource.query(`
|
||||
SELECT view_name
|
||||
FROM timescaledb_information.continuous_aggregates
|
||||
WHERE view_name IN ('accepted_share_10m', 'accepted_share_1h', 'accepted_share_1d', 'accepted_share_block_10m')
|
||||
WHERE view_name IN ('accepted_share_10m', 'accepted_share_1h', 'accepted_share_1d', 'accepted_share_block_10m', 'accepted_share_pool_10m')
|
||||
ORDER BY view_name
|
||||
`);
|
||||
expect(aggregates.map(row => row.view_name)).toEqual([
|
||||
@@ -84,6 +84,7 @@ describe('TimescaleDB and Redis integration', () => {
|
||||
'accepted_share_1d',
|
||||
'accepted_share_1h',
|
||||
'accepted_share_block_10m',
|
||||
'accepted_share_pool_10m',
|
||||
]);
|
||||
|
||||
const legacyTables = await dataSource.query(`
|
||||
@@ -228,6 +229,7 @@ describe('TimescaleDB and Redis integration', () => {
|
||||
|
||||
await dataSource.query(`CALL refresh_continuous_aggregate('accepted_share_10m', NULL, NULL)`);
|
||||
await dataSource.query(`CALL refresh_continuous_aggregate('accepted_share_block_10m', NULL, NULL)`);
|
||||
await dataSource.query(`CALL refresh_continuous_aggregate('accepted_share_pool_10m', NULL, NULL)`);
|
||||
const aggregateRows = await dataSource.query(`
|
||||
SELECT "shares"::float AS shares, "acceptedCount"::int AS "acceptedCount"
|
||||
FROM accepted_share_10m
|
||||
@@ -250,6 +252,16 @@ describe('TimescaleDB and Redis integration', () => {
|
||||
expect(blockAggregateRows).toEqual(expect.arrayContaining([
|
||||
expect.objectContaining({ shares: 96, acceptedCount: 2, networkDifficulty: 100000 }),
|
||||
]));
|
||||
|
||||
const poolAggregateRows = await dataSource.query(`
|
||||
SELECT "shares"::float AS shares, "acceptedCount"::int AS "acceptedCount"
|
||||
FROM accepted_share_pool_10m
|
||||
WHERE "payoutMode" = $1
|
||||
`, ['solo']);
|
||||
|
||||
expect(poolAggregateRows).toEqual(expect.arrayContaining([
|
||||
expect.objectContaining({ shares: 96, acceptedCount: 2 }),
|
||||
]));
|
||||
});
|
||||
|
||||
it('should retain daily and all-time best share from completed rollup buckets', async () => {
|
||||
|
||||
Reference in New Issue
Block a user