diff --git a/pnpm-lock.yaml b/pnpm-lock.yaml index 88574182..08887baa 100644 --- a/pnpm-lock.yaml +++ b/pnpm-lock.yaml @@ -189,6 +189,9 @@ importers: typescript: specifier: ^5.3.3 version: 5.3.3 + vitest: + specifier: ^1.6.0 + version: 1.6.0(@types/node@18.16.3) test/e2e: dependencies: @@ -815,9 +818,6 @@ packages: '@jridgewell/sourcemap-codec@1.4.14': resolution: {integrity: sha512-XPSJHWmi394fuUuzDnGz1wiKqWfo1yXecHQMRf2l6hztTO+nPru658AyDngaBe7isIxEkRsPR3FZh+s7iVa4Uw==} - '@jridgewell/sourcemap-codec@1.4.15': - resolution: {integrity: sha512-eF2rxCRulEKXHTRiDrDy6erMYWqNw4LPdQ8UQA4huuxaQsVeRPFl2oM8oDGxMFhJUWZf9McpLtJasDDZb/Bpeg==} - '@jridgewell/sourcemap-codec@1.5.0': resolution: {integrity: sha512-gv3ZRaISU3fjPAgNsriBRqGWQL6quFx04YMPW/zD8XMLsU32mhCCbfbO6KZFLjvYpCZ8zyDEgqsgf+PwPaM7GQ==} @@ -2464,15 +2464,16 @@ packages: glob@10.4.5: resolution: {integrity: sha512-7Bv8RF0k6xjo7d4A/PxYLbUCfb6c+Vpd2/mB2yRDlew7Jb5hEXiCD9ibfO7wpk8i4sevK6DFny9h7EYbM3/sHg==} + deprecated: Old versions of glob are not supported, and contain widely publicized security vulnerabilities, which have been fixed in the current version. Please update. Support for old versions may be purchased (at exorbitant rates) by contacting i@izs.me hasBin: true glob@7.2.0: resolution: {integrity: sha512-lmLf6gtyrPq8tTjSmrO94wBeQbFR3HbLHbuyD69wuyQkImp2hWqMGB47OX65FBkPffO641IP9jWa1z4ivqG26Q==} - deprecated: Glob versions prior to v9 are no longer supported + deprecated: Old versions of glob are not supported, and contain widely publicized security vulnerabilities, which have been fixed in the current version. Please update. Support for old versions may be purchased (at exorbitant rates) by contacting i@izs.me glob@7.2.3: resolution: {integrity: sha512-nFR0zLpU2YCaRxwoCJvL6UvCH2JFyFVIvwTLsIf21AuHlMskA1hhTdk+LlYJtOlYt9v6dvszD2BGRqBL+iQK9Q==} - deprecated: Glob versions prior to v9 are no longer supported + deprecated: Old versions of glob are not supported, and contain widely publicized security vulnerabilities, which have been fixed in the current version. Please update. Support for old versions may be purchased (at exorbitant rates) by contacting i@izs.me globals@11.12.0: resolution: {integrity: sha512-WOBp/EEGUiIsJSp7wcv/y6MO+lV9UoncWqxuFfm8eBwzWNgyfBd6Gz+IeKQ9jCmyhoH99g15M3T+QaVHFjizVA==} @@ -3220,9 +3221,6 @@ packages: resolution: {integrity: sha512-qTAAlrEsl8s4OiEQY69wDvcMIdQN6wdz5ojQiOy6YRMuynxenON0O5oCpJI6lshc6scgAY8qvJ2On/p+CXY0GA==} engines: {node: '>=4'} - picocolors@1.0.0: - resolution: {integrity: sha512-1fygroTLlHu66zi26VoTDv8yRgm0Fccecssto+MhsZ0D/DGW2sm8E8AjW7NU5VVTRt5GxbeZ5qBuJr+HyLYkjQ==} - picocolors@1.1.1: resolution: {integrity: sha512-xceH2snhtb5M9liqDsmEw56le376mTZkEX/jEb/RxNFyegNul7eNslCXP9FDj/Lcu0X8KEyMceP2ntpaHrDEVA==} @@ -3754,6 +3752,7 @@ packages: tar@7.2.0: resolution: {integrity: sha512-hctwP0Nb4AB60bj8WQgRYaMOuJYRAPMGiQUAotms5igN8ppfQM+IvjQ5HcKu1MaZh2Wy2KWVTe563Yj8dfc14w==} engines: {node: '>=18'} + deprecated: Old versions of tar are not supported, and contain widely publicized security vulnerabilities, which have been fixed in the current version. Please update. Support for old versions may be purchased (at exorbitant rates) by contacting i@izs.me tdigest@0.1.2: resolution: {integrity: sha512-+G0LLgjjo9BZX2MfdvPfH+MKLCrxlXSYec5DaPYP1fe6Iyhf0/fSmJ0bFiZ1F8BT6cGXl2LpltQptzjXKWEkKA==} @@ -3931,6 +3930,7 @@ packages: uuid@8.3.2: resolution: {integrity: sha512-+NYs2QeMWy+GWFOEm9xnn6HCDp0l7QBD7ml8zLUmJ+93Q5NF0NocErnwkTkXVFNiX3/fpC6afS8Dhb/gz7R7eg==} + deprecated: uuid@10 and below is no longer supported. For ESM codebases, update to uuid@latest. For CommonJS codebases, use uuid@11 (but be aware this version will likely be deprecated in 2028). hasBin: true v8-compile-cache-lib@3.0.1: @@ -4230,7 +4230,7 @@ snapshots: '@babel/traverse': 7.21.5 '@babel/types': 7.21.5 convert-source-map: 1.9.0 - debug: 4.3.4 + debug: 4.3.6 gensync: 1.0.0-beta.2 json5: 2.2.3 semver: 6.3.1 @@ -4777,7 +4777,7 @@ snapshots: '@jridgewell/gen-mapping@0.3.3': dependencies: '@jridgewell/set-array': 1.1.2 - '@jridgewell/sourcemap-codec': 1.4.15 + '@jridgewell/sourcemap-codec': 1.5.0 '@jridgewell/trace-mapping': 0.3.18 '@jridgewell/resolve-uri@3.1.0': {} @@ -4786,8 +4786,6 @@ snapshots: '@jridgewell/sourcemap-codec@1.4.14': {} - '@jridgewell/sourcemap-codec@1.4.15': {} - '@jridgewell/sourcemap-codec@1.5.0': {} '@jridgewell/trace-mapping@0.3.18': @@ -4798,7 +4796,7 @@ snapshots: '@jridgewell/trace-mapping@0.3.9': dependencies: '@jridgewell/resolve-uri': 3.1.0 - '@jridgewell/sourcemap-codec': 1.4.15 + '@jridgewell/sourcemap-codec': 1.5.0 '@js-sdsl/ordered-map@4.4.2': {} @@ -5912,7 +5910,7 @@ snapshots: dependencies: '@fastify/error': 3.4.1 archy: 1.0.0 - debug: 4.3.4 + debug: 4.3.6 fastq: 1.17.1 transitivePeerDependencies: - supports-color @@ -6953,7 +6951,7 @@ snapshots: istanbul-lib-source-maps@4.0.1: dependencies: - debug: 4.3.4 + debug: 4.3.6 istanbul-lib-coverage: 3.2.0 source-map: 0.6.1 transitivePeerDependencies: @@ -7437,8 +7435,6 @@ snapshots: postgres-date: 1.0.7 postgres-interval: 1.2.0 - picocolors@1.0.0: {} - picocolors@1.1.1: {} picomatch@2.3.1: {} @@ -8233,7 +8229,7 @@ snapshots: dependencies: browserslist: 4.21.5 escalade: 3.1.1 - picocolors: 1.0.0 + picocolors: 1.1.1 uri-js@4.4.1: dependencies: @@ -8306,7 +8302,7 @@ snapshots: cac: 6.7.14 debug: 4.3.6 pathe: 1.1.2 - picocolors: 1.0.0 + picocolors: 1.1.1 vite: 5.4.11(@types/node@18.16.3) transitivePeerDependencies: - '@types/node' @@ -8342,7 +8338,7 @@ snapshots: local-pkg: 0.5.1 magic-string: 0.30.17 pathe: 1.1.2 - picocolors: 1.0.0 + picocolors: 1.1.1 std-env: 3.8.0 strip-literal: 2.1.1 tinybench: 2.9.0 diff --git a/src/cli/config/redisKeys.ts b/src/cli/config/redisKeys.ts index 78b28e4b..a4788649 100644 --- a/src/cli/config/redisKeys.ts +++ b/src/cli/config/redisKeys.ts @@ -17,6 +17,9 @@ export const getRedisKeys = (config: AltoConfig) => { gasPriceQueue: `${prefix}:gas-price`, // Sender manager queue - senderManagerQueue: `${prefix}:sender-manager` + senderManagerQueue: `${prefix}:sender-manager`, + + // Wallets currently checked out of the sender manager queue + senderManagerCheckedOut: `${prefix}:sender-manager:checked-out` } } diff --git a/src/executor/senderManager/createRedisSenderManager.test.ts b/src/executor/senderManager/createRedisSenderManager.test.ts new file mode 100644 index 00000000..875d43a4 --- /dev/null +++ b/src/executor/senderManager/createRedisSenderManager.test.ts @@ -0,0 +1,143 @@ +// Integration tests for the Redis sender manager against a real Redis. +// Skipped unless REDIS_URL is set, e.g.: +// docker run --rm -d -p 6379:6379 redis:7-alpine +// REDIS_URL=redis://localhost:6379 pnpm --filter @pimlico/alto test +import Redis from "ioredis" +import { generatePrivateKey, privateKeyToAccount } from "viem/accounts" +import { afterAll, afterEach, describe, expect, it } from "vitest" +import { getRedisKeys } from "../../cli/config/redisKeys" +import type { AltoConfig } from "../../createConfig" +import { + RELEASE_WALLET_LUA, + RESTORE_WALLET_LUA, + createRedisSenderManager +} from "./createRedisSenderManager" + +const redisUrl = process.env.REDIS_URL + +const STALE_MS = 30 * 60 * 1000 + +describe.skipIf(!redisUrl)("createRedisSenderManager", () => { + const walletA = privateKeyToAccount(generatePrivateKey()) + const walletB = privateKeyToAccount(generatePrivateKey()) + const stranger = privateKeyToAccount(generatePrivateKey()) + + const noopLogger = { + trace: () => {}, + info: () => {}, + warn: () => {}, + error: () => {} + } + const metrics = { + walletsTotal: { set: () => {} }, + walletsAvailable: { set: () => {} } + } as unknown as Parameters[0]["metrics"] + const config = { + executorPrivateKeys: [walletA, walletB], + maxExecutors: undefined, + redisKeyPrefix: `sender-manager-test-${process.pid}`, + chainId: 999, + logLevel: "silent", + executorLogLevel: "silent", + getLogger: () => noopLogger + } as unknown as AltoConfig + + const { senderManagerQueue: queueKey, senderManagerCheckedOut: markerKey } = + getRedisKeys(config) + + const redis = new Redis(redisUrl as string) + + const release = (address: string) => + redis.eval(RELEASE_WALLET_LUA, 2, queueKey, markerKey, address) + const restore = (address: string, staleCutoff: number) => + redis.eval( + RESTORE_WALLET_LUA, + 2, + queueKey, + markerKey, + address, + staleCutoff + ) + + afterEach(async () => { + await redis.del(queueKey, markerKey) + }) + + afterAll(() => { + redis.disconnect() + }) + + it("seeds on first boot and round-trips checkout/release", async () => { + const manager = await createRedisSenderManager({ + config, + metrics, + redisEndpoint: redisUrl as string + }) + expect(await redis.llen(queueKey)).toBe(2) + + const wallet = await manager.getWallet() + expect(await redis.llen(queueKey)).toBe(1) + expect(await redis.zscore(markerKey, wallet.address)).not.toBeNull() + + await manager.markWalletProcessed(wallet) + expect(await redis.llen(queueKey)).toBe(2) + expect(await redis.zscore(markerKey, wallet.address)).toBeNull() + }) + + it("boot prunes unconfigured wallets and keeps a non-empty pool unseeded", async () => { + await redis.rpush(queueKey, walletA.address, stranger.address) + await redis.zadd(markerKey, Date.now(), walletB.address) + + await createRedisSenderManager({ + config, + metrics, + redisEndpoint: redisUrl as string + }) + + expect(await redis.lrange(queueKey, 0, -1)).toEqual([walletA.address]) + expect(await redis.zscore(markerKey, walletB.address)).not.toBeNull() + }) + + it("release without ownership marker is a no-op", async () => { + await redis.rpush(queueKey, walletA.address) + expect(await release(walletA.address)).toBe(0) + expect(await redis.lrange(queueKey, 0, -1)).toEqual([walletA.address]) + }) + + it("release does not duplicate an already-queued wallet", async () => { + await redis.rpush(queueKey, walletA.address) + await redis.zadd(markerKey, Date.now(), walletA.address) + expect(await release(walletA.address)).toBe(1) + expect(await redis.lrange(queueKey, 0, -1)).toEqual([walletA.address]) + expect(await redis.zscore(markerKey, walletA.address)).toBeNull() + }) + + it("restore skips a queued wallet", async () => { + await redis.rpush(queueKey, walletA.address) + expect(await restore(walletA.address, Date.now() - STALE_MS)).toBe(0) + expect(await redis.lrange(queueKey, 0, -1)).toEqual([walletA.address]) + }) + + it("restore skips a wallet with a fresh marker", async () => { + await redis.zadd(markerKey, Date.now(), walletA.address) + expect(await restore(walletA.address, Date.now() - STALE_MS)).toBe(0) + expect(await redis.llen(queueKey)).toBe(0) + expect(await redis.zscore(markerKey, walletA.address)).not.toBeNull() + }) + + it("restore re-queues a wallet with a stale marker", async () => { + await redis.zadd( + markerKey, + Date.now() - STALE_MS - 1000, + walletA.address + ) + expect(await restore(walletA.address, Date.now() - STALE_MS)).toBe(1) + expect(await redis.lrange(queueKey, 0, -1)).toEqual([walletA.address]) + expect(await redis.zscore(markerKey, walletA.address)).toBeNull() + }) + + it("restore re-queues a wallet absent from both structures", async () => { + expect(await restore(walletA.address, Date.now() - STALE_MS)).toBe(1) + expect(await redis.lrange(queueKey, 0, -1)).toEqual([walletA.address]) + }) +}) diff --git a/src/executor/senderManager/createRedisSenderManager.ts b/src/executor/senderManager/createRedisSenderManager.ts index 3902a0fd..05a70560 100644 --- a/src/executor/senderManager/createRedisSenderManager.ts +++ b/src/executor/senderManager/createRedisSenderManager.ts @@ -6,37 +6,63 @@ import { getRedisKeys } from "../../cli/config/redisKeys" import type { AltoConfig } from "../../createConfig" import type { SenderManager } from "../senderManager" -async function createRedisQueue({ - redis, - name, - entries -}: { - redis: Redis - name: string - entries: string[] -}) { - const hasElements = await redis.llen(name) - - // Ensure queue is populated on startup - // Avoids race case where queue is populated twice due to multi (atomic txs) - if (hasElements === 0) { - const multi = redis.multi() - multi.del(name) - multi.rpush(name, ...entries) - await multi.exec() - } - - return { - llen: () => redis.llen(name), - pop: () => redis.rpop(name), - push: (entry: string) => redis.lpush(name, entry) - } -} - const delay = async (delay: number) => { await new Promise((resolve) => setTimeout(resolve, delay)) } +// Atomically pops a wallet and stamps its checked-out marker, so a held +// wallet is never visible in neither structure (which the restore scan on +// another pod would misread as a leak). +export const CHECKOUT_WALLET_LUA = ` +local address = redis.call('RPOP', KEYS[1]) +if address then + redis.call('ZADD', KEYS[2], ARGV[1], address) +end +return address +` + +// Ownership-conditional release: the marker is the ownership token (checkout +// always stamps it), so a release only pushes if this holder's marker still +// exists. Retrying after a client-side timeout whose eval actually reached +// the server is then a no-op (ZREM == 0) instead of double-queueing a wallet +// another pod has since checked out. +export const RELEASE_WALLET_LUA = ` +if redis.call('ZREM', KEYS[2], ARGV[1]) == 1 then + if redis.call('LPOS', KEYS[1], ARGV[1]) == false then + redis.call('LPUSH', KEYS[1], ARGV[1]) + end + return 1 +end +return 0 +` + +// Leak restore: re-validates inside the script (atomically vs concurrent +// checkouts/restores) that the wallet is not queued and has no fresh marker +// before pushing. ARGV[2] is the stale cutoff timestamp. +export const RESTORE_WALLET_LUA = ` +if redis.call('LPOS', KEYS[1], ARGV[1]) ~= false then + return 0 +end +local score = redis.call('ZSCORE', KEYS[2], ARGV[1]) +if score and tonumber(score) >= tonumber(ARGV[2]) then + return 0 +end +redis.call('ZREM', KEYS[2], ARGV[1]) +redis.call('LPUSH', KEYS[1], ARGV[1]) +return 1 +` + +// Held wallets re-stamp their marker every refresh tick (covers arbitrarily +// long holds, e.g. quarantine), so a marker older than the staleness window +// means its holder is gone — crashed pod, or a release that never went +// through. A wallet must additionally be observed continuously missing for +// the grace period before restore, so holders whose re-stamps were blocked +// (e.g. a Redis OOM window) or old-code pods during a rolling deploy get +// time to reassert ownership before anyone re-queues their wallet. +const CHECKED_OUT_REFRESH_MS = 60 * 1000 +const CHECKED_OUT_STALE_MS = 30 * 60 * 1000 +const RESTORE_GRACE_MS = 10 * 60 * 1000 + export const createRedisSenderManager = async ({ config, metrics, @@ -57,18 +83,213 @@ export const createRedisSenderManager = async ({ ) const redis = new Redis(redisEndpoint) - const redisQueueName = getRedisKeys(config).senderManagerQueue - const redisQueue = await createRedisQueue({ - redis, - name: redisQueueName, - entries: wallets.map((w) => w.address) - }) - - // Track active wallets for this instance + const { senderManagerQueue, senderManagerCheckedOut } = getRedisKeys(config) + const configured = new Set(wallets.map((w) => w.address)) + + const checkoutWallet = async () => + (await redis.eval( + CHECKOUT_WALLET_LUA, + 2, + senderManagerQueue, + senderManagerCheckedOut, + Date.now() + )) as string | null + + const releaseWallet = async (address: string) => + (await redis.eval( + RELEASE_WALLET_LUA, + 2, + senderManagerQueue, + senderManagerCheckedOut, + address + )) as number + + const restoreWallet = async (address: string) => + (await redis.eval( + RESTORE_WALLET_LUA, + 2, + senderManagerQueue, + senderManagerCheckedOut, + address, + Date.now() - CHECKED_OUT_STALE_MS + )) as number + + const updateAvailableMetric = async () => { + try { + metrics.walletsAvailable.set(await redis.llen(senderManagerQueue)) + } catch (err) { + logger.warn({ err }, "failed to update walletsAvailable metric") + } + } + + const readPoolState = async () => { + const [queued, markers] = await Promise.all([ + redis.lrange(senderManagerQueue, 0, -1), + redis.zrange(senderManagerCheckedOut, 0, -1, "WITHSCORES") + ]) + const markerScores = new Map() + for (let i = 0; i < markers.length; i += 2) { + markerScores.set(markers[i], Number(markers[i + 1])) + } + return { queued, markerScores } + } + + // Boot: prune leftovers of wallets that are no longer configured, and + // seed the pool on genuine first boot (both structures empty). Restores + // of leaked wallets are deliberately NOT done at boot — they run in the + // maintenance loop behind the grace period. Never fail boot over Redis: + // the maintenance loop heals whatever this missed. + try { + const { queued, markerScores } = await readPoolState() + for (const address of queued) { + if (!configured.has(address)) { + await redis.lrem(senderManagerQueue, 0, address) + logger.warn( + { executor: address }, + "pruned unconfigured wallet from sender manager queue" + ) + } + } + for (const address of markerScores.keys()) { + if (!configured.has(address)) { + await redis.zrem(senderManagerCheckedOut, address) + } + } + const anyConfiguredPresent = + queued.some((address) => configured.has(address)) || + [...markerScores.keys()].some((address) => configured.has(address)) + if (!anyConfiguredPresent) { + for (const wallet of wallets) { + await restoreWallet(wallet.address) + } + logger.info("seeded sender manager queue") + } + } catch (err) { + logger.error( + { err }, + "sender manager boot reconciliation failed, maintenance loop will retry" + ) + } + + // Wallets held by this instance, and held wallets whose release failed + // and is awaiting retry by the maintenance loop. const activeWallets = new Set() + const pendingRelease = new Set() + // First time each configured wallet was observed missing from both the + // queue and fresh markers; restore happens only after the grace period. + const missingSince = new Map() + + const maintain = async () => { + // Retry releases that failed (a Redis outage can outlast + // markWalletProcessed's retries). + for (const wallet of [...pendingRelease]) { + try { + const released = await releaseWallet(wallet.address) + pendingRelease.delete(wallet) + activeWallets.delete(wallet) + await updateAvailableMetric() + if (released === 1) { + logger.info( + { executor: wallet.address }, + "returned wallet to sender manager queue after earlier failure" + ) + } else { + logger.warn( + { executor: wallet.address }, + "wallet was no longer owned at release retry" + ) + } + } catch (err) { + logger.error( + { err, executor: wallet.address }, + "failed to return wallet to sender manager queue" + ) + } + } + + // Re-stamp markers for held wallets so other pods' restore scans + // never mistake a live hold for a leak. Skip pending releases: their + // marker may already belong to a new holder. + const now = Date.now() + for (const wallet of [...activeWallets]) { + if (pendingRelease.has(wallet)) { + continue + } + try { + await redis.zadd(senderManagerCheckedOut, now, wallet.address) + } catch (err) { + logger.error( + { err, executor: wallet.address }, + "failed to refresh wallet checkout marker" + ) + } + } + + // Only trust marker staleness while Redis writes work: during a + // write outage (e.g. OOM) holders cannot re-stamp their markers, so + // stale markers say nothing about liveness. The probe shares the + // writes' fate; on failure the missing-clock resets, giving holders + // a full grace period after recovery to reassert ownership. + try { + await redis.set( + `${senderManagerQueue}:write-probe`, + Date.now(), + "PX", + CHECKED_OUT_REFRESH_MS * 2 + ) + } catch (err) { + missingSince.clear() + logger.warn( + { err }, + "redis writes unhealthy, skipping wallet restore scan" + ) + return + } + + // Restore wallets that have been continuously missing from both the + // queue and fresh markers for the grace period — leaked by a crashed + // pod or a release that never completed. + try { + const { queued, markerScores } = await readPoolState() + const staleCutoff = Date.now() - CHECKED_OUT_STALE_MS + const queuedSet = new Set(queued) + for (const wallet of wallets) { + const address = wallet.address + const markerScore = markerScores.get(address) + const missing = + !queuedSet.has(address) && + (markerScore === undefined || markerScore < staleCutoff) && + !activeWallets.has(wallet) + if (!missing) { + missingSince.delete(address) + continue + } + const firstSeen = missingSince.get(address) ?? Date.now() + missingSince.set(address, firstSeen) + if (Date.now() - firstSeen < RESTORE_GRACE_MS) { + continue + } + if ((await restoreWallet(address)) === 1) { + logger.warn( + { executor: address }, + "restored missing wallet to sender manager queue" + ) + await updateAvailableMetric() + } + missingSince.delete(address) + } + } catch (err) { + logger.error({ err }, "wallet restore scan failed") + } + } + setInterval(() => { + maintain().catch((err) => + logger.error({ err }, "wallet maintenance loop failed") + ) + }, CHECKED_OUT_REFRESH_MS).unref() logger.info( - `Created redis sender manager with queueName: ${redisQueueName}` + `Created redis sender manager with queueName: ${senderManagerQueue}` ) return { getAllWallets: () => [...wallets], @@ -78,8 +299,10 @@ export const createRedisSenderManager = async ({ let walletAddress: string | null = null while (!walletAddress) { - walletAddress = await redisQueue.pop() - await delay(100) + walletAddress = await checkoutWallet() + if (!walletAddress) { + await delay(100) + } } const wallet = wallets.find((w) => w.address === walletAddress) @@ -96,23 +319,46 @@ export const createRedisSenderManager = async ({ "got wallet from sender manager" ) - await redisQueue.llen().then((len) => { - metrics.walletsAvailable.set(len) - }) + await updateAvailableMetric() return wallet }, markWalletProcessed: async (wallet: Account) => { - if (activeWallets.delete(wallet)) { - await redisQueue.push(wallet.address) - const len = await redisQueue.llen() - metrics.walletsAvailable.set(len) - } else { + if (!activeWallets.has(wallet)) { logger.warn( { executor: wallet.address }, "Attempted to mark a wallet as processed that wasn't active" ) + return + } + + // Return the wallet to Redis BEFORE dropping local bookkeeping + // and never throw — if the push fails the wallet must stay + // tracked so the maintenance loop, shutdown sweep and the + // restore scan can recover it instead of leaking it from the + // pool. + for (let attempt = 1; attempt <= 3; attempt++) { + try { + const released = await releaseWallet(wallet.address) + activeWallets.delete(wallet) + pendingRelease.delete(wallet) + await updateAvailableMetric() + if (released === 0) { + logger.warn( + { executor: wallet.address }, + "wallet was no longer owned at release" + ) + } + return + } catch (err) { + logger.error( + { err, executor: wallet.address, attempt }, + "failed to return wallet to sender manager queue" + ) + await delay(1000) + } } + pendingRelease.add(wallet) }, getActiveWallets: () => { return [...activeWallets] diff --git a/src/executor/senderManager/validateAndRefill.ts b/src/executor/senderManager/validateAndRefill.ts index 9fc733cf..15283ffb 100644 --- a/src/executor/senderManager/validateAndRefill.ts +++ b/src/executor/senderManager/validateAndRefill.ts @@ -64,19 +64,10 @@ export const validateAndRefillWallets = async ({ { balancesMissing, totalBalanceMissing }, "balances missing" ) + // Executor wallets are funded externally now, so an underfunded + // utility wallet is expected — skip refilling without logging an + // error every cycle. metrics.utilityWalletInsufficientBalance.set(1) - logger.error( - { - minBalance: formatEther(minBalance), - utilityWalletBalance: formatEther(utilityWalletBalance), - totalBalanceMissing: formatEther(totalBalanceMissing), - minRefillAmount: formatEther( - totalBalanceMissing - utilityWalletBalance - ), - utilityAccount: utilityAccount.address - }, - "utility wallet has insufficient balance to refill wallets" - ) return } diff --git a/src/package.json b/src/package.json index 6e168b86..bb40dde0 100644 --- a/src/package.json +++ b/src/package.json @@ -32,7 +32,8 @@ "build:esm": "tsc -p ./tsconfig.esm.json && tsc-alias -p tsconfig.esm.json", "dev": "nodemon --exec DOTENV_CONFIG_PATH=$(pwd)/../.env tsx -r tsconfig-paths/register cli/alto.ts run", "lint": "eslint src/**/*.ts", - "lint:fix": "eslint src/**/*.ts --fix" + "lint:fix": "eslint src/**/*.ts --fix", + "test": "vitest run" }, "dependencies": { "@fastify/websocket": "^10.0.1", @@ -77,6 +78,7 @@ "ts-node": "^10.9.2", "tsc-alias": "^1.8.8", "tsconfig-paths": "^4.2.0", - "typescript": "^5.3.3" + "typescript": "^5.3.3", + "vitest": "^1.6.0" } } diff --git a/src/vitest.config.ts b/src/vitest.config.ts new file mode 100644 index 00000000..325676fb --- /dev/null +++ b/src/vitest.config.ts @@ -0,0 +1,11 @@ +import { defineConfig } from "vitest/config" + +export default defineConfig({ + test: { + environment: "node", + // forks pool exits cleanly despite the sender manager's long-lived + // ioredis connection (SenderManager exposes no close()) + pool: "forks", + testTimeout: 15_000 + } +})