Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
9 changes: 9 additions & 0 deletions src/routes/webhooks.js
Original file line number Diff line number Diff line change
Expand Up @@ -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");
Expand Down Expand Up @@ -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({
Expand Down
25 changes: 18 additions & 7 deletions src/services/airdrops.js
Original file line number Diff line number Diff line change
Expand Up @@ -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 = [];
Expand Down
42 changes: 40 additions & 2 deletions src/services/alerts.js
Original file line number Diff line number Diff line change
Expand Up @@ -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)}`;
}
Expand Down Expand Up @@ -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;
Expand Down Expand Up @@ -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;
}
Expand Down Expand Up @@ -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;
Expand Down Expand Up @@ -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,
};
60 changes: 47 additions & 13 deletions src/services/webhook.js
Original file line number Diff line number Diff line change
Expand Up @@ -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';
}
}

Expand Down