diff --git a/Makefile b/Makefile index 97d3a3c3..6b735fb8 100644 --- a/Makefile +++ b/Makefile @@ -34,5 +34,5 @@ src/apps/relay/RelayWebsocket.o: build/StrfryTemplates.h test-subid: build/subid_tests build/subid_tests -build/subid_tests: test/SubIdTests.cpp build/golpe.h +build/subid_tests: test/tests/SubIdTests.cpp build/golpe.h $(CXX) $(CXXFLAGS) $(INCS) $< -o $@ diff --git a/docs/architecture.md b/docs/architecture.md index 570963b7..54aad9fb 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -154,7 +154,7 @@ The query engine is the most complicated part of the relay, so there is a differ To bootstrap the tests, we load in a set of [real-world nostr events](https://wiki.wellorder.net/wiki/nostr-datasets/). -There is a simple but inefficient filter implementation in `test/dumbFilter.pl` that can be used to check if an event matches a filter. In a loop, we randomly generate a complicated filter group and pipe the entire DB's worth of events through the dumb filter and record which events it matched. Next, we perform the query using strfry's query engine (using a `strfry scan`) and ensure it matches. This gives us confidence that querying for "old" records in the DB will be performed correctly. +There is a simple but inefficient filter implementation in `test/utils/dumbFilter.pl` that can be used to check if an event matches a filter. In a loop, we randomly generate a complicated filter group and pipe the entire DB's worth of events through the dumb filter and record which events it matched. Next, we perform the query using strfry's query engine (using a `strfry scan`) and ensure it matches. This gives us confidence that querying for "old" records in the DB will be performed correctly. Next, we need to verify that monitoring for "new" records will function also. For this, in a loop we create a set of hundreds of random filters and install them in the monitoring engine. One of which is selected as a sample. The entire DB's worth of events is "posted to the relay" (actually just iterated over in the DB using `strfry monitor`), and we record which events were matched. This is then compared against a full-DB scan using the same query. diff --git a/package.json b/package.json new file mode 100644 index 00000000..121e94cb --- /dev/null +++ b/package.json @@ -0,0 +1,23 @@ +{ + "name": "strfry", + "version": "1.0.0", + "description": "", + "main": "index.js", + "scripts": { + "test": "echo \"Error: no test specified\" && exit 1" + }, + "keywords": [], + "author": "", + "license": "ISC", + "devEngines": { + "packageManager": { + "name": "pnpm", + "version": "^11.22.0", + "onFail": "download" + } + }, + "type": "module", + "dependencies": { + "@nostr/tools": "jsr:^2.24.2" + } +} diff --git a/pnpm-lock.yaml b/pnpm-lock.yaml new file mode 100644 index 00000000..2ede3dde --- /dev/null +++ b/pnpm-lock.yaml @@ -0,0 +1,276 @@ +--- +lockfileVersion: '9.0' + +importers: + + .: + configDependencies: {} + packageManagerDependencies: + '@pnpm/exe': + specifier: ^11.22.0 + version: 11.22.0 + pnpm: + specifier: ^11.22.0 + version: 11.22.0 + +packages: + + '@pnpm/exe@11.22.0': + resolution: {integrity: sha512-B1SGeKm+v9pX9YkzMmrnO2FbgBd8TDwzZ3jSj6J6ThxdyGvI4TsOfMeHARW4Wb25visPpmMDfIXrU76EZdJM4g==} + hasBin: true + + '@pnpm/linux-arm64@11.22.0': + resolution: {integrity: sha512-xzzn3jYG9QaiFZPaHcWM3yX4Fm9UGz5E42mzpbvUEvXYf5O9fwNolenOhGLpRl2qS5u35fjZSvDZbWlOwHMcug==} + cpu: [arm64] + os: [linux] + + '@pnpm/linux-x64@11.22.0': + resolution: {integrity: sha512-isvaPctGinbsM2hsTtRsMarN8Sr5QhXDTNmn8Xv9Lp1PjincCvH2RBbhpo+xYIOxzgls1dtQSdgoORUbfRVALQ==} + cpu: [x64] + os: [linux] + + '@pnpm/linuxstatic-arm64@11.22.0': + resolution: {integrity: sha512-i4J+AQWW0T3JBdXaLsvxu8ZmMG1HLS6JZQESdOR9uBv8Zumb4yg8GYT83eWRLacr6ngVdhEBiC3W4eOG64MFbA==} + cpu: [arm64] + os: [linux] + libc: [musl] + + '@pnpm/linuxstatic-x64@11.22.0': + resolution: {integrity: sha512-QYzk8jhSuSbVthW/OxEOhU3f3zhjpw4HqgedrlwbVNb7btCWn5del2Hb5PsYYka2PbmKq5EtF5IzYN8Yp1CfxQ==} + cpu: [x64] + os: [linux] + libc: [musl] + + '@pnpm/macos-arm64@11.22.0': + resolution: {integrity: sha512-Io8Axk5kutPgMuAfOg1QGj3J0/KpLvH5iiPcmp7up8D4q7BNWT02ndQo3RyGWqrlkf5nwspM41GFpWTShrZ4Aw==} + cpu: [arm64] + os: [darwin] + + '@pnpm/win-arm64@11.22.0': + resolution: {integrity: sha512-QgaRuKGQKov7xW2utPCgDT83fn/PU5cD6HFgib48oTslz7wv26E89d4jMCxL4GRmEL2YCfkYuj5ITj1oVWg5WQ==} + cpu: [arm64] + os: [win32] + + '@pnpm/win-x64@11.22.0': + resolution: {integrity: sha512-iWYsiSwpgxqur+TwsnSoPdOQaEfd4ygB84mb5tJUOgim6PHr1xhKJBndaQzH+At/WkLxLdYN+K8dYc4M7naezg==} + cpu: [x64] + os: [win32] + + '@reflink/reflink-darwin-arm64@0.1.19': + resolution: {integrity: sha512-ruy44Lpepdk1FqDz38vExBY/PVUsjxZA+chd9wozjUH9JjuDT/HEaQYA6wYN9mf041l0yLVar6BCZuWABJvHSA==} + engines: {node: '>= 10'} + cpu: [arm64] + os: [darwin] + + '@reflink/reflink-darwin-x64@0.1.19': + resolution: {integrity: sha512-By85MSWrMZa+c26TcnAy8SDk0sTUkYlNnwknSchkhHpGXOtjNDUOxJE9oByBnGbeuIE1PiQsxDG3Ud+IVV9yuA==} + engines: {node: '>= 10'} + cpu: [x64] + os: [darwin] + + '@reflink/reflink-linux-arm64-gnu@0.1.19': + resolution: {integrity: sha512-7P+er8+rP9iNeN+bfmccM4hTAaLP6PQJPKWSA4iSk2bNvo6KU6RyPgYeHxXmzNKzPVRcypZQTpFgstHam6maVg==} + engines: {node: '>= 10'} + cpu: [arm64] + os: [linux] + libc: [glibc] + + '@reflink/reflink-linux-arm64-musl@0.1.19': + resolution: {integrity: sha512-37iO/Dp6m5DDaC2sf3zPtx/hl9FV3Xze4xoYidrxxS9bgP3S8ALroxRK6xBG/1TtfXKTvolvp+IjrUU6ujIGmA==} + engines: {node: '>= 10'} + cpu: [arm64] + os: [linux] + libc: [musl] + + '@reflink/reflink-linux-x64-gnu@0.1.19': + resolution: {integrity: sha512-jbI8jvuYCaA3MVUdu8vLoLAFqC+iNMpiSuLbxlAgg7x3K5bsS8nOpTRnkLF7vISJ+rVR8W+7ThXlXlUQ93ulkw==} + engines: {node: '>= 10'} + cpu: [x64] + os: [linux] + libc: [glibc] + + '@reflink/reflink-linux-x64-musl@0.1.19': + resolution: {integrity: sha512-e9FBWDe+lv7QKAwtKOt6A2W/fyy/aEEfr0g6j/hWzvQcrzHCsz07BNQYlNOjTfeytrtLU7k449H1PI95jA4OjQ==} + engines: {node: '>= 10'} + cpu: [x64] + os: [linux] + libc: [musl] + + '@reflink/reflink-win32-arm64-msvc@0.1.19': + resolution: {integrity: sha512-09PxnVIQcd+UOn4WAW73WU6PXL7DwGS6wPlkMhMg2zlHHG65F3vHepOw06HFCq+N42qkaNAc8AKIabWvtk6cIQ==} + engines: {node: '>= 10'} + cpu: [arm64] + os: [win32] + + '@reflink/reflink-win32-x64-msvc@0.1.19': + resolution: {integrity: sha512-E//yT4ni2SyhwP8JRjVGWr3cbnhWDiPLgnQ66qqaanjjnMiu3O/2tjCPQXlcGc/DEYofpDc9fvhv6tALQsMV9w==} + engines: {node: '>= 10'} + cpu: [x64] + os: [win32] + + '@reflink/reflink@0.1.19': + resolution: {integrity: sha512-DmCG8GzysnCZ15bres3N5AHCmwBwYgp0As6xjhQ47rAUTUXxJiK+lLUxaGsX3hd/30qUpVElh05PbGuxRPgJwA==} + engines: {node: '>= 10'} + + detect-libc@2.1.2: + resolution: {integrity: sha512-Btj2BOOO83o3WyH59e8MgXsxEQVcarkUOpEYrubB0urwnN10yQ364rsiByU11nZlqWYZm05i/of7io4mzihBtQ==} + engines: {node: '>=8'} + + pnpm@11.22.0: + resolution: {integrity: sha512-H/hwxMYTPf2I+yr8Rt0T1H8JyXlLQ4xv20fKmMrzvBY4HuC+k6CRuOOCTPAfiJ9G19niCRD7C+GrD7W6qA3WIQ==} + engines: {node: '>=22.13'} + hasBin: true + +snapshots: + + '@pnpm/exe@11.22.0': + dependencies: + '@reflink/reflink': 0.1.19 + detect-libc: 2.1.2 + optionalDependencies: + '@pnpm/linux-arm64': 11.22.0 + '@pnpm/linux-x64': 11.22.0 + '@pnpm/linuxstatic-arm64': 11.22.0 + '@pnpm/linuxstatic-x64': 11.22.0 + '@pnpm/macos-arm64': 11.22.0 + '@pnpm/win-arm64': 11.22.0 + '@pnpm/win-x64': 11.22.0 + + '@pnpm/linux-arm64@11.22.0': + optional: true + + '@pnpm/linux-x64@11.22.0': + optional: true + + '@pnpm/linuxstatic-arm64@11.22.0': + optional: true + + '@pnpm/linuxstatic-x64@11.22.0': + optional: true + + '@pnpm/macos-arm64@11.22.0': + optional: true + + '@pnpm/win-arm64@11.22.0': + optional: true + + '@pnpm/win-x64@11.22.0': + optional: true + + '@reflink/reflink-darwin-arm64@0.1.19': + optional: true + + '@reflink/reflink-darwin-x64@0.1.19': + optional: true + + '@reflink/reflink-linux-arm64-gnu@0.1.19': + optional: true + + '@reflink/reflink-linux-arm64-musl@0.1.19': + optional: true + + '@reflink/reflink-linux-x64-gnu@0.1.19': + optional: true + + '@reflink/reflink-linux-x64-musl@0.1.19': + optional: true + + '@reflink/reflink-win32-arm64-msvc@0.1.19': + optional: true + + '@reflink/reflink-win32-x64-msvc@0.1.19': + optional: true + + '@reflink/reflink@0.1.19': + optionalDependencies: + '@reflink/reflink-darwin-arm64': 0.1.19 + '@reflink/reflink-darwin-x64': 0.1.19 + '@reflink/reflink-linux-arm64-gnu': 0.1.19 + '@reflink/reflink-linux-arm64-musl': 0.1.19 + '@reflink/reflink-linux-x64-gnu': 0.1.19 + '@reflink/reflink-linux-x64-musl': 0.1.19 + '@reflink/reflink-win32-arm64-msvc': 0.1.19 + '@reflink/reflink-win32-x64-msvc': 0.1.19 + + detect-libc@2.1.2: {} + + pnpm@11.22.0: {} + +--- +lockfileVersion: '9.0' + +settings: + autoInstallPeers: true + excludeLinksFromLockfile: false + +importers: + + .: + dependencies: + '@nostr/tools': + specifier: jsr:^2.24.2 + version: '@jsr/nostr__tools@2.24.2' + +packages: + + '@jsr/nostr__tools@2.24.2': + resolution: {integrity: sha512-oEIsZtycernGHCDFqKGjuq6OHpjOGL8zrExfS2YlsFam8b2MhjShKC8qeV+Tp18puw+FAHWcvnXYx0wtx7IHeQ==, tarball: https://npm.jsr.io/~/11/@jsr/nostr__tools/2.24.2.tgz} + + '@noble/ciphers@2.1.1': + resolution: {integrity: sha512-bysYuiVfhxNJuldNXlFEitTVdNnYUc+XNJZd7Qm2a5j1vZHgY+fazadNFWFaMK/2vye0JVlxV3gHmC0WDfAOQw==} + engines: {node: '>= 20.19.0'} + + '@noble/curves@2.0.1': + resolution: {integrity: sha512-vs1Az2OOTBiP4q0pwjW5aF0xp9n4MxVrmkFBxc6EKZc6ddYx5gaZiAsZoq0uRRXWbi3AT/sBqn05eRPtn1JCPw==} + engines: {node: '>= 20.19.0'} + + '@noble/hashes@2.0.1': + resolution: {integrity: sha512-XlOlEbQcE9fmuXxrVTXCTlG2nlRXa9Rj3rr5Ue/+tX+nmkgbX720YHh0VR3hBF9xDvwnb8D2shVGOwNx+ulArw==} + engines: {node: '>= 20.19.0'} + + '@scure/base@2.0.0': + resolution: {integrity: sha512-3E1kpuZginKkek01ovG8krQ0Z44E3DHPjc5S2rjJw9lZn3KSQOs8S7wqikF/AH7iRanHypj85uGyxk0XAyC37w==} + + '@scure/bip32@2.0.1': + resolution: {integrity: sha512-4Md1NI5BzoVP+bhyJaY3K6yMesEFzNS1sE/cP+9nuvE7p/b0kx9XbpDHHFl8dHtufcbdHRUUQdRqLIPHN/s7yA==} + + '@scure/bip39@2.0.1': + resolution: {integrity: sha512-PsxdFj/d2AcJcZDX1FXN3dDgitDDTmwf78rKZq1a6c1P1Nan1X/Sxc7667zU3U+AN60g7SxxP0YCVw2H/hBycg==} + + nostr-wasm@0.1.0: + resolution: {integrity: sha512-78BTryCLcLYv96ONU8Ws3Q1JzjlAt+43pWQhIl86xZmWeegYCNLPml7yQ+gG3vR6V5h4XGj+TxO+SS5dsThQIA==} + +snapshots: + + '@jsr/nostr__tools@2.24.2': + dependencies: + '@noble/ciphers': 2.1.1 + '@noble/curves': 2.0.1 + '@noble/hashes': 2.0.1 + '@scure/base': 2.0.0 + '@scure/bip32': 2.0.1 + '@scure/bip39': 2.0.1 + nostr-wasm: 0.1.0 + + '@noble/ciphers@2.1.1': {} + + '@noble/curves@2.0.1': + dependencies: + '@noble/hashes': 2.0.1 + + '@noble/hashes@2.0.1': {} + + '@scure/base@2.0.0': {} + + '@scure/bip32@2.0.1': + dependencies: + '@noble/curves': 2.0.1 + '@noble/hashes': 2.0.1 + '@scure/base': 2.0.0 + + '@scure/bip39@2.0.1': + dependencies: + '@noble/hashes': 2.0.1 + '@scure/base': 2.0.0 + + nostr-wasm@0.1.0: {} diff --git a/src/ReadRestrictor.h b/src/ReadRestrictor.h index 1e6dc73d..7045716a 100644 --- a/src/ReadRestrictor.h +++ b/src/ReadRestrictor.h @@ -49,10 +49,13 @@ struct ReadRestrictor { static bool isFilterAllowedToCount(const NostrFilterGroup &fg, Bytes32 pubkey) { if (restrictedKinds().empty()) return true; + + bool pubkeyIsNull = pubkey.isNull(); + for (const auto &f: fg.filters) { if (!f.kinds) continue; bool hasSomeRestrictedKind = false; - for(size_t i = 0; isize(); ++i) { + for (size_t i = 0; i < f.kinds->size(); ++i) { uint64_t kind = f.kinds->at(i); if (restrictedKinds().contains(kind)) { hasSomeRestrictedKind = true; @@ -60,9 +63,12 @@ struct ReadRestrictor { } } if (hasSomeRestrictedKind) { - if (pubkey.isNull()) { + if (pubkeyIsNull) { return false; } + if (!cfg().relay__auth__restrictReadToInvolvedPubkey) { + continue; + } bool authorScoped = f.authors && allPubkeysMatch(*f.authors, pubkey); bool pScoped = false; if (auto it = f.tags.find('p'); it != f.tags.end()) { @@ -86,7 +92,7 @@ struct ReadRestrictor { // Returns true if the event should be sent to the subscriber static bool shouldSendToSubscriber(const PackedEventView &packed, const Bytes32 &subscriberAuthedPubkey) { - if (!(restrictedKinds().contains(packed.kind()) && cfg().relay__auth__restrictReadToInvolvedPubkey)) { + if (!restrictedKinds().contains(packed.kind())) { return true; } @@ -94,22 +100,21 @@ struct ReadRestrictor { return false; } - Bytes32 recipientPubkey; - bool foundRecipient = false; + if(!cfg().relay__auth__restrictReadToInvolvedPubkey) return true; + + bool involved = subscriberAuthedPubkey == packed.pubkey(); packed.foreachTag([&](char tagName, std::string_view tagVal) { if (tagName == 'p' && tagVal.size() == 32) { - recipientPubkey = Bytes32(tagVal); - foundRecipient = true; - return false; + if (subscriberAuthedPubkey == Bytes32(tagVal)) { + involved = true; + return false; + } } + return true; }); - if (!foundRecipient) { - return false; - } - - return subscriberAuthedPubkey == recipientPubkey || subscriberAuthedPubkey == packed.pubkey(); + return involved; } }; diff --git a/src/apps/relay/RelayIngester.cpp b/src/apps/relay/RelayIngester.cpp index 686e0a5c..2509a8ce 100644 --- a/src/apps/relay/RelayIngester.cpp +++ b/src/apps/relay/RelayIngester.cpp @@ -225,27 +225,34 @@ void RelayServer::ingesterProcessReq(lmdb::txn &txn, RelayServerCtx &rsctx, uint } auto it = rsctx.connIdToAuthSession.find(connId); - bool isAuthed = (it != rsctx.connIdToAuthSession.end()); - bool requiresAuth = false; + bool hasSession = it != rsctx.connIdToAuthSession.end(); + bool isAuthed = hasSession && !it->second.authed.isNull(); + bool shouldRejectReq = false; if (countOnly) { // COUNT can't be filtered per event, so a restricted-kind filter must be // scoped to the client's own pubkeys via authors/#p. - requiresAuth = !ReadRestrictor::isFilterAllowedToCount(filterGroup, isAuthed ? it->second.authed : Bytes32()); + shouldRejectReq = !ReadRestrictor::isFilterAllowedToCount(filterGroup, isAuthed ? it->second.authed : Bytes32()); } else { // if the filter group contains no filter that has a kind that is not restricted, // that means an unauthenticated client won't be able to see anything if (ReadRestrictor::isFilterGroupFullyRestricted(filterGroup)) { - requiresAuth = !isAuthed; + shouldRejectReq = !isAuthed; } } - if (requiresAuth) { - auto challenge = rsctx.challengeGenerator.get(); - rsctx.connIdToAuthSession.emplace(connId, challenge); - LI << "[" << connId << "] Requesting initial AUTH"; - sendAuthChallenge(connId, challenge); - sendClosedError(connId, outSubIdStr, "auth-required: requested filter requires authentication"); + if (shouldRejectReq) { + if (!hasSession) { + auto challenge = rsctx.challengeGenerator.get(); + rsctx.connIdToAuthSession.emplace(connId, challenge); + LI << "[" << connId << "] Requesting initial AUTH"; + sendAuthChallenge(connId, challenge); + sendClosedError(connId, outSubIdStr, "auth-required: requested filter requires authentication"); + } else if (countOnly && isAuthed) { + sendClosedError(connId, outSubIdStr, "count-failed: can only count events you are involved in"); + } else { + sendClosedError(connId, outSubIdStr, "auth-required: requested filter requires authentication"); + } return; } @@ -342,11 +349,15 @@ void RelayServer::ingesterProcessNegentropy(lmdb::txn &txn, RelayServerCtx &rsct // that means an unauthenticated client won't be able to see anything if (ReadRestrictor::isFilterGroupFullyRestricted(filter)) { auto it = rsctx.connIdToAuthSession.find(connId); - if (it == rsctx.connIdToAuthSession.end()) { - auto challenge = rsctx.challengeGenerator.get(); - rsctx.connIdToAuthSession.emplace(connId, challenge); - LI << "[" << connId << "] Requesting initial AUTH"; - sendAuthChallenge(connId, challenge); + bool hasSession = it != rsctx.connIdToAuthSession.end(); + bool isAuthed = hasSession && !it->second.authed.isNull(); + if (!isAuthed) { + if (!hasSession) { + auto challenge = rsctx.challengeGenerator.get(); + rsctx.connIdToAuthSession.emplace(connId, challenge); + LI << "[" << connId << "] Requesting initial AUTH"; + sendAuthChallenge(connId, challenge); + } PROM_INC_RELAY_MSG("NEG-ERR"); sendToConn(connId, tao::json::to_string(tao::json::value::array({ "NEG-ERR", diff --git a/src/apps/relay/golpe.yaml b/src/apps/relay/golpe.yaml index 33ce2051..553d0b44 100644 --- a/src/apps/relay/golpe.yaml +++ b/src/apps/relay/golpe.yaml @@ -22,7 +22,7 @@ config: desc: "External relay URL (beginning with wss://). Required in order to validate challenge responses." default: "" - name: relay__auth__restrictedReadKinds - desc: "Comma-separated list of event kinds that require NIP-42 AUTH to read via REQ/COUNT/NEG-OPEN. A filter with no 'kinds' field is treated as restricted. Example for DMs: '4,1059'." + desc: "Comma-separated list of event kinds that require NIP-42 AUTH to read via REQ/COUNT/NEG-OPEN. Example for DMs: '4,1059'." default: "4, 1059" - name: relay__auth__restrictReadToInvolvedPubkey desc: "When a filter matches restrictedReadKinds, also require the authenticated pubkey to appear in its 'authors' or '#p' set." diff --git a/strfry.conf b/strfry.conf index 5333f354..81da7afa 100644 --- a/strfry.conf +++ b/strfry.conf @@ -59,7 +59,7 @@ relay { # External relay URL (beginning with wss://). Required in order to validate challenge responses. serviceUrl = "" - # Comma-separated list of event kinds that require NIP-42 AUTH to read via REQ/COUNT/NEG-OPEN. A filter with no 'kinds' field is treated as restricted. Example for DMs: '4,1059'. + # Comma-separated list of event kinds that require NIP-42 AUTH to read via REQ/COUNT/NEG-OPEN. Example for DMs: '4,1059'. restrictedReadKinds = "4, 1059" # When a filter matches restrictedReadKinds, also require the authenticated pubkey to appear in its 'authors' or '#p' set. diff --git a/test/README.md b/test/README.md index b2c0fdb6..c1d52eac 100644 --- a/test/README.md +++ b/test/README.md @@ -4,9 +4,11 @@ Tests should be run from the *root* of the project. ## Tests for event writing, including replacements, deletions, etc: - perl test/writeTest.pl + node test/writeTest.js -Note that this script relies on [`nostril`](https://github.com/jb55/nostril) being installed in your path. +## Restricted read tests (REQ/COUNT/negentropy + ReadRestrictor logic): + + node test/readRestrictTest.js ## Fuzz tests diff --git a/test/runTests.sh b/test/runTests.sh index 67fe614b..500a95c5 100755 --- a/test/runTests.sh +++ b/test/runTests.sh @@ -13,36 +13,41 @@ info() { echo -e "${YELLOW}[INFO]${NC} $*"; } info "running write tests..." -node "./test/writeTest.js" \ - && pass "./test/writeTest.js" \ - || fail "./test/writeTest.js failed" +node "./test/tests/writeTest.js" \ + && pass "./test/tests/writeTest.js" \ + || fail "./test/tests/writeTest.js failed" +info "running restricted read tests..." + +node "./test/tests/readRestrictTest.js" \ + && pass "./test/tests/readRestrictTest.js" \ + || fail "./test/tests/readRestrictTest.js failed" info "Seeding events..." -perl "./test/generate-seed-data.pl" -o - | ./strfry --config ./test/cfgs/test.conf import --no-verify +perl "./test/utils/generate-seed-data.pl" -o - | ./strfry --config ./test/cfgs/test.conf import --no-verify info "running filterFuzzTest..." -perl "./test/filterFuzzTest.pl" scan \ -&& pass "./test/filterFuzzTest.pl scan" \ -|| fail "./test/filterFuzzTest.pl scan failed" +perl "./test/tests/filterFuzzTest.pl" scan \ +&& pass "./test/tests/filterFuzzTest.pl scan" \ +|| fail "./test/tests/filterFuzzTest.pl scan failed" -perl "./test/filterFuzzTest.pl" scan-limit \ -&& pass "./test/filterFuzzTest.pl scan-limit" \ -|| fail "./test/filterFuzzTest.pl scan-limit failed" +perl "./test/tests/filterFuzzTest.pl" scan-limit \ +&& pass "./test/tests/filterFuzzTest.pl scan-limit" \ +|| fail "./test/tests/filterFuzzTest.pl scan-limit failed" -perl "./test/filterFuzzTest.pl" monitor \ -&& pass "./test/filterFuzzTest.pl monitor" \ -|| fail "./test/filterFuzzTest.pl monitor failed" +perl "./test/tests/filterFuzzTest.pl" monitor \ +&& pass "./test/tests/filterFuzzTest.pl monitor" \ +|| fail "./test/tests/filterFuzzTest.pl monitor failed" # sync tests info "running sync tests..." -perl "./test/runSyncTests.pl" \ - && pass "./test/syncTests.pl" \ - || fail "./test/runSyncTests.pl failed" +perl "./test/tests/runSyncTests.pl" \ + && pass "./test/tests/syncTests.pl" \ + || fail "./test/tests/runSyncTests.pl failed" pass "All tests passed." diff --git a/test/SubIdTests.cpp b/test/tests/SubIdTests.cpp similarity index 100% rename from test/SubIdTests.cpp rename to test/tests/SubIdTests.cpp diff --git a/test/filterFuzzTest.pl b/test/tests/filterFuzzTest.pl similarity index 98% rename from test/filterFuzzTest.pl rename to test/tests/filterFuzzTest.pl index f26f1679..5e27447a 100755 --- a/test/filterFuzzTest.pl +++ b/test/tests/filterFuzzTest.pl @@ -24,7 +24,7 @@ } close $fh; -push @topics, "nosuchtopic"; +push @$topics, "nosuchtopic"; sub genRandomFilterGroup { my $useLimit = shift; @@ -136,7 +136,7 @@ sub testScan { my $headCmd = @$fg == 1 && $fg->[0]->{limit} ? "| head -n $fg->[0]->{limit}" : ""; - my $resA = `./strfry export --reverse 2>/dev/null | perl test/dumbFilter.pl '$fge' $headCmd | jq -r .id | sort | sha256sum`; + my $resA = `./strfry export --reverse 2>/dev/null | perl test/utils/dumbFilter.pl '$fge' $headCmd | jq -r .id | sort | sha256sum`; my $resB = `./strfry scan --pause 1 --metrics '$fge' | jq -r .id | sort | sha256sum`; print "$resA\n$resB\n"; diff --git a/test/tests/readRestrictTest.js b/test/tests/readRestrictTest.js new file mode 100644 index 00000000..bc69bc14 --- /dev/null +++ b/test/tests/readRestrictTest.js @@ -0,0 +1,394 @@ +import os from "node:os"; +import path from "node:path"; +import { openWebSocket, WsClient } from "../utils/websocketClient.js"; +import { + writeConfig, + addEvent, + cleanDb, + runStrfry, + runRelaySuite, +} from "../utils/relay.js"; +import ids from "../utils/dummyIds.json" with { type: "json" }; +import { signEvent } from "../utils/events.js"; + +const workDir = path.join(os.tmpdir(), "strfry-read-restrict-tests"); +const relayDbDir = path.join(workDir, "relay-db"); +const syncDbDir = path.join(workDir, "sync-db"); +const relayCfgPath = path.join(workDir, "readRestrictRelay.conf"); +const syncCfgPath = path.join(workDir, "readRestrictSync.conf"); +const relayPort = 40553; + +let authChallengeString = null; + +const pass = (msg) => console.log(`Pass: ${msg}`); + +// pre auth +async function testRestrictedReqAndCountRequireAuth({ wsUrl, client }) { + client.send(["REQ", "restricted-req", { kinds: [4] }]); + const authReq = await client.waitFor((m) => m[0] === "AUTH"); + expect( + typeof authReq[1] === "string" && authReq[1].length > 0, + "REQ must return AUTH challenge", + ); + + authChallengeString = authReq[1]; + + const closedReq = await client.waitFor( + (m) => m[0] === "CLOSED" && m[1] === "restricted-req", + ); + expect( + String(closedReq[2]).includes("auth-required"), + "REQ must be closed with auth-required", + ); + + client.send(["COUNT", "restricted-count", { kinds: [4] }]); + + const closedCount = await client.waitFor( + (m) => m[0] === "CLOSED" && m[1] === "restricted-count", + ); + expect( + String(closedCount[2]).includes("auth-required"), + "COUNT must be closed with auth-required", + ); + pass("testRestrictedReqAndCountRequireAuth"); +} + +async function testCountUnrestrictedAllowed({ wsUrl, client }) { + client.send(["COUNT", "count-open", { kinds: [1] }]); + const count = await client.waitFor( + (m) => m[0] === "COUNT" && m[1] === "count-open", + 4_000, + ); + expect( + typeof count[2]?.count === "number", + "COUNT on unrestricted kinds should return a count body", + ); + pass("testCountUnrestrictedAllowed"); +} + +async function testReqWorkerFiltersRestrictedInitialScan({ wsUrl, client }) { + client.send(["REQ", "mixed-initial", { kinds: [1, 4] }]); + let msgs = await client.collectUntil( + (m) => m[0] === "EOSE" && m[1] === "mixed-initial", + 4_000, + ); + + let events = msgs + .filter((m) => m[0] === "EVENT" && m[1] === "mixed-initial") + .map((m) => m[2]); + let kinds = events.map((ev) => ev.kind); + expect( + kinds.length >= 1, + "mixed initial REQ should return at least one event", + ); + expect( + kinds.every((k) => k === 1), + "mixed initial REQ should not return restricted kind 4", + ); + + client.send(["REQ", "omitted-kinds", {}]); + msgs = await client.collectUntil( + (m) => m[0] === "EOSE" && m[1] === "omitted-kinds", + 4_000, + ); + + events = msgs + .filter((m) => m[0] === "EVENT" && m[1] === "omitted-kinds") + .map((m) => m[2]); + kinds = events.map((ev) => ev.kind); + expect( + kinds.length >= 1, + "REQ with omitted kinds should return at least one event", + ); + expect( + kinds.every((k) => k !== 4), + "REQ with omitted kinds should not return restricted kind 4", + ); + pass("testReqWorkerFiltersRestrictedInitialScan"); +} + +async function testReqMonitorFiltersRestrictedLiveEvents({ wsUrl, client }) { + client.send(["REQ", "mixed-live", { kinds: [1, 4] }]); + await client.waitFor((m) => m[0] === "EOSE" && m[1] === "mixed-live", 4_000); + + addEvent(relayCfgPath, { + kind: 4, + from: 0, + tags: [["p", ids[1].pub]], + content: "live-restricted", + }); + addEvent(relayCfgPath, { kind: 1, from: 0, content: "live-public" }); + + const firstLiveEvent = await client.waitFor( + (m) => + m[0] === "EVENT" && + m[1] === "mixed-live" && + m[2].content === "live-public", + 5_000, + ); + expect( + firstLiveEvent[2].kind === 1, + "live mixed REQ should deliver kind 1 event", + ); + + pass("testReqMonitorFiltersRestrictedLiveEvents"); +} + +function testNegentropyMixedFilterBlocksRestrictedWithoutAuth({ wsUrl }) { + cleanDb(syncDbDir); + writeConfig(`db = "${syncDbDir}/"\n`, syncCfgPath); + + const syncRes = runStrfry( + [ + "--config", + syncCfgPath, + "sync", + wsUrl, + "--filter", + '{"kinds":[1,4]}', + "--print-missing", + "--timeout", + "10", + ], + { stdio: ["ignore", "pipe", "pipe"] }, + ); + + if (syncRes.status !== 0) { + throw new Error( + `mixed negentropy sync failed: ${syncRes.stderr || syncRes.stdout}`, + ); + } + + const needLines = syncRes.stdout + .trim() + .split("\n") + .filter((line) => line.startsWith("need,")); + expect( + needLines.length === 2, + "mixed negentropy sync should only expose non-restricted event ids", + ); + pass("testNegentropyMixedFilterBlocksRestrictedWithoutAuth"); +} + +function testNegentropyRestrictedFilterRequiresAuth({ wsUrl }) { + cleanDb(syncDbDir); + writeConfig(`db = "${syncDbDir}/"\n`, syncCfgPath); + + const syncRes = runStrfry( + [ + "--config", + syncCfgPath, + "sync", + wsUrl, + "--filter", + '{"kinds":[4]}', + "--print-missing", + "--timeout", + "10", + ], + { stdio: ["ignore", "pipe", "pipe"] }, + ); + + expect( + syncRes.status !== 0, + "restricted negentropy sync should fail without auth", + ); + const logs = `${syncRes.stdout}\n${syncRes.stderr}`; + expect( + logs.includes("NEG-ERR") && logs.includes("auth-required"), + "restricted negentropy sync should fail with auth-required NEG-ERR", + ); + pass("testNegentropyRestrictedFilterRequiresAuth"); +} + +// post auth + +async function testRestrictedFilterReturnsAllIfAuthenticatedAndInvolvementNotRequired({ + wsUrl, + client, +}) { + // send AUTH message with challenge string + const authEvent = signEvent( + authChallengeString, + "wss://relay.test", + ids[0].sec, + ); + client.send(["AUTH", authEvent]); + + const authOk = await client.waitFor( + (m) => m[3] === "successfully authenticated", + 10000, + ); + + expect( + authOk[0] === "OK", + "AUTH should return 'successfully authenticated' with 'ok' status", + ); + + // send REQ with mixed filter + client.send(["REQ", "restricted-auth", { kinds: [4] }]); + const msgs = await client.collectUntil( + (m) => m[0] === "EOSE" && m[1] === "restricted-auth", + 4_000, + ); + + const events = msgs + .filter((m) => m[0] === "EVENT" && m[1] === "restricted-auth") + .map((m) => m[2]); + + const kinds = events.map((ev) => ev.kind); + expect( + kinds.length >= 1, + "restricted authenticated REQ should return at least one event", + ); + expect( + kinds.includes(4), + "restricted authenticated REQ should return restricted kind 4 events", + ); + pass( + "testRestrictedFilterReturnsAllIfAuthenticatedAndInvolvementNotRequired", + ); +} + +async function testRestrictedFilterCountAuthenticatedNotScoped({ + wsUrl, + client, +}) { + client.send(["COUNT", "count-restricted-unscoped", { kinds: [4] }]); + const count = await client.waitFor( + (m) => m[1] === "count-restricted-unscoped", + 4_000, + ); + expect( + count[0] === "COUNT", + "COUNT must be successful when authenticated and restrictReadToInvolvedPubkey is false", + ); + pass("testRestrictedFilterCountAuthenticatedNotScoped"); +} + +async function testOnlyReturnsMyOwnEventsIfInvolvedRequired({ wsUrl, client }) { + client.send(["REQ", "only-my-own", {}]); + const msgs = await client.collectUntil( + (m) => m[0] === "EOSE" && m[1] === "only-my-own", + 4_000, + ); + + const events = msgs + .filter((m) => m[0] === "EVENT" && m[1] === "only-my-own") + .map((m) => m[2]); + + expect( + events.every((ev) => ev.pubkey === ids[0].pub), + "Must only return my own events when restrictReadToInvolvedPubkey is set", + ); + + pass("testOnlyReturnsMyOwnEventsIfInvolvedRequired"); +} + +async function testCountFailsWhenFilterNotFullyScopedIfInvolvedRequired({ + wsUrl, + client, +}) { + client.send(["COUNT", "failing-count", { kinds: [4] }]); + + const msg = await client.waitFor((m) => m[1] === "failing-count"); + + expect(String(msg[2]).includes("count-failed")); + + pass("testCountFailsWhenFilterNotFullyScopedIfInvolvedRequired"); +} + +async function testCountSuccessfulWhenFilterFullyScoped({ wsUrl, client }) { + client.send([ + "COUNT", + "succeeding-count", + { kinds: [4], authors: [ids[0].pub] }, + ]); + + const msg = await client.waitFor((m) => m[1] === "succeeding-count"); + + expect( + Number.isFinite(msg[2].count), + "returned count should be a valid number", + ); + + pass("testCountSuccessFulWhenFilterFullyScoped"); +} + +function config(relayDb, relayPort, restrictReadToInvolvedPubkey) { + return ` +db = "${relayDb}/" + +relay { + bind = "127.0.0.1" + port = ${relayPort} + nofiles = 0 + autoPingSeconds = 0 + + auth { + enabled = true + serviceUrl = "wss://relay.test" + restrictedReadKinds = "4, 1059" + restrictReadToInvolvedPubkey = ${restrictReadToInvolvedPubkey ? "true" : "false"} + } + + numThreads { + ingester = 1 + reqWorker = 1 + reqMonitor = 1 + negentropy = 1 + } + + negentropy { + enabled = true + maxSyncEvents = 100000 + } +} +`; +} + +function expect(cond, msg) { + if (!cond) throw new Error(msg); +} + +async function main() { + console.log("* read restriction relay integration tests"); + + cleanDb(syncDbDir); + writeConfig(`db = "${syncDbDir}/"\n`, syncCfgPath); + + await runRelaySuite({ + config: config(relayDbDir, relayPort, true), + relayConfigPath: relayCfgPath, + relayPort: relayPort, + relayDbPath: relayDbDir, + tests: async ({ wsUrl, client }) => { + writeConfig(config(relayDbDir, relayPort, false), relayCfgPath); + await testRestrictedReqAndCountRequireAuth({ wsUrl, client }); + await testCountUnrestrictedAllowed({ wsUrl, client }); + await testReqWorkerFiltersRestrictedInitialScan({ wsUrl, client }); + await testReqMonitorFiltersRestrictedLiveEvents({ wsUrl, client }); + testNegentropyMixedFilterBlocksRestrictedWithoutAuth({ wsUrl }); + testNegentropyRestrictedFilterRequiresAuth({ wsUrl }); + await testRestrictedFilterReturnsAllIfAuthenticatedAndInvolvementNotRequired( + { wsUrl, client }, + ); + await testRestrictedFilterCountAuthenticatedNotScoped({ wsUrl, client }); + writeConfig(config(relayDbDir, relayPort, true), relayCfgPath); + console.log("Writing new config, wait..."); + await new Promise((resolve) => setTimeout(resolve, 3000)); + await testOnlyReturnsMyOwnEventsIfInvolvedRequired({ wsUrl, client }); + await testCountFailsWhenFilterNotFullyScopedIfInvolvedRequired({ + wsUrl, + client, + }); + await testCountSuccessfulWhenFilterFullyScoped({ wsUrl, client }); + }, + }); + console.log("All read restriction tests passed!"); +} + +main().catch((err) => { + console.error(err && err.stack ? err.stack : String(err)); + process.exit(1); +}); diff --git a/test/runSyncTests.pl b/test/tests/runSyncTests.pl similarity index 95% rename from test/runSyncTests.pl rename to test/tests/runSyncTests.pl index 45f3a322..33d78969 100755 --- a/test/runSyncTests.pl +++ b/test/tests/runSyncTests.pl @@ -48,7 +48,7 @@ sub test { my $cmd = qq{ ./strfry --config test/cfgs/test.conf export 2>/dev/null | head -$num | - perl test/syncTest.pl $params $redir + perl test/tests/syncTest.pl $params $redir }; print "CMD: $cmd\n"; diff --git a/test/syncTest.pl b/test/tests/syncTest.pl similarity index 93% rename from test/syncTest.pl rename to test/tests/syncTest.pl index c2943cec..f33ab269 100755 --- a/test/syncTest.pl +++ b/test/tests/syncTest.pl @@ -55,8 +55,8 @@ system("./strfry --config test/cfgs/syncTest2.conf sync ws://127.0.0.1:40551 --dir both --filter '$filter'"); }); -my $hash1 = `./strfry --config test/cfgs/syncTest1.conf export | perl test/dumbFilter.pl '$filter' | sort | sha256sum`; -my $hash2 = `./strfry --config test/cfgs/syncTest2.conf export | perl test/dumbFilter.pl '$filter' | sort | sha256sum`; +my $hash1 = `./strfry --config test/cfgs/syncTest1.conf export | perl test/utils/dumbFilter.pl '$filter' | sort | sha256sum`; +my $hash2 = `./strfry --config test/cfgs/syncTest2.conf export | perl test/utils/dumbFilter.pl '$filter' | sort | sha256sum`; die "hashes differ" unless $hash1 eq $hash2; diff --git a/test/writeTest.js b/test/tests/writeTest.js similarity index 86% rename from test/writeTest.js rename to test/tests/writeTest.js index aef9df6c..1b21986a 100644 --- a/test/writeTest.js +++ b/test/tests/writeTest.js @@ -1,64 +1,18 @@ import { spawnSync } from "node:child_process"; import { mkdirSync, rmSync } from "node:fs"; import path from "node:path"; -import { buildEvent } from "./events.js"; - -let ids = [ - { - sec: "c1eee22f68dc218d98263cfecb350db6fc6b3e836b47423b66c62af7ae3e32bb", - pub: "003ba9b2c5bd8afeed41a4ce362a8b7fc3ab59c25b6a1359cae9093f296dac01", - }, - { - sec: "a0b459d9ff90e30dc9d1749b34c4401dfe80ac2617c7732925ff994e8d5203ff", - pub: "cc49e2a58373abc226eee84bee9ba954615aa2ef1563c4f955a74c4606a3b1fa", - }, - { - sec: "3bf1d347477a4f8d2e6a3b1c9d5e7f0a2b4c6d8e0f1a3b5c7d9e1f2a4b6c8d0e", - pub: "e83f7c58291be6a4d05c13f78a2e94b61d37c085f4a2960db51c8ef73a41d9f2", - }, -]; - -function addEvent(evInput) { - const event = buildEvent(evInput); - - spawnSync( - "./strfry", - ["--config", "test/cfgs/writeTest.conf", "import", "--no-verify"], - { - input: JSON.stringify(event) + "\n", - encoding: "utf-8", - stdio: ["pipe", "ignore", "ignore"], - }, - ); - - if (process.env.DUMP_EVENTS) { - console.log(event); - } - - return event.id; -} - -function cleanDb() { - const dir = path.join(process.cwd(), "strfry-db-test"); - const file = path.join(dir, "data.mdb"); - mkdirSync(dir, { recursive: true }); - rmSync(file, { force: true }); -} +import { buildEvent } from "../utils/events.js"; +import { cleanDb, addEvent, runStrfry } from "../utils/relay.js"; +import ids from "../utils/dummyIds.json" with { type: "json" }; function doTest(spec) { console.log("*", spec.desc || "unnamed"); - cleanDb(); + cleanDb(path.join(process.cwd(), "strfry-db-test")); const eventIds = []; for (const ev of spec.events) { - ev.pub = ids[ev.from || 0].pub; - ev.sec = ids[ev.from || 0].sec; - - // deep clone - const e = JSON.parse(JSON.stringify(ev)); - const replaceEV = (v) => { if (typeof v === "string") { return v.replace(/EV_(\d+)/g, (_, i) => eventIds[Number(i)]); @@ -70,9 +24,9 @@ function doTest(spec) { return v; }; - replaceEV(e); + replaceEV(ev); - const id = addEvent(e); + const id = addEvent("test/cfgs/writeTest.conf", ev).id; eventIds.push(id); } @@ -84,11 +38,7 @@ function doTest(spec) { } } - const result = spawnSync( - "./strfry", - ["--config", "test/cfgs/writeTest.conf", "export"], - { encoding: "utf-8" }, - ); + const result = runStrfry(["--config", "test/cfgs/writeTest.conf", "export"]); if (result.error) throw result.error; diff --git a/test/dumbFilter.pl b/test/utils/dumbFilter.pl similarity index 100% rename from test/dumbFilter.pl rename to test/utils/dumbFilter.pl diff --git a/test/utils/dummyIds.json b/test/utils/dummyIds.json new file mode 100644 index 00000000..b3db3bb0 --- /dev/null +++ b/test/utils/dummyIds.json @@ -0,0 +1,14 @@ +[ + { + "sec": "c1eee22f68dc218d98263cfecb350db6fc6b3e836b47423b66c62af7ae3e32bb", + "pub": "003ba9b2c5bd8afeed41a4ce362a8b7fc3ab59c25b6a1359cae9093f296dac01" + }, + { + "sec": "a0b459d9ff90e30dc9d1749b34c4401dfe80ac2617c7732925ff994e8d5203ff", + "pub": "cc49e2a58373abc226eee84bee9ba954615aa2ef1563c4f955a74c4606a3b1fa" + }, + { + "sec": "3bf1d347477a4f8d2e6a3b1c9d5e7f0a2b4c6d8e0f1a3b5c7d9e1f2a4b6c8d0e", + "pub": "e83f7c58291be6a4d05c13f78a2e94b61d37c085f4a2960db51c8ef73a41d9f2" + } +] \ No newline at end of file diff --git a/test/events.js b/test/utils/events.js similarity index 59% rename from test/events.js rename to test/utils/events.js index 0ac28a87..e87f3dfb 100644 --- a/test/events.js +++ b/test/utils/events.js @@ -1,4 +1,5 @@ import { createHash } from "node:crypto"; +import { finalizeEvent } from "@nostr/tools"; export function sha256(message) { return createHash("sha256").update(message).digest("hex"); @@ -8,6 +9,10 @@ export function bytesToHex(bytes) { return Buffer.from(bytes).toString("hex"); } +export function hexToBytes(hex) { + return Uint8Array.from(hex.match(/.{2}/g), (byte) => parseInt(byte, 16)); +} + function serializeEvent(evt) { return JSON.stringify([ 0, @@ -46,3 +51,21 @@ export function buildEvent({ return { ...evt, id, sig }; } + +export function signEvent(authChallengeString, relayUrl, sec) { + const secBytes = hexToBytes(sec); + const authEvent = finalizeEvent( + { + kind: 22242, + created_at: Math.floor(Date.now() / 1000), + tags: [ + ["relay", relayUrl], + ["challenge", authChallengeString], + ], + content: "", + }, + secBytes, + ); + + return authEvent; +} diff --git a/test/generate-seed-data.pl b/test/utils/generate-seed-data.pl similarity index 100% rename from test/generate-seed-data.pl rename to test/utils/generate-seed-data.pl diff --git a/test/generate-seed-data.sh b/test/utils/generate-seed-data.sh similarity index 100% rename from test/generate-seed-data.sh rename to test/utils/generate-seed-data.sh diff --git a/test/utils/relay.js b/test/utils/relay.js new file mode 100644 index 00000000..e2b1f83a --- /dev/null +++ b/test/utils/relay.js @@ -0,0 +1,104 @@ +import { spawn, spawnSync } from "node:child_process"; +import { mkdirSync, rmSync, writeFileSync } from "node:fs"; +import { setTimeout as delay } from "node:timers/promises"; +import path from "node:path"; +import { buildEvent } from "./events.js"; +import { waitForRelay, openWebSocket, WsClient } from "./websocketClient.js"; +import ids from "./dummyIds.json" with { type: "json" }; + +let ts = 1_700_000_000; + +export function runStrfry(args, opts = {}) { + const res = spawnSync("./strfry", args, { + encoding: "utf-8", + ...opts, + }); + if (res.error) throw res.error; + return res; +} + +export function cleanDb(dir) { + rmSync(dir, { recursive: true, force: true }); + mkdirSync(dir, { recursive: true }); +} + +export function writeConfig(config, configPath) { + mkdirSync(path.dirname(configPath), { recursive: true }); + writeFileSync(configPath, config.trim() + "\n", "utf-8"); +} + +export function addEvent(configPath, evInput) { + const event = buildEvent({ + sec: ids[evInput.from ?? 0].sec, + pub: ids[evInput.from ?? 0].pub, + content: evInput.content ?? "", + kind: evInput.kind ?? 1, + tags: evInput.tags ?? [], + created_at: evInput.created_at ?? ts++, + }); + + const res = runStrfry(["--config", configPath, "import", "--no-verify"], { + input: JSON.stringify(event) + "\n", + stdio: ["pipe", "ignore", "pipe"], + }); + + if (res.status !== 0) { + throw new Error(`import failed: ${res.stderr}`); + } + + if (process.env.DUMP_EVENTS) { + console.log(event); + } + + return event; +} + +export async function runRelaySuite({ + relayConfigPath, + relayPort, + relayDbPath, + tests, +}) { + const wsUrl = `ws://127.0.0.1:${relayPort}`; + + cleanDb(relayDbPath); + addEvent(relayConfigPath, { kind: 1, from: 0, content: "seed-public" }); + addEvent(relayConfigPath, { + kind: 4, + from: 0, + tags: [["p", ids[1].pub]], + content: "seed-restricted", + }); + + const relayProc = spawn("./strfry", ["--config", relayConfigPath, "relay"], { + stdio: ["ignore", "pipe", "pipe"], + }); + + let relayLogs = ""; + relayProc.stdout.on("data", (d) => { + relayLogs += d.toString(); + }); + relayProc.stderr.on("data", (d) => { + relayLogs += d.toString(); + }); + + let client; + try { + await waitForRelay(wsUrl); + const ws = await openWebSocket(wsUrl); + client = new WsClient(ws); + await tests({ wsUrl, client }); + } catch (e) { + throw new Error( + `${String(e && e.message ? e.message : e)}\n\nRelay logs:\n${relayLogs}`, + ); + } finally { + if (client) await client.close(); + relayProc.kill("SIGTERM"); + await Promise.race([ + new Promise((resolve) => relayProc.once("exit", resolve)), + delay(2_000), + ]); + if (!relayProc.killed) relayProc.kill("SIGKILL"); + } +} diff --git a/test/utils/websocketClient.js b/test/utils/websocketClient.js new file mode 100644 index 00000000..31ddb99a --- /dev/null +++ b/test/utils/websocketClient.js @@ -0,0 +1,102 @@ +import { setTimeout as delay } from "node:timers/promises"; + +export async function openWebSocket(url, timeoutMs = 4_000) { + return await new Promise((resolve, reject) => { + const ws = new WebSocket(url); + const timer = setTimeout(() => { + ws.close(); + reject(new Error("websocket open timeout")); + }, timeoutMs); + + ws.addEventListener("open", () => { + clearTimeout(timer); + resolve(ws); + }); + + ws.addEventListener("error", () => { + clearTimeout(timer); + reject(new Error("websocket connection failed")); + }); + }); +} + +export class WsClient { + constructor(ws) { + this.ws = ws; + this.queue = []; + this.waiters = []; + + this.ws.addEventListener("message", (ev) => { + const msg = JSON.parse(String(ev.data)); + if (this.waiters.length > 0) { + const waiter = this.waiters.shift(); + clearTimeout(waiter.timer); + waiter.resolve(msg); + } else { + this.queue.push(msg); + } + }); + } + + send(msg) { + this.ws.send(JSON.stringify(msg)); + } + + async nextMessage(timeoutMs = 3_000) { + if (this.queue.length > 0) return this.queue.shift(); + + return await new Promise((resolve, reject) => { + const timer = setTimeout(() => { + this.waiters = this.waiters.filter((w) => w.resolve !== resolve); + reject(new Error("message timeout")); + }, timeoutMs); + + this.waiters.push({ resolve, reject, timer }); + }); + } + + async waitFor(predicate, timeoutMs = 3_000) { + const deadline = Date.now() + timeoutMs; + while (Date.now() < deadline) { + const remaining = Math.max(1, deadline - Date.now()); + const msg = await this.nextMessage(remaining); + if (predicate(msg)) return msg; + } + throw new Error("waitFor timeout"); + } + + async collectUntil(predicate, timeoutMs = 3_000) { + const out = []; + while (true) { + const msg = await this.nextMessage(timeoutMs); + out.push(msg); + if (predicate(msg)) return out; + } + } + + async close() { + if (this.ws.readyState === WebSocket.CLOSED) return; + await new Promise((resolve) => { + this.ws.addEventListener("close", () => resolve(), { once: true }); + this.ws.close(); + }); + } +} + +export async function waitForRelay(wsUrl) { + const deadline = Date.now() + 12_000; + let lastErr = ""; + + while (Date.now() < deadline) { + try { + const ws = await openWebSocket(wsUrl, 500); + ws.close(); + return; + } catch (e) { + lastErr = String(e.message || e); + await delay(100); + } + } + + throw new Error(`relay did not start in time: ${lastErr}`); +}