mirror of
https://github.com/benjamin-wilson/public-pool.git
synced 2026-09-29 09:05:06 -07:00
Compare commits
15
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
789b529bd5 | ||
|
|
5889e3af5b | ||
|
|
f10cef61ce | ||
|
|
07a2b173d0 | ||
|
|
4c56839c02 | ||
|
|
1cdf3fbf7e | ||
|
|
d1ebc55871 | ||
|
|
7d96c06bbf | ||
|
|
dd39c6edcf | ||
|
|
86ad296189 | ||
|
|
f49e7b5b8f | ||
|
|
3cb4cbfbdf | ||
|
|
f31cdef7f0 | ||
|
|
943dd60eef | ||
|
|
4150d1e425 |
@@ -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
|
||||
|
||||
@@ -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}
|
||||
|
||||
@@ -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}
|
||||
|
||||
@@ -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}
|
||||
|
||||
Executable
+33
@@ -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,
|
||||
|
||||
@@ -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 {}
|
||||
@@ -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,
|
||||
|
||||
@@ -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
|
||||
}),
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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
@@ -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;
|
||||
|
||||
@@ -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));
|
||||
|
||||
|
||||
@@ -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';
|
||||
|
||||
@@ -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', () => {
|
||||
|
||||
@@ -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;
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -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);
|
||||
|
||||
@@ -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 {
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -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);
|
||||
|
||||
@@ -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;
|
||||
|
||||
@@ -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,
|
||||
};
|
||||
}
|
||||
}
|
||||
@@ -76,7 +76,7 @@ describe('TimescaleDB and Redis integration', () => {
|
||||
const aggregates = await dataSource.query(`
|
||||
SELECT view_name
|
||||
FROM timescaledb_information.continuous_aggregates
|
||||
WHERE view_name IN ('accepted_share_10m', 'accepted_share_1h', 'accepted_share_1d', 'accepted_share_block_10m')
|
||||
WHERE view_name IN ('accepted_share_10m', 'accepted_share_1h', 'accepted_share_1d', 'accepted_share_block_10m', 'accepted_share_pool_10m')
|
||||
ORDER BY view_name
|
||||
`);
|
||||
expect(aggregates.map(row => row.view_name)).toEqual([
|
||||
@@ -84,6 +84,7 @@ describe('TimescaleDB and Redis integration', () => {
|
||||
'accepted_share_1d',
|
||||
'accepted_share_1h',
|
||||
'accepted_share_block_10m',
|
||||
'accepted_share_pool_10m',
|
||||
]);
|
||||
|
||||
const legacyTables = await dataSource.query(`
|
||||
@@ -228,6 +229,7 @@ describe('TimescaleDB and Redis integration', () => {
|
||||
|
||||
await dataSource.query(`CALL refresh_continuous_aggregate('accepted_share_10m', NULL, NULL)`);
|
||||
await dataSource.query(`CALL refresh_continuous_aggregate('accepted_share_block_10m', NULL, NULL)`);
|
||||
await dataSource.query(`CALL refresh_continuous_aggregate('accepted_share_pool_10m', NULL, NULL)`);
|
||||
const aggregateRows = await dataSource.query(`
|
||||
SELECT "shares"::float AS shares, "acceptedCount"::int AS "acceptedCount"
|
||||
FROM accepted_share_10m
|
||||
@@ -250,6 +252,16 @@ describe('TimescaleDB and Redis integration', () => {
|
||||
expect(blockAggregateRows).toEqual(expect.arrayContaining([
|
||||
expect.objectContaining({ shares: 96, acceptedCount: 2, networkDifficulty: 100000 }),
|
||||
]));
|
||||
|
||||
const poolAggregateRows = await dataSource.query(`
|
||||
SELECT "shares"::float AS shares, "acceptedCount"::int AS "acceptedCount"
|
||||
FROM accepted_share_pool_10m
|
||||
WHERE "payoutMode" = $1
|
||||
`, ['solo']);
|
||||
|
||||
expect(poolAggregateRows).toEqual(expect.arrayContaining([
|
||||
expect.objectContaining({ shares: 96, acceptedCount: 2 }),
|
||||
]));
|
||||
});
|
||||
|
||||
it('should retain daily and all-time best share from completed rollup buckets', async () => {
|
||||
|
||||
Reference in New Issue
Block a user