mirror of
https://github.com/logos-messaging/logos-delivery.git
synced 2026-07-25 22:13:12 +00:00
feat: HTTP REST API: Filter support v2 (#1890)
Filter v2 rest api support implemented Filter rest api documentation updated with v1 and v2 interface support. Separated legacy filter rest interface Fix code and tests of v2 Filter rest api Filter v2 message push test added Applied autoshard to Filter V2 Redesigned FilterPushHandling, code style, catch up apps and tests with filter v2 interface changes Rename of FilterV1SubscriptionsRequest to FilterLegacySubscribeRequest, fix broken chat2 app, fix tests Changed Filter v2 push handler subscription to simple register Separate node's filterUnsubscribe and filterUnsubscribeAll
This commit is contained in:
+12
-6
@@ -262,10 +262,12 @@ proc writeAndPrint(c: Chat) {.async.} =
|
||||
echo "You are now known as " & c.nick
|
||||
|
||||
elif line.startsWith("/exit"):
|
||||
if not c.node.wakuFilter.isNil():
|
||||
if not c.node.wakuFilterLegacy.isNil():
|
||||
echo "unsubscribing from content filters..."
|
||||
|
||||
await c.node.unsubscribe(pubsubTopic=some(DefaultPubsubTopic), contentTopics=c.contentTopic)
|
||||
let peerOpt = c.node.peerManager.selectPeer(WakuLegacyFilterCodec)
|
||||
if peerOpt.isSome():
|
||||
await c.node.legacyFilterUnsubscribe(pubsubTopic=some(DefaultPubsubTopic), contentTopics=c.contentTopic, peer=peerOpt.get())
|
||||
|
||||
echo "quitting..."
|
||||
|
||||
@@ -464,14 +466,18 @@ proc processInput(rfd: AsyncFD, rng: ref HmacDrbgContext) {.async.} =
|
||||
if peerInfo.isOk():
|
||||
await node.mountFilter()
|
||||
await node.mountFilterClient()
|
||||
node.peerManager.addServicePeer(peerInfo.value, WakuFilterCodec)
|
||||
node.peerManager.addServicePeer(peerInfo.value, WakuLegacyFilterCodec)
|
||||
|
||||
proc filterHandler(pubsubTopic: PubsubTopic, msg: WakuMessage) {.async, gcsafe, closure.} =
|
||||
trace "Hit filter handler", contentTopic=msg.contentTopic
|
||||
chat.printReceivedMessage(msg)
|
||||
|
||||
await node.subscribe(pubsubTopic=some(DefaultPubsubTopic), contentTopics=chat.contentTopic, filterHandler)
|
||||
|
||||
await node.legacyFilterSubscribe(pubsubTopic=some(DefaultPubsubTopic),
|
||||
contentTopics=chat.contentTopic,
|
||||
filterHandler,
|
||||
peerInfo.value)
|
||||
# TODO: Here to support FilterV2 relevant subscription, but still
|
||||
# Legacy Filter is concurrent to V2 untill legacy filter will be removed
|
||||
else:
|
||||
error "Filter not mounted. Couldn't parse conf.filternode",
|
||||
error = peerInfo.error
|
||||
@@ -485,7 +491,7 @@ proc processInput(rfd: AsyncFD, rng: ref HmacDrbgContext) {.async.} =
|
||||
chat.printReceivedMessage(msg)
|
||||
|
||||
let topic = DefaultPubsubTopic
|
||||
await node.subscribe(some(topic), @[ContentTopic("")], handler)
|
||||
node.subscribe(topic, handler)
|
||||
|
||||
if conf.rlnRelay:
|
||||
info "WakuRLNRelay is enabled"
|
||||
|
||||
@@ -18,6 +18,7 @@ import
|
||||
../../../waku/waku_node,
|
||||
../../../waku/node/peer_manager,
|
||||
../../waku/waku_filter,
|
||||
../../waku/waku_filter_v2,
|
||||
../../waku/waku_store,
|
||||
# Chat 2 imports
|
||||
../chat2/chat2,
|
||||
@@ -297,7 +298,8 @@ when isMainModule:
|
||||
if conf.filternode != "":
|
||||
let filterPeer = parsePeerInfo(conf.filternode)
|
||||
if filterPeer.isOk():
|
||||
bridge.nodev2.peerManager.addServicePeer(filterPeer.value, WakuFilterCodec)
|
||||
bridge.nodev2.peerManager.addServicePeer(filterPeer.value, WakuLegacyFilterCodec)
|
||||
bridge.nodev2.peerManager.addServicePeer(filterPeer.value, WakuFilterSubscribeCodec)
|
||||
else:
|
||||
error "Error parsing conf.filternode", error = filterPeer.error
|
||||
|
||||
|
||||
+11
-4
@@ -37,6 +37,8 @@ import
|
||||
../../waku/waku_store,
|
||||
../../waku/waku_lightpush,
|
||||
../../waku/waku_filter,
|
||||
../../waku/waku_filter_v2,
|
||||
../../waku/waku_filter_v2/client as waku_filter_client,
|
||||
./wakunode2_validator_signed,
|
||||
./internal_config,
|
||||
./external_config
|
||||
@@ -46,6 +48,7 @@ import
|
||||
../../waku/node/rest/debug/handlers as rest_debug_api,
|
||||
../../waku/node/rest/relay/handlers as rest_relay_api,
|
||||
../../waku/node/rest/relay/topic_cache,
|
||||
../../waku/node/rest/filter/legacy_handlers as rest_legacy_filter_api,
|
||||
../../waku/node/rest/filter/handlers as rest_filter_api,
|
||||
../../waku/node/rest/store/handlers as rest_store_api,
|
||||
../../waku/node/rest/health/handlers as rest_health_api,
|
||||
@@ -470,8 +473,9 @@ proc setupProtocols(node: WakuNode,
|
||||
if conf.filternode != "":
|
||||
let filterNode = parsePeerInfo(conf.filternode)
|
||||
if filterNode.isOk():
|
||||
await mountFilterClient(node)
|
||||
node.peerManager.addServicePeer(filterNode.value, WakuFilterCodec)
|
||||
await node.mountFilterClient()
|
||||
node.peerManager.addServicePeer(filterNode.value, WakuLegacyFilterCodec)
|
||||
node.peerManager.addServicePeer(filterNode.value, WakuFilterSubscribeCodec)
|
||||
else:
|
||||
return err("failed to set node waku filter peer: " & filterNode.error)
|
||||
|
||||
@@ -577,8 +581,11 @@ proc startRestServer(app: App, address: ValidIpAddress, port: Port, conf: WakuNo
|
||||
|
||||
## Filter REST API
|
||||
if conf.filter:
|
||||
let filterCache = rest_filter_api.MessageCache.init(capacity=rest_filter_api.filterMessageCacheDefaultCapacity)
|
||||
installFilterApiHandlers(server.router, app.node, filterCache)
|
||||
let legacyFilterCache = rest_legacy_filter_api.MessageCache.init()
|
||||
rest_legacy_filter_api.installLegacyFilterRestApiHandlers(server.router, app.node, legacyFilterCache)
|
||||
|
||||
let filterCache = rest_filter_api.MessageCache.init()
|
||||
rest_filter_api.installFilterRestApiHandlers(server.router, app.node, filterCache)
|
||||
|
||||
## Store REST API
|
||||
installStoreApiHandlers(server.router, app.node)
|
||||
|
||||
Reference in New Issue
Block a user