Prevent accounting cache refresh stampedes

This commit is contained in:
Ben
2026-08-04 23:47:41 -04:00
parent 1cdf3fbf7e
commit 4c56839c02
8 changed files with 306 additions and 8 deletions
+4
View File
@@ -170,6 +170,10 @@ SHARE_ACCOUNTING_FLUSH_INTERVAL_MS=25
SHARE_ACCOUNTING_MAX_QUEUE_SIZE=50000 SHARE_ACCOUNTING_MAX_QUEUE_SIZE=50000
SHARE_ACCOUNTING_SUMMARY_CACHE_MS=2500 SHARE_ACCOUNTING_SUMMARY_CACHE_MS=2500
SHARE_ACCOUNTING_SUMMARY_CACHE_MAX=10000 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 CLIENT_HASHRATE_PERSIST_INTERVAL_MS=60000
SHARE_ROLLUP_ENABLED=true SHARE_ROLLUP_ENABLED=true
SHARE_ROLLUP_INTERVAL_MS=60000 SHARE_ROLLUP_INTERVAL_MS=60000
+4
View File
@@ -106,6 +106,10 @@ services:
SHARE_ACCOUNTING_MAX_QUEUE_SIZE: ${SHARE_ACCOUNTING_MAX_QUEUE_SIZE:-50000} 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_MS: ${SHARE_ACCOUNTING_SUMMARY_CACHE_MS:-2500}
SHARE_ACCOUNTING_SUMMARY_CACHE_MAX: ${SHARE_ACCOUNTING_SUMMARY_CACHE_MAX:-10000} 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_ENABLED: ${SHARE_ROLLUP_ENABLED:-true}
SHARE_ROLLUP_INTERVAL_MS: ${SHARE_ROLLUP_INTERVAL_MS:-60000} SHARE_ROLLUP_INTERVAL_MS: ${SHARE_ROLLUP_INTERVAL_MS:-60000}
SHARE_ROLLUP_SAFETY_LAG_SECONDS: ${SHARE_ROLLUP_SAFETY_LAG_SECONDS:-30} SHARE_ROLLUP_SAFETY_LAG_SECONDS: ${SHARE_ROLLUP_SAFETY_LAG_SECONDS:-30}
+4
View File
@@ -110,6 +110,10 @@ services:
API_HEADERS_TIMEOUT_MS: ${API_HEADERS_TIMEOUT_MS:-10000} API_HEADERS_TIMEOUT_MS: ${API_HEADERS_TIMEOUT_MS:-10000}
API_MAX_REQUESTS_PER_SOCKET: ${API_MAX_REQUESTS_PER_SOCKET:-100} API_MAX_REQUESTS_PER_SOCKET: ${API_MAX_REQUESTS_PER_SOCKET:-100}
API_MAX_CONNECTIONS_PER_WORKER: ${API_MAX_CONNECTIONS_PER_WORKER:-500} 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} PM2_ENABLED: ${PM2_ENABLED:-true}
STRATUM_WORKERS: ${STRATUM_WORKERS:-auto} STRATUM_WORKERS: ${STRATUM_WORKERS:-auto}
STRATUM_PORTS: ${STRATUM_PORTS:-3333,3332,3331,3330} STRATUM_PORTS: ${STRATUM_PORTS:-3333,3332,3331,3330}
+4
View File
@@ -114,6 +114,10 @@ services:
API_HEADERS_TIMEOUT_MS: ${API_HEADERS_TIMEOUT_MS:-10000} API_HEADERS_TIMEOUT_MS: ${API_HEADERS_TIMEOUT_MS:-10000}
API_MAX_REQUESTS_PER_SOCKET: ${API_MAX_REQUESTS_PER_SOCKET:-100} API_MAX_REQUESTS_PER_SOCKET: ${API_MAX_REQUESTS_PER_SOCKET:-100}
API_MAX_CONNECTIONS_PER_WORKER: ${API_MAX_CONNECTIONS_PER_WORKER:-500} 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" PM2_ENABLED: "true"
STRATUM_WORKERS: ${STRATUM_WORKERS:-2} STRATUM_WORKERS: ${STRATUM_WORKERS:-2}
STRATUM_MIN_DIFFICULTY: ${STRATUM_MIN_DIFFICULTY:-1} STRATUM_MIN_DIFFICULTY: ${STRATUM_MIN_DIFFICULTY:-1}
@@ -365,11 +365,65 @@ describe('ShareAccountingService', () => {
expect(repository.query).toHaveBeenCalledTimes(1); expect(repository.query).toHaveBeenCalledTimes(1);
expect(redis.setJsonCache).toHaveBeenCalledWith( expect(redis.setJsonCache).toHaveBeenCalledWith(
expect.stringContaining('accounting:summary:'), expect.stringContaining('accounting:summary:'),
expect.objectContaining({ totalAcceptedShares: 1 }), expect.objectContaining({
30000, 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 () => { it('should refresh pool summaries from completed rollup buckets and current round rollups', async () => {
const redis = { const redis = {
getJsonCache: jest.fn().mockResolvedValue(null), getJsonCache: jest.fn().mockResolvedValue(null),
@@ -1,5 +1,6 @@
import { Injectable, OnModuleDestroy, OnModuleInit } from '@nestjs/common'; import { Injectable, OnModuleDestroy, OnModuleInit } from '@nestjs/common';
import { InjectRepository } from '@nestjs/typeorm'; import { InjectRepository } from '@nestjs/typeorm';
import { randomUUID } from 'crypto';
import { Repository } from 'typeorm'; import { Repository } from 'typeorm';
import { AcceptedShareEntity } from '../accepted-share/accepted-share.entity'; 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_MAX_QUEUE_SIZE = 50000;
const DEFAULT_SUMMARY_CACHE_MS = 2500; const DEFAULT_SUMMARY_CACHE_MS = 2500;
const DEFAULT_SUMMARY_CACHE_MAX = 10000; 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_INTERVAL_MS = 60000;
const DEFAULT_ROLLUP_SAFETY_LAG_SECONDS = 10; const DEFAULT_ROLLUP_SAFETY_LAG_SECONDS = 10;
const DEFAULT_ROLLUP_MAX_SHARES_PER_BATCH = 5000000; 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 activeRollup: Promise<void> | null = null;
private summaryCache = new Map<string, SummaryCacheEntry>(); private summaryCache = new Map<string, SummaryCacheEntry>();
private summaryInFlight = new Map<string, Promise<ShareAccountingSummary>>(); private summaryInFlight = new Map<string, Promise<ShareAccountingSummary>>();
private redisSummaryRefreshInFlight = new Map<string, Promise<void>>();
private readonly poolSummaryCacheKey = 'accounting:pool-summary'; private readonly poolSummaryCacheKey = 'accounting:pool-summary';
private readonly batchSize = this.readPositiveInt('SHARE_ACCOUNTING_BATCH_SIZE', DEFAULT_BATCH_SIZE); 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); 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 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 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 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 shareRollupEnabled = this.readBoolean('SHARE_ROLLUP_ENABLED', true);
private readonly shareRollupIntervalMs = this.readPositiveInt('SHARE_ROLLUP_INTERVAL_MS', DEFAULT_ROLLUP_INTERVAL_MS); 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); 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> { private async loadCachedSummary(filter: AccountingFilter, cacheKey: string): Promise<ShareAccountingSummary> {
const redisCacheKey = `accounting:summary:${cacheKey}`; 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 const cached = await this.redisMessagingService
?.getJsonCache<ShareAccountingSummary>(redisCacheKey) ?.getJsonCache<RedisSummaryCacheEntry | ShareAccountingSummary>(redisCacheKey)
.catch(error => { .catch(error => {
console.error(`Share accounting summary cache read failed: ${error.message}`); console.error(`Share accounting summary cache read failed: ${error.message}`);
return null; return null;
}); });
if (cached != null) { if (cached == null) {
return null;
}
if (this.isRedisSummaryCacheEntry(cached)) {
return 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); 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 await this.redisMessagingService
?.setJsonCache(redisCacheKey, summary, this.redisSummaryCacheMs) ?.setJsonCache(redisCacheKey, entry, this.redisSummaryStaleCacheMs)
.catch(error => { .catch(error => {
console.error(`Share accounting summary cache write failed: ${error.message}`); 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> { private async loadSummary(filter: AccountingFilter): Promise<ShareAccountingSummary> {
@@ -924,3 +1085,9 @@ interface SummaryCacheEntry {
expiresAt: number; expiresAt: number;
value: ShareAccountingSummary; value: ShareAccountingSummary;
} }
interface RedisSummaryCacheEntry {
schemaVersion: 1;
refreshedAtMs: number;
value: ShareAccountingSummary;
}
+26 -1
View File
@@ -120,6 +120,19 @@ describe('RedisMessagingService', () => {
consoleSpy.mockRestore(); 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 () => { it('stores payout variants and same-height reorg templates under distinct tip keys', async () => {
await service.connect(); await service.connect();
const firstHash = '11'.repeat(32); const firstHash = '11'.repeat(32);
@@ -493,7 +506,10 @@ function createRedisClient() {
connect: jest.fn().mockResolvedValue(undefined), connect: jest.fn().mockResolvedValue(undefined),
quit: jest.fn().mockResolvedValue(undefined), quit: jest.fn().mockResolvedValue(undefined),
on: jest.fn(), 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); store.set(key, value);
return Promise.resolve('OK'); return Promise.resolve('OK');
}), }),
@@ -512,6 +528,15 @@ function createRedisClient() {
}); });
return Promise.resolve(deleted); 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) => { sAdd: jest.fn((key: string, value: string) => {
const set = sets.get(key) ?? new Set<string>(); const set = sets.get(key) ?? new Set<string>();
set.add(value); set.add(value);
+36
View File
@@ -32,6 +32,7 @@ const blockTemplateHeightPointerKey = (height: number, payoutMode: PayoutMode) =
const blockTemplateLatestPointerKey = (payoutMode: PayoutMode) => const blockTemplateLatestPointerKey = (payoutMode: PayoutMode) =>
`${BLOCK_TEMPLATE_LATEST_KEY}:${payoutMode}`; `${BLOCK_TEMPLATE_LATEST_KEY}:${payoutMode}`;
const jsonCacheKey = (key: string) => `json-cache:${key}`; const jsonCacheKey = (key: string) => `json-cache:${key}`;
const jsonCacheLockKey = (key: string) => `json-cache-lock:${key}`;
export interface BlockTemplateUpdate { export interface BlockTemplateUpdate {
schemaVersion: 1; 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> { private async readBlockTemplatePointer(value: string | null): Promise<IBlockTemplate | null> {
if (value == null) { if (value == null) {
return null; return null;