From 8e2abcd0282e21cc7f31352c45671e1b942a19fd Mon Sep 17 00:00:00 2001 From: ThunderBlaze Date: Sat, 18 Apr 2026 09:59:41 +0530 Subject: [PATCH 1/2] feat(relay): add maxTotalEventsPerReq to cap aggregate events across filters --- src/DBQuery.h | 7 +++++++ src/QueryScheduler.h | 1 + src/apps/relay/golpe.yaml | 3 +++ strfry.conf | 3 +++ 4 files changed, 14 insertions(+) diff --git a/src/DBQuery.h b/src/DBQuery.h index 7e91f4bd..2828c900 100644 --- a/src/DBQuery.h +++ b/src/DBQuery.h @@ -290,6 +290,7 @@ struct DBQuery : NonCopyable { flat_hash_set sentEventsCurr; uint64_t lastWorkChecked = 0; + uint64_t maxTotalEvents = 0; uint64_t currScanTime = 0; uint64_t currScanSaveRestores = 0; uint64_t totalTime = 0; @@ -319,6 +320,7 @@ struct DBQuery : NonCopyable { } sentEventsCurr.insert(levId); + if (maxTotalEvents > 0 && sentEventsFull.size() >= maxTotalEvents) return true; return sentEventsCurr.size() >= f.limit; }, [&](uint64_t approxWork){ if (approxWork > lastWorkChecked + 2'000) { @@ -355,6 +357,11 @@ struct DBQuery : NonCopyable { currScanTime = 0; currScanSaveRestores = 0; + + if (maxTotalEvents > 0 && sentEventsFull.size() >= maxTotalEvents) { + LW << "[" << sub.connId << "] REQ='" << sub.subId.sv() << "' hit maxTotalEventsPerReq limit (" << maxTotalEvents << ")"; + break; + } } if (logMetrics) { diff --git a/src/QueryScheduler.h b/src/QueryScheduler.h index 72722c82..47626861 100644 --- a/src/QueryScheduler.h +++ b/src/QueryScheduler.h @@ -33,6 +33,7 @@ struct QueryScheduler : NonCopyable { } DBQuery *q = new DBQuery(sub); + q->maxTotalEvents = cfg().relay__maxTotalEventsPerReq; connQueries.try_emplace(q->sub.subId, q); running.push_front(q); diff --git a/src/apps/relay/golpe.yaml b/src/apps/relay/golpe.yaml index 864d117b..dece6bbe 100644 --- a/src/apps/relay/golpe.yaml +++ b/src/apps/relay/golpe.yaml @@ -73,6 +73,9 @@ config: - name: relay__maxFilterLimit desc: "Maximum records that can be returned per filter" default: 500 + - name: relay__maxTotalEventsPerReq + desc: "Maximum total events returned across all filters in a REQ (0 = unlimited)" + default: 0 - name: relay__maxTagsPerFilter desc: "Maximum number of tag filters allowed per filter (O(N^2) in matching, so keep it small)" default: 3 diff --git a/strfry.conf b/strfry.conf index 6f4cdcc9..6de8fa42 100644 --- a/strfry.conf +++ b/strfry.conf @@ -110,6 +110,9 @@ relay { # Maximum records that can be returned per filter maxFilterLimit = 500 + # Maximum total events returned across all filters in a REQ (0 = unlimited) + maxTotalEventsPerReq = 0 + # Maximum number of tag filters allowed per filter maxTagsPerFilter = 3 From 72b138c77ca90e075ebb9d2fa4bd40043148df5e Mon Sep 17 00:00:00 2001 From: ThunderBlaze Date: Thu, 23 Apr 2026 16:22:12 +0530 Subject: [PATCH 2/2] relay: refactor to send CLOSED and set default to 2000 --- src/DBQuery.h | 2 ++ src/QueryScheduler.h | 4 ++-- src/apps/relay/RelayNegentropy.cpp | 2 +- src/apps/relay/RelayReqWorker.cpp | 7 ++++++- src/apps/relay/RelayServer.h | 10 +++++++--- src/apps/relay/golpe.yaml | 2 +- strfry.conf | 2 +- 7 files changed, 20 insertions(+), 9 deletions(-) diff --git a/src/DBQuery.h b/src/DBQuery.h index 2828c900..65366f8d 100644 --- a/src/DBQuery.h +++ b/src/DBQuery.h @@ -291,6 +291,7 @@ struct DBQuery : NonCopyable { uint64_t lastWorkChecked = 0; uint64_t maxTotalEvents = 0; + bool hitMaxTotalEvents = false; uint64_t currScanTime = 0; uint64_t currScanSaveRestores = 0; uint64_t totalTime = 0; @@ -360,6 +361,7 @@ struct DBQuery : NonCopyable { if (maxTotalEvents > 0 && sentEventsFull.size() >= maxTotalEvents) { LW << "[" << sub.connId << "] REQ='" << sub.subId.sv() << "' hit maxTotalEventsPerReq limit (" << maxTotalEvents << ")"; + hitMaxTotalEvents = true; break; } } diff --git a/src/QueryScheduler.h b/src/QueryScheduler.h index 47626861..99ba2d69 100644 --- a/src/QueryScheduler.h +++ b/src/QueryScheduler.h @@ -6,7 +6,7 @@ struct QueryScheduler : NonCopyable { std::function onEvent; std::function &levIds)> onEventBatch; - std::function onComplete; + std::function onComplete; // If false, then levIds returned to above callbacks can be stale (because they were deleted) // If false, then onEvent's eventPayload will always be "" @@ -102,7 +102,7 @@ struct QueryScheduler : NonCopyable { auto connId = q->sub.connId; removeSub(connId, q->sub.subId); - if (onComplete) onComplete(txn, q->sub, q->sentEventsFull.size()); + if (onComplete) onComplete(txn, q->sub, q->sentEventsFull.size(), q->hitMaxTotalEvents); delete q; } else { diff --git a/src/apps/relay/RelayNegentropy.cpp b/src/apps/relay/RelayNegentropy.cpp index 65dbf8de..d35e5773 100644 --- a/src/apps/relay/RelayNegentropy.cpp +++ b/src/apps/relay/RelayNegentropy.cpp @@ -137,7 +137,7 @@ void RelayServer::runNegentropy(ThreadPool::Thread &thr) { } }; - queries.onComplete = [&](lmdb::txn &txn, Subscription &sub, uint64_t){ + queries.onComplete = [&](lmdb::txn &txn, Subscription &sub, uint64_t, bool){ auto *userView = views.findView(sub.connId, sub.subId); if (!userView) return; diff --git a/src/apps/relay/RelayReqWorker.cpp b/src/apps/relay/RelayReqWorker.cpp index 5478c522..e185ac7f 100644 --- a/src/apps/relay/RelayReqWorker.cpp +++ b/src/apps/relay/RelayReqWorker.cpp @@ -11,7 +11,7 @@ void RelayServer::runReqWorker(ThreadPool::Thread &thr) { sendEvent(sub.connId, sub.subId, decodeEventPayload(txn, decomp, eventPayload, nullptr, nullptr)); }; - queries.onComplete = [&](lmdb::txn &, Subscription &sub, uint64_t total){ + queries.onComplete = [&](lmdb::txn &, Subscription &sub, uint64_t total, bool hitMaxTotalEvents){ if (sub.countOnly) { bool limited = false; @@ -27,6 +27,11 @@ void RelayServer::runReqWorker(ThreadPool::Thread &thr) { if (limited) countBody["limited"] = true; sendToConn(sub.connId, tao::json::to_string(tao::json::value::array({ "COUNT", sub.subId.str(), countBody }))); + } else if (hitMaxTotalEvents) { + // Per NIP-01, terminate with CLOSED + a machine-readable "rate-limited:" + // reason so historical-sync clients don't mistake a truncated result + // for a complete one. + sendClosed(sub.connId, sub.subId.str(), "rate-limited: max aggregate events reached for request"); } else { PROM_INC_RELAY_MSG("EOSE"); sendToConn(sub.connId, tao::json::to_string(tao::json::value::array({ "EOSE", sub.subId.str() }))); diff --git a/src/apps/relay/RelayServer.h b/src/apps/relay/RelayServer.h index d8685ca5..e2a1c3cd 100644 --- a/src/apps/relay/RelayServer.h +++ b/src/apps/relay/RelayServer.h @@ -254,14 +254,18 @@ struct RelayServer { hubTrigger->send(); } - void sendClosedError(uint64_t connId, const std::string &subId, std::string &&payload) { + void sendClosed(uint64_t connId, const std::string &subId, std::string &&reason) { PROM_INC_RELAY_MSG("CLOSED"); - LI << "sending closed to [" << connId << "]: " << payload; - auto reply = tao::json::value::array({ "CLOSED", subId, std::string("ERROR: ") + payload }); + LI << "sending closed to [" << connId << "]: " << reason; + auto reply = tao::json::value::array({ "CLOSED", subId, std::move(reason) }); tpWebsocket.dispatch(0, MsgWebsocket{MsgWebsocket::Send{connId, std::move(tao::json::to_string(reply))}}); hubTrigger->send(); } + void sendClosedError(uint64_t connId, const std::string &subId, std::string &&payload) { + sendClosed(connId, subId, std::string("ERROR: ") + payload); + } + void sendOKResponse(uint64_t connId, std::string_view eventIdHex, bool written, std::string_view message) { PROM_INC_RELAY_MSG("OK"); auto reply = tao::json::value::array({ "OK", eventIdHex, written, message }); diff --git a/src/apps/relay/golpe.yaml b/src/apps/relay/golpe.yaml index dece6bbe..c6b632cf 100644 --- a/src/apps/relay/golpe.yaml +++ b/src/apps/relay/golpe.yaml @@ -75,7 +75,7 @@ config: default: 500 - name: relay__maxTotalEventsPerReq desc: "Maximum total events returned across all filters in a REQ (0 = unlimited)" - default: 0 + default: 2000 - name: relay__maxTagsPerFilter desc: "Maximum number of tag filters allowed per filter (O(N^2) in matching, so keep it small)" default: 3 diff --git a/strfry.conf b/strfry.conf index 6de8fa42..2db14f0e 100644 --- a/strfry.conf +++ b/strfry.conf @@ -111,7 +111,7 @@ relay { maxFilterLimit = 500 # Maximum total events returned across all filters in a REQ (0 = unlimited) - maxTotalEventsPerReq = 0 + maxTotalEventsPerReq = 2000 # Maximum number of tag filters allowed per filter maxTagsPerFilter = 3