7 Commits
Author SHA1 Message Date
Ben 789b529bd5 Fix one-off hashrate calculation 2026-09-15 21:54:06 -04:00
Ben 5889e3af5b constrain suggest difficulty 2026-08-27 22:41:48 -04:00
Ben f10cef61ce Automate deployment of renewed TLS certificates 2026-08-18 08:52:25 -04:00
Ben 07a2b173d0 Harden Stratum V1 share validation 2026-08-05 16:34:43 -04:00
Ben 4c56839c02 Prevent accounting cache refresh stampedes 2026-08-04 23:47:41 -04:00
Ben 1cdf3fbf7e Cap concurrent API connections 2026-08-04 19:24:29 -04:00
Ben d1ebc55871 Bound public API connection lifetimes 2026-08-04 19:20:48 -04:00
16 changed files with 714 additions and 116 deletions
+13
View File
@@ -32,6 +32,14 @@ API_PORT=3334
API_BIND_HOST=127.0.0.1
API_PUBLIC_PORT=3334
API_WORKERS=4
# Bound abandoned or slow API sockets. These settings apply only to the
# HTTP(S) listener and do not affect Stratum connections.
API_CONNECTION_TIMEOUT_MS=15000
API_KEEP_ALIVE_TIMEOUT_MS=5000
API_REQUEST_TIMEOUT_MS=15000
API_HEADERS_TIMEOUT_MS=10000
API_MAX_REQUESTS_PER_SOCKET=100
API_MAX_CONNECTIONS_PER_WORKER=500
# Docker json-file log rotation. Applies when using the compose files.
DOCKER_LOG_MAX_SIZE=100m
@@ -64,6 +72,7 @@ STRATUM_TCP_KEEPALIVE_INITIAL_DELAY_MS=60000
# Disconnect slow/non-reading SV1 clients before their per-socket write queue
# can grow without bound. Includes bytes already buffered by Node.js.
STRATUM_MAX_SOCKET_BUFFER_BYTES=262144
STRATUM_MAX_INBOUND_LINE_BYTES=65536
# Immediately publish a consensus-valid, subsidy-only solo job before the full
# transaction template. SV2 uses the same activation to select its pre-staged
# future job; the full job then follows on the active tip.
@@ -162,6 +171,10 @@ SHARE_ACCOUNTING_FLUSH_INTERVAL_MS=25
SHARE_ACCOUNTING_MAX_QUEUE_SIZE=50000
SHARE_ACCOUNTING_SUMMARY_CACHE_MS=2500
SHARE_ACCOUNTING_SUMMARY_CACHE_MAX=10000
SHARE_ACCOUNTING_REDIS_SUMMARY_CACHE_MS=300000
SHARE_ACCOUNTING_REDIS_SUMMARY_STALE_CACHE_MS=3600000
SHARE_ACCOUNTING_REDIS_SUMMARY_LOCK_MS=30000
SHARE_ACCOUNTING_REDIS_SUMMARY_LOCK_WAIT_MS=5000
CLIENT_HASHRATE_PERSIST_INTERVAL_MS=60000
SHARE_ROLLUP_ENABLED=true
SHARE_ROLLUP_INTERVAL_MS=60000
+10
View File
@@ -74,6 +74,12 @@ services:
API_PORT: ${API_PORT:-3334}
API_SECURE: ${API_SECURE:-false}
API_WORKERS: ${API_WORKERS:-4}
API_CONNECTION_TIMEOUT_MS: ${API_CONNECTION_TIMEOUT_MS:-15000}
API_KEEP_ALIVE_TIMEOUT_MS: ${API_KEEP_ALIVE_TIMEOUT_MS:-5000}
API_REQUEST_TIMEOUT_MS: ${API_REQUEST_TIMEOUT_MS:-15000}
API_HEADERS_TIMEOUT_MS: ${API_HEADERS_TIMEOUT_MS:-10000}
API_MAX_REQUESTS_PER_SOCKET: ${API_MAX_REQUESTS_PER_SOCKET:-100}
API_MAX_CONNECTIONS_PER_WORKER: ${API_MAX_CONNECTIONS_PER_WORKER:-500}
PM2_ENABLED: ${PM2_ENABLED:-true}
STRATUM_WORKERS: ${STRATUM_WORKERS:-auto}
STRATUM_PORTS: ${STRATUM_PORTS:-3333,3332,3331,3330}
@@ -100,6 +106,10 @@ services:
SHARE_ACCOUNTING_MAX_QUEUE_SIZE: ${SHARE_ACCOUNTING_MAX_QUEUE_SIZE:-50000}
SHARE_ACCOUNTING_SUMMARY_CACHE_MS: ${SHARE_ACCOUNTING_SUMMARY_CACHE_MS:-2500}
SHARE_ACCOUNTING_SUMMARY_CACHE_MAX: ${SHARE_ACCOUNTING_SUMMARY_CACHE_MAX:-10000}
SHARE_ACCOUNTING_REDIS_SUMMARY_CACHE_MS: ${SHARE_ACCOUNTING_REDIS_SUMMARY_CACHE_MS:-300000}
SHARE_ACCOUNTING_REDIS_SUMMARY_STALE_CACHE_MS: ${SHARE_ACCOUNTING_REDIS_SUMMARY_STALE_CACHE_MS:-3600000}
SHARE_ACCOUNTING_REDIS_SUMMARY_LOCK_MS: ${SHARE_ACCOUNTING_REDIS_SUMMARY_LOCK_MS:-30000}
SHARE_ACCOUNTING_REDIS_SUMMARY_LOCK_WAIT_MS: ${SHARE_ACCOUNTING_REDIS_SUMMARY_LOCK_WAIT_MS:-5000}
SHARE_ROLLUP_ENABLED: ${SHARE_ROLLUP_ENABLED:-true}
SHARE_ROLLUP_INTERVAL_MS: ${SHARE_ROLLUP_INTERVAL_MS:-60000}
SHARE_ROLLUP_SAFETY_LAG_SECONDS: ${SHARE_ROLLUP_SAFETY_LAG_SECONDS:-30}
+10
View File
@@ -104,6 +104,16 @@ services:
API_PORT: ${API_PORT:-3334}
API_SECURE: ${API_SECURE:-false}
API_WORKERS: ${API_WORKERS:-4}
API_CONNECTION_TIMEOUT_MS: ${API_CONNECTION_TIMEOUT_MS:-15000}
API_KEEP_ALIVE_TIMEOUT_MS: ${API_KEEP_ALIVE_TIMEOUT_MS:-5000}
API_REQUEST_TIMEOUT_MS: ${API_REQUEST_TIMEOUT_MS:-15000}
API_HEADERS_TIMEOUT_MS: ${API_HEADERS_TIMEOUT_MS:-10000}
API_MAX_REQUESTS_PER_SOCKET: ${API_MAX_REQUESTS_PER_SOCKET:-100}
API_MAX_CONNECTIONS_PER_WORKER: ${API_MAX_CONNECTIONS_PER_WORKER:-500}
SHARE_ACCOUNTING_REDIS_SUMMARY_CACHE_MS: ${SHARE_ACCOUNTING_REDIS_SUMMARY_CACHE_MS:-300000}
SHARE_ACCOUNTING_REDIS_SUMMARY_STALE_CACHE_MS: ${SHARE_ACCOUNTING_REDIS_SUMMARY_STALE_CACHE_MS:-3600000}
SHARE_ACCOUNTING_REDIS_SUMMARY_LOCK_MS: ${SHARE_ACCOUNTING_REDIS_SUMMARY_LOCK_MS:-30000}
SHARE_ACCOUNTING_REDIS_SUMMARY_LOCK_WAIT_MS: ${SHARE_ACCOUNTING_REDIS_SUMMARY_LOCK_WAIT_MS:-5000}
PM2_ENABLED: ${PM2_ENABLED:-true}
STRATUM_WORKERS: ${STRATUM_WORKERS:-auto}
STRATUM_PORTS: ${STRATUM_PORTS:-3333,3332,3331,3330}
+10
View File
@@ -108,6 +108,16 @@ services:
DB_DATABASE: public_pool_mainnet
REDIS_URL: redis://redis:6379
API_WORKERS: ${API_WORKERS:-4}
API_CONNECTION_TIMEOUT_MS: ${API_CONNECTION_TIMEOUT_MS:-15000}
API_KEEP_ALIVE_TIMEOUT_MS: ${API_KEEP_ALIVE_TIMEOUT_MS:-5000}
API_REQUEST_TIMEOUT_MS: ${API_REQUEST_TIMEOUT_MS:-15000}
API_HEADERS_TIMEOUT_MS: ${API_HEADERS_TIMEOUT_MS:-10000}
API_MAX_REQUESTS_PER_SOCKET: ${API_MAX_REQUESTS_PER_SOCKET:-100}
API_MAX_CONNECTIONS_PER_WORKER: ${API_MAX_CONNECTIONS_PER_WORKER:-500}
SHARE_ACCOUNTING_REDIS_SUMMARY_CACHE_MS: ${SHARE_ACCOUNTING_REDIS_SUMMARY_CACHE_MS:-300000}
SHARE_ACCOUNTING_REDIS_SUMMARY_STALE_CACHE_MS: ${SHARE_ACCOUNTING_REDIS_SUMMARY_STALE_CACHE_MS:-3600000}
SHARE_ACCOUNTING_REDIS_SUMMARY_LOCK_MS: ${SHARE_ACCOUNTING_REDIS_SUMMARY_LOCK_MS:-30000}
SHARE_ACCOUNTING_REDIS_SUMMARY_LOCK_WAIT_MS: ${SHARE_ACCOUNTING_REDIS_SUMMARY_LOCK_WAIT_MS:-5000}
PM2_ENABLED: "true"
STRATUM_WORKERS: ${STRATUM_WORKERS:-2}
STRATUM_MIN_DIFFICULTY: ${STRATUM_MIN_DIFFICULTY:-1}
+33
View File
@@ -0,0 +1,33 @@
#!/bin/sh
set -eu
lineage=${RENEWED_LINEAGE:-/etc/letsencrypt/live/public-pool.io}
secrets_dir=${PUBLIC_POOL_SECRETS_DIR:-/home/ben/public-pool-timescaledb-test/secrets}
container=${PUBLIC_POOL_CONTAINER:-public-pool}
owner=${PUBLIC_POOL_CERT_OWNER:-ben}
group=${PUBLIC_POOL_CERT_GROUP:-ben}
cert_source="$lineage/fullchain.pem"
key_source="$lineage/privkey.pem"
test -r "$cert_source"
test -r "$key_source"
openssl x509 -in "$cert_source" -noout >/dev/null
openssl pkey -in "$key_source" -noout >/dev/null
cert_temp=$(mktemp "$secrets_dir/.cert.pem.XXXXXX")
key_temp=$(mktemp "$secrets_dir/.key.pem.XXXXXX")
trap 'rm -f "$cert_temp" "$key_temp"' EXIT HUP INT TERM
cat "$cert_source" > "$cert_temp"
cat "$key_source" > "$key_temp"
chown "$owner:$group" "$cert_temp" "$key_temp"
chmod 0644 "$cert_temp"
chmod 0600 "$key_temp"
mv -f "$cert_temp" "$secrets_dir/cert.pem"
mv -f "$key_temp" "$secrets_dir/key.pem"
# HTTPS can reload its secure context, but the Stratum TLS listeners currently
# read their certificate only when they start.
docker restart --time 30 "$container" >/dev/null
@@ -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),
@@ -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;
@@ -103,6 +107,7 @@ export class ShareAccountingService implements OnModuleInit, OnModuleDestroy {
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);
@@ -110,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);
@@ -597,25 +608,175 @@ export class ShareAccountingService implements OnModuleInit, OnModuleDestroy {
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.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, summary, this.redisSummaryCacheMs)
?.setJsonCache(redisCacheKey, entry, this.redisSummaryStaleCacheMs)
.catch(error => {
console.error(`Share accounting summary cache write failed: ${error.message}`);
});
}
return summary;
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> {
@@ -924,3 +1085,9 @@ interface SummaryCacheEntry {
expiresAt: number;
value: ShareAccountingSummary;
}
interface RedisSummaryCacheEntry {
schemaVersion: 1;
refreshedAtMs: number;
value: ShareAccountingSummary;
}
+29 -6
View File
@@ -10,6 +10,18 @@ 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.');
@@ -23,13 +35,18 @@ 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: {
options.https = {
key: readFileSync(keyPath),
cert: readFileSync(certPath),
}
};
}
@@ -75,11 +92,17 @@ async function bootstrap() {
console.log(`API listening on ${address}`);
});
const server: any = app.getHttpServer();
server.headersTimeout = readPositiveInt('API_HEADERS_TIMEOUT_MS', DEFAULT_API_HEADERS_TIMEOUT_MS);
server.maxConnections = readPositiveInt(
'API_MAX_CONNECTIONS_PER_WORKER',
DEFAULT_API_MAX_CONNECTIONS_PER_WORKER,
);
server.dropMaxConnection = true;
// --- Live-reload TLS certs/keys when they change on disk ---
if (secure) {
// Fastify's underlying Node https server
const server: any = app.getHttpServer();
// Guard: only HTTPS servers expose setSecureContext
if (typeof server?.setSecureContext === 'function') {
let reloadTimer: NodeJS.Timeout | null = null;
+146 -36
View File
@@ -169,6 +169,59 @@ describe('StratumV1Client', () => {
expect(socket.on).toHaveBeenCalled();
});
it('disconnects when an unterminated inbound message exceeds the configured limit', async () => {
(configService.get as jest.Mock).mockImplementation((key: string) => {
if (key === 'STRATUM_MAX_INBOUND_LINE_BYTES') return '16';
if (key === 'NETWORK') return 'testnet';
return null;
});
client = new StratumV1Client(
socket,
stratumV1JobsService,
bitcoinRpcService,
clientService,
notificationService,
blocksService,
configService,
addressSettings,
shareAccountingService as any,
redisMessagingService as any,
);
socketEmitter(Buffer.from('x'.repeat(17)));
await Promise.resolve();
expect(socket.destroy).toHaveBeenCalled();
expect((client as any).buffer).toBe('');
});
it('accepts multiple complete messages when each line is within the inbound limit', async () => {
(configService.get as jest.Mock).mockImplementation((key: string) => {
if (key === 'STRATUM_MAX_INBOUND_LINE_BYTES') return '64';
if (key === 'NETWORK') return 'testnet';
return null;
});
client = new StratumV1Client(
socket,
stratumV1JobsService,
bitcoinRpcService,
clientService,
notificationService,
blocksService,
configService,
addressSettings,
shareAccountingService as any,
redisMessagingService as any,
);
jest.spyOn(client as any, 'handleMessage').mockResolvedValue(undefined);
socketEmitter(Buffer.from('{"id":1}\n{"id":2}\n'));
await Promise.resolve();
expect((client as any).handleMessage).toHaveBeenCalledTimes(2);
expect(socket.destroy).not.toHaveBeenCalled();
});
it('should clean up socket state only once when destroyed repeatedly', async () => {
const timer = setInterval(() => undefined, 1000);
const removeListenerSpy = jest.spyOn(socket, 'removeListener');
@@ -409,24 +462,47 @@ describe('StratumV1Client', () => {
expect(socket.write).toHaveBeenCalledWith(`{"id":null,"method":"mining.set_difficulty","params":[512]}\n`, expect.any(Function));
});
it('should clamp suggested difficulty to the configured minimum', async () => {
(configService.get as jest.Mock).mockImplementation((key: string) => {
switch (key) {
case 'STRATUM_MIN_DIFFICULTY':
return '1';
case 'DEV_FEE_ADDRESS':
return 'tb1qumezefzdeqqwn5zfvgdrhxjzc5ylr39uhuxcz4';
case 'NETWORK':
return 'testnet';
}
return null;
});
it('should reject suggested difficulty below the protocol minimum', async () => {
jest.spyOn(socket, 'write').mockImplementation((data) => true);
emitMessage(`{"id":4,"method":"mining.suggest_difficulty","params":[0]}`);
await new Promise((r) => setTimeout(r, 1));
expect(socket.write).toHaveBeenCalledWith(`{"id":null,"method":"mining.set_difficulty","params":[1]}\n`, expect.any(Function));
expect(socket.write).toHaveBeenCalledWith(
expect.stringContaining('Suggest difficulty validation error'),
expect.any(Function),
);
expect((client as any).sessionDifficulty).toBe(100000);
});
it.each([12884901888, 1e303])(
'should reject excessive suggested difficulty %s',
async (suggestedDifficulty) => {
jest.spyOn(socket, 'write').mockImplementation((data) => true);
emitMessage(`{"id":4,"method":"mining.suggest_difficulty","params":[${suggestedDifficulty}]}`);
await new Promise((r) => setTimeout(r, 1));
expect(socket.write).toHaveBeenCalledWith(
expect.stringContaining('Suggest difficulty validation error'),
expect.any(Function),
);
expect((client as any).sessionDifficulty).toBe(100000);
},
);
it('should reject excessive password-provided starting difficulty', async () => {
jest.spyOn(socket, 'write').mockImplementation((data) => true);
emitMessage(MockRecording1.MINING_SUBSCRIBE);
emitMessage(`{"id":3,"method":"mining.authorize","params":["tb1qumezefzdeqqwn5zfvgdrhxjzc5ylr39uhuxcz4.worker","d=12884901888"]}`);
await new Promise((r) => setTimeout(r, 1));
expect(socket.write).toHaveBeenCalledWith(
expect.stringContaining('Authorization validation error'),
expect.any(Function),
);
expect((client as any).sessionDifficulty).toBe(100000);
});
it('should set difficulty', async () => {
@@ -482,7 +558,7 @@ describe('StratumV1Client', () => {
jest.spyOn(client as any, 'write').mockImplementation((data) => Promise.resolve(true));
emitMessage(MockRecording1.MINING_SUBSCRIBE);
emitMessage(`{"id": 4, "method": "mining.suggest_difficulty", "params": [0]}`);
emitMessage(`{"id": 4, "method": "mining.suggest_difficulty", "params": [1e-9]}`);
emitMessage(MockRecording1.MINING_AUTHORIZE);
await new Promise((r) => setTimeout(r, 100));
emitMessage(MockRecording1.MINING_SUBMIT);
@@ -510,7 +586,7 @@ describe('StratumV1Client', () => {
emitMessage(MockRecording1.MINING_SUBSCRIBE);
emitMessage(`{"id": 4, "method": "mining.suggest_difficulty", "params": [0]}`);
emitMessage(`{"id": 4, "method": "mining.suggest_difficulty", "params": [1e-9]}`);
emitMessage(MockRecording1.MINING_AUTHORIZE);
@@ -536,7 +612,7 @@ describe('StratumV1Client', () => {
sessionId: MockRecording1.EXTRA_NONCE,
jobId: '1',
jobTemplateId: '1',
creditedDifficulty: 0,
creditedDifficulty: 1e-9,
isBlockCandidate: false,
}));
});
@@ -547,7 +623,7 @@ describe('StratumV1Client', () => {
const fullBlockSpy = jest.spyOn(MiningJob.prototype, 'copyAndUpdateBlock');
emitMessage(MockRecording1.MINING_SUBSCRIBE);
emitMessage(`{"id": 4, "method": "mining.suggest_difficulty", "params": [0]}`);
emitMessage(`{"id": 4, "method": "mining.suggest_difficulty", "params": [1e-9]}`);
emitMessage(MockRecording1.MINING_AUTHORIZE);
await new Promise((r) => setTimeout(r, 100));
@@ -563,7 +639,7 @@ describe('StratumV1Client', () => {
jest.spyOn(client as any, 'write').mockImplementation((data) => Promise.resolve(true));
emitMessage(MockRecording1.MINING_SUBSCRIBE);
emitMessage(`{"id": 4, "method": "mining.suggest_difficulty", "params": [0]}`);
emitMessage(`{"id": 4, "method": "mining.suggest_difficulty", "params": [1e-9]}`);
emitMessage(MockRecording1.MINING_AUTHORIZE);
await new Promise((r) => setTimeout(r, 100));
@@ -597,7 +673,7 @@ describe('StratumV1Client', () => {
const clientUpdateIfHigherSpy = jest.spyOn(clientService as any, 'updateBestDifficultyIfHigher').mockResolvedValue({ affected: 1 });
emitMessage(MockRecording1.MINING_SUBSCRIBE);
emitMessage(`{"id": 4, "method": "mining.suggest_difficulty", "params": [0]}`);
emitMessage(`{"id": 4, "method": "mining.suggest_difficulty", "params": [1e-9]}`);
emitMessage(MockRecording1.MINING_AUTHORIZE);
await new Promise((r) => setTimeout(r, 100));
@@ -671,7 +747,7 @@ describe('StratumV1Client', () => {
jest.spyOn(client as any, 'write').mockImplementation((data) => Promise.resolve(true));
emitMessage(MockRecording1.MINING_SUBSCRIBE);
emitMessage(`{"id": 4, "method": "mining.suggest_difficulty", "params": [0]}`);
emitMessage(`{"id": 4, "method": "mining.suggest_difficulty", "params": [1e-9]}`);
emitMessage(MockRecording1.MINING_AUTHORIZE);
await new Promise((r) => setTimeout(r, 100));
@@ -693,7 +769,7 @@ describe('StratumV1Client', () => {
});
emitMessage(MockRecording1.MINING_SUBSCRIBE);
emitMessage(`{"id": 4, "method": "mining.suggest_difficulty", "params": [0]}`);
emitMessage(`{"id": 4, "method": "mining.suggest_difficulty", "params": [1e-9]}`);
emitMessage(MockRecording1.MINING_AUTHORIZE);
await new Promise((r) => setTimeout(r, 100));
@@ -712,7 +788,7 @@ describe('StratumV1Client', () => {
jest.spyOn(client as any, 'write').mockImplementation((data) => Promise.resolve(true));
emitMessage(MockRecording1.MINING_SUBSCRIBE);
emitMessage(`{"id": 4, "method": "mining.suggest_difficulty", "params": [0]}`);
emitMessage(`{"id": 4, "method": "mining.suggest_difficulty", "params": [1e-9]}`);
emitMessage(MockRecording1.MINING_AUTHORIZE);
await new Promise((r) => setTimeout(r, 100));
@@ -737,7 +813,7 @@ describe('StratumV1Client', () => {
const calculateDifficultySpy = jest.spyOn(client as any, 'calculateDifficulty');
emitMessage(MockRecording1.MINING_SUBSCRIBE);
emitMessage(`{"id": 4, "method": "mining.suggest_difficulty", "params": [0]}`);
emitMessage(`{"id": 4, "method": "mining.suggest_difficulty", "params": [1e-9]}`);
emitMessage(MockRecording1.MINING_AUTHORIZE);
await new Promise((r) => setTimeout(r, 100));
@@ -775,7 +851,7 @@ describe('StratumV1Client', () => {
});
emitMessage(MockRecording1.MINING_SUBSCRIBE);
emitMessage(`{"id": 4, "method": "mining.suggest_difficulty", "params": [0]}`);
emitMessage(`{"id": 4, "method": "mining.suggest_difficulty", "params": [1e-9]}`);
emitMessage(MockRecording1.MINING_AUTHORIZE);
await new Promise((r) => setTimeout(r, 100));
@@ -800,11 +876,12 @@ describe('StratumV1Client', () => {
jest.spyOn(socket, 'write').mockImplementation(() => true);
emitMessage(MockRecording1.MINING_SUBSCRIBE);
emitMessage(`{"id": 4, "method": "mining.suggest_difficulty", "params": [0]}`);
emitMessage(`{"id": 4, "method": "mining.suggest_difficulty", "params": [1e-9]}`);
emitMessage(MockRecording1.MINING_AUTHORIZE);
await new Promise((r) => setTimeout(r, 100));
expect((client as any).isDuplicateSubmission('old-tip-share')).toBe(false);
expect((client as any).rememberSubmission('old-tip-share', '1')).toBe(true);
const nextTip = {
...MockRecording1.BLOCK_TEMPLATE,
previousblockhash: '11'.repeat(32),
@@ -817,7 +894,9 @@ describe('StratumV1Client', () => {
expect((client as any).isDuplicateSubmission('old-tip-share')).toBe(true);
});
it('should bound duplicate tracking and expire entries by TTL', () => {
it('should preserve live accepted-share dedup entries when capacity is reached', () => {
const getSubmissionContext = jest.spyOn(stratumV1JobsService, 'getSubmissionContext')
.mockReturnValue({} as any);
(configService.get as jest.Mock).mockImplementation((key: string) => {
if (key === 'STRATUM_SUBMISSION_DEDUP_TTL_MS') return '1000';
if (key === 'STRATUM_SUBMISSION_DEDUP_MAX_ENTRIES') return '2';
@@ -826,13 +905,18 @@ describe('StratumV1Client', () => {
});
expect((client as any).isDuplicateSubmission('one')).toBe(false);
expect((client as any).rememberSubmission('one', '1')).toBe(true);
expect((client as any).isDuplicateSubmission('two')).toBe(false);
expect((client as any).isDuplicateSubmission('three')).toBe(false);
expect([...(client as any).miningSubmissionHashes.keys()]).toEqual(['two', 'three']);
expect((client as any).rememberSubmission('two', '1')).toBe(true);
expect((client as any).rememberSubmission('three', '1')).toBe(false);
expect([...(client as any).miningSubmissionHashes.keys()]).toEqual(['one', 'two']);
jest.advanceTimersByTime(1001);
expect((client as any).isDuplicateSubmission('two')).toBe(true);
getSubmissionContext.mockReturnValue(null);
expect((client as any).isDuplicateSubmission('two')).toBe(false);
expect([...(client as any).miningSubmissionHashes.keys()]).toEqual(['two']);
expect((client as any).rememberSubmission('three', '1')).toBe(true);
expect([...(client as any).miningSubmissionHashes.keys()]).toEqual(['three']);
});
it('should hash and reject stale non-block shares without accounting them', async () => {
@@ -846,7 +930,7 @@ describe('StratumV1Client', () => {
const buildHeaderSpy = jest.spyOn(MiningJob.prototype, 'buildHeaderBuffer');
emitMessage(MockRecording1.MINING_SUBSCRIBE);
emitMessage(`{"id": 4, "method": "mining.suggest_difficulty", "params": [0]}`);
emitMessage(`{"id": 4, "method": "mining.suggest_difficulty", "params": [1e-9]}`);
emitMessage(MockRecording1.MINING_AUTHORIZE);
await new Promise((r) => setTimeout(r, 100));
@@ -876,7 +960,7 @@ describe('StratumV1Client', () => {
});
emitMessage(MockRecording1.MINING_SUBSCRIBE);
emitMessage(`{"id": 4, "method": "mining.suggest_difficulty", "params": [0]}`);
emitMessage(`{"id": 4, "method": "mining.suggest_difficulty", "params": [1e-9]}`);
emitMessage(MockRecording1.MINING_AUTHORIZE);
await new Promise((r) => setTimeout(r, 100));
@@ -913,7 +997,7 @@ describe('StratumV1Client', () => {
});
emitMessage(MockRecording1.MINING_SUBSCRIBE);
emitMessage(`{"id": 4, "method": "mining.suggest_difficulty", "params": [0]}`);
emitMessage(`{"id": 4, "method": "mining.suggest_difficulty", "params": [1e-9]}`);
emitMessage(MockRecording1.MINING_AUTHORIZE);
await new Promise((r) => setTimeout(r, 100));
@@ -939,7 +1023,7 @@ describe('StratumV1Client', () => {
const hashSpy = jest.spyOn(MiningSubmitMessage.prototype, 'hash');
emitMessage(MockRecording1.MINING_SUBSCRIBE);
emitMessage(`{"id": 4, "method": "mining.suggest_difficulty", "params": [0]}`);
emitMessage(`{"id": 4, "method": "mining.suggest_difficulty", "params": [1e-9]}`);
emitMessage(MockRecording1.MINING_AUTHORIZE);
await new Promise((r) => setTimeout(r, 100));
@@ -956,7 +1040,7 @@ describe('StratumV1Client', () => {
jest.spyOn(client as any, 'write').mockImplementation((data) => Promise.resolve(true));
emitMessage(MockRecording1.MINING_SUBSCRIBE);
emitMessage(`{"id": 4, "method": "mining.suggest_difficulty", "params": [0]}`);
emitMessage(`{"id": 4, "method": "mining.suggest_difficulty", "params": [1e-9]}`);
emitMessage(MockRecording1.MINING_AUTHORIZE);
await new Promise((r) => setTimeout(r, 100));
stratumV1JobsService.blocks = {};
@@ -981,13 +1065,39 @@ describe('StratumV1Client', () => {
expect((client as any).write).lastCalledWith(`{"id":5,"result":null,"error":[23,"Difficulty too low",""]}\n`);
expect(await clientService.connectedClientCount()).toBe(1);
expect(shareAccountingService.recordAcceptedShare).not.toHaveBeenCalled();
expect((client as any).miningSubmissionHashes.size).toBe(0);
});
it.each([
{ label: 'before the advertised job time', offsetSeconds: -1 },
{ label: 'more than two hours in the future', offsetSeconds: (2 * 60 * 60) + 1 },
])('rejects ntime $label before hashing or accounting', async ({ offsetSeconds }) => {
jest.spyOn(client as any, 'write').mockResolvedValue(true);
const calculateDifficultySpy = jest.spyOn(client as any, 'calculateDifficulty');
emitMessage(MockRecording1.MINING_SUBSCRIBE);
emitMessage(`{"id": 4, "method": "mining.suggest_difficulty", "params": [1e-9]}`);
emitMessage(MockRecording1.MINING_AUTHORIZE);
await new Promise((resolve) => setTimeout(resolve, 100));
const submission = JSON.parse(MockRecording1.MINING_SUBMIT);
const baseTime = parseInt(MockRecording1.TIME, 16);
submission.params[3] = (baseTime + offsetSeconds).toString(16).padStart(8, '0');
emitMessage(JSON.stringify(submission));
await new Promise((resolve) => setTimeout(resolve, 100));
expect((client as any).write).lastCalledWith(
`{"id":5,"result":null,"error":[20,"Invalid ntime",""]}\n`,
);
expect(calculateDifficultySpy).not.toHaveBeenCalled();
expect(shareAccountingService.recordAcceptedShare).not.toHaveBeenCalled();
});
it('rejects version rolling masks outside the negotiated BIP320 range', async () => {
jest.spyOn(client as any, 'write').mockImplementation(() => Promise.resolve(true));
emitMessage(MockRecording1.MINING_SUBSCRIBE);
emitMessage(`{"id": 4, "method": "mining.suggest_difficulty", "params": [0]}`);
emitMessage(`{"id": 4, "method": "mining.suggest_difficulty", "params": [1e-9]}`);
emitMessage(MockRecording1.MINING_AUTHORIZE);
await new Promise((resolve) => setTimeout(resolve, 100));
@@ -1114,7 +1224,7 @@ describe('StratumV1Client', () => {
jest.spyOn(addressSettings, 'resetBestDifficultyAndShares').mockResolvedValue(undefined);
emitMessage(MockRecording1.MINING_SUBSCRIBE);
emitMessage(`{"id": 4, "method": "mining.suggest_difficulty", "params": [0]}`);
emitMessage(`{"id": 4, "method": "mining.suggest_difficulty", "params": [1e-9]}`);
emitMessage(MockRecording1.MINING_AUTHORIZE);
await new Promise((r) => setTimeout(r, 100));
+72 -28
View File
@@ -46,6 +46,8 @@ const DEFAULT_CLIENT_HASHRATE_PERSIST_INTERVAL_MS = 60 * 1000;
const DEFAULT_SUBMISSION_DEDUP_TTL_MS = 5 * 60 * 1000;
const DEFAULT_SUBMISSION_DEDUP_MAX_ENTRIES = 10_000;
const DEFAULT_MAX_SOCKET_BUFFER_BYTES = 256 * 1024;
const DEFAULT_MAX_INBOUND_LINE_BYTES = 64 * 1024;
const MAX_NTIME_FUTURE_SECONDS = 2 * 60 * 60;
const VERSION_ROLLING_MASK = 0x1fffe000;
export interface MiningJobBroadcastResult {
@@ -55,6 +57,11 @@ export interface MiningJobBroadcastResult {
preStaged?: boolean;
}
interface SubmissionDedupEntry {
expiresAt: number;
jobId: string;
}
export class StratumV1Client {
private static blockedUserAgentLogState = new Map<string, { nextLogAt: number, suppressed: number }>();
private static validationErrorLogState = new Map<string, { nextLogAt: number, suppressed: number, sample: string }>();
@@ -90,8 +97,9 @@ export class StratumV1Client {
private lastHashRatePersistedAt = 0;
private readonly network: bitcoinjs.Network;
private readonly maxSocketBufferBytes: number;
private readonly maxInboundLineBytes: number;
private miningSubmissionHashes = new Map<string, number>();
private miningSubmissionHashes = new Map<string, SubmissionDedupEntry>();
constructor(
public readonly socket: Socket,
@@ -115,6 +123,7 @@ export class StratumV1Client {
this.socket.on('data', this.socketDataHandler);
this.network = this.getNetwork();
this.maxSocketBufferBytes = this.readMaxSocketBufferBytes();
this.maxInboundLineBytes = this.readMaxInboundLineBytes();
}
@@ -150,9 +159,15 @@ export class StratumV1Client {
return;
}
this.buffer += data.toString();
const lines = this.buffer.split('\n');
this.buffer = lines.pop() || ''; // Save the last part of the data (incomplete line) to the buffer
const lines = `${this.buffer}${data.toString()}`.split('\n');
const incompleteLine = lines.pop() || '';
if (Buffer.byteLength(incompleteLine) > this.maxInboundLineBytes
|| lines.some(line => Buffer.byteLength(line) > this.maxInboundLineBytes)) {
this.buffer = '';
this.closeSocket();
return;
}
this.buffer = incompleteLine;
for (const m of lines.filter(l => l.length > 0)) {
if (this.connectionClosed || this.socket.destroyed || this.socket.writableEnded) {
@@ -449,8 +464,11 @@ export class StratumV1Client {
}
this.backgroundWork.push(
setInterval(async () => {
await this.checkDifficulty();
setInterval(() => {
void this.checkDifficulty().catch((error) => {
console.error('Stratum difficulty check failed; closing client connection', error);
this.closeSocket();
});
}, 60 * 1000)
);
@@ -679,6 +697,15 @@ export class StratumV1Client {
);
return false;
}
const maximumNtime = Math.floor(Date.now() / 1000) + MAX_NTIME_FUTURE_SECONDS;
if (timestamp < jobTemplate.block.timestamp || timestamp > maximumNtime) {
await this.writeSubmissionError(
submission,
eStratumErrorCode.OtherUnknown,
'Invalid ntime',
);
return false;
}
// The optional BIP310 field contains replacement bits, not an XOR
// delta. A legacy five-field submission leaves the advertised version
// unchanged, including any bits already set inside the rolling mask.
@@ -757,6 +784,15 @@ export class StratumV1Client {
const creditedDifficulty = isBlockCandidate && !meetsSessionTarget
? Math.min(this.sessionDifficulty, submissionDifficulty)
: this.sessionDifficulty;
if (!this.rememberSubmission(submissionHash, job.jobId)) {
await this.writeSubmissionError(
submission,
eStratumErrorCode.OtherUnknown,
'Submission dedup capacity exceeded',
);
this.closeSocket();
return false;
}
let blockSubmissionResult: string = null;
if (status === 'stale') {
@@ -1052,33 +1088,31 @@ export class StratumV1Client {
private isDuplicateSubmission(submissionHash: string): boolean {
const now = Date.now();
const existingExpiry = this.miningSubmissionHashes.get(submissionHash);
if (existingExpiry != null && existingExpiry > now) {
return true;
}
if (existingExpiry != null) {
this.miningSubmissionHashes.delete(submissionHash);
this.pruneExpiredSubmissions(now);
return this.miningSubmissionHashes.has(submissionHash);
}
for (const [hash, expiresAt] of this.miningSubmissionHashes) {
if (expiresAt <= now) {
private rememberSubmission(submissionHash: string, jobId: string): boolean {
const now = Date.now();
this.pruneExpiredSubmissions(now);
if (this.miningSubmissionHashes.has(submissionHash)
|| this.miningSubmissionHashes.size >= this.getSubmissionDedupMaxEntries()) {
return false;
}
this.miningSubmissionHashes.set(submissionHash, {
expiresAt: now + this.getSubmissionDedupTtlMs(),
jobId,
});
return true;
}
private pruneExpiredSubmissions(now: number): void {
for (const [hash, entry] of this.miningSubmissionHashes) {
if (entry.expiresAt <= now
&& this.stratumV1JobsService.getSubmissionContext(entry.jobId) == null) {
this.miningSubmissionHashes.delete(hash);
}
}
this.miningSubmissionHashes.set(
submissionHash,
now + this.getSubmissionDedupTtlMs(),
);
const maxEntries = this.getSubmissionDedupMaxEntries();
while (this.miningSubmissionHashes.size > maxEntries) {
const oldestHash = this.miningSubmissionHashes.keys().next().value;
if (oldestHash == null) {
break;
}
this.miningSubmissionHashes.delete(oldestHash);
}
return false;
}
private getSubmissionDedupTtlMs(): number {
@@ -1111,6 +1145,16 @@ export class StratumV1Client {
: DEFAULT_MAX_SOCKET_BUFFER_BYTES;
}
private readMaxInboundLineBytes(): number {
const configured = Number(
this.configService.get<string>('STRATUM_MAX_INBOUND_LINE_BYTES')
?? process.env.STRATUM_MAX_INBOUND_LINE_BYTES,
);
return Number.isSafeInteger(configured) && configured > 0
? configured
: DEFAULT_MAX_INBOUND_LINE_BYTES;
}
private getValidationErrorSignature(errors: ValidationError[]): string {
if (errors.length === 0) {
return 'unknown';
+42 -3
View File
@@ -32,7 +32,30 @@ describe('StratumV1ClientStatistics', () => {
await statistics.addShares(client, 64);
}
expect(statistics.hashRate).toBeGreaterThan(0);
expect(statistics.hashRate).toBeCloseTo((64 * 4294967296) / 62);
});
it('keeps exactly the configured number of share samples', async () => {
for (let i = 0; i < 31; i++) {
jest.setSystemTime(new Date(Date.parse('2026-05-06T12:00:00Z') + (i * 31000)));
await statistics.addShares(client, 64);
}
expect((statistics as any).submissionCache).toHaveLength(30);
expect((statistics as any).submissionCacheDifficultySum).toBe(30 * 64);
});
it('excludes pre-window work and corrects share-terminated sampling bias', async () => {
for (let i = 0; i < 30; i++) {
jest.setSystemTime(new Date(Date.parse('2026-05-06T12:00:00Z') + (i * 31000)));
await statistics.addShares(client, 64);
}
const elapsedSeconds = 29 * 31;
const expectedUnbiasedDifficulty = 28 * 64;
expect(statistics.hashRate).toBeCloseTo(
(expectedUnbiasedDifficulty * 4294967296) / elapsedSeconds,
);
});
it('should not suggest a difficulty change before enough time or shares have passed', () => {
@@ -51,7 +74,23 @@ describe('StratumV1ClientStatistics', () => {
await statistics.addShares(client, 64);
}
expect(statistics.getSuggestedDifficulty(64)).toBe(2048);
expect(statistics.getSuggestedDifficulty(64)).toBe(1024);
});
it('does not retarget when accepted shares have a zero-duration sample window', async () => {
for (let i = 0; i < 5; i++) {
await statistics.addShares(client, 64);
}
expect(() => statistics.getSuggestedDifficulty(64)).not.toThrow();
expect(statistics.getSuggestedDifficulty(64)).toBeNull();
});
it('returns a finite power-of-two difficulty for values above 32-bit range', () => {
const result = (statistics as any).nearestPowerOfTwo(2 ** 40 + 1);
expect(result).toBe(2 ** 40);
expect(Number.isFinite(result)).toBe(true);
});
it('should decrease difficulty for slow submissions', async () => {
@@ -60,7 +99,7 @@ describe('StratumV1ClientStatistics', () => {
await statistics.addShares(client, 64);
}
expect(statistics.getSuggestedDifficulty(128)).toBe(16);
expect(statistics.getSuggestedDifficulty(128)).toBe(8);
});
it('should not suggest a difficulty below the configured minimum', () => {
+46 -28
View File
@@ -16,9 +16,9 @@ export class StratumV1ClientStatistics {
}
public async addShares(_client: ClientEntity, targetDifficulty: number) {
var date = new Date();
const date = new Date();
if (this.submissionCache.length > CACHE_SIZE) {
if (this.submissionCache.length >= CACHE_SIZE) {
this.submissionCacheDifficultySum -= this.submissionCache[0].difficulty;
this.submissionCache.shift();
}
@@ -28,9 +28,12 @@ export class StratumV1ClientStatistics {
});
this.submissionCacheDifficultySum += targetDifficulty;
const time = new Date().getTime() - this.submissionCache[0].time.getTime();
if(time > 60000 && this.submissionCache.length > 2) {
this.hashRate = (this.submissionCacheDifficultySum * 4294967296) / (time / 1000);
const elapsedSeconds = (date.getTime() - this.submissionCache[0].time.getTime()) / 1000;
if (elapsedSeconds > 60) {
const difficultyPerSecond = this.getDifficultyPerSecond(elapsedSeconds);
if (difficultyPerSecond != null) {
this.hashRate = difficultyPerSecond * 4294967296;
}
}
}
@@ -46,15 +49,16 @@ export class StratumV1ClientStatistics {
}
}
const sum = this.submissionCache.reduce((pre, cur) => {
pre += cur.difficulty;
return pre;
}, 0);
const diffSeconds = (this.submissionCache[this.submissionCache.length - 1].time.getTime() - this.submissionCache[0].time.getTime()) / 1000;
const difficultyPerSecond = sum / diffSeconds;
const difficultyPerSecond = this.getDifficultyPerSecond(diffSeconds);
if (difficultyPerSecond == null) {
return null;
}
const targetDifficulty = difficultyPerSecond * this.targetSubmitShareEveryNSeconds;
if (!Number.isFinite(targetDifficulty) || targetDifficulty <= 0) {
return null;
}
if ((clientDifficulty * 2) < targetDifficulty || (clientDifficulty / 2) > targetDifficulty) {
return this.nearestPowerOfTwo(targetDifficulty)
@@ -63,27 +67,41 @@ export class StratumV1ClientStatistics {
return null;
}
private nearestPowerOfTwo(val): number {
if (val === 0) {
/**
* Estimate work rate from a share-terminated sample window.
*
* A cache of N shares contains N - 1 observed inter-share intervals. The
* first share's work predates the window and must not be counted. Because
* the window closes on a share arrival, the reciprocal elapsed time also
* has the usual finite-sample Poisson bias; multiplying by
* (intervalCount - 1) / intervalCount removes it.
*/
private getDifficultyPerSecond(elapsedSeconds: number): number | null {
const sampleCount = this.submissionCache.length;
if (sampleCount <= 2 || !Number.isFinite(elapsedSeconds) || elapsedSeconds <= 0) {
return null;
}
if (val < this.minDifficulty) {
const intervalCount = sampleCount - 1;
const observedDifficulty = this.submissionCacheDifficultySum
- this.submissionCache[0].difficulty;
const unbiasedDifficulty = observedDifficulty * (intervalCount - 1) / intervalCount;
const difficultyPerSecond = unbiasedDifficulty / elapsedSeconds;
return Number.isFinite(difficultyPerSecond) && difficultyPerSecond > 0
? difficultyPerSecond
: null;
}
private nearestPowerOfTwo(val: number): number {
if (!Number.isFinite(val) || val <= 0) {
return null;
}
if (val <= this.minDifficulty) {
return this.minDifficulty;
}
let x = val | (val >> 1);
x = x | (x >> 2);
x = x | (x >> 4);
x = x | (x >> 8);
x = x | (x >> 16);
x = x | (x >> 32);
const res = x - (x >> 1);
if (res == 0 && val * 100 < this.minDifficulty) {
return this.minDifficulty;
}
if (res == 0) {
return this.nearestPowerOfTwo(val * 100) / 100;
}
return res;
const result = 2 ** Math.floor(Math.log2(val));
return Number.isFinite(result) ? Math.max(this.minDifficulty, result) : null;
}
}
@@ -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=')) {
@@ -1,10 +1,12 @@
import { Expose, Transform } from 'class-transformer';
import { ArrayMaxSize, ArrayMinSize, IsArray, IsNumber } from 'class-validator';
import { ArrayMaxSize, ArrayMinSize, IsArray, IsNumber, IsPositive, Max } from 'class-validator';
import { eRequestMethod } from '../enums/eRequestMethod';
import { eResponseMethod } from '../enums/eResponseMethod';
import { StratumBaseMessage } from './StratumBaseMessage';
export const MAX_STRATUM_DIFFICULTY = 2 ** 32;
export class SuggestDifficulty extends StratumBaseMessage {
@IsArray()
@ArrayMinSize(1)
@@ -16,6 +18,8 @@ export class SuggestDifficulty extends StratumBaseMessage {
@Expose()
@IsNumber()
@IsPositive()
@Max(MAX_STRATUM_DIFFICULTY)
@Transform(({ value, key, obj, type }) => {
return Number(obj.params[0]);
})
@@ -37,4 +41,3 @@ export class SuggestDifficulty extends StratumBaseMessage {
}
}
}
+26 -1
View File
@@ -120,6 +120,19 @@ describe('RedisMessagingService', () => {
consoleSpy.mockRestore();
});
it('uses an owner token when locking JSON cache refreshes', async () => {
await service.connect();
await expect(service.tryAcquireJsonCacheLock('summary', 'owner-a', 1000)).resolves.toBe(true);
await expect(service.tryAcquireJsonCacheLock('summary', 'owner-b', 1000)).resolves.toBe(false);
await service.releaseJsonCacheLock('summary', 'owner-b');
await expect(service.tryAcquireJsonCacheLock('summary', 'owner-b', 1000)).resolves.toBe(false);
await service.releaseJsonCacheLock('summary', 'owner-a');
await expect(service.tryAcquireJsonCacheLock('summary', 'owner-b', 1000)).resolves.toBe(true);
});
it('stores payout variants and same-height reorg templates under distinct tip keys', async () => {
await service.connect();
const firstHash = '11'.repeat(32);
@@ -493,7 +506,10 @@ function createRedisClient() {
connect: jest.fn().mockResolvedValue(undefined),
quit: jest.fn().mockResolvedValue(undefined),
on: jest.fn(),
set: jest.fn((key: string, value: string) => {
set: jest.fn((key: string, value: string, options?: { NX?: boolean }) => {
if (options?.NX && store.has(key)) {
return Promise.resolve(null);
}
store.set(key, value);
return Promise.resolve('OK');
}),
@@ -512,6 +528,15 @@ function createRedisClient() {
});
return Promise.resolve(deleted);
}),
eval: jest.fn((_script: string, options: { keys: string[], arguments: string[] }) => {
const [key] = options.keys;
const [owner] = options.arguments;
if (store.get(key) !== owner) {
return Promise.resolve(0);
}
store.delete(key);
return Promise.resolve(1);
}),
sAdd: jest.fn((key: string, value: string) => {
const set = sets.get(key) ?? new Set<string>();
set.add(value);
+36
View File
@@ -32,6 +32,7 @@ const blockTemplateHeightPointerKey = (height: number, payoutMode: PayoutMode) =
const blockTemplateLatestPointerKey = (payoutMode: PayoutMode) =>
`${BLOCK_TEMPLATE_LATEST_KEY}:${payoutMode}`;
const jsonCacheKey = (key: string) => `json-cache:${key}`;
const jsonCacheLockKey = (key: string) => `json-cache-lock:${key}`;
export interface BlockTemplateUpdate {
schemaVersion: 1;
@@ -524,6 +525,41 @@ export class RedisMessagingService implements OnModuleInit, OnModuleDestroy {
);
}
public async tryAcquireJsonCacheLock(key: string, owner: string, ttlMs: number): Promise<boolean | null> {
if (!await this.ensureConnected()) {
return null;
}
if (owner.length === 0 || ttlMs <= 0) {
return false;
}
const result = await this.publisher.set(
jsonCacheLockKey(key),
owner,
{ NX: true, PX: Math.max(1, Math.ceil(ttlMs)) },
);
return result === 'OK';
}
public async releaseJsonCacheLock(key: string, owner: string): Promise<void> {
if (!await this.ensureConnected() || owner.length === 0) {
return;
}
await this.publisher.eval(
`
if redis.call('GET', KEYS[1]) == ARGV[1] then
return redis.call('DEL', KEYS[1])
end
return 0
`,
{
keys: [jsonCacheLockKey(key)],
arguments: [owner],
},
);
}
private async readBlockTemplatePointer(value: string | null): Promise<IBlockTemplate | null> {
if (value == null) {
return null;