From df24c8e81c69676d11c734c57d11930182346e0c Mon Sep 17 00:00:00 2001 From: Sergei Tikhomirov Date: Thu, 18 Jun 2026 14:55:22 +0200 Subject: [PATCH] feat(store): add liblogosdelivery eligibility hooks and store query (Step 15) Expose verifier/provider C callbacks, N8 canonical bytes, inbound wrapper, logosdelivery_store_query, Nim parity tests, and C ABI smoke. --- Makefile | 19 ++ library/kernel_api/protocols/store_api.nim | 58 +----- library/liblogosdelivery.h | 57 ++++++ library/liblogosdelivery.nim | 1 + library/logos_delivery_api/node_api.nim | 2 + library/store_eligibility/eligibility_api.nim | 170 ++++++++++++++++++ .../store_eligibility/store_query_json.nim | 65 +++++++ library/tests/test_eligibility_hooks.c | 43 +++++ .../waku/waku_store/eligibility_canonical.nim | 97 ++++++++++ .../waku/waku_store/eligibility_hooks.nim | 94 ++++++++++ logos_delivery/waku/waku_store/protocol.nim | 10 +- tests/all_tests_waku.nim | 2 + .../test_store_eligibility_canonical.nim | 35 ++++ .../test_store_eligibility_hooks.nim | 119 ++++++++++++ 14 files changed, 715 insertions(+), 57 deletions(-) create mode 100644 library/store_eligibility/eligibility_api.nim create mode 100644 library/store_eligibility/store_query_json.nim create mode 100644 library/tests/test_eligibility_hooks.c create mode 100644 logos_delivery/waku/waku_store/eligibility_canonical.nim create mode 100644 logos_delivery/waku/waku_store/eligibility_hooks.nim create mode 100644 tests/waku_store/test_store_eligibility_canonical.nim create mode 100644 tests/waku_store/test_store_eligibility_hooks.nim diff --git a/Makefile b/Makefile index daa8a8ad7..19c4bee99 100644 --- a/Makefile +++ b/Makefile @@ -478,6 +478,25 @@ else ifeq ($(detected_OS),Windows) -lws2_32 endif +logosdelivery_eligibility_smoke: | build liblogosdelivery + @echo -e $(BUILD_MSG) "build/$@" +ifeq ($(detected_OS),Darwin) + gcc -o build/logosdelivery_eligibility_smoke \ + library/tests/test_eligibility_hooks.c \ + -I./library \ + -L./build \ + -llogosdelivery \ + -Wl,-rpath,./build +else ifeq ($(detected_OS),Linux) + gcc -o build/logosdelivery_eligibility_smoke \ + library/tests/test_eligibility_hooks.c \ + -I./library \ + -L./build \ + -llogosdelivery \ + -Wl,-rpath,'$$ORIGIN' +endif + ./build/logosdelivery_eligibility_smoke + cwaku_example: | build liblogosdelivery echo -e $(BUILD_MSG) "build/$@" && \ cc -o "build/$@" \ diff --git a/library/kernel_api/protocols/store_api.nim b/library/kernel_api/protocols/store_api.nim index 75d43fee1..325c5a79b 100644 --- a/library/kernel_api/protocols/store_api.nim +++ b/library/kernel_api/protocols/store_api.nim @@ -6,63 +6,11 @@ import logos_delivery/waku/waku_core/message/digest, logos_delivery/waku/waku_store/common, logos_delivery/waku/common/paging, - library/declare_lib + library/declare_lib, + ../../store_eligibility/store_query_json func fromJsonNode(jsonContent: JsonNode): Result[StoreQueryRequest, string] = - var contentTopics: seq[string] - if jsonContent.contains("contentTopics"): - contentTopics = collect(newSeq): - for cTopic in jsonContent["contentTopics"].getElems(): - cTopic.getStr() - - var msgHashes: seq[WakuMessageHash] - if jsonContent.contains("messageHashes"): - for hashJsonObj in jsonContent["messageHashes"].getElems(): - let hash = hashJsonObj.getStr().hexToHash().valueOr: - return err("Failed converting message hash hex string to bytes: " & error) - msgHashes.add(hash) - - let pubsubTopic = - if jsonContent.contains("pubsubTopic"): - some(jsonContent["pubsubTopic"].getStr()) - else: - none(string) - - let paginationCursor = - if jsonContent.contains("paginationCursor"): - let hash = jsonContent["paginationCursor"].getStr().hexToHash().valueOr: - return err("Failed converting paginationCursor hex string to bytes: " & error) - some(hash) - else: - none(WakuMessageHash) - - let paginationForwardBool = jsonContent["paginationForward"].getBool() - let paginationForward = - if paginationForwardBool: PagingDirection.FORWARD else: PagingDirection.BACKWARD - - let paginationLimit = - if jsonContent.contains("paginationLimit"): - some(uint64(jsonContent["paginationLimit"].getInt())) - else: - none(uint64) - - let startTime = ?jsonContent.getProtoInt64("timeStart") - let endTime = ?jsonContent.getProtoInt64("timeEnd") - - return ok( - StoreQueryRequest( - requestId: jsonContent["requestId"].getStr(), - includeData: jsonContent["includeData"].getBool(), - pubsubTopic: pubsubTopic, - contentTopics: contentTopics, - startTime: startTime, - endTime: endTime, - messageHashes: msgHashes, - paginationCursor: paginationCursor, - paginationForward: paginationForward, - paginationLimit: paginationLimit, - ) - ) + storeQueryRequestFromJson(jsonContent) proc waku_store_query( ctx: ptr FFIContext[LogosDelivery], diff --git a/library/liblogosdelivery.h b/library/liblogosdelivery.h index 7b4bbd5ab..8ce133f65 100644 --- a/library/liblogosdelivery.h +++ b/library/liblogosdelivery.h @@ -127,6 +127,63 @@ extern "C" // header liblogosdelivery_kernel.h. It is intentionally not declared here so // this header only promises the stable Messaging / Reliable Channels surface. + /* + * Inbound (provider-side) eligibility verifier. + * + * Called by liblogosdelivery for every inbound Store request. + * Runs on the liblogosdelivery async handler thread; + * implementation must not re-enter the library. + */ + typedef int (*EligibilityVerifierCb)( + const char *proof_hex, + const char *canonical_hex, + const char *requester_peer_id, + char *out_desc, + size_t out_desc_len, + void *user_data); + + /* + * Outbound (user-side) eligibility provider. + * + * Called by liblogosdelivery before sending an outgoing Store query + * when a provider callback is registered. + */ + typedef int (*EligibilityProviderCb)( + const char *canonical_hex, + const char *provider_peer_id, + char *out_proof_hex, + size_t out_buf_len, + void *user_data); + + /* Register or replace the inbound eligibility verifier. + * Pass NULL to clear a previous registration. */ + int logosdelivery_set_eligibility_verifier( + void *ctx, + EligibilityVerifierCb cb, + void *user_data); + + /* Register or replace the outbound eligibility provider. + * Pass NULL to clear a previous registration. */ + int logosdelivery_set_eligibility_provider( + void *ctx, + EligibilityProviderCb cb, + void *user_data); + + /* + * Issue a Store query to the given provider. + * + * queryJson – JSON object (see store_api / integration docs for field names) + * providerAddr – multiaddr string of the target Store provider peer + * + * Returns StoreQueryResponse JSON via callback when RET_OK. + */ + int logosdelivery_store_query( + void *ctx, + FFICallBack callback, + void *userData, + const char *queryJson, + const char *providerAddr); + #ifdef __cplusplus } #endif diff --git a/library/liblogosdelivery.nim b/library/liblogosdelivery.nim index 755503e7c..faa458aec 100644 --- a/library/liblogosdelivery.nim +++ b/library/liblogosdelivery.nim @@ -20,6 +20,7 @@ import ## Include different APIs, i.e. all procs with {.ffi.} pragma include + ./store_eligibility/eligibility_api, ./logos_delivery_api/node_api, ./logos_delivery_api/messaging_api, ./logos_delivery_api/debug_api, diff --git a/library/logos_delivery_api/node_api.nim b/library/logos_delivery_api/node_api.nim index cfee2f5e2..9b733263c 100644 --- a/library/logos_delivery_api/node_api.nim +++ b/library/logos_delivery_api/node_api.nim @@ -39,6 +39,8 @@ proc logosdelivery_destroy( callback(RET_ERR, unsafeAddr msg[0], cast[csize_t](len(msg)), userData) return RET_ERR + dropState(ctx) + ## always need to invoke the callback although we don't retrieve value to the caller callback(RET_OK, nil, 0, userData) diff --git a/library/store_eligibility/eligibility_api.nim b/library/store_eligibility/eligibility_api.nim new file mode 100644 index 000000000..dae14ce87 --- /dev/null +++ b/library/store_eligibility/eligibility_api.nim @@ -0,0 +1,170 @@ +import std/[json, options, strutils, tables, sugar, locks] +import chronos, results, stew/byteutils, ffi +import + logos_delivery/waku/factory/waku, + logos_delivery/waku/waku_core/peers, + logos_delivery/waku/waku_store/[ + common, protocol, client, eligibility_canonical, eligibility_hooks + ], + ../declare_lib, + ./store_query_json + +const EligibilityOutProofHexMinLen = 4096 + +type + EligibilityProviderCb* = proc ( + canonical_hex, provider_peer_id: cstring, out_proof_hex: cstring, + out_buf_len: csize_t, user_data: pointer + ): cint {.cdecl, gcsafe, raises: [].} + + StoreEligibilityState* = ref object + verifierCb: EligibilityVerifierCb + verifierUserData: pointer + providerCb: EligibilityProviderCb + providerUserData: pointer + innerHandler: StoreQueryRequestHandler + handlerWrapped: bool + +export EligibilityVerifierCb + +var + storeEligibilityByCtx = initTable[pointer, StoreEligibilityState]() + eligibilityHookLock: Lock + +initLock(eligibilityHookLock) + +proc getOrCreateStateUnlocked(key: pointer): StoreEligibilityState = + result = storeEligibilityByCtx.getOrDefault(key) + if result.isNil: + result = StoreEligibilityState.new() + storeEligibilityByCtx[key] = result + +proc dropState*(ctx: ptr FFIContext[Waku]) = + eligibilityHookLock.acquire() + defer: + eligibilityHookLock.release() + storeEligibilityByCtx.del cast[pointer](ctx) + +proc applyVerifierWrapper*(ctx: ptr FFIContext[Waku], state: StoreEligibilityState) = + let store = ctx.myLib[].node.wakuStore + if store.isNil: + return + if not state.handlerWrapped: + state.innerHandler = store.requestHandler + state.handlerWrapped = true + store.requestHandler = makeEligibilityWrappedHandler( + state.verifierCb, state.verifierUserData, state.innerHandler, store + ) + +proc clearVerifierWrapper*(ctx: ptr FFIContext[Waku], state: StoreEligibilityState) = + let store = ctx.myLib[].node.wakuStore + if store.isNil or not state.handlerWrapped: + return + store.requestHandler = state.innerHandler + state.handlerWrapped = false + +proc logosdelivery_set_eligibility_verifier( + ctx: ptr FFIContext[Waku], cb: EligibilityVerifierCb, userData: pointer +): cint {.dynlib, exportc, cdecl.} = + initializeLibrary() + if isNil(ctx): + return RET_ERR + eligibilityHookLock.acquire() + defer: + eligibilityHookLock.release() + let state = getOrCreateStateUnlocked(cast[pointer](ctx)) + state.verifierCb = cb + state.verifierUserData = userData + if cb.isNil: + clearVerifierWrapper(ctx, state) + else: + applyVerifierWrapper(ctx, state) + RET_OK + +proc logosdelivery_set_eligibility_provider( + ctx: ptr FFIContext[Waku], cb: EligibilityProviderCb, userData: pointer +): cint {.dynlib, exportc, cdecl.} = + initializeLibrary() + if isNil(ctx): + return RET_ERR + eligibilityHookLock.acquire() + defer: + eligibilityHookLock.release() + let state = getOrCreateStateUnlocked(cast[pointer](ctx)) + state.providerCb = cb + state.providerUserData = userData + RET_OK + +proc snapshotProviderCb(ctx: ptr FFIContext[Waku]): (EligibilityProviderCb, pointer) {. + gcsafe +.} = + eligibilityHookLock.acquire() + defer: + eligibilityHookLock.release() + {.cast(gcsafe).}: + let state = storeEligibilityByCtx.getOrDefault(cast[pointer](ctx)) + if state.isNil: + return (nil, nil) + return (state.providerCb, state.providerUserData) + +proc logosdelivery_store_query( + ctx: ptr FFIContext[Waku], + callback: FFICallBack, + userData: pointer, + queryJson: cstring, + providerAddr: cstring, +) {.ffi.} = + requireInitializedNode(ctx, "STORE_QUERY"): + return err(errMsg) + + let jsonContentRes = catch: + parseJson($queryJson) + + if jsonContentRes.isErr(): + return err("StoreRequest failed parsing store request: " & jsonContentRes.error.msg) + + var storeQueryRequest = storeQueryRequestFromJson(jsonContentRes.get()).valueOr: + return err("StoreRequest invalid query: " & error) + + storeQueryRequest.eligibilityProof = none(seq[byte]) + + let (providerCb, providerUserData) = snapshotProviderCb(ctx) + + let peer = peers.parsePeerInfo(($providerAddr).split(",")).valueOr: + return err("StoreRequest failed to parse peer addr: " & $error) + + if not providerCb.isNil: + let canonicalHex = storeEligibilityCanonicalHex(storeQueryRequest) + var outProof = newString(EligibilityOutProofHexMinLen) + outProof.setLen(EligibilityOutProofHexMinLen) + let providerRc = providerCb( + cstring(canonicalHex), + cstring($peer.peerId), + cstring(outProof), + csize_t(EligibilityOutProofHexMinLen), + providerUserData, + ) + if providerRc < 0: + return err("eligibility provider callback failed") + + var proofHex = outProof.strip(chars = {'\0'}) + if proofHex.startsWith("0x"): + proofHex = proofHex[2 .. ^1] + let proofBytes = + try: + proofHex.hexToSeqByte() + except ValueError: + return err("eligibility provider returned invalid proof hex") + storeQueryRequest.eligibilityProof = some(proofBytes) + + let queryResponse = ( + await ctx.myLib[].node.wakuStoreClient.query(storeQueryRequest, peer) + ).valueOr: + return err("StoreRequest failed store query: " & $error) + + let res = $(%*(queryResponse.toHex())) + return ok(res) + +export + logosdelivery_set_eligibility_verifier, logosdelivery_set_eligibility_provider, + logosdelivery_store_query diff --git a/library/store_eligibility/store_query_json.nim b/library/store_eligibility/store_query_json.nim new file mode 100644 index 000000000..d8281f9c7 --- /dev/null +++ b/library/store_eligibility/store_query_json.nim @@ -0,0 +1,65 @@ +import std/[json, options, strutils, sugar] +import results +import + logos_delivery/waku/waku_core/message/digest, + logos_delivery/waku/waku_store/common, + logos_delivery/waku/common/paging, + ../utils + +func storeQueryRequestFromJson*( + jsonContent: JsonNode +): Result[StoreQueryRequest, string] = + var contentTopics: seq[string] + if jsonContent.contains("contentTopics"): + contentTopics = collect(newSeq): + for cTopic in jsonContent["contentTopics"].getElems(): + cTopic.getStr() + + var msgHashes: seq[WakuMessageHash] + if jsonContent.contains("messageHashes"): + for hashJsonObj in jsonContent["messageHashes"].getElems(): + let hash = hashJsonObj.getStr().hexToHash().valueOr: + return err("Failed converting message hash hex string to bytes: " & error) + msgHashes.add(hash) + + let pubsubTopic = + if jsonContent.contains("pubsubTopic"): + some(jsonContent["pubsubTopic"].getStr()) + else: + none(string) + + let paginationCursor = + if jsonContent.contains("paginationCursor"): + let hash = jsonContent["paginationCursor"].getStr().hexToHash().valueOr: + return err("Failed converting paginationCursor hex string to bytes: " & error) + some(hash) + else: + none(WakuMessageHash) + + let paginationForwardBool = jsonContent["paginationForward"].getBool() + let paginationForward = + if paginationForwardBool: PagingDirection.FORWARD else: PagingDirection.BACKWARD + + let paginationLimit = + if jsonContent.contains("paginationLimit"): + some(uint64(jsonContent["paginationLimit"].getInt())) + else: + none(uint64) + + let startTime = ?jsonContent.getProtoInt64("timeStart") + let endTime = ?jsonContent.getProtoInt64("timeEnd") + + ok( + StoreQueryRequest( + requestId: jsonContent["requestId"].getStr(), + includeData: jsonContent["includeData"].getBool(), + pubsubTopic: pubsubTopic, + contentTopics: contentTopics, + startTime: startTime, + endTime: endTime, + messageHashes: msgHashes, + paginationCursor: paginationCursor, + paginationForward: paginationForward, + paginationLimit: paginationLimit, + ) + ) diff --git a/library/tests/test_eligibility_hooks.c b/library/tests/test_eligibility_hooks.c new file mode 100644 index 000000000..37ee01c80 --- /dev/null +++ b/library/tests/test_eligibility_hooks.c @@ -0,0 +1,43 @@ +#include "../liblogosdelivery.h" +#include +#include + +static int verifier_invoked = 0; + +static int test_verifier_cb( + const char *proof_hex, + const char *canonical_hex, + const char *requester_peer_id, + char *out_desc, + size_t out_desc_len, + void *user_data) { + (void)canonical_hex; + (void)requester_peer_id; + (void)out_desc; + (void)out_desc_len; + (void)user_data; + verifier_invoked++; + if (proof_hex != NULL) { + return -1; + } + return 0; +} + +int main(void) { + int rc; + + rc = logosdelivery_set_eligibility_verifier(NULL, test_verifier_cb, NULL); + if (rc != RET_ERR) { + fprintf(stderr, "expected RET_ERR for NULL ctx, got %d\n", rc); + return 1; + } + + rc = logosdelivery_set_eligibility_provider(NULL, NULL, NULL); + if (rc != RET_ERR) { + fprintf(stderr, "expected RET_ERR for NULL ctx on provider clear, got %d\n", rc); + return 1; + } + + printf("eligibility ABI smoke: registration entry points linked OK\n"); + return 0; +} diff --git a/logos_delivery/waku/waku_store/eligibility_canonical.nim b/logos_delivery/waku/waku_store/eligibility_canonical.nim new file mode 100644 index 000000000..4c528c935 --- /dev/null +++ b/logos_delivery/waku/waku_store/eligibility_canonical.nim @@ -0,0 +1,97 @@ +{.push raises: [].} + +import std/[options, strutils], stew/[endians2, byteutils] +import ../common/paging +import ./common + +const StoreEligibilityDomainPrefixHead* = "/LEZ/v0.1/StoreEligibility/" + +static: + doAssert StoreEligibilityDomainPrefixHead.len == 27 + +proc storeEligibilityDomainPrefixBytes*(): seq[byte] = + result = newSeq[byte](32) + for i in 0 ..< StoreEligibilityDomainPrefixHead.len: + result[i] = byte(StoreEligibilityDomainPrefixHead[i]) + +proc appendLeU32(buf: var seq[byte], value: uint32) = + buf.add toBytes(value, Endianness.littleEndian) + +proc appendLeU64(buf: var seq[byte], value: uint64) = + buf.add toBytes(value, Endianness.littleEndian) + +proc appendLeI64(buf: var seq[byte], value: int64) = + appendLeU64(buf, cast[uint64](value)) + +proc appendBorshString*(buf: var seq[byte], value: string) = + appendLeU32(buf, uint32(value.len)) + for ch in value: + buf.add byte(uint8(ch)) + +proc appendOptionalString*(buf: var seq[byte], value: Option[string]) = + if value.isSome(): + buf.add byte(1) + buf.appendBorshString(value.get()) + else: + buf.add byte(0) + +proc appendOptionalInt64*(buf: var seq[byte], value: Option[int64]) = + if value.isSome(): + buf.add byte(1) + appendLeI64(buf, value.get()) + else: + buf.add byte(0) + +proc appendOptionalUint64*(buf: var seq[byte], value: Option[uint64]) = + if value.isSome(): + buf.add byte(1) + appendLeU64(buf, value.get()) + else: + buf.add byte(0) + +proc storeEligibilityCanonicalBody*(req: StoreQueryRequest): seq[byte] = + var buf = newSeqOfCap[byte](256) + buf.appendBorshString(req.requestId) + buf.add byte(if req.includeData: 1 else: 0) + let pubsub = + if req.pubsubTopic.isSome(): some($req.pubsubTopic.get()) else: none(string) + buf.appendOptionalString(pubsub) + appendLeU32(buf, uint32(req.contentTopics.len)) + for topic in req.contentTopics: + buf.appendBorshString($topic) + buf.appendOptionalInt64( + if req.startTime.isSome(): some(int64(req.startTime.get())) else: none(int64) + ) + buf.appendOptionalInt64( + if req.endTime.isSome(): some(int64(req.endTime.get())) else: none(int64) + ) + appendLeU32(buf, uint32(req.messageHashes.len)) + for h in req.messageHashes: + buf.add h.toOpenArray(0, h.high) + if req.paginationCursor.isSome(): + buf.add byte(1) + buf.add req.paginationCursor.get().toOpenArray(0, 31) + else: + buf.add byte(0) + buf.add byte(if req.paginationForward.into(): 1 else: 0) + buf.appendOptionalUint64(req.paginationLimit) + buf + +proc storeEligibilityCanonicalPayload*(req: StoreQueryRequest): seq[byte] = + var req = req + req.eligibilityProof = none(seq[byte]) + result = newSeqOfCap[byte](32 + 256) + result.add storeEligibilityDomainPrefixBytes() + result.add storeEligibilityCanonicalBody(req) + +proc bytesToLowerHex*(data: openArray[byte]): string = + const hexChars = "0123456789abcdef" + result = newString(data.len * 2) + var j = 0 + for b in data: + result[j] = hexChars[int(b shr 4) and 0xF] + result[j + 1] = hexChars[int(b) and 0xF] + j += 2 + +proc storeEligibilityCanonicalHex*(req: StoreQueryRequest): string = + bytesToLowerHex(storeEligibilityCanonicalPayload(req)) diff --git a/logos_delivery/waku/waku_store/eligibility_hooks.nim b/logos_delivery/waku/waku_store/eligibility_hooks.nim new file mode 100644 index 000000000..50eef7eff --- /dev/null +++ b/logos_delivery/waku/waku_store/eligibility_hooks.nim @@ -0,0 +1,94 @@ +{.push raises: [].} + +import std/[options, strutils], chronos +import ./[common, protocol, eligibility_canonical] + +const EligibilityOutDescMinLen* = 512 + +type + EligibilityVerifierCb* = proc ( + proof_hex, canonical_hex, requester_peer_id: cstring, + out_desc: cstring, out_desc_len: csize_t, user_data: pointer + ): cint {.cdecl, gcsafe, raises: [].} + +proc defaultEligibilityDesc*(code: EligibilityStatusCode): string = + case code + of EligibilityStatusCode.OK: + "ok" + of EligibilityStatusCode.PARAMS_REJECTED: + "params rejected" + of EligibilityStatusCode.PROOF_INVALID: + "proof invalid" + of EligibilityStatusCode.STREAM_NOT_ACTIVE: + "stream not active" + +proc eligibilityCodeFromCint*(raw: cint): EligibilityStatusCode = + if raw < 0: + return EligibilityStatusCode.PROOF_INVALID + case raw + of 0: + EligibilityStatusCode.OK + of 1: + EligibilityStatusCode.PARAMS_REJECTED + of 2: + EligibilityStatusCode.PROOF_INVALID + of 3: + EligibilityStatusCode.STREAM_NOT_ACTIVE + else: + EligibilityStatusCode.PROOF_INVALID + +proc makeEligibilityWrappedHandler*( + verifierCb: EligibilityVerifierCb, + verifierUserData: pointer, + innerHandler: StoreQueryRequestHandler, + store: WakuStore, +): StoreQueryRequestHandler = + proc ( + req: StoreQueryRequest + ): Future[StoreQueryResult] {.async, gcsafe.} = + if verifierCb.isNil: + return await innerHandler(req) + + var reqClean = req + let proofBytes = + if req.eligibilityProof.isSome(): req.eligibilityProof.get() else: @[] + reqClean.eligibilityProof = none(seq[byte]) + + let canonicalHex = storeEligibilityCanonicalHex(reqClean) + var proofHexStorage = "" + let proofHexPtr: cstring = + if proofBytes.len > 0: + proofHexStorage = bytesToLowerHex(proofBytes) + cstring(proofHexStorage) + else: + nil + + let requesterPeerId = + if store.inboundRequestPeerId.isSome(): $store.inboundRequestPeerId.get() + else: "" + + var outDescBuf = newString(EligibilityOutDescMinLen) + outDescBuf.setLen(EligibilityOutDescMinLen) + let verdictRaw = verifierCb( + proofHexPtr, + cstring(canonicalHex), + cstring(requesterPeerId), + cstring(outDescBuf), + csize_t(EligibilityOutDescMinLen), + verifierUserData, + ) + + let code = eligibilityCodeFromCint(verdictRaw) + if code != EligibilityStatusCode.OK: + var desc = outDescBuf.strip(chars = {'\0'}) + if desc.len == 0: + desc = defaultEligibilityDesc(code) + let res = StoreQueryResponse( + requestId: req.requestId, + statusCode: uint32(StatusCode.BAD_REQUEST), + statusDesc: "BAD_REQUEST", + eligibilityStatus: some(EligibilityStatus(code: code, desc: desc)), + ) + return ok(res) + + return await innerHandler(req) diff --git a/logos_delivery/waku/waku_store/protocol.nim b/logos_delivery/waku/waku_store/protocol.nim index d09532d43..cbb8edbdc 100644 --- a/logos_delivery/waku/waku_store/protocol.nim +++ b/logos_delivery/waku/waku_store/protocol.nim @@ -34,6 +34,7 @@ type WakuStore* = ref object of LPProtocol rng: crypto.Rng requestHandler*: StoreQueryRequestHandler requestRateLimiter*: RequestRateLimiter + inboundRequestPeerId*: Option[PeerId] ## Protocol @@ -59,6 +60,10 @@ proc handleQueryRequest( peerId = requestor, requestId = requestId, request = req waku_store_queries.inc() + self.inboundRequestPeerId = some(requestor) + defer: + self.inboundRequestPeerId = none(PeerId) + let queryResult = await self.requestHandler(req) res = queryResult.valueOr: @@ -71,8 +76,9 @@ proc handleQueryRequest( return (res.encode().buffer, "not_parsed_requestId") res.requestId = requestId - res.statusCode = 200 - res.statusDesc = "OK" + if res.statusCode == 0: + res.statusCode = 200 + res.statusDesc = "OK" info "sending store query response", peerId = requestor, requestId = requestId, messages = res.messages.len diff --git a/tests/all_tests_waku.nim b/tests/all_tests_waku.nim index 5498db307..7cae38ddb 100644 --- a/tests/all_tests_waku.nim +++ b/tests/all_tests_waku.nim @@ -38,6 +38,8 @@ when os == "Linux" and import ./waku_store/test_client, ./waku_store/test_rpc_codec, + ./waku_store/test_store_eligibility_canonical, + ./waku_store/test_store_eligibility_hooks, ./waku_store/test_waku_store, ./waku_store/test_wakunode_store diff --git a/tests/waku_store/test_store_eligibility_canonical.nim b/tests/waku_store/test_store_eligibility_canonical.nim new file mode 100644 index 000000000..6d2c84d7e --- /dev/null +++ b/tests/waku_store/test_store_eligibility_canonical.nim @@ -0,0 +1,35 @@ +{.used.} + +import std/options, testutils/unittests +import logos_delivery/waku/[common/paging, waku_core, waku_store/common, waku_store/eligibility_canonical] + +const N8ReferenceWireHex = + "2f4c455a2f76302e312f53746f7265456c69676962696c6974792f0000000000050000007265712d3101010d0000002f77616b752f322f746f70696301000000140000002f6d792d6170702f312f636861742f70726f746f010a000000000000000002000000010101010101010101010101010101010101010101010101010101010101010102020202020202020202020202020202020202020202020202020202020202020001016400000000000000" + +procSuite "Waku Store - eligibility canonical (N8)": + test "matches lez-payment-streams-core n8_canonical_wire_hex reference": + let hash1 = block: + var h: WakuMessageHash + for i in 0 .. 31: + h[i] = byte(1) + h + let hash2 = block: + var h: WakuMessageHash + for i in 0 .. 31: + h[i] = byte(2) + h + + let query = StoreQueryRequest( + requestId: "req-1", + includeData: true, + pubsubTopic: some("/waku/2/topic"), + contentTopics: @["/my-app/1/chat/proto"], + startTime: some(Timestamp(10)), + endTime: none(Timestamp), + messageHashes: @[hash1, hash2], + paginationCursor: none(WakuMessageHash), + paginationForward: PagingDirection.FORWARD, + paginationLimit: some(uint64(100)), + ) + + check storeEligibilityCanonicalHex(query) == N8ReferenceWireHex diff --git a/tests/waku_store/test_store_eligibility_hooks.nim b/tests/waku_store/test_store_eligibility_hooks.nim new file mode 100644 index 000000000..281113674 --- /dev/null +++ b/tests/waku_store/test_store_eligibility_hooks.nim @@ -0,0 +1,119 @@ +{.used.} + +import std/options, testutils/unittests, chronos, libp2p/crypto/crypto +import + logos_delivery/waku/[common/paging, node/peer_manager, waku_core, waku_store, waku_store/common], + logos_delivery/waku/waku_store/[eligibility_hooks, self_req_handler], + ../testlib/[wakucore, testasync, futures], + ./store_utils + +procSuite "Store eligibility verifier wrapper": + var serverSwitch {.threadvar.}: Switch + var server {.threadvar.}: WakuStore + var innerCalled {.threadvar.}: int + var verifierCalled {.threadvar.}: int + var lastProofHex {.threadvar.}: string + var innerHandler {.threadvar.}: StoreQueryRequestHandler + + asyncSetup: + innerCalled = 0 + verifierCalled = 0 + lastProofHex = "" + + innerHandler = proc( + req: StoreQueryRequest + ): Future[StoreQueryResult] {.async, gcsafe.} = + innerCalled.inc() + return ok(StoreQueryResponse(requestId: req.requestId, statusCode: 200)) + + serverSwitch = newTestSwitch() + server = await newTestWakuStore( + serverSwitch, + handler = innerHandler, + ) + await serverSwitch.start() + await sleepAsync(100.millis) + + asyncTeardown: + await serverSwitch.stop() + + asyncTest "verifier reject skips inner handler": + let verifierCb = proc ( + proof_hex, canonical_hex, requester_peer_id: cstring, + out_desc: cstring, out_desc_len: csize_t, user_data: pointer + ): cint {.cdecl, gcsafe, raises: [].} = + verifierCalled.inc() + if not isNil(proof_hex): + lastProofHex = $proof_hex + cint(ord(EligibilityStatusCode.PROOF_INVALID)) + + server.requestHandler = makeEligibilityWrappedHandler( + verifierCb, nil, innerHandler, server + ) + server.inboundRequestPeerId = some(serverSwitch.peerInfo.peerId) + + let req = StoreQueryRequest( + requestId: "r1", + paginationForward: PagingDirection.FORWARD, + eligibilityProof: some(@[byte(0xAB)]), + ) + let res = (await server.handleSelfStoreRequest(req)).get() + + check: + verifierCalled == 1 + innerCalled == 0 + res.statusCode == uint32(StatusCode.BAD_REQUEST) + res.eligibilityStatus.isSome() + res.eligibilityStatus.get().code == EligibilityStatusCode.PROOF_INVALID + lastProofHex.len > 0 + + asyncTest "verifier OK delegates to inner handler": + innerCalled = 0 + verifierCalled = 0 + let verifierCb = proc ( + proof_hex, canonical_hex, requester_peer_id: cstring, + out_desc: cstring, out_desc_len: csize_t, user_data: pointer + ): cint {.cdecl, gcsafe, raises: [].} = + verifierCalled.inc() + cint(ord(EligibilityStatusCode.OK)) + + server.requestHandler = makeEligibilityWrappedHandler( + verifierCb, nil, innerHandler, server + ) + + let req = StoreQueryRequest( + requestId: "r2", paginationForward: PagingDirection.FORWARD + ) + discard (await server.handleSelfStoreRequest(req)).get() + + check: + verifierCalled == 1 + innerCalled == 1 + + asyncTest "NULL proof passes NULL proof_hex to verifier": + innerCalled = 0 + verifierCalled = 0 + lastProofHex = "unset" + let verifierCb = proc ( + proof_hex, canonical_hex, requester_peer_id: cstring, + out_desc: cstring, out_desc_len: csize_t, user_data: pointer + ): cint {.cdecl, gcsafe, raises: [].} = + verifierCalled.inc() + if isNil(proof_hex): + lastProofHex = "" + else: + lastProofHex = $proof_hex + cint(ord(EligibilityStatusCode.OK)) + + server.requestHandler = makeEligibilityWrappedHandler( + verifierCb, nil, innerHandler, server + ) + + let req = StoreQueryRequest( + requestId: "r3", paginationForward: PagingDirection.FORWARD + ) + discard (await server.handleSelfStoreRequest(req)).get() + + check: + verifierCalled == 1 + lastProofHex == ""