15 Commits
31 changed files with 1547 additions and 150 deletions
+13
View File
@@ -32,6 +32,14 @@ API_PORT=3334
API_BIND_HOST=127.0.0.1
API_PUBLIC_PORT=3334
API_WORKERS=4
# Bound abandoned or slow API sockets. These settings apply only to the
# HTTP(S) listener and do not affect Stratum connections.
API_CONNECTION_TIMEOUT_MS=15000
API_KEEP_ALIVE_TIMEOUT_MS=5000
API_REQUEST_TIMEOUT_MS=15000
API_HEADERS_TIMEOUT_MS=10000
API_MAX_REQUESTS_PER_SOCKET=100
API_MAX_CONNECTIONS_PER_WORKER=500
# Docker json-file log rotation. Applies when using the compose files.
DOCKER_LOG_MAX_SIZE=100m
@@ -64,6 +72,7 @@ STRATUM_TCP_KEEPALIVE_INITIAL_DELAY_MS=60000
# Disconnect slow/non-reading SV1 clients before their per-socket write queue
# can grow without bound. Includes bytes already buffered by Node.js.
STRATUM_MAX_SOCKET_BUFFER_BYTES=262144
STRATUM_MAX_INBOUND_LINE_BYTES=65536
# Immediately publish a consensus-valid, subsidy-only solo job before the full
# transaction template. SV2 uses the same activation to select its pre-staged
# future job; the full job then follows on the active tip.
@@ -162,6 +171,10 @@ SHARE_ACCOUNTING_FLUSH_INTERVAL_MS=25
SHARE_ACCOUNTING_MAX_QUEUE_SIZE=50000
SHARE_ACCOUNTING_SUMMARY_CACHE_MS=2500
SHARE_ACCOUNTING_SUMMARY_CACHE_MAX=10000
SHARE_ACCOUNTING_REDIS_SUMMARY_CACHE_MS=300000
SHARE_ACCOUNTING_REDIS_SUMMARY_STALE_CACHE_MS=3600000
SHARE_ACCOUNTING_REDIS_SUMMARY_LOCK_MS=30000
SHARE_ACCOUNTING_REDIS_SUMMARY_LOCK_WAIT_MS=5000
CLIENT_HASHRATE_PERSIST_INTERVAL_MS=60000
SHARE_ROLLUP_ENABLED=true
SHARE_ROLLUP_INTERVAL_MS=60000
+10
View File
@@ -74,6 +74,12 @@ services:
API_PORT: ${API_PORT:-3334}
API_SECURE: ${API_SECURE:-false}
API_WORKERS: ${API_WORKERS:-4}
API_CONNECTION_TIMEOUT_MS: ${API_CONNECTION_TIMEOUT_MS:-15000}
API_KEEP_ALIVE_TIMEOUT_MS: ${API_KEEP_ALIVE_TIMEOUT_MS:-5000}
API_REQUEST_TIMEOUT_MS: ${API_REQUEST_TIMEOUT_MS:-15000}
API_HEADERS_TIMEOUT_MS: ${API_HEADERS_TIMEOUT_MS:-10000}
API_MAX_REQUESTS_PER_SOCKET: ${API_MAX_REQUESTS_PER_SOCKET:-100}
API_MAX_CONNECTIONS_PER_WORKER: ${API_MAX_CONNECTIONS_PER_WORKER:-500}
PM2_ENABLED: ${PM2_ENABLED:-true}
STRATUM_WORKERS: ${STRATUM_WORKERS:-auto}
STRATUM_PORTS: ${STRATUM_PORTS:-3333,3332,3331,3330}
@@ -100,6 +106,10 @@ services:
SHARE_ACCOUNTING_MAX_QUEUE_SIZE: ${SHARE_ACCOUNTING_MAX_QUEUE_SIZE:-50000}
SHARE_ACCOUNTING_SUMMARY_CACHE_MS: ${SHARE_ACCOUNTING_SUMMARY_CACHE_MS:-2500}
SHARE_ACCOUNTING_SUMMARY_CACHE_MAX: ${SHARE_ACCOUNTING_SUMMARY_CACHE_MAX:-10000}
SHARE_ACCOUNTING_REDIS_SUMMARY_CACHE_MS: ${SHARE_ACCOUNTING_REDIS_SUMMARY_CACHE_MS:-300000}
SHARE_ACCOUNTING_REDIS_SUMMARY_STALE_CACHE_MS: ${SHARE_ACCOUNTING_REDIS_SUMMARY_STALE_CACHE_MS:-3600000}
SHARE_ACCOUNTING_REDIS_SUMMARY_LOCK_MS: ${SHARE_ACCOUNTING_REDIS_SUMMARY_LOCK_MS:-30000}
SHARE_ACCOUNTING_REDIS_SUMMARY_LOCK_WAIT_MS: ${SHARE_ACCOUNTING_REDIS_SUMMARY_LOCK_WAIT_MS:-5000}
SHARE_ROLLUP_ENABLED: ${SHARE_ROLLUP_ENABLED:-true}
SHARE_ROLLUP_INTERVAL_MS: ${SHARE_ROLLUP_INTERVAL_MS:-60000}
SHARE_ROLLUP_SAFETY_LAG_SECONDS: ${SHARE_ROLLUP_SAFETY_LAG_SECONDS:-30}
+10
View File
@@ -104,6 +104,16 @@ services:
API_PORT: ${API_PORT:-3334}
API_SECURE: ${API_SECURE:-false}
API_WORKERS: ${API_WORKERS:-4}
API_CONNECTION_TIMEOUT_MS: ${API_CONNECTION_TIMEOUT_MS:-15000}
API_KEEP_ALIVE_TIMEOUT_MS: ${API_KEEP_ALIVE_TIMEOUT_MS:-5000}
API_REQUEST_TIMEOUT_MS: ${API_REQUEST_TIMEOUT_MS:-15000}
API_HEADERS_TIMEOUT_MS: ${API_HEADERS_TIMEOUT_MS:-10000}
API_MAX_REQUESTS_PER_SOCKET: ${API_MAX_REQUESTS_PER_SOCKET:-100}
API_MAX_CONNECTIONS_PER_WORKER: ${API_MAX_CONNECTIONS_PER_WORKER:-500}
SHARE_ACCOUNTING_REDIS_SUMMARY_CACHE_MS: ${SHARE_ACCOUNTING_REDIS_SUMMARY_CACHE_MS:-300000}
SHARE_ACCOUNTING_REDIS_SUMMARY_STALE_CACHE_MS: ${SHARE_ACCOUNTING_REDIS_SUMMARY_STALE_CACHE_MS:-3600000}
SHARE_ACCOUNTING_REDIS_SUMMARY_LOCK_MS: ${SHARE_ACCOUNTING_REDIS_SUMMARY_LOCK_MS:-30000}
SHARE_ACCOUNTING_REDIS_SUMMARY_LOCK_WAIT_MS: ${SHARE_ACCOUNTING_REDIS_SUMMARY_LOCK_WAIT_MS:-5000}
PM2_ENABLED: ${PM2_ENABLED:-true}
STRATUM_WORKERS: ${STRATUM_WORKERS:-auto}
STRATUM_PORTS: ${STRATUM_PORTS:-3333,3332,3331,3330}
+10
View File
@@ -108,6 +108,16 @@ services:
DB_DATABASE: public_pool_mainnet
REDIS_URL: redis://redis:6379
API_WORKERS: ${API_WORKERS:-4}
API_CONNECTION_TIMEOUT_MS: ${API_CONNECTION_TIMEOUT_MS:-15000}
API_KEEP_ALIVE_TIMEOUT_MS: ${API_KEEP_ALIVE_TIMEOUT_MS:-5000}
API_REQUEST_TIMEOUT_MS: ${API_REQUEST_TIMEOUT_MS:-15000}
API_HEADERS_TIMEOUT_MS: ${API_HEADERS_TIMEOUT_MS:-10000}
API_MAX_REQUESTS_PER_SOCKET: ${API_MAX_REQUESTS_PER_SOCKET:-100}
API_MAX_CONNECTIONS_PER_WORKER: ${API_MAX_CONNECTIONS_PER_WORKER:-500}
SHARE_ACCOUNTING_REDIS_SUMMARY_CACHE_MS: ${SHARE_ACCOUNTING_REDIS_SUMMARY_CACHE_MS:-300000}
SHARE_ACCOUNTING_REDIS_SUMMARY_STALE_CACHE_MS: ${SHARE_ACCOUNTING_REDIS_SUMMARY_STALE_CACHE_MS:-3600000}
SHARE_ACCOUNTING_REDIS_SUMMARY_LOCK_MS: ${SHARE_ACCOUNTING_REDIS_SUMMARY_LOCK_MS:-30000}
SHARE_ACCOUNTING_REDIS_SUMMARY_LOCK_WAIT_MS: ${SHARE_ACCOUNTING_REDIS_SUMMARY_LOCK_WAIT_MS:-5000}
PM2_ENABLED: "true"
STRATUM_WORKERS: ${STRATUM_WORKERS:-2}
STRATUM_MIN_DIFFICULTY: ${STRATUM_MIN_DIFFICULTY:-1}
+33
View File
@@ -0,0 +1,33 @@
#!/bin/sh
set -eu
lineage=${RENEWED_LINEAGE:-/etc/letsencrypt/live/public-pool.io}
secrets_dir=${PUBLIC_POOL_SECRETS_DIR:-/home/ben/public-pool-timescaledb-test/secrets}
container=${PUBLIC_POOL_CONTAINER:-public-pool}
owner=${PUBLIC_POOL_CERT_OWNER:-ben}
group=${PUBLIC_POOL_CERT_GROUP:-ben}
cert_source="$lineage/fullchain.pem"
key_source="$lineage/privkey.pem"
test -r "$cert_source"
test -r "$key_source"
openssl x509 -in "$cert_source" -noout >/dev/null
openssl pkey -in "$key_source" -noout >/dev/null
cert_temp=$(mktemp "$secrets_dir/.cert.pem.XXXXXX")
key_temp=$(mktemp "$secrets_dir/.key.pem.XXXXXX")
trap 'rm -f "$cert_temp" "$key_temp"' EXIT HUP INT TERM
cat "$cert_source" > "$cert_temp"
cat "$key_source" > "$key_temp"
chown "$owner:$group" "$cert_temp" "$key_temp"
chmod 0644 "$cert_temp"
chmod 0600 "$key_temp"
mv -f "$cert_temp" "$secrets_dir/cert.pem"
mv -f "$key_temp" "$secrets_dir/key.pem"
# HTTPS can reload its secure context, but the Stratum TLS listeners currently
# read their certificate only when they start.
docker restart --time 30 "$container" >/dev/null
@@ -0,0 +1,46 @@
import { MigrationInterface, QueryRunner } from 'typeorm';
export class AcceptedShare10mLookupIndexes1781403000000 implements MigrationInterface {
public name = 'AcceptedShare10mLookupIndexes1781403000000';
public async up(queryRunner: QueryRunner): Promise<void> {
await queryRunner.query(`
CREATE INDEX IF NOT EXISTS "IDX_accepted_share_10m_bucket"
ON "accepted_share_10m" ("bucket" DESC)
`);
await queryRunner.query(`
CREATE INDEX IF NOT EXISTS "IDX_accepted_share_10m_address_bucket"
ON "accepted_share_10m" ("address", "bucket" DESC)
`);
await queryRunner.query(`
CREATE INDEX IF NOT EXISTS "IDX_accepted_share_10m_address_mode_bucket"
ON "accepted_share_10m" ("address", "payoutMode", "bucket" DESC)
`);
await queryRunner.query(`
CREATE INDEX IF NOT EXISTS "IDX_accepted_share_10m_group_bucket"
ON "accepted_share_10m" ("address", "clientName", "bucket" DESC)
`);
await queryRunner.query(`
CREATE INDEX IF NOT EXISTS "IDX_accepted_share_10m_group_mode_bucket"
ON "accepted_share_10m" ("address", "clientName", "payoutMode", "bucket" DESC)
`);
await queryRunner.query(`
CREATE INDEX IF NOT EXISTS "IDX_accepted_share_10m_client_bucket"
ON "accepted_share_10m" ("clientId", "bucket" DESC)
`);
await queryRunner.query(`
CREATE INDEX IF NOT EXISTS "IDX_accepted_share_10m_client_mode_bucket"
ON "accepted_share_10m" ("clientId", "payoutMode", "bucket" DESC)
`);
}
public async down(queryRunner: QueryRunner): Promise<void> {
await queryRunner.query(`DROP INDEX IF EXISTS "IDX_accepted_share_10m_client_mode_bucket"`);
await queryRunner.query(`DROP INDEX IF EXISTS "IDX_accepted_share_10m_group_mode_bucket"`);
await queryRunner.query(`DROP INDEX IF EXISTS "IDX_accepted_share_10m_address_mode_bucket"`);
await queryRunner.query(`DROP INDEX IF EXISTS "IDX_accepted_share_10m_client_bucket"`);
await queryRunner.query(`DROP INDEX IF EXISTS "IDX_accepted_share_10m_group_bucket"`);
await queryRunner.query(`DROP INDEX IF EXISTS "IDX_accepted_share_10m_address_bucket"`);
await queryRunner.query(`DROP INDEX IF EXISTS "IDX_accepted_share_10m_bucket"`);
}
}
@@ -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(
@@ -365,11 +365,65 @@ describe('ShareAccountingService', () => {
expect(repository.query).toHaveBeenCalledTimes(1);
expect(redis.setJsonCache).toHaveBeenCalledWith(
expect.stringContaining('accounting:summary:'),
expect.objectContaining({ totalAcceptedShares: 1 }),
30000,
expect.objectContaining({
schemaVersion: 1,
refreshedAtMs: expect.any(Number),
value: expect.objectContaining({ totalAcceptedShares: 1 }),
}),
3600000,
);
});
it('serves a stale shared summary while one API worker refreshes it', async () => {
process.env.SHARE_ACCOUNTING_SUMMARY_CACHE_MS = '0';
process.env.SHARE_ACCOUNTING_REDIS_SUMMARY_CACHE_MS = '100';
const staleSummary = {
...new ShareAccountingService({} as any).emptySummary(),
totalAcceptedShares: 7,
};
const repository = {
query: jest.fn().mockResolvedValueOnce([{
totalAcceptedShares: '8',
totalCreditedDifficulty: '256',
acceptedSharesLast10Minutes: '1',
creditedDifficultyLast10Minutes: '32',
acceptedSharesLastHour: '8',
creditedDifficultyLastHour: '256',
acceptedSharesLastDay: '8',
creditedDifficultyLastDay: '256',
hashRateLast10Minutes: '1',
hashRateLastHour: '1',
latestShareAt: null,
}]),
};
const redis = {
getJsonCache: jest.fn().mockResolvedValue({
schemaVersion: 1,
refreshedAtMs: Date.now() - 1000,
value: staleSummary,
}),
setJsonCache: jest.fn().mockResolvedValue(undefined),
tryAcquireJsonCacheLock: jest.fn().mockResolvedValue(true),
releaseJsonCacheLock: jest.fn().mockResolvedValue(undefined),
};
const service = new ShareAccountingService(repository as any, redis as any);
await expect(service.getAddressSummary('bc1qstale')).resolves.toEqual(staleSummary);
await new Promise(resolve => setImmediate(resolve));
expect(repository.query).toHaveBeenCalledTimes(1);
expect(redis.tryAcquireJsonCacheLock).toHaveBeenCalledTimes(1);
expect(redis.setJsonCache).toHaveBeenCalledWith(
expect.stringContaining('accounting:summary:'),
expect.objectContaining({
schemaVersion: 1,
value: expect.objectContaining({ totalAcceptedShares: 8 }),
}),
3600000,
);
expect(redis.releaseJsonCacheLock).toHaveBeenCalledTimes(1);
});
it('should refresh pool summaries from completed rollup buckets and current round rollups', async () => {
const redis = {
getJsonCache: jest.fn().mockResolvedValue(null),
@@ -413,6 +467,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"'),
@@ -1,5 +1,6 @@
import { Injectable, OnModuleDestroy, OnModuleInit } from '@nestjs/common';
import { InjectRepository } from '@nestjs/typeorm';
import { randomUUID } from 'crypto';
import { Repository } from 'typeorm';
import { AcceptedShareEntity } from '../accepted-share/accepted-share.entity';
@@ -88,7 +89,10 @@ const DEFAULT_FLUSH_INTERVAL_MS = 25;
const DEFAULT_MAX_QUEUE_SIZE = 50000;
const DEFAULT_SUMMARY_CACHE_MS = 2500;
const DEFAULT_SUMMARY_CACHE_MAX = 10000;
const DEFAULT_REDIS_SUMMARY_CACHE_MS = 30000;
const DEFAULT_REDIS_SUMMARY_CACHE_MS = 5 * 60 * 1000;
const DEFAULT_REDIS_SUMMARY_STALE_CACHE_MS = 60 * 60 * 1000;
const DEFAULT_REDIS_SUMMARY_LOCK_MS = 30 * 1000;
const DEFAULT_REDIS_SUMMARY_LOCK_WAIT_MS = 5 * 1000;
const DEFAULT_ROLLUP_INTERVAL_MS = 60000;
const DEFAULT_ROLLUP_SAFETY_LAG_SECONDS = 10;
const DEFAULT_ROLLUP_MAX_SHARES_PER_BATCH = 5000000;
@@ -102,6 +106,8 @@ export class ShareAccountingService implements OnModuleInit, OnModuleDestroy {
private rollupTimer: NodeJS.Timeout | null = null;
private activeRollup: Promise<void> | null = null;
private summaryCache = new Map<string, SummaryCacheEntry>();
private summaryInFlight = new Map<string, Promise<ShareAccountingSummary>>();
private redisSummaryRefreshInFlight = new Map<string, Promise<void>>();
private readonly poolSummaryCacheKey = 'accounting:pool-summary';
private readonly batchSize = this.readPositiveInt('SHARE_ACCOUNTING_BATCH_SIZE', DEFAULT_BATCH_SIZE);
private readonly flushIntervalMs = this.readPositiveInt('SHARE_ACCOUNTING_FLUSH_INTERVAL_MS', DEFAULT_FLUSH_INTERVAL_MS);
@@ -109,6 +115,12 @@ export class ShareAccountingService implements OnModuleInit, OnModuleDestroy {
private readonly summaryCacheMs = this.readNonNegativeInt('SHARE_ACCOUNTING_SUMMARY_CACHE_MS', DEFAULT_SUMMARY_CACHE_MS);
private readonly summaryCacheMax = this.readPositiveInt('SHARE_ACCOUNTING_SUMMARY_CACHE_MAX', DEFAULT_SUMMARY_CACHE_MAX);
private readonly redisSummaryCacheMs = this.readNonNegativeInt('SHARE_ACCOUNTING_REDIS_SUMMARY_CACHE_MS', DEFAULT_REDIS_SUMMARY_CACHE_MS);
private readonly redisSummaryStaleCacheMs = Math.max(
this.redisSummaryCacheMs,
this.readPositiveInt('SHARE_ACCOUNTING_REDIS_SUMMARY_STALE_CACHE_MS', DEFAULT_REDIS_SUMMARY_STALE_CACHE_MS),
);
private readonly redisSummaryLockMs = this.readPositiveInt('SHARE_ACCOUNTING_REDIS_SUMMARY_LOCK_MS', DEFAULT_REDIS_SUMMARY_LOCK_MS);
private readonly redisSummaryLockWaitMs = this.readPositiveInt('SHARE_ACCOUNTING_REDIS_SUMMARY_LOCK_WAIT_MS', DEFAULT_REDIS_SUMMARY_LOCK_WAIT_MS);
private readonly shareRollupEnabled = this.readBoolean('SHARE_ROLLUP_ENABLED', true);
private readonly shareRollupIntervalMs = this.readPositiveInt('SHARE_ROLLUP_INTERVAL_MS', DEFAULT_ROLLUP_INTERVAL_MS);
private readonly shareRollupSafetyLagSeconds = this.readNonNegativeInt('SHARE_ROLLUP_SAFETY_LAG_SECONDS', DEFAULT_ROLLUP_SAFETY_LAG_SECONDS);
@@ -565,46 +577,219 @@ export class ShareAccountingService implements OnModuleInit, OnModuleDestroy {
return cached.value;
}
const value = this.loadCachedSummary(filter, cacheKey).catch(error => {
if (cached != null) {
this.summaryCache.delete(cacheKey);
throw error;
});
if (this.summaryCacheMs > 0) {
this.summaryCache.set(cacheKey, {
expiresAt: now + this.summaryCacheMs,
value,
});
this.trimSummaryCache();
}
const inFlight = this.summaryInFlight.get(cacheKey);
if (inFlight != null) {
return inFlight;
}
const value = this.loadCachedSummary(filter, cacheKey)
.then(summary => {
if (this.summaryCacheMs > 0) {
this.summaryCache.set(cacheKey, {
expiresAt: Date.now() + this.summaryCacheMs,
value: summary,
});
this.trimSummaryCache();
}
return summary;
})
.finally(() => {
this.summaryInFlight.delete(cacheKey);
});
this.summaryInFlight.set(cacheKey, value);
return value;
}
private async loadCachedSummary(filter: AccountingFilter, cacheKey: string): Promise<ShareAccountingSummary> {
const redisCacheKey = `accounting:summary:${cacheKey}`;
if (this.redisMessagingService == null || this.redisSummaryCacheMs <= 0) {
return this.loadSummary(filter);
}
const cached = await this.readRedisSummaryCache(redisCacheKey);
if (cached != null) {
if (Date.now() - cached.refreshedAtMs >= this.redisSummaryCacheMs) {
this.refreshRedisSummaryCacheInBackground(filter, redisCacheKey);
}
return cached.value;
}
return this.loadColdRedisSummaryCache(filter, redisCacheKey);
}
private async readRedisSummaryCache(redisCacheKey: string): Promise<RedisSummaryCacheEntry | null> {
const cached = await this.redisMessagingService
?.getJsonCache<ShareAccountingSummary>(redisCacheKey)
?.getJsonCache<RedisSummaryCacheEntry | ShareAccountingSummary>(redisCacheKey)
.catch(error => {
console.error(`Share accounting summary cache read failed: ${error.message}`);
return null;
});
if (cached != null) {
if (cached == null) {
return null;
}
if (this.isRedisSummaryCacheEntry(cached)) {
return cached;
}
// Cache entries written by older workers have no timestamp. Their old
// short Redis TTL still bounds how long they can be treated as fresh.
return {
schemaVersion: 1,
refreshedAtMs: Date.now(),
value: cached,
};
}
private async loadColdRedisSummaryCache(
filter: AccountingFilter,
redisCacheKey: string,
): Promise<ShareAccountingSummary> {
if (!this.supportsRedisSummaryLock()) {
const summary = await this.loadSummary(filter);
await this.storeRedisSummaryCache(redisCacheKey, summary);
return summary;
}
const owner = randomUUID();
const acquired = await this.tryAcquireRedisSummaryLock(redisCacheKey, owner);
if (acquired === true) {
return this.loadAndStoreRedisSummary(filter, redisCacheKey, owner);
}
if (acquired == null) {
const summary = await this.loadSummary(filter);
await this.storeRedisSummaryCache(redisCacheKey, summary);
return summary;
}
const deadline = Date.now() + this.redisSummaryLockWaitMs;
while (Date.now() < deadline) {
await new Promise(resolve => setTimeout(resolve, 50));
const cached = await this.readRedisSummaryCache(redisCacheKey);
if (cached != null) {
return cached.value;
}
}
// Redis locking is an optimization, not an availability dependency.
// If a lock holder died or a refresh exceeded its budget, fail open.
const summary = await this.loadSummary(filter);
await this.redisMessagingService
?.setJsonCache(redisCacheKey, summary, this.redisSummaryCacheMs)
.catch(error => {
console.error(`Share accounting summary cache write failed: ${error.message}`);
});
await this.storeRedisSummaryCache(redisCacheKey, summary);
return summary;
}
private refreshRedisSummaryCacheInBackground(filter: AccountingFilter, redisCacheKey: string): void {
if (this.redisSummaryRefreshInFlight.has(redisCacheKey)) {
return;
}
const refresh = this.refreshRedisSummaryCache(filter, redisCacheKey)
.catch(error => {
console.error(`Share accounting summary background refresh failed: ${error.message}`);
})
.finally(() => {
this.redisSummaryRefreshInFlight.delete(redisCacheKey);
});
this.redisSummaryRefreshInFlight.set(redisCacheKey, refresh);
}
private async refreshRedisSummaryCache(filter: AccountingFilter, redisCacheKey: string): Promise<void> {
if (!this.supportsRedisSummaryLock()) {
const summary = await this.loadSummary(filter);
await this.storeRedisSummaryCache(redisCacheKey, summary);
return;
}
const owner = randomUUID();
if (await this.tryAcquireRedisSummaryLock(redisCacheKey, owner) !== true) {
return;
}
await this.loadAndStoreRedisSummary(filter, redisCacheKey, owner);
}
private async loadAndStoreRedisSummary(
filter: AccountingFilter,
redisCacheKey: string,
owner: string,
): Promise<ShareAccountingSummary> {
try {
const summary = await this.loadSummary(filter);
await this.storeRedisSummaryCache(redisCacheKey, summary);
return summary;
} finally {
if (this.supportsRedisSummaryLock()) {
await this.redisMessagingService
.releaseJsonCacheLock(redisCacheKey, owner)
.catch(error => {
console.error(`Share accounting summary cache lock release failed: ${error.message}`);
});
}
}
}
private async storeRedisSummaryCache(
redisCacheKey: string,
summary: ShareAccountingSummary,
): Promise<void> {
const entry: RedisSummaryCacheEntry = {
schemaVersion: 1,
refreshedAtMs: Date.now(),
value: summary,
};
await this.redisMessagingService
?.setJsonCache(redisCacheKey, entry, this.redisSummaryStaleCacheMs)
.catch(error => {
console.error(`Share accounting summary cache write failed: ${error.message}`);
});
}
private async tryAcquireRedisSummaryLock(redisCacheKey: string, owner: string): Promise<boolean | null> {
if (!this.supportsRedisSummaryLock()) {
return null;
}
try {
return await this.redisMessagingService
.tryAcquireJsonCacheLock(redisCacheKey, owner, this.redisSummaryLockMs);
} catch (error) {
console.error(`Share accounting summary cache lock failed: ${error.message}`);
return null;
}
}
private supportsRedisSummaryLock(): boolean {
return typeof this.redisMessagingService?.tryAcquireJsonCacheLock === 'function'
&& typeof this.redisMessagingService?.releaseJsonCacheLock === 'function';
}
private isRedisSummaryCacheEntry(
cached: RedisSummaryCacheEntry | ShareAccountingSummary,
): cached is RedisSummaryCacheEntry {
const candidate = cached as Partial<RedisSummaryCacheEntry>;
return candidate.schemaVersion === 1
&& Number.isFinite(candidate.refreshedAtMs)
&& candidate.value != null;
}
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 +834,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 +907,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;
@@ -796,6 +1040,13 @@ export class ShareAccountingService implements OnModuleInit, OnModuleDestroy {
}
private trimSummaryCache(): void {
const now = Date.now();
for (const [key, entry] of this.summaryCache) {
if (entry.expiresAt <= now) {
this.summaryCache.delete(key);
}
}
while (this.summaryCache.size > this.summaryCacheMax) {
const firstKey = this.summaryCache.keys().next().value;
if (firstKey == null) {
@@ -832,5 +1083,11 @@ interface PendingShare {
interface SummaryCacheEntry {
expiresAt: number;
value: Promise<ShareAccountingSummary>;
value: ShareAccountingSummary;
}
interface RedisSummaryCacheEntry {
schemaVersion: 1;
refreshedAtMs: number;
value: ShareAccountingSummary;
}
@@ -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,
+62
View File
@@ -0,0 +1,62 @@
import { CacheModule } from '@nestjs/cache-manager';
import { Module } from '@nestjs/common';
import { ConfigModule, ConfigService } from '@nestjs/config';
import { TypeOrmModule } from '@nestjs/typeorm';
import { AppController } from './app.controller';
import { AddressController } from './controllers/address/address.controller';
import { ClientController } from './controllers/client/client.controller';
import { UserAgentReportModule } from './ORM/_views/user-agent-report/user-agent-report.module';
import { AddressSettingsModule } from './ORM/address-settings/address-settings.module';
import { BlocksModule } from './ORM/blocks/blocks.module';
import { ClientStatisticsModule } from './ORM/client-statistics/client-statistics.module';
import { ClientModule } from './ORM/client/client.module';
import { PayoutSnapshotModule } from './ORM/payout-snapshot/payout-snapshot.module';
import { RpcBlocksModule } from './ORM/rpc-block/rpc-block.module';
import { ShareAccountingModule } from './ORM/share-accounting/share-accounting.module';
import { createDatabaseOptions } from './database.config';
import { BitcoinAddressValidator } from './models/validators/bitcoin-address.validator';
import { BitcoinRpcService } from './services/bitcoin-rpc.service';
import { RedisMessagingModule } from './services/redis-messaging.module';
import { Sv2AuthorityService } from './services/sv2-authority.service';
@Module({
imports: [
ConfigModule.forRoot(),
TypeOrmModule.forRootAsync({
useFactory: (configService: ConfigService) => createDatabaseOptions({
...process.env,
DB_HOST: configService.get('DB_HOST'),
DB_PORT: configService.get('DB_PORT'),
DB_USERNAME: configService.get('DB_USERNAME'),
DB_PASSWORD: configService.get('DB_PASSWORD'),
DB_DATABASE: configService.get('DB_DATABASE'),
DB_LOGGING: configService.get('DB_LOGGING'),
DB_POOL_SIZE: configService.get('DB_POOL_SIZE'),
}),
imports: [ConfigModule],
inject: [ConfigService],
}),
CacheModule.register(),
RedisMessagingModule,
ClientStatisticsModule,
ClientModule,
AddressSettingsModule,
BlocksModule,
RpcBlocksModule,
UserAgentReportModule,
ShareAccountingModule,
PayoutSnapshotModule,
],
controllers: [
AppController,
ClientController,
AddressController,
],
providers: [
BitcoinRpcService,
BitcoinAddressValidator,
Sv2AuthorityService,
],
})
export class ApiModule {}
+2 -2
View File
@@ -6,7 +6,7 @@ describe('AppController', () => {
blocksService?: any;
addressSettingsService?: any;
userAgentReportService?: any;
stratumV2Service?: any;
sv2AuthorityService?: any;
redisMessagingService?: any;
} = {}) => {
return new AppController(
@@ -26,7 +26,7 @@ describe('AppController', () => {
overrides.userAgentReportService ?? {
getReport: jest.fn().mockResolvedValue([]),
},
overrides.stratumV2Service ?? {
overrides.sv2AuthorityService ?? {
getPoolAuthorityPublicKey: jest.fn().mockResolvedValue({
publicKey: '',
configured: false,
+3 -3
View File
@@ -10,9 +10,9 @@ import { ClientStatisticsService } from './ORM/client-statistics/client-statisti
import { ClientService } from './ORM/client/client.service';
import { BitcoinRpcService } from './services/bitcoin-rpc.service';
import { UserAgentReportView } from './ORM/_views/user-agent-report/user-agent-report.view';
import { StratumV2Service } from './services/stratum-v2.service';
import { ShareAccountingService } from './ORM/share-accounting/share-accounting.service';
import { RedisMessagingService } from './services/redis-messaging.service';
import { Sv2AuthorityService } from './services/sv2-authority.service';
import { normalizePayoutMode } from './types/payout-mode';
@Controller()
@@ -30,7 +30,7 @@ export class AppController {
private readonly bitcoinRpcService: BitcoinRpcService,
private readonly addressSettingsService: AddressSettingsService,
private readonly userAgentReportService: UserAgentReportService,
private readonly stratumV2Service: StratumV2Service,
private readonly sv2AuthorityService: Sv2AuthorityService,
private readonly shareAccountingService: ShareAccountingService,
private readonly redisMessagingService: RedisMessagingService
) { }
@@ -78,7 +78,7 @@ export class AppController {
const [blockData, highScores, poolAuthority, userAgentReport] = await Promise.all([
withInfoTimeout('found blocks', this.blocksService.getFoundBlocks(), staleInfo?.blockData ?? []),
withInfoTimeout('high scores', this.addressSettingsService.getHighScores(), staleInfo?.highScores ?? []),
withInfoTimeout('SV2 authority', this.stratumV2Service.getPoolAuthorityPublicKey(), {
withInfoTimeout('SV2 authority', this.sv2AuthorityService.getPoolAuthorityPublicKey(), {
publicKey: staleInfo?.sv2?.poolAuthorityPublicKey ?? '',
configured: staleInfo?.sv2?.authorityKeyConfigured ?? false
}),
+2
View File
@@ -33,6 +33,7 @@ import { StratumV1Service } from './services/stratum-v1.service';
import { Sv2JobDeclarationRegistryService } from './services/sv2-job-declaration-registry.service';
import { Sv2JobDeclarationService } from './services/sv2-job-declaration.service';
import { Sv2TemplateDistributionService } from './services/sv2-template-distribution.service';
import { Sv2AuthorityService } from './services/sv2-authority.service';
import { StratumV2Service } from './services/stratum-v2.service';
import { TelegramService } from './services/telegram.service';
import { TemplateProviderService } from './services/template-provider.service';
@@ -94,6 +95,7 @@ const ORMModules = [
Sv2JobDeclarationRegistryService,
Sv2JobDeclarationService,
Sv2TemplateDistributionService,
Sv2AuthorityService,
CustomWorkService,
DatumService,
BTCPayService,
+8
View File
@@ -18,6 +18,10 @@ 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 { AcceptedShare10mLookupIndexes1781403000000 } from './ORM/_migrations/AcceptedShare10mLookupIndexes1781403000000';
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 +66,10 @@ export const databaseMigrations = [
ShareRollupStoragePolicy1781305000000,
AcceptedShareHighScores1781309000000,
UserAgentReportNonzeroHashrate1781313000000,
PoolSummaryContinuousAggregate1781400000000,
PoolSummaryRefreshWindow1781401000000,
AcceptedShareHighScoreNumericRetention1781402000000,
AcceptedShare10mLookupIndexes1781403000000,
];
export function createDatabaseOptions(env: NodeJS.ProcessEnv): TypeOrmModuleOptions & DataSourceOptions {
+35 -10
View File
@@ -7,8 +7,21 @@ import { readFileSync, watch } from 'fs';
import * as path from 'path';
import * as ecc from 'tiny-secp256k1';
import { ApiModule } from './api.module';
import { AppModule } from './app.module';
const DEFAULT_API_CONNECTION_TIMEOUT_MS = 15_000;
const DEFAULT_API_KEEP_ALIVE_TIMEOUT_MS = 5_000;
const DEFAULT_API_REQUEST_TIMEOUT_MS = 15_000;
const DEFAULT_API_HEADERS_TIMEOUT_MS = 10_000;
const DEFAULT_API_MAX_REQUESTS_PER_SOCKET = 100;
const DEFAULT_API_MAX_CONNECTIONS_PER_WORKER = 500;
function readPositiveInt(name: string, fallback: number): number {
const value = Number(process.env[name]);
return Number.isInteger(value) && value > 0 ? value : fallback;
}
async function bootstrap() {
if (process.env.API_PORT == null) {
console.error('It appears your environment is not configured, create and populate an .env file.');
@@ -22,17 +35,23 @@ async function bootstrap() {
const keyPath = path.join(currentDirectory, 'secrets', 'key.pem');
const certPath = path.join(currentDirectory, 'secrets', 'cert.pem');
let options: any = {};
let options: any = serveApi
? {
connectionTimeout: readPositiveInt('API_CONNECTION_TIMEOUT_MS', DEFAULT_API_CONNECTION_TIMEOUT_MS),
keepAliveTimeout: readPositiveInt('API_KEEP_ALIVE_TIMEOUT_MS', DEFAULT_API_KEEP_ALIVE_TIMEOUT_MS),
requestTimeout: readPositiveInt('API_REQUEST_TIMEOUT_MS', DEFAULT_API_REQUEST_TIMEOUT_MS),
maxRequestsPerSocket: readPositiveInt('API_MAX_REQUESTS_PER_SOCKET', DEFAULT_API_MAX_REQUESTS_PER_SOCKET),
}
: {};
if (secure) {
options = {
https: {
key: readFileSync(keyPath),
cert: readFileSync(certPath),
}
options.https = {
key: readFileSync(keyPath),
cert: readFileSync(certPath),
};
}
const app = await NestFactory.create<NestFastifyApplication>(AppModule, new FastifyAdapter(options));
const appModule = process.env.API_ONLY === 'true' ? ApiModule : AppModule;
const app = await NestFactory.create<NestFastifyApplication>(appModule, new FastifyAdapter(options));
app.setGlobalPrefix('api');
app.useGlobalPipes(
new ValidationPipe({
@@ -54,7 +73,7 @@ async function bootstrap() {
});
app.enableCors();
useContainer(app.select(AppModule), { fallbackOnErrors: true });
useContainer(app.select(appModule), { fallbackOnErrors: true });
// Taproot
bitcoinjs.initEccLib(ecc);
@@ -73,11 +92,17 @@ async function bootstrap() {
console.log(`API listening on ${address}`);
});
const server: any = app.getHttpServer();
server.headersTimeout = readPositiveInt('API_HEADERS_TIMEOUT_MS', DEFAULT_API_HEADERS_TIMEOUT_MS);
server.maxConnections = readPositiveInt(
'API_MAX_CONNECTIONS_PER_WORKER',
DEFAULT_API_MAX_CONNECTIONS_PER_WORKER,
);
server.dropMaxConnection = true;
// --- Live-reload TLS certs/keys when they change on disk ---
if (secure) {
// Fastify's underlying Node https server
const server: any = app.getHttpServer();
// Guard: only HTTPS servers expose setSecureContext
if (typeof server?.setSecureContext === 'function') {
let reloadTimer: NodeJS.Timeout | null = null;
+146 -36
View File
@@ -169,6 +169,59 @@ describe('StratumV1Client', () => {
expect(socket.on).toHaveBeenCalled();
});
it('disconnects when an unterminated inbound message exceeds the configured limit', async () => {
(configService.get as jest.Mock).mockImplementation((key: string) => {
if (key === 'STRATUM_MAX_INBOUND_LINE_BYTES') return '16';
if (key === 'NETWORK') return 'testnet';
return null;
});
client = new StratumV1Client(
socket,
stratumV1JobsService,
bitcoinRpcService,
clientService,
notificationService,
blocksService,
configService,
addressSettings,
shareAccountingService as any,
redisMessagingService as any,
);
socketEmitter(Buffer.from('x'.repeat(17)));
await Promise.resolve();
expect(socket.destroy).toHaveBeenCalled();
expect((client as any).buffer).toBe('');
});
it('accepts multiple complete messages when each line is within the inbound limit', async () => {
(configService.get as jest.Mock).mockImplementation((key: string) => {
if (key === 'STRATUM_MAX_INBOUND_LINE_BYTES') return '64';
if (key === 'NETWORK') return 'testnet';
return null;
});
client = new StratumV1Client(
socket,
stratumV1JobsService,
bitcoinRpcService,
clientService,
notificationService,
blocksService,
configService,
addressSettings,
shareAccountingService as any,
redisMessagingService as any,
);
jest.spyOn(client as any, 'handleMessage').mockResolvedValue(undefined);
socketEmitter(Buffer.from('{"id":1}\n{"id":2}\n'));
await Promise.resolve();
expect((client as any).handleMessage).toHaveBeenCalledTimes(2);
expect(socket.destroy).not.toHaveBeenCalled();
});
it('should clean up socket state only once when destroyed repeatedly', async () => {
const timer = setInterval(() => undefined, 1000);
const removeListenerSpy = jest.spyOn(socket, 'removeListener');
@@ -409,24 +462,47 @@ describe('StratumV1Client', () => {
expect(socket.write).toHaveBeenCalledWith(`{"id":null,"method":"mining.set_difficulty","params":[512]}\n`, expect.any(Function));
});
it('should clamp suggested difficulty to the configured minimum', async () => {
(configService.get as jest.Mock).mockImplementation((key: string) => {
switch (key) {
case 'STRATUM_MIN_DIFFICULTY':
return '1';
case 'DEV_FEE_ADDRESS':
return 'tb1qumezefzdeqqwn5zfvgdrhxjzc5ylr39uhuxcz4';
case 'NETWORK':
return 'testnet';
}
return null;
});
it('should reject suggested difficulty below the protocol minimum', async () => {
jest.spyOn(socket, 'write').mockImplementation((data) => true);
emitMessage(`{"id":4,"method":"mining.suggest_difficulty","params":[0]}`);
await new Promise((r) => setTimeout(r, 1));
expect(socket.write).toHaveBeenCalledWith(`{"id":null,"method":"mining.set_difficulty","params":[1]}\n`, expect.any(Function));
expect(socket.write).toHaveBeenCalledWith(
expect.stringContaining('Suggest difficulty validation error'),
expect.any(Function),
);
expect((client as any).sessionDifficulty).toBe(100000);
});
it.each([12884901888, 1e303])(
'should reject excessive suggested difficulty %s',
async (suggestedDifficulty) => {
jest.spyOn(socket, 'write').mockImplementation((data) => true);
emitMessage(`{"id":4,"method":"mining.suggest_difficulty","params":[${suggestedDifficulty}]}`);
await new Promise((r) => setTimeout(r, 1));
expect(socket.write).toHaveBeenCalledWith(
expect.stringContaining('Suggest difficulty validation error'),
expect.any(Function),
);
expect((client as any).sessionDifficulty).toBe(100000);
},
);
it('should reject excessive password-provided starting difficulty', async () => {
jest.spyOn(socket, 'write').mockImplementation((data) => true);
emitMessage(MockRecording1.MINING_SUBSCRIBE);
emitMessage(`{"id":3,"method":"mining.authorize","params":["tb1qumezefzdeqqwn5zfvgdrhxjzc5ylr39uhuxcz4.worker","d=12884901888"]}`);
await new Promise((r) => setTimeout(r, 1));
expect(socket.write).toHaveBeenCalledWith(
expect.stringContaining('Authorization validation error'),
expect.any(Function),
);
expect((client as any).sessionDifficulty).toBe(100000);
});
it('should set difficulty', async () => {
@@ -482,7 +558,7 @@ describe('StratumV1Client', () => {
jest.spyOn(client as any, 'write').mockImplementation((data) => Promise.resolve(true));
emitMessage(MockRecording1.MINING_SUBSCRIBE);
emitMessage(`{"id": 4, "method": "mining.suggest_difficulty", "params": [0]}`);
emitMessage(`{"id": 4, "method": "mining.suggest_difficulty", "params": [1e-9]}`);
emitMessage(MockRecording1.MINING_AUTHORIZE);
await new Promise((r) => setTimeout(r, 100));
emitMessage(MockRecording1.MINING_SUBMIT);
@@ -510,7 +586,7 @@ describe('StratumV1Client', () => {
emitMessage(MockRecording1.MINING_SUBSCRIBE);
emitMessage(`{"id": 4, "method": "mining.suggest_difficulty", "params": [0]}`);
emitMessage(`{"id": 4, "method": "mining.suggest_difficulty", "params": [1e-9]}`);
emitMessage(MockRecording1.MINING_AUTHORIZE);
@@ -536,7 +612,7 @@ describe('StratumV1Client', () => {
sessionId: MockRecording1.EXTRA_NONCE,
jobId: '1',
jobTemplateId: '1',
creditedDifficulty: 0,
creditedDifficulty: 1e-9,
isBlockCandidate: false,
}));
});
@@ -547,7 +623,7 @@ describe('StratumV1Client', () => {
const fullBlockSpy = jest.spyOn(MiningJob.prototype, 'copyAndUpdateBlock');
emitMessage(MockRecording1.MINING_SUBSCRIBE);
emitMessage(`{"id": 4, "method": "mining.suggest_difficulty", "params": [0]}`);
emitMessage(`{"id": 4, "method": "mining.suggest_difficulty", "params": [1e-9]}`);
emitMessage(MockRecording1.MINING_AUTHORIZE);
await new Promise((r) => setTimeout(r, 100));
@@ -563,7 +639,7 @@ describe('StratumV1Client', () => {
jest.spyOn(client as any, 'write').mockImplementation((data) => Promise.resolve(true));
emitMessage(MockRecording1.MINING_SUBSCRIBE);
emitMessage(`{"id": 4, "method": "mining.suggest_difficulty", "params": [0]}`);
emitMessage(`{"id": 4, "method": "mining.suggest_difficulty", "params": [1e-9]}`);
emitMessage(MockRecording1.MINING_AUTHORIZE);
await new Promise((r) => setTimeout(r, 100));
@@ -597,7 +673,7 @@ describe('StratumV1Client', () => {
const clientUpdateIfHigherSpy = jest.spyOn(clientService as any, 'updateBestDifficultyIfHigher').mockResolvedValue({ affected: 1 });
emitMessage(MockRecording1.MINING_SUBSCRIBE);
emitMessage(`{"id": 4, "method": "mining.suggest_difficulty", "params": [0]}`);
emitMessage(`{"id": 4, "method": "mining.suggest_difficulty", "params": [1e-9]}`);
emitMessage(MockRecording1.MINING_AUTHORIZE);
await new Promise((r) => setTimeout(r, 100));
@@ -671,7 +747,7 @@ describe('StratumV1Client', () => {
jest.spyOn(client as any, 'write').mockImplementation((data) => Promise.resolve(true));
emitMessage(MockRecording1.MINING_SUBSCRIBE);
emitMessage(`{"id": 4, "method": "mining.suggest_difficulty", "params": [0]}`);
emitMessage(`{"id": 4, "method": "mining.suggest_difficulty", "params": [1e-9]}`);
emitMessage(MockRecording1.MINING_AUTHORIZE);
await new Promise((r) => setTimeout(r, 100));
@@ -693,7 +769,7 @@ describe('StratumV1Client', () => {
});
emitMessage(MockRecording1.MINING_SUBSCRIBE);
emitMessage(`{"id": 4, "method": "mining.suggest_difficulty", "params": [0]}`);
emitMessage(`{"id": 4, "method": "mining.suggest_difficulty", "params": [1e-9]}`);
emitMessage(MockRecording1.MINING_AUTHORIZE);
await new Promise((r) => setTimeout(r, 100));
@@ -712,7 +788,7 @@ describe('StratumV1Client', () => {
jest.spyOn(client as any, 'write').mockImplementation((data) => Promise.resolve(true));
emitMessage(MockRecording1.MINING_SUBSCRIBE);
emitMessage(`{"id": 4, "method": "mining.suggest_difficulty", "params": [0]}`);
emitMessage(`{"id": 4, "method": "mining.suggest_difficulty", "params": [1e-9]}`);
emitMessage(MockRecording1.MINING_AUTHORIZE);
await new Promise((r) => setTimeout(r, 100));
@@ -737,7 +813,7 @@ describe('StratumV1Client', () => {
const calculateDifficultySpy = jest.spyOn(client as any, 'calculateDifficulty');
emitMessage(MockRecording1.MINING_SUBSCRIBE);
emitMessage(`{"id": 4, "method": "mining.suggest_difficulty", "params": [0]}`);
emitMessage(`{"id": 4, "method": "mining.suggest_difficulty", "params": [1e-9]}`);
emitMessage(MockRecording1.MINING_AUTHORIZE);
await new Promise((r) => setTimeout(r, 100));
@@ -775,7 +851,7 @@ describe('StratumV1Client', () => {
});
emitMessage(MockRecording1.MINING_SUBSCRIBE);
emitMessage(`{"id": 4, "method": "mining.suggest_difficulty", "params": [0]}`);
emitMessage(`{"id": 4, "method": "mining.suggest_difficulty", "params": [1e-9]}`);
emitMessage(MockRecording1.MINING_AUTHORIZE);
await new Promise((r) => setTimeout(r, 100));
@@ -800,11 +876,12 @@ describe('StratumV1Client', () => {
jest.spyOn(socket, 'write').mockImplementation(() => true);
emitMessage(MockRecording1.MINING_SUBSCRIBE);
emitMessage(`{"id": 4, "method": "mining.suggest_difficulty", "params": [0]}`);
emitMessage(`{"id": 4, "method": "mining.suggest_difficulty", "params": [1e-9]}`);
emitMessage(MockRecording1.MINING_AUTHORIZE);
await new Promise((r) => setTimeout(r, 100));
expect((client as any).isDuplicateSubmission('old-tip-share')).toBe(false);
expect((client as any).rememberSubmission('old-tip-share', '1')).toBe(true);
const nextTip = {
...MockRecording1.BLOCK_TEMPLATE,
previousblockhash: '11'.repeat(32),
@@ -817,7 +894,9 @@ describe('StratumV1Client', () => {
expect((client as any).isDuplicateSubmission('old-tip-share')).toBe(true);
});
it('should bound duplicate tracking and expire entries by TTL', () => {
it('should preserve live accepted-share dedup entries when capacity is reached', () => {
const getSubmissionContext = jest.spyOn(stratumV1JobsService, 'getSubmissionContext')
.mockReturnValue({} as any);
(configService.get as jest.Mock).mockImplementation((key: string) => {
if (key === 'STRATUM_SUBMISSION_DEDUP_TTL_MS') return '1000';
if (key === 'STRATUM_SUBMISSION_DEDUP_MAX_ENTRIES') return '2';
@@ -826,13 +905,18 @@ describe('StratumV1Client', () => {
});
expect((client as any).isDuplicateSubmission('one')).toBe(false);
expect((client as any).rememberSubmission('one', '1')).toBe(true);
expect((client as any).isDuplicateSubmission('two')).toBe(false);
expect((client as any).isDuplicateSubmission('three')).toBe(false);
expect([...(client as any).miningSubmissionHashes.keys()]).toEqual(['two', 'three']);
expect((client as any).rememberSubmission('two', '1')).toBe(true);
expect((client as any).rememberSubmission('three', '1')).toBe(false);
expect([...(client as any).miningSubmissionHashes.keys()]).toEqual(['one', 'two']);
jest.advanceTimersByTime(1001);
expect((client as any).isDuplicateSubmission('two')).toBe(true);
getSubmissionContext.mockReturnValue(null);
expect((client as any).isDuplicateSubmission('two')).toBe(false);
expect([...(client as any).miningSubmissionHashes.keys()]).toEqual(['two']);
expect((client as any).rememberSubmission('three', '1')).toBe(true);
expect([...(client as any).miningSubmissionHashes.keys()]).toEqual(['three']);
});
it('should hash and reject stale non-block shares without accounting them', async () => {
@@ -846,7 +930,7 @@ describe('StratumV1Client', () => {
const buildHeaderSpy = jest.spyOn(MiningJob.prototype, 'buildHeaderBuffer');
emitMessage(MockRecording1.MINING_SUBSCRIBE);
emitMessage(`{"id": 4, "method": "mining.suggest_difficulty", "params": [0]}`);
emitMessage(`{"id": 4, "method": "mining.suggest_difficulty", "params": [1e-9]}`);
emitMessage(MockRecording1.MINING_AUTHORIZE);
await new Promise((r) => setTimeout(r, 100));
@@ -876,7 +960,7 @@ describe('StratumV1Client', () => {
});
emitMessage(MockRecording1.MINING_SUBSCRIBE);
emitMessage(`{"id": 4, "method": "mining.suggest_difficulty", "params": [0]}`);
emitMessage(`{"id": 4, "method": "mining.suggest_difficulty", "params": [1e-9]}`);
emitMessage(MockRecording1.MINING_AUTHORIZE);
await new Promise((r) => setTimeout(r, 100));
@@ -913,7 +997,7 @@ describe('StratumV1Client', () => {
});
emitMessage(MockRecording1.MINING_SUBSCRIBE);
emitMessage(`{"id": 4, "method": "mining.suggest_difficulty", "params": [0]}`);
emitMessage(`{"id": 4, "method": "mining.suggest_difficulty", "params": [1e-9]}`);
emitMessage(MockRecording1.MINING_AUTHORIZE);
await new Promise((r) => setTimeout(r, 100));
@@ -939,7 +1023,7 @@ describe('StratumV1Client', () => {
const hashSpy = jest.spyOn(MiningSubmitMessage.prototype, 'hash');
emitMessage(MockRecording1.MINING_SUBSCRIBE);
emitMessage(`{"id": 4, "method": "mining.suggest_difficulty", "params": [0]}`);
emitMessage(`{"id": 4, "method": "mining.suggest_difficulty", "params": [1e-9]}`);
emitMessage(MockRecording1.MINING_AUTHORIZE);
await new Promise((r) => setTimeout(r, 100));
@@ -956,7 +1040,7 @@ describe('StratumV1Client', () => {
jest.spyOn(client as any, 'write').mockImplementation((data) => Promise.resolve(true));
emitMessage(MockRecording1.MINING_SUBSCRIBE);
emitMessage(`{"id": 4, "method": "mining.suggest_difficulty", "params": [0]}`);
emitMessage(`{"id": 4, "method": "mining.suggest_difficulty", "params": [1e-9]}`);
emitMessage(MockRecording1.MINING_AUTHORIZE);
await new Promise((r) => setTimeout(r, 100));
stratumV1JobsService.blocks = {};
@@ -981,13 +1065,39 @@ describe('StratumV1Client', () => {
expect((client as any).write).lastCalledWith(`{"id":5,"result":null,"error":[23,"Difficulty too low",""]}\n`);
expect(await clientService.connectedClientCount()).toBe(1);
expect(shareAccountingService.recordAcceptedShare).not.toHaveBeenCalled();
expect((client as any).miningSubmissionHashes.size).toBe(0);
});
it.each([
{ label: 'before the advertised job time', offsetSeconds: -1 },
{ label: 'more than two hours in the future', offsetSeconds: (2 * 60 * 60) + 1 },
])('rejects ntime $label before hashing or accounting', async ({ offsetSeconds }) => {
jest.spyOn(client as any, 'write').mockResolvedValue(true);
const calculateDifficultySpy = jest.spyOn(client as any, 'calculateDifficulty');
emitMessage(MockRecording1.MINING_SUBSCRIBE);
emitMessage(`{"id": 4, "method": "mining.suggest_difficulty", "params": [1e-9]}`);
emitMessage(MockRecording1.MINING_AUTHORIZE);
await new Promise((resolve) => setTimeout(resolve, 100));
const submission = JSON.parse(MockRecording1.MINING_SUBMIT);
const baseTime = parseInt(MockRecording1.TIME, 16);
submission.params[3] = (baseTime + offsetSeconds).toString(16).padStart(8, '0');
emitMessage(JSON.stringify(submission));
await new Promise((resolve) => setTimeout(resolve, 100));
expect((client as any).write).lastCalledWith(
`{"id":5,"result":null,"error":[20,"Invalid ntime",""]}\n`,
);
expect(calculateDifficultySpy).not.toHaveBeenCalled();
expect(shareAccountingService.recordAcceptedShare).not.toHaveBeenCalled();
});
it('rejects version rolling masks outside the negotiated BIP320 range', async () => {
jest.spyOn(client as any, 'write').mockImplementation(() => Promise.resolve(true));
emitMessage(MockRecording1.MINING_SUBSCRIBE);
emitMessage(`{"id": 4, "method": "mining.suggest_difficulty", "params": [0]}`);
emitMessage(`{"id": 4, "method": "mining.suggest_difficulty", "params": [1e-9]}`);
emitMessage(MockRecording1.MINING_AUTHORIZE);
await new Promise((resolve) => setTimeout(resolve, 100));
@@ -1114,7 +1224,7 @@ describe('StratumV1Client', () => {
jest.spyOn(addressSettings, 'resetBestDifficultyAndShares').mockResolvedValue(undefined);
emitMessage(MockRecording1.MINING_SUBSCRIBE);
emitMessage(`{"id": 4, "method": "mining.suggest_difficulty", "params": [0]}`);
emitMessage(`{"id": 4, "method": "mining.suggest_difficulty", "params": [1e-9]}`);
emitMessage(MockRecording1.MINING_AUTHORIZE);
await new Promise((r) => setTimeout(r, 100));
+73 -29
View File
@@ -46,6 +46,8 @@ const DEFAULT_CLIENT_HASHRATE_PERSIST_INTERVAL_MS = 60 * 1000;
const DEFAULT_SUBMISSION_DEDUP_TTL_MS = 5 * 60 * 1000;
const DEFAULT_SUBMISSION_DEDUP_MAX_ENTRIES = 10_000;
const DEFAULT_MAX_SOCKET_BUFFER_BYTES = 256 * 1024;
const DEFAULT_MAX_INBOUND_LINE_BYTES = 64 * 1024;
const MAX_NTIME_FUTURE_SECONDS = 2 * 60 * 60;
const VERSION_ROLLING_MASK = 0x1fffe000;
export interface MiningJobBroadcastResult {
@@ -55,6 +57,11 @@ export interface MiningJobBroadcastResult {
preStaged?: boolean;
}
interface SubmissionDedupEntry {
expiresAt: number;
jobId: string;
}
export class StratumV1Client {
private static blockedUserAgentLogState = new Map<string, { nextLogAt: number, suppressed: number }>();
private static validationErrorLogState = new Map<string, { nextLogAt: number, suppressed: number, sample: string }>();
@@ -90,8 +97,9 @@ export class StratumV1Client {
private lastHashRatePersistedAt = 0;
private readonly network: bitcoinjs.Network;
private readonly maxSocketBufferBytes: number;
private readonly maxInboundLineBytes: number;
private miningSubmissionHashes = new Map<string, number>();
private miningSubmissionHashes = new Map<string, SubmissionDedupEntry>();
constructor(
public readonly socket: Socket,
@@ -115,6 +123,7 @@ export class StratumV1Client {
this.socket.on('data', this.socketDataHandler);
this.network = this.getNetwork();
this.maxSocketBufferBytes = this.readMaxSocketBufferBytes();
this.maxInboundLineBytes = this.readMaxInboundLineBytes();
}
@@ -150,9 +159,15 @@ export class StratumV1Client {
return;
}
this.buffer += data.toString();
const lines = this.buffer.split('\n');
this.buffer = lines.pop() || ''; // Save the last part of the data (incomplete line) to the buffer
const lines = `${this.buffer}${data.toString()}`.split('\n');
const incompleteLine = lines.pop() || '';
if (Buffer.byteLength(incompleteLine) > this.maxInboundLineBytes
|| lines.some(line => Buffer.byteLength(line) > this.maxInboundLineBytes)) {
this.buffer = '';
this.closeSocket();
return;
}
this.buffer = incompleteLine;
for (const m of lines.filter(l => l.length > 0)) {
if (this.connectionClosed || this.socket.destroyed || this.socket.writableEnded) {
@@ -449,8 +464,11 @@ export class StratumV1Client {
}
this.backgroundWork.push(
setInterval(async () => {
await this.checkDifficulty();
setInterval(() => {
void this.checkDifficulty().catch((error) => {
console.error('Stratum difficulty check failed; closing client connection', error);
this.closeSocket();
});
}, 60 * 1000)
);
@@ -679,6 +697,15 @@ export class StratumV1Client {
);
return false;
}
const maximumNtime = Math.floor(Date.now() / 1000) + MAX_NTIME_FUTURE_SECONDS;
if (timestamp < jobTemplate.block.timestamp || timestamp > maximumNtime) {
await this.writeSubmissionError(
submission,
eStratumErrorCode.OtherUnknown,
'Invalid ntime',
);
return false;
}
// The optional BIP310 field contains replacement bits, not an XOR
// delta. A legacy five-field submission leaves the advertised version
// unchanged, including any bits already set inside the rolling mask.
@@ -757,6 +784,15 @@ export class StratumV1Client {
const creditedDifficulty = isBlockCandidate && !meetsSessionTarget
? Math.min(this.sessionDifficulty, submissionDifficulty)
: this.sessionDifficulty;
if (!this.rememberSubmission(submissionHash, job.jobId)) {
await this.writeSubmissionError(
submission,
eStratumErrorCode.OtherUnknown,
'Submission dedup capacity exceeded',
);
this.closeSocket();
return false;
}
let blockSubmissionResult: string = null;
if (status === 'stale') {
@@ -1052,33 +1088,31 @@ export class StratumV1Client {
private isDuplicateSubmission(submissionHash: string): boolean {
const now = Date.now();
const existingExpiry = this.miningSubmissionHashes.get(submissionHash);
if (existingExpiry != null && existingExpiry > now) {
return true;
}
if (existingExpiry != null) {
this.miningSubmissionHashes.delete(submissionHash);
}
this.pruneExpiredSubmissions(now);
return this.miningSubmissionHashes.has(submissionHash);
}
for (const [hash, expiresAt] of this.miningSubmissionHashes) {
if (expiresAt <= now) {
private rememberSubmission(submissionHash: string, jobId: string): boolean {
const now = Date.now();
this.pruneExpiredSubmissions(now);
if (this.miningSubmissionHashes.has(submissionHash)
|| this.miningSubmissionHashes.size >= this.getSubmissionDedupMaxEntries()) {
return false;
}
this.miningSubmissionHashes.set(submissionHash, {
expiresAt: now + this.getSubmissionDedupTtlMs(),
jobId,
});
return true;
}
private pruneExpiredSubmissions(now: number): void {
for (const [hash, entry] of this.miningSubmissionHashes) {
if (entry.expiresAt <= now
&& this.stratumV1JobsService.getSubmissionContext(entry.jobId) == null) {
this.miningSubmissionHashes.delete(hash);
}
}
this.miningSubmissionHashes.set(
submissionHash,
now + this.getSubmissionDedupTtlMs(),
);
const maxEntries = this.getSubmissionDedupMaxEntries();
while (this.miningSubmissionHashes.size > maxEntries) {
const oldestHash = this.miningSubmissionHashes.keys().next().value;
if (oldestHash == null) {
break;
}
this.miningSubmissionHashes.delete(oldestHash);
}
return false;
}
private getSubmissionDedupTtlMs(): number {
@@ -1111,6 +1145,16 @@ export class StratumV1Client {
: DEFAULT_MAX_SOCKET_BUFFER_BYTES;
}
private readMaxInboundLineBytes(): number {
const configured = Number(
this.configService.get<string>('STRATUM_MAX_INBOUND_LINE_BYTES')
?? process.env.STRATUM_MAX_INBOUND_LINE_BYTES,
);
return Number.isSafeInteger(configured) && configured > 0
? configured
: DEFAULT_MAX_INBOUND_LINE_BYTES;
}
private getValidationErrorSignature(errors: ValidationError[]): string {
if (errors.length === 0) {
return 'unknown';
+42 -3
View File
@@ -32,7 +32,30 @@ describe('StratumV1ClientStatistics', () => {
await statistics.addShares(client, 64);
}
expect(statistics.hashRate).toBeGreaterThan(0);
expect(statistics.hashRate).toBeCloseTo((64 * 4294967296) / 62);
});
it('keeps exactly the configured number of share samples', async () => {
for (let i = 0; i < 31; i++) {
jest.setSystemTime(new Date(Date.parse('2026-05-06T12:00:00Z') + (i * 31000)));
await statistics.addShares(client, 64);
}
expect((statistics as any).submissionCache).toHaveLength(30);
expect((statistics as any).submissionCacheDifficultySum).toBe(30 * 64);
});
it('excludes pre-window work and corrects share-terminated sampling bias', async () => {
for (let i = 0; i < 30; i++) {
jest.setSystemTime(new Date(Date.parse('2026-05-06T12:00:00Z') + (i * 31000)));
await statistics.addShares(client, 64);
}
const elapsedSeconds = 29 * 31;
const expectedUnbiasedDifficulty = 28 * 64;
expect(statistics.hashRate).toBeCloseTo(
(expectedUnbiasedDifficulty * 4294967296) / elapsedSeconds,
);
});
it('should not suggest a difficulty change before enough time or shares have passed', () => {
@@ -51,7 +74,23 @@ describe('StratumV1ClientStatistics', () => {
await statistics.addShares(client, 64);
}
expect(statistics.getSuggestedDifficulty(64)).toBe(2048);
expect(statistics.getSuggestedDifficulty(64)).toBe(1024);
});
it('does not retarget when accepted shares have a zero-duration sample window', async () => {
for (let i = 0; i < 5; i++) {
await statistics.addShares(client, 64);
}
expect(() => statistics.getSuggestedDifficulty(64)).not.toThrow();
expect(statistics.getSuggestedDifficulty(64)).toBeNull();
});
it('returns a finite power-of-two difficulty for values above 32-bit range', () => {
const result = (statistics as any).nearestPowerOfTwo(2 ** 40 + 1);
expect(result).toBe(2 ** 40);
expect(Number.isFinite(result)).toBe(true);
});
it('should decrease difficulty for slow submissions', async () => {
@@ -60,7 +99,7 @@ describe('StratumV1ClientStatistics', () => {
await statistics.addShares(client, 64);
}
expect(statistics.getSuggestedDifficulty(128)).toBe(16);
expect(statistics.getSuggestedDifficulty(128)).toBe(8);
});
it('should not suggest a difficulty below the configured minimum', () => {
+46 -28
View File
@@ -16,9 +16,9 @@ export class StratumV1ClientStatistics {
}
public async addShares(_client: ClientEntity, targetDifficulty: number) {
var date = new Date();
const date = new Date();
if (this.submissionCache.length > CACHE_SIZE) {
if (this.submissionCache.length >= CACHE_SIZE) {
this.submissionCacheDifficultySum -= this.submissionCache[0].difficulty;
this.submissionCache.shift();
}
@@ -28,9 +28,12 @@ export class StratumV1ClientStatistics {
});
this.submissionCacheDifficultySum += targetDifficulty;
const time = new Date().getTime() - this.submissionCache[0].time.getTime();
if(time > 60000 && this.submissionCache.length > 2) {
this.hashRate = (this.submissionCacheDifficultySum * 4294967296) / (time / 1000);
const elapsedSeconds = (date.getTime() - this.submissionCache[0].time.getTime()) / 1000;
if (elapsedSeconds > 60) {
const difficultyPerSecond = this.getDifficultyPerSecond(elapsedSeconds);
if (difficultyPerSecond != null) {
this.hashRate = difficultyPerSecond * 4294967296;
}
}
}
@@ -46,15 +49,16 @@ export class StratumV1ClientStatistics {
}
}
const sum = this.submissionCache.reduce((pre, cur) => {
pre += cur.difficulty;
return pre;
}, 0);
const diffSeconds = (this.submissionCache[this.submissionCache.length - 1].time.getTime() - this.submissionCache[0].time.getTime()) / 1000;
const difficultyPerSecond = sum / diffSeconds;
const difficultyPerSecond = this.getDifficultyPerSecond(diffSeconds);
if (difficultyPerSecond == null) {
return null;
}
const targetDifficulty = difficultyPerSecond * this.targetSubmitShareEveryNSeconds;
if (!Number.isFinite(targetDifficulty) || targetDifficulty <= 0) {
return null;
}
if ((clientDifficulty * 2) < targetDifficulty || (clientDifficulty / 2) > targetDifficulty) {
return this.nearestPowerOfTwo(targetDifficulty)
@@ -63,27 +67,41 @@ export class StratumV1ClientStatistics {
return null;
}
private nearestPowerOfTwo(val): number {
if (val === 0) {
/**
* Estimate work rate from a share-terminated sample window.
*
* A cache of N shares contains N - 1 observed inter-share intervals. The
* first share's work predates the window and must not be counted. Because
* the window closes on a share arrival, the reciprocal elapsed time also
* has the usual finite-sample Poisson bias; multiplying by
* (intervalCount - 1) / intervalCount removes it.
*/
private getDifficultyPerSecond(elapsedSeconds: number): number | null {
const sampleCount = this.submissionCache.length;
if (sampleCount <= 2 || !Number.isFinite(elapsedSeconds) || elapsedSeconds <= 0) {
return null;
}
if (val < this.minDifficulty) {
const intervalCount = sampleCount - 1;
const observedDifficulty = this.submissionCacheDifficultySum
- this.submissionCache[0].difficulty;
const unbiasedDifficulty = observedDifficulty * (intervalCount - 1) / intervalCount;
const difficultyPerSecond = unbiasedDifficulty / elapsedSeconds;
return Number.isFinite(difficultyPerSecond) && difficultyPerSecond > 0
? difficultyPerSecond
: null;
}
private nearestPowerOfTwo(val: number): number {
if (!Number.isFinite(val) || val <= 0) {
return null;
}
if (val <= this.minDifficulty) {
return this.minDifficulty;
}
let x = val | (val >> 1);
x = x | (x >> 2);
x = x | (x >> 4);
x = x | (x >> 8);
x = x | (x >> 16);
x = x | (x >> 32);
const res = x - (x >> 1);
if (res == 0 && val * 100 < this.minDifficulty) {
return this.minDifficulty;
}
if (res == 0) {
return this.nearestPowerOfTwo(val * 100) / 100;
}
return res;
const result = 2 ** Math.floor(Math.log2(val));
return Number.isFinite(result) ? Math.max(this.minDifficulty, result) : null;
}
}
+3 -2
View File
@@ -991,7 +991,7 @@ describe('StratumV2Client extended channels', () => {
.toBe(parseInt(activation.bits, 16));
});
it('validates fixed/rolling versions, required bits, and the advertised minimum nTime', async () => {
it('validates BIP320 version rolling, required bits, and the advertised minimum nTime', async () => {
const { client } = await createClient();
const now = process.hrtime.bigint();
const context = {
@@ -1007,11 +1007,12 @@ describe('StratumV2Client extended channels', () => {
};
expect((client as any).isSubmissionHeaderValid(context, 0x20000004, 100)).toBe(true);
expect((client as any).isSubmissionHeaderValid(context, 0x20002004, 100)).toBe(false);
expect((client as any).isSubmissionHeaderValid(context, 0x20002004, 100)).toBe(true);
expect((client as any).isSubmissionHeaderValid(context, 0x20000004, 99)).toBe(false);
expect((client as any).isSubmissionHeaderValid(context, 0x20000004, 101)).toBe(true);
expect((client as any).isSubmissionHeaderValid(context, 0x20000004, 0xffffffff)).toBe(true);
expect((client as any).isSubmissionHeaderValid(context, 0x20000004, 0x1_0000_0000)).toBe(false);
expect((client as any).isSubmissionHeaderValid(context, 0x00000004, 100)).toBe(false);
(client as any).versionRollingEnabled = true;
expect((client as any).isSubmissionHeaderValid(context, 0x20002004, 100)).toBe(true);
+5 -5
View File
@@ -2088,11 +2088,11 @@ export class StratumV2Client {
): boolean {
const submitted = submittedVersion >>> 0;
const base = headerContext.baseVersion >>> 0;
if (this.versionRollingEnabled) {
if (((submitted ^ base) & BIP320_CONSENSUS_VERSION_MASK) !== 0) {
return false;
}
} else if (submitted !== base) {
// Some SV2 clients roll BIP320 general-purpose version bits even when
// they did not explicitly require version rolling during setup. Those
// bits are non-consensus signalling space; reject only changes outside
// that mask and continue enforcing required bits below.
if (((submitted ^ base) & BIP320_CONSENSUS_VERSION_MASK) !== 0) {
return false;
}
if ((submitted & headerContext.requiredVersionBits) !== headerContext.requiredVersionBits) {
@@ -1,9 +1,10 @@
import { Expose, Transform } from 'class-transformer';
import { ArrayMaxSize, ArrayMinSize, IsArray, IsNumber, IsOptional, IsString, MaxLength } from 'class-validator';
import { ArrayMaxSize, ArrayMinSize, IsArray, IsNumber, IsOptional, IsPositive, IsString, Max, MaxLength } from 'class-validator';
import { eRequestMethod } from '../enums/eRequestMethod';
import { IsBitcoinAddress } from '../validators/bitcoin-address.validator';
import { StratumBaseMessage } from './StratumBaseMessage';
import { MAX_STRATUM_DIFFICULTY } from './SuggestDifficultyMessage';
export class AuthorizationMessage extends StratumBaseMessage {
@@ -35,6 +36,8 @@ export class AuthorizationMessage extends StratumBaseMessage {
@Expose()
@IsNumber()
@IsPositive()
@Max(MAX_STRATUM_DIFFICULTY)
@Transform(({ value, key, obj, type }) => {
const password: string | null = obj.params[1];
if (password?.includes('d=')) {
@@ -69,4 +72,4 @@ export class AuthorizationMessage extends StratumBaseMessage {
result: true
};
}
}
}
@@ -1,10 +1,12 @@
import { Expose, Transform } from 'class-transformer';
import { ArrayMaxSize, ArrayMinSize, IsArray, IsNumber } from 'class-validator';
import { ArrayMaxSize, ArrayMinSize, IsArray, IsNumber, IsPositive, Max } from 'class-validator';
import { eRequestMethod } from '../enums/eRequestMethod';
import { eResponseMethod } from '../enums/eResponseMethod';
import { StratumBaseMessage } from './StratumBaseMessage';
export const MAX_STRATUM_DIFFICULTY = 2 ** 32;
export class SuggestDifficulty extends StratumBaseMessage {
@IsArray()
@ArrayMinSize(1)
@@ -16,6 +18,8 @@ export class SuggestDifficulty extends StratumBaseMessage {
@Expose()
@IsNumber()
@IsPositive()
@Max(MAX_STRATUM_DIFFICULTY)
@Transform(({ value, key, obj, type }) => {
return Number(obj.params[0]);
})
@@ -37,4 +41,3 @@ export class SuggestDifficulty extends StratumBaseMessage {
}
}
}
+26 -1
View File
@@ -120,6 +120,19 @@ describe('RedisMessagingService', () => {
consoleSpy.mockRestore();
});
it('uses an owner token when locking JSON cache refreshes', async () => {
await service.connect();
await expect(service.tryAcquireJsonCacheLock('summary', 'owner-a', 1000)).resolves.toBe(true);
await expect(service.tryAcquireJsonCacheLock('summary', 'owner-b', 1000)).resolves.toBe(false);
await service.releaseJsonCacheLock('summary', 'owner-b');
await expect(service.tryAcquireJsonCacheLock('summary', 'owner-b', 1000)).resolves.toBe(false);
await service.releaseJsonCacheLock('summary', 'owner-a');
await expect(service.tryAcquireJsonCacheLock('summary', 'owner-b', 1000)).resolves.toBe(true);
});
it('stores payout variants and same-height reorg templates under distinct tip keys', async () => {
await service.connect();
const firstHash = '11'.repeat(32);
@@ -493,7 +506,10 @@ function createRedisClient() {
connect: jest.fn().mockResolvedValue(undefined),
quit: jest.fn().mockResolvedValue(undefined),
on: jest.fn(),
set: jest.fn((key: string, value: string) => {
set: jest.fn((key: string, value: string, options?: { NX?: boolean }) => {
if (options?.NX && store.has(key)) {
return Promise.resolve(null);
}
store.set(key, value);
return Promise.resolve('OK');
}),
@@ -512,6 +528,15 @@ function createRedisClient() {
});
return Promise.resolve(deleted);
}),
eval: jest.fn((_script: string, options: { keys: string[], arguments: string[] }) => {
const [key] = options.keys;
const [owner] = options.arguments;
if (store.get(key) !== owner) {
return Promise.resolve(0);
}
store.delete(key);
return Promise.resolve(1);
}),
sAdd: jest.fn((key: string, value: string) => {
const set = sets.get(key) ?? new Set<string>();
set.add(value);
+36
View File
@@ -32,6 +32,7 @@ const blockTemplateHeightPointerKey = (height: number, payoutMode: PayoutMode) =
const blockTemplateLatestPointerKey = (payoutMode: PayoutMode) =>
`${BLOCK_TEMPLATE_LATEST_KEY}:${payoutMode}`;
const jsonCacheKey = (key: string) => `json-cache:${key}`;
const jsonCacheLockKey = (key: string) => `json-cache-lock:${key}`;
export interface BlockTemplateUpdate {
schemaVersion: 1;
@@ -524,6 +525,41 @@ export class RedisMessagingService implements OnModuleInit, OnModuleDestroy {
);
}
public async tryAcquireJsonCacheLock(key: string, owner: string, ttlMs: number): Promise<boolean | null> {
if (!await this.ensureConnected()) {
return null;
}
if (owner.length === 0 || ttlMs <= 0) {
return false;
}
const result = await this.publisher.set(
jsonCacheLockKey(key),
owner,
{ NX: true, PX: Math.max(1, Math.ceil(ttlMs)) },
);
return result === 'OK';
}
public async releaseJsonCacheLock(key: string, owner: string): Promise<void> {
if (!await this.ensureConnected() || owner.length === 0) {
return;
}
await this.publisher.eval(
`
if redis.call('GET', KEYS[1]) == ARGV[1] then
return redis.call('DEL', KEYS[1])
end
return 0
`,
{
keys: [jsonCacheLockKey(key)],
arguments: [owner],
},
);
}
private async readBlockTemplatePointer(value: string | null): Promise<IBlockTemplate | null> {
if (value == null) {
return null;
+37
View File
@@ -0,0 +1,37 @@
import { Injectable } from '@nestjs/common';
import { ConfigService } from '@nestjs/config';
import * as crypto from 'crypto';
import { encodeSv2AuthorityPublicKey } from '../models/sv2/sv2-authority-key';
import { xOnlyPubKeyFromPriv } from '../models/sv2/sv2-noise';
@Injectable()
export class Sv2AuthorityService {
private authorityPublicKey: string | null = null;
private authorityKeyConfigured = false;
constructor(
private readonly configService: ConfigService,
) {}
public async getPoolAuthorityPublicKey(): Promise<{ publicKey: string; configured: boolean }> {
if (this.authorityPublicKey != null) {
return {
publicKey: this.authorityPublicKey,
configured: this.authorityKeyConfigured,
};
}
const configuredAuthorityKey = this.configService.get<string>('SV2_AUTHORITY_PRIVKEY');
this.authorityKeyConfigured = configuredAuthorityKey?.length === 64;
const authorityPrivKey = configuredAuthorityKey?.length === 64
? Buffer.from(configuredAuthorityKey, 'hex')
: crypto.randomBytes(32);
this.authorityPublicKey = encodeSv2AuthorityPublicKey(xOnlyPubKeyFromPriv(authorityPrivKey));
return {
publicKey: this.authorityPublicKey,
configured: this.authorityKeyConfigured,
};
}
}
+13 -1
View File
@@ -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 () => {