mirror of
https://github.com/logos-messaging/logos-delivery.git
synced 2026-07-26 06:23:14 +00:00
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.
This commit is contained in:
parent
a081681977
commit
df24c8e81c
19
Makefile
19
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/$@" \
|
||||
|
||||
@ -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],
|
||||
|
||||
@ -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
|
||||
|
||||
@ -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,
|
||||
|
||||
@ -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)
|
||||
|
||||
|
||||
170
library/store_eligibility/eligibility_api.nim
Normal file
170
library/store_eligibility/eligibility_api.nim
Normal file
@ -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
|
||||
65
library/store_eligibility/store_query_json.nim
Normal file
65
library/store_eligibility/store_query_json.nim
Normal file
@ -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,
|
||||
)
|
||||
)
|
||||
43
library/tests/test_eligibility_hooks.c
Normal file
43
library/tests/test_eligibility_hooks.c
Normal file
@ -0,0 +1,43 @@
|
||||
#include "../liblogosdelivery.h"
|
||||
#include <stdio.h>
|
||||
#include <string.h>
|
||||
|
||||
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;
|
||||
}
|
||||
97
logos_delivery/waku/waku_store/eligibility_canonical.nim
Normal file
97
logos_delivery/waku/waku_store/eligibility_canonical.nim
Normal file
@ -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))
|
||||
94
logos_delivery/waku/waku_store/eligibility_hooks.nim
Normal file
94
logos_delivery/waku/waku_store/eligibility_hooks.nim
Normal file
@ -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)
|
||||
@ -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
|
||||
|
||||
@ -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
|
||||
|
||||
|
||||
35
tests/waku_store/test_store_eligibility_canonical.nim
Normal file
35
tests/waku_store/test_store_eligibility_canonical.nim
Normal file
@ -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
|
||||
119
tests/waku_store/test_store_eligibility_hooks.nim
Normal file
119
tests/waku_store/test_store_eligibility_hooks.nim
Normal file
@ -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 == ""
|
||||
Loading…
x
Reference in New Issue
Block a user