16 Commits
32 changed files with 1585 additions and 153 deletions
+14 -1
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.
@@ -102,7 +111,7 @@ STRATUM_SUBMISSION_DEDUP_MAX_ENTRIES=10000
# Per-client lifecycle bounds. Obsolete queued jobs are coalesced; clients are
# disconnected instead of dropping still-recoverable candidates or buffering forever.
SV2_JOB_RETENTION_MS=300000
SV2_MAX_RETAINED_JOBS_PER_CHANNEL=16
SV2_MAX_RETAINED_JOBS_PER_CHANNEL=64
SV2_MAX_QUEUED_JOB_OPERATIONS=4
SV2_MAX_SOCKET_BUFFER_BYTES=262144
SV2_SOCKET_WRITE_TIMEOUT_MS=2000
@@ -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
+1 -1
View File
@@ -125,7 +125,7 @@ Standard and extended candidates retain exact header/body reconstruction, and
late network-target candidates remain recoverable without crediting stale shares.
Pending SV2 canonical jobs are coalesced per client, while a new-tip activation
is moved ahead of any not-yet-started canonical work for that tip. The finite
defaults are 16 retained jobs per channel, four queued operations, 256 KiB of
defaults are 64 retained jobs per channel, four queued operations, 256 KiB of
outstanding socket writes, and a two-second write-callback deadline. A client is
disconnected if `SV2_MAX_RETAINED_JOBS_PER_CHANNEL`,
`SV2_MAX_QUEUED_JOB_OPERATIONS`, `SV2_MAX_SOCKET_BUFFER_BYTES`, or
+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;
}
}
+38 -2
View File
@@ -927,6 +927,41 @@ describe('StratumV2Client extended channels', () => {
)).toBe(true);
});
it('keeps enough retained SV2 jobs for repeated fast activation and full-template followup cycles', async () => {
const { client, jobTemplate } = await createClient();
await (client as any).handleOpenExtendedMiningChannel(serializeOpenExtendedMiningChannel({
requestId: 1,
userIdentity: 'tb1qumezefzdeqqwn5zfvgdrhxjzc5ylr39uhuxcz4.worker',
nominalHashRate: 0,
maxTarget: Buffer.alloc(32, 0xff),
minExtranonceSize: 8,
}));
let currentTemplate = jobTemplate;
for (let index = 0; index < 20; index++) {
const activation = {
...createActivationTemplate(currentTemplate),
previousblockhash: (index + 1).toString(16).padStart(64, '0'),
mintime: parseInt(MockRecording1.TIME, 16) + index + 1,
curtime: parseInt(MockRecording1.TIME, 16) + index + 1,
};
await client.enqueueWorkActivation(activation as any);
currentTemplate = {
...createCanonicalFollowup(currentTemplate, activation),
blockData: {
...createCanonicalFollowup(currentTemplate, activation).blockData,
id: `canonical-followup-${index}`,
},
};
await client.enqueueCanonicalJob(currentTemplate);
}
expect((client as any).socket.destroyed).toBe(false);
const channel = (client as any).channels.get(1);
expect(channel.extendedJobs.size).toBeLessThanOrEqual((client as any).maxRetainedJobsPerChannel);
expect(channel.stagedFutureJobId).toBeDefined();
});
it('activates a staged standard future and rejects incompatible activation versions', async () => {
const { client, sentFrames, jobTemplate } = await createClient();
await (client as any).handleOpenStandardMiningChannel(serializeOpenStandardMiningChannel({
@@ -956,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 = {
@@ -972,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);
+6 -6
View File
@@ -78,7 +78,7 @@ const DEFAULT_DIFFICULTY_CHECK_INTERVAL_MS = 60 * 1000;
const DEFAULT_CLIENT_HASHRATE_PERSIST_INTERVAL_MS = 60 * 1000;
const FIXED_STANDARD_EXTRANONCE2 = '0000000000000000';
const DEFAULT_SV2_JOB_RETENTION_MS = 5 * 60 * 1000;
const DEFAULT_SV2_MAX_RETAINED_JOBS_PER_CHANNEL = 16;
const DEFAULT_SV2_MAX_RETAINED_JOBS_PER_CHANNEL = 64;
const DEFAULT_SV2_MAX_QUEUED_JOB_OPERATIONS = 4;
const DEFAULT_SV2_MAX_SOCKET_BUFFER_BYTES = 256 * 1024;
const DEFAULT_SV2_SOCKET_WRITE_TIMEOUT_MS = 2 * 1000;
@@ -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 () => {