From 470ea9c1508099b618dd673c9130a7feb223e4d8 Mon Sep 17 00:00:00 2001 From: codemaster1104 Date: Tue, 14 Apr 2026 00:13:12 +0530 Subject: [PATCH] Add IP address logging in relay components --- TODO | 1 - src/apps/relay/RelayIngester.cpp | 26 +++++++++++++------------- src/apps/relay/RelayNegentropy.cpp | 14 +++++++------- src/apps/relay/RelayServer.h | 17 +++++++++++++++-- src/apps/relay/RelayWebsocket.cpp | 10 ++++++++++ src/apps/relay/RelayWriter.cpp | 10 +++++----- 6 files changed, 50 insertions(+), 28 deletions(-) diff --git a/TODO b/TODO index e0e5a19b..a80a403a 100644 --- a/TODO +++ b/TODO @@ -14,7 +14,6 @@ rate limits (maybe not needed now that we have plugins?) event writes per second per ip max connections per ip (nginx?) max bandwidth up/down (nginx?) - log IP address in sendNoticeError and elsewhere where it makes sense ? events that contain IP/pubkey/etc block-lists in their contents ? limit on total number of events from a DBScan, not just per filter ? time limit on DBScan diff --git a/src/apps/relay/RelayIngester.cpp b/src/apps/relay/RelayIngester.cpp index d31c72b5..e337aa71 100644 --- a/src/apps/relay/RelayIngester.cpp +++ b/src/apps/relay/RelayIngester.cpp @@ -18,7 +18,7 @@ void RelayServer::runIngester(ThreadPool::Thread &thr) { if (msg->payload.starts_with('[')) { auto payload = tao::json::from_string(msg->payload); - if (cfg().relay__logging__dumpInAll) LI << "[" << msg->connId << "] dumpInAll: " << msg->payload; + if (cfg().relay__logging__dumpInAll) LI << "[" << msg->connId << " " << renderIP(msg->ipAddr) << "] dumpInAll: " << msg->payload; auto &arr = jsonGetArray(payload, "message is not an array"); if (arr.size() < 2) throw herr("too few array elements"); @@ -27,18 +27,18 @@ void RelayServer::runIngester(ThreadPool::Thread &thr) { if (cmd == "EVENT") { PROM_INC_CLIENT_MSG(cmd); - if (cfg().relay__logging__dumpInEvents) LI << "[" << msg->connId << "] dumpInEvent: " << msg->payload; + if (cfg().relay__logging__dumpInEvents) LI << "[" << msg->connId << " " << renderIP(msg->ipAddr) << "] dumpInEvent: " << msg->payload; try { ingesterProcessEvent(txn, rsctx, msg->connId, msg->ipAddr, arr[1], writerMsgs); } catch (std::exception &e) { sendOKResponse(msg->connId, arr[1].is_object() && arr[1].at("id").is_string() ? arr[1].at("id").get_string() : "?", false, std::string("invalid: ") + e.what()); - if (cfg().relay__logging__invalidEvents) LI << "Rejected invalid event: " << e.what(); + if (cfg().relay__logging__invalidEvents) LI << "[" << msg->connId << " " << renderIP(msg->ipAddr) << "] Rejected invalid event: " << e.what(); } } else if (cmd == "AUTH") { PROM_INC_CLIENT_MSG(cmd); - if (cfg().relay__logging__dumpInAll) LI << "[" << msg->connId << "] dumpInAuth: " << msg->payload; + if (cfg().relay__logging__dumpInAll) LI << "[" << msg->connId << " " << renderIP(msg->ipAddr) << "] dumpInAuth: " << msg->payload; try { ingesterProcessAuth(rsctx, msg->connId, arr[1]); @@ -47,7 +47,7 @@ void RelayServer::runIngester(ThreadPool::Thread &thr) { } } else if (cmd == "REQ" || cmd == "COUNT") { PROM_INC_CLIENT_MSG(cmd); - if (cfg().relay__logging__dumpInReqs) LI << "[" << msg->connId << "] dumpInReq: " << msg->payload; + if (cfg().relay__logging__dumpInReqs) LI << "[" << msg->connId << " " << renderIP(msg->ipAddr) << "] dumpInReq: " << msg->payload; std::string subIdStr; @@ -59,7 +59,7 @@ void RelayServer::runIngester(ThreadPool::Thread &thr) { } } else if (cmd == "CLOSE") { PROM_INC_CLIENT_MSG(cmd); - if (cfg().relay__logging__dumpInReqs) LI << "[" << msg->connId << "] dumpInReq: " << msg->payload; + if (cfg().relay__logging__dumpInReqs) LI << "[" << msg->connId << " " << renderIP(msg->ipAddr) << "] dumpInReq: " << msg->payload; try { ingesterProcessClose(txn, msg->connId, arr); @@ -120,7 +120,7 @@ void RelayServer::ingesterProcessEvent(lmdb::txn &txn, RelayServerCtx &rsctx, ui if (packed.kind() == 6 || packed.kind() == 16) { if (origJson.at("content").get_string().find("[\"-\"]") != std::string::npos) { auto idHex = to_hex(packed.id()); - LI << "Repost embedded a protected event, blocking: " << idHex; + LI << "[" << connId << " " << renderIP(ipAddr) << "] Repost embedded a protected event, blocking: " << idHex; sendOKResponse(connId, idHex, false, "blocked: reposts can't embed protected events"); return; } @@ -143,14 +143,14 @@ void RelayServer::ingesterProcessEvent(lmdb::txn &txn, RelayServerCtx &rsctx, ui auto idHex = to_hex(packed.id()); if (!cfg().relay__auth__enabled) { - LI << "[" << connId << "] Protected event and auth disabled, rejecting: " << idHex; + LI << "[" << connId << " " << renderIP(ipAddr) << "] Protected event and auth disabled, rejecting: " << idHex; sendOKResponse(connId, idHex, false, "blocked: event marked as protected"); return; } if (cfg().relay__auth__serviceUrl.empty()) { // If we don't have a serviceUrl, just fail - LI << "[" << connId << "] Protected event and no serviceUrl configured, rejecting: " << idHex; + LI << "[" << connId << " " << renderIP(ipAddr) << "] Protected event and no serviceUrl configured, rejecting: " << idHex; sendOKResponse(connId, idHex, false, "blocked: event marked as protected"); return; } @@ -161,7 +161,7 @@ void RelayServer::ingesterProcessEvent(lmdb::txn &txn, RelayServerCtx &rsctx, ui auto challenge = rsctx.challengeGenerator.get(); rsctx.connIdToAuthStatus.emplace(connId, challenge); - LI << "[" << connId << "] Protected event, requesting AUTH: " << idHex; + LI << "[" << connId << " " << renderIP(ipAddr) << "] Protected event, requesting AUTH: " << idHex; sendAuthChallenge(connId, challenge); sendOKResponse(connId, idHex, false, "auth-required: event marked as protected"); return; @@ -171,7 +171,7 @@ void RelayServer::ingesterProcessEvent(lmdb::txn &txn, RelayServerCtx &rsctx, ui if (!as.isAuthed()) { // not authenticated - LI << "[" << connId << "] Protected event, AUTH already requested: " << idHex; + LI << "[" << connId << " " << renderIP(ipAddr) << "] Protected event, AUTH already requested: " << idHex; sendOKResponse(connId, idHex, false, "auth-required: event marked as protected"); return; } else if (as.authed != packed.pubkey()) { @@ -189,7 +189,7 @@ void RelayServer::ingesterProcessEvent(lmdb::txn &txn, RelayServerCtx &rsctx, ui auto existing = lookupEventById(txn, packed.id()); if (existing) { auto hexId = to_hex(packed.id()); - LI << "[" << connId << "] Duplicate event, skipping: " << hexId; + LI << "[" << connId << " " << renderIP(ipAddr) << "] Duplicate event, skipping: " << hexId; sendOKResponse(connId, hexId, true, "duplicate: have this event"); return; } @@ -270,7 +270,7 @@ void RelayServer::ingesterProcessAuth(RelayServerCtx &rsctx, uint64_t connId, co // set the connection as authenticated with this pubkey as.authed = packed.pubkey(); - LI << "[" << connId << "] AUTHed as " << to_hex(packed.pubkey()); + LI << "[" << connId << " " << getIpForConn(connId) << "] AUTHed as " << to_hex(packed.pubkey()); sendOKResponse(connId, to_hex(packed.id()), true, "successfully authenticated"); } diff --git a/src/apps/relay/RelayNegentropy.cpp b/src/apps/relay/RelayNegentropy.cpp index 65dbf8de..4d56c28b 100644 --- a/src/apps/relay/RelayNegentropy.cpp +++ b/src/apps/relay/RelayNegentropy.cpp @@ -99,7 +99,7 @@ void RelayServer::runNegentropy(ThreadPool::Thread &thr) { Negentropy ne(storage, 500'000); resp = ne.reconcile(msg); } catch (std::exception &e) { - LI << "[" << connId << "] Error parsing negentropy message: " << e.what(); + LI << "[" << connId << " " << getIpForConn(connId) << "] Error parsing negentropy message: " << e.what(); PROM_INC_RELAY_MSG("NEG-ERR"); sendToConn(connId, tao::json::to_string(tao::json::value::array({ @@ -112,7 +112,7 @@ void RelayServer::runNegentropy(ThreadPool::Thread &thr) { return; } - LI << "[" << connId << "] negentropy SEND session=" << subId.sv() << " bytesOut=" << resp.size(); + LI << "[" << connId << " " << getIpForConn(connId) << "] negentropy SEND session=" << subId.sv() << " bytesOut=" << resp.size(); PROM_INC_RELAY_MSG("NEG-MSG"); sendToConn(connId, tao::json::to_string(tao::json::value::array({ @@ -144,11 +144,11 @@ void RelayServer::runNegentropy(ThreadPool::Thread &thr) { auto *view = std::get_if(userView); if (!view) throw herr("bad variant, expected memory view"); - LI << "[" << sub.connId << "] negentropy QUERY matched " << view->levIds.size() << " events in " + LI << "[" << sub.connId << " " << getIpForConn(sub.connId) << "] negentropy QUERY matched " << view->levIds.size() << " events in " << (hoytech::curr_time_us() - view->startTime) << "us"; if (view->levIds.size() > cfg().relay__negentropy__maxSyncEvents) { - LI << "[" << sub.connId << "] negentropy QUERY size exceeded " << cfg().relay__negentropy__maxSyncEvents; + LI << "[" << sub.connId << " " << getIpForConn(sub.connId) << "] negentropy QUERY size exceeded " << cfg().relay__negentropy__maxSyncEvents; PROM_INC_RELAY_MSG("NEG-ERR"); sendToConn(sub.connId, tao::json::to_string(tao::json::value::array({ @@ -204,7 +204,7 @@ void RelayServer::runNegentropy(ThreadPool::Thread &thr) { return true; }); - LI << "[" << connId << "] negentropy NEW session=" << subId.sv() << " bytesIn=" << msg->negPayload.size() << " tree=" << (treeId ? std::to_string(*treeId) : "NONE"); + LI << "[" << connId << " " << getIpForConn(connId) << "] negentropy NEW session=" << subId.sv() << " bytesIn=" << msg->negPayload.size() << " tree=" << (treeId ? std::to_string(*treeId) : "NONE"); if (treeId) { negentropy::storage::BTreeLMDB storage(txn, negentropyDbi, *treeId); @@ -242,7 +242,7 @@ void RelayServer::runNegentropy(ThreadPool::Thread &thr) { continue; } - LI << "[" << msg->connId << "] negentropy RECV session=" << msg->subId.sv() << " bytesIn=" << msg->negPayload.size(); + LI << "[" << msg->connId << " " << getIpForConn(msg->connId) << "] negentropy RECV session=" << msg->subId.sv() << " bytesIn=" << msg->negPayload.size(); if (auto *view = std::get_if(userView)) { if (!view->storageVector.sealed) { @@ -259,7 +259,7 @@ void RelayServer::runNegentropy(ThreadPool::Thread &thr) { handleReconcile(msg->connId, msg->subId, subStorage, msg->negPayload); } } else if (auto msg = std::get_if(&newMsg.msg)) { - LI << "[" << msg->connId << "] negentropy CLOSE session=" << msg->subId.sv(); + LI << "[" << msg->connId << " " << getIpForConn(msg->connId) << "] negentropy CLOSE session=" << msg->subId.sv(); queries.removeSub(msg->connId, msg->subId); views.removeView(msg->connId, msg->subId); diff --git a/src/apps/relay/RelayServer.h b/src/apps/relay/RelayServer.h index d8685ca5..94c96fbc 100644 --- a/src/apps/relay/RelayServer.h +++ b/src/apps/relay/RelayServer.h @@ -3,6 +3,7 @@ #include #include #include +#include #include #include @@ -213,6 +214,18 @@ struct RelayServer { void runSignalHandler(); + // Connection tracking + + mutable std::shared_mutex connIdToIpMutex; + flat_hash_map connIdToIp; + + std::string getIpForConn(uint64_t connId) const { + std::shared_lock lock(connIdToIpMutex); + auto it = connIdToIp.find(connId); + if (it != connIdToIp.end()) return renderIP(it->second); + return "unknown"; + } + // Utils (can be called by any thread) void sendToConn(uint64_t connId, std::string &&payload) { @@ -248,7 +261,7 @@ struct RelayServer { void sendNoticeError(uint64_t connId, std::string &&payload) { PROM_INC_RELAY_MSG("NOTICE"); - LI << "sending error to [" << connId << "]: " << payload; + LI << "sending error to [" << connId << " " << getIpForConn(connId) << "]: " << payload; auto reply = tao::json::value::array({ "NOTICE", std::string("ERROR: ") + payload }); tpWebsocket.dispatch(0, MsgWebsocket{MsgWebsocket::Send{connId, std::move(tao::json::to_string(reply))}}); hubTrigger->send(); @@ -256,7 +269,7 @@ struct RelayServer { void sendClosedError(uint64_t connId, const std::string &subId, std::string &&payload) { PROM_INC_RELAY_MSG("CLOSED"); - LI << "sending closed to [" << connId << "]: " << payload; + LI << "sending closed to [" << connId << " " << getIpForConn(connId) << "]: " << payload; auto reply = tao::json::value::array({ "CLOSED", subId, std::string("ERROR: ") + payload }); tpWebsocket.dispatch(0, MsgWebsocket{MsgWebsocket::Send{connId, std::move(tao::json::to_string(reply))}}); hubTrigger->send(); diff --git a/src/apps/relay/RelayWebsocket.cpp b/src/apps/relay/RelayWebsocket.cpp index 8cf7820f..5cbb77eb 100644 --- a/src/apps/relay/RelayWebsocket.cpp +++ b/src/apps/relay/RelayWebsocket.cpp @@ -267,6 +267,11 @@ void RelayServer::runWebsocket(ThreadPool::Thread &thr) { if (c->ipAddr.size() == 0) c->ipAddr = ws->getAddressBytes(); + { + std::unique_lock lock(connIdToIpMutex); + connIdToIp[connId] = c->ipAddr; + } + ws->setUserData((void*)c); connIdToConnection.emplace(connId, c); @@ -300,6 +305,11 @@ void RelayServer::runWebsocket(ThreadPool::Thread &thr) { tpIngester.dispatch(connId, MsgIngester{MsgIngester::CloseConn{connId}}); + { + std::unique_lock lock(connIdToIpMutex); + connIdToIp.erase(connId); + } + connIdToConnection.erase(connId); delete c; diff --git a/src/apps/relay/RelayWriter.cpp b/src/apps/relay/RelayWriter.cpp index 61e20b10..8b71e0f5 100644 --- a/src/apps/relay/RelayWriter.cpp +++ b/src/apps/relay/RelayWriter.cpp @@ -49,7 +49,7 @@ void RelayServer::runWriter(ThreadPool::Thread &thr) { PackedEventView packed(msg->packedStr); auto eventIdHex = to_hex(packed.id()); - if (okMsg.size()) LI << "[" << msg->connId << "] write policy blocked event " << eventIdHex << ": " << okMsg; + if (okMsg.size()) LI << "[" << msg->connId << " " << renderIP(msg->ipAddr) << "] write policy blocked event " << eventIdHex << ": " << okMsg; sendOKResponse(msg->connId, eventIdHex, res == PluginEventSifterResult::ShadowReject, okMsg); } @@ -89,8 +89,10 @@ void RelayServer::runWriter(ThreadPool::Thread &thr) { std::string message; bool written = false; + MsgWriter::AddEvent *addEventMsg = static_cast(newEvent.userData); + if (newEvent.status == EventWriteStatus::Written) { - LI << "Inserted event. id=" << eventIdHex << " levId=" << newEvent.levId; + LI << "[" << addEventMsg->connId << " " << renderIP(addEventMsg->ipAddr) << "] Inserted event. id=" << eventIdHex << " levId=" << newEvent.levId; written = true; } else if (newEvent.status == EventWriteStatus::Duplicate) { message = "duplicate: have this event"; @@ -102,11 +104,9 @@ void RelayServer::runWriter(ThreadPool::Thread &thr) { } if (newEvent.status != EventWriteStatus::Written) { - LI << "Rejected event. " << message << ", id=" << eventIdHex; + LI << "[" << addEventMsg->connId << " " << renderIP(addEventMsg->ipAddr) << "] Rejected event. " << message << ", id=" << eventIdHex; } - MsgWriter::AddEvent *addEventMsg = static_cast(newEvent.userData); - sendOKResponse(addEventMsg->connId, eventIdHex, written, message); } }