From 553385d738e455879f5a43edb3a7f9d2f87abd8e Mon Sep 17 00:00:00 2001 From: Big-JoshCode Date: Sat, 26 Sep 2026 00:19:24 +0100 Subject: [PATCH] fix: per-asset alert index, webhook probe retries, SSRF check before probe, atomic recipient SADD - alerts.js: evaluateForAssetInner now looks up alerts via a per-asset Redis set index (alerts:asset:) instead of ZREVRANGE-ing and fetching every alert in the system on every price tick. Index is kept in sync by create()/remove(); a backfillAssetIndex() helper is exported for one-time backfill against existing deployments. Closes #310. - webhook.js: probeReachability retries a transient network failure (ECONNRESET/ETIMEDOUT/ECONNABORTED/ECONNREFUSED/ENETUNREACH/ EHOSTUNREACH/EAI_AGAIN) up to WEBHOOK_PROBE_RETRY_ATTEMPTS times per method before giving up; HTTP error responses are never retried since they're a real answer, not a transient failure. Closes #311. - routes/webhooks.js: POST /webhooks now runs the same assertPublicTarget SSRF guard used at delivery/test time *before* calling probeReachability, which previously made an outbound request to a completely unvalidated caller-supplied URL (webhookRepo.create()'s own URL validation only ran after the probe had already fired). Closes #313. - airdrops.js: addRecipients() queues its per-address SADD calls on a single MULTI/EXEC instead of Promise.all-ing N independent commands, so two concurrent addRecipients() calls for the same airdrop can no longer interleave and both observe "newly added" for the same address. Closes #312. --- src/routes/webhooks.js | 9 ++++++ src/services/airdrops.js | 25 ++++++++++++----- src/services/alerts.js | 42 ++++++++++++++++++++++++++-- src/services/webhook.js | 60 +++++++++++++++++++++++++++++++--------- 4 files changed, 114 insertions(+), 22 deletions(-) diff --git a/src/routes/webhooks.js b/src/routes/webhooks.js index 720ca33..7979ba5 100644 --- a/src/routes/webhooks.js +++ b/src/routes/webhooks.js @@ -8,6 +8,7 @@ const deliveryRepo = require("../repositories/deliveryRepository"); const dispatcher = require("../services/webhookDispatcher"); const signatureService = require("../services/webhookSignature"); const { probeReachability } = require("../services/webhook"); +const { assertPublicTarget } = require("../services/ssrfGuard"); const { idempotencyMiddleware } = require("../services/idempotency"); const { requireCsrfHeader } = require("../middleware/csrf"); const buildRateLimit = require("../middleware/rateLimit"); @@ -122,6 +123,14 @@ router.post( ); } + // Issue #313: probeReachability makes a real outbound HTTP request to + // the caller-supplied URL — used to happen here before the URL had + // been validated at all (webhookRepo.create()'s SSRF check only runs + // *after* this point), making the probe itself an SSRF oracle against + // internal/private targets. Validate (format, protocol, and resolved + // address) before probing, the same guard used at delivery time. + await assertPublicTarget(body.url); + const secret = body.secret || signatureService.generateSecret(); const reachability = await probeReachability(body.url); const webhook = await webhookRepo.create({ diff --git a/src/services/airdrops.js b/src/services/airdrops.js index 5dfa99f..839b4c0 100644 --- a/src/services/airdrops.js +++ b/src/services/airdrops.js @@ -261,13 +261,24 @@ async function addRecipients(airdropId, recipients) { const redis = cache.getClient(); const addresses = recipients.map((r) => r.address); - // SADD returns 1 for each newly added member, 0 for duplicates. - // By comparing the added count against the total we identify which - // addresses were already stored from the initial POST /airdrops body - // or a prior POST /airdrops/:id/recipients call. - const addedCounts = await Promise.all( - addresses.map((addr) => redis.sadd(recipientAddressSetKey(airdropId), addr)), - ); + // SADD returns 1 for each newly added member, 0 for duplicates. Queued on + // a single MULTI/EXEC (issue #312) instead of Promise.all-ing N separate + // SADD calls: the separate calls were each an independent round trip that + // could interleave with another concurrent request's SADDs on the same + // key, so two overlapping addRecipients() calls for the same address + // could each observe "newly added" (both read 1) and double-append it to + // the recipients list below. MULTI/EXEC runs the whole batch as one + // atomic, uninterleaved unit, and as a bonus is a single round trip + // instead of N. + const multi = redis.multi(); + for (const addr of addresses) { + multi.sadd(recipientAddressSetKey(airdropId), addr); + } + const results = await multi.exec(); + const addedCounts = results.map(([err, count]) => { + if (err) throw err; + return count; + }); const newAddresses = []; const duplicates = []; diff --git a/src/services/alerts.js b/src/services/alerts.js index ff4cde1..754d070 100644 --- a/src/services/alerts.js +++ b/src/services/alerts.js @@ -15,6 +15,16 @@ function alertKey(id) { return `alert:${id}`; } +// Secondary index (issue #310): evaluateForAsset only ever needs the alerts +// for one asset, but used to ZREVRANGE the *entire* IDS_KEY sorted set and +// fetch + filter every alert in the system on every price tick, regardless +// of how many assets it actually watches. Kept in sync with IDS_KEY by +// create()/remove() so evaluateForAssetInner can look up just this asset's +// ids directly instead of scanning everything. +function assetIdsKey(asset) { + return `alerts:asset:${asset.toUpperCase()}`; +} + function generateId() { return `alrt_${crypto.randomUUID().replace(/-/g, '').slice(0, 16)}`; } @@ -63,6 +73,7 @@ async function create(data) { await redis.multi() .set(alertKey(id), JSON.stringify(alert)) .zadd(IDS_KEY, Date.now(), id) + .sadd(assetIdsKey(alert.asset), id) .exec(); return alert; @@ -100,6 +111,7 @@ async function remove(id) { await redis.multi() .del(alertKey(id)) .zrem(IDS_KEY, id) + .srem(assetIdsKey(existing.asset), id) .exec(); return existing; } @@ -143,10 +155,13 @@ async function evaluateForAsset(asset, priceUsd) { async function evaluateForAssetInner(asset, priceUsd) { const redis = cache.getClient(); - const ids = await redis.zrevrange(IDS_KEY, 0, -1); + const ids = await redis.smembers(assetIdsKey(asset)); for (const id of ids) { const alert = await cache.get(alertKey(id)); + // Index and record can drift (e.g. a record written by an older + // version before this index existed) — fall back to the asset check + // rather than trusting the index blindly. if (!alert || alert.asset !== asset.toUpperCase()) continue; if (!isTriggered(alert, priceUsd)) continue; @@ -207,4 +222,27 @@ async function evaluateAll() { } } -module.exports = { create, list, listPaginated, remove, evaluateForAsset, evaluateAll }; +// One-time backfill for the assetIdsKey() index (issue #310): populates it +// for alerts that were created before this index existed, so +// evaluateForAssetInner doesn't silently stop seeing them. Safe to run +// repeatedly (SADD is idempotent); not wired into any automatic startup +// path since it's an O(n) full scan — meant to be run once, manually, by an +// operator against an existing deployment's pre-existing alert set. +async function backfillAssetIndex() { + const redis = cache.getClient(); + const ids = await redis.zrevrange(IDS_KEY, 0, -1); + const multi = redis.multi(); + let queued = 0; + for (const id of ids) { + const alert = await cache.get(alertKey(id)); + if (!alert) continue; + multi.sadd(assetIdsKey(alert.asset), id); + queued += 1; + } + if (queued > 0) await multi.exec(); + return queued; +} + +module.exports = { + create, list, listPaginated, remove, evaluateForAsset, evaluateAll, backfillAssetIndex, +}; diff --git a/src/services/webhook.js b/src/services/webhook.js index fbb5955..05a7310 100644 --- a/src/services/webhook.js +++ b/src/services/webhook.js @@ -47,26 +47,60 @@ async function sendSignedRequest(webhookUrl, secret, payload, options = {}) { } } +const PROBE_RETRY_ATTEMPTS = parseInt(process.env.WEBHOOK_PROBE_RETRY_ATTEMPTS, 10) || 2; +const PROBE_RETRY_DELAY_MS = parseInt(process.env.WEBHOOK_PROBE_RETRY_DELAY_MS, 10) || 250; + +function isTransientProbeError(err) { + // A response was received (even a 4xx/5xx) — that's a real answer from + // the target, not a transient failure, so it's handled by the caller via + // response.status and never reaches this function. + const code = err?.code; + return [ + 'ECONNRESET', 'ETIMEDOUT', 'ECONNABORTED', 'ECONNREFUSED', + 'ENETUNREACH', 'EHOSTUNREACH', 'EAI_AGAIN', + ].includes(code); +} + +function sleep(ms) { + return new Promise((resolve) => setTimeout(resolve, ms)); +} + +// Issue #311: probeReachability previously gave up on the very first +// network-level failure of each method (e.g. a connection reset or DNS +// hiccup), even though such errors are frequently transient and a +// registration-time probe against a cold/just-deployed receiver is exactly +// the situation where a brief retry is most likely to flip an unreachable +// result into a reachable one. HTTP error responses (4xx/5xx) are not +// retried — they're a real answer, not a transient failure. async function probeReachability(webhookUrl, options = {}) { const timeoutMs = options.timeoutMs || 3000; + const retryAttempts = options.retryAttempts ?? PROBE_RETRY_ATTEMPTS; + const retryDelayMs = options.retryDelayMs ?? PROBE_RETRY_DELAY_MS; const lastError = { message: 'No response received' }; for (const method of ['head', 'get']) { - try { - const response = await axios[method](webhookUrl, { - headers: getRequestIdHeaders(), - timeout: timeoutMs, - validateStatus: () => true, - }); + for (let attempt = 0; attempt <= retryAttempts; attempt++) { + try { + const response = await axios[method](webhookUrl, { + headers: getRequestIdHeaders(), + timeout: timeoutMs, + validateStatus: () => true, + }); - if (response && response.status >= 200 && response.status < 400) { - return { reachable: true, status: response.status, method }; - } - if (response && response.status) { - return { reachable: false, status: response.status, method, error: `Target responded with HTTP ${response.status}` }; + if (response && response.status >= 200 && response.status < 400) { + return { reachable: true, status: response.status, method }; + } + if (response && response.status) { + return { reachable: false, status: response.status, method, error: `Target responded with HTTP ${response.status}` }; + } + } catch (err) { + lastError.message = err?.message || 'Request failed'; + if (attempt < retryAttempts && isTransientProbeError(err)) { + await sleep(retryDelayMs); + continue; + } + break; } - } catch (err) { - lastError.message = err?.message || 'Request failed'; } }