mirror of
https://github.com/logos-messaging/logos-delivery.git
synced 2026-08-11 06:13:16 +00:00
201 lines
6.3 KiB
Nim
201 lines
6.3 KiB
Nim
{.push raises: [Defect].}
|
|
|
|
import
|
|
std/[options, times, random],
|
|
stew/results,
|
|
chronicles,
|
|
chronos,
|
|
metrics,
|
|
bearssl/rand,
|
|
libp2p/crypto/crypto
|
|
|
|
import
|
|
../waku_message,
|
|
../waku_relay,
|
|
../../node/peer_manager/peer_manager,
|
|
../../utils/requests,
|
|
./rpc,
|
|
./rpc_codec
|
|
|
|
|
|
logScope:
|
|
topics = "wakulightpush"
|
|
|
|
declarePublicGauge waku_lightpush_peers, "number of lightpush peers"
|
|
declarePublicGauge waku_lightpush_errors, "number of lightpush protocol errors", ["type"]
|
|
declarePublicGauge waku_lightpush_messages, "number of lightpush messages received", ["type"]
|
|
|
|
|
|
const
|
|
WakuLightPushCodec* = "/vac/waku/lightpush/2.0.0-beta1"
|
|
|
|
const
|
|
MaxRpcSize* = MaxWakuMessageSize + 64 * 1024 # We add a 64kB safety buffer for protocol overhead
|
|
|
|
const
|
|
DandelionQ = 0.2 # Dandelion q paramter
|
|
EpochDuration = chronos.minutes(10)
|
|
|
|
# Error types (metric label values)
|
|
const
|
|
dialFailure = "dial_failure"
|
|
decodeRpcFailure = "decode_rpc_failure"
|
|
|
|
|
|
type
|
|
PushResponseHandler* = proc(response: PushResponse) {.gcsafe, closure.}
|
|
|
|
PushRequestHandler* = proc(requestId: string, msg: PushRequest) {.gcsafe, closure.}
|
|
|
|
WakuLightPushResult*[T] = Result[T, string]
|
|
|
|
WakuDandelionStem = ref object
|
|
isStemState: bool
|
|
dandelionRelay: RemotePeerInfo
|
|
|
|
WakuLightPush* = ref object of LPProtocol
|
|
rng*: ref rand.HmacDrbgContext
|
|
peerManager*: PeerManager
|
|
requestHandler*: PushRequestHandler
|
|
relayReference*: WakuRelay
|
|
wakuDandelionStem: Option[WakuDandelionStem]
|
|
|
|
proc request(wl: WakuLightPush, req: PushRequest, peer: RemotePeerInfo): Future[WakuLightPushResult[PushResponse]] {.async, gcsafe.} =
|
|
let connOpt = await wl.peerManager.dialPeer(peer, WakuLightPushCodec)
|
|
if connOpt.isNone():
|
|
waku_lightpush_errors.inc(labelValues = [dialFailure])
|
|
return err(dialFailure)
|
|
|
|
let connection = connOpt.get()
|
|
|
|
let rpc = PushRPC(requestId: generateRequestId(wl.rng), request: req)
|
|
await connection.writeLP(rpc.encode().buffer)
|
|
|
|
var message = await connection.readLp(MaxRpcSize.int)
|
|
let res = PushRPC.init(message)
|
|
|
|
if res.isErr():
|
|
waku_lightpush_errors.inc(labelValues = [decodeRpcFailure])
|
|
return err(decodeRpcFailure)
|
|
|
|
let rpcRes = res.get()
|
|
if rpcRes.response == PushResponse():
|
|
return err("empty response body")
|
|
|
|
return ok(rpcRes.response)
|
|
|
|
proc request*(wl: WakuLightPush, req: PushRequest): Future[WakuLightPushResult[PushResponse]] {.async, gcsafe.} =
|
|
let peerOpt = wl.peerManager.selectPeer(WakuLightPushCodec)
|
|
if peerOpt.isNone():
|
|
waku_lightpush_errors.inc(labelValues = [dialFailure])
|
|
return err(dialFailure)
|
|
|
|
return await wl.request(req, peerOpt.get())
|
|
|
|
proc updateEpoch(wl: WakuLightPush) =
|
|
randomize()
|
|
let wd = wl.wakuDandelionStem.get()
|
|
if rand(1.0) < DandelionQ:
|
|
wd.isStemState = false
|
|
else:
|
|
wd.isStemState = true
|
|
let peerOpt = wl.peerManager.selectPeer(WakuLightPushCodec) # TODO: select random peer from pubSubTopic mesh; all Dandelion supporting nodes have to support the WakuLightPushCodec
|
|
# todo: if peerOpt.isNone: ; retry until we get working peer
|
|
wd.dandelionRelay = peerOpt.get()
|
|
|
|
proc startDandelionEpochLoop(wl: WakuLightPush) =
|
|
updateEpoch(wl)
|
|
|
|
let currentTime = getTime().toUnix()
|
|
let timeToNextEpochBoundry = chronos.seconds(
|
|
EpochDuration.seconds - (currentTime mod EpochDuration.seconds))
|
|
|
|
# https://github.com/nim-lang/Nim/issues/17369
|
|
var executeUpdateEpoch: proc(data: pointer) {.gcsafe, raises: [Defect].}
|
|
executeUpdateEpoch = proc(udata: pointer) {.gcsafe.} =
|
|
updateEpoch(wl)
|
|
discard setTimer(Moment.fromNow(EpochDuration), executeUpdateEpoch)
|
|
|
|
discard setTimer(Moment.fromNow(timeToNextEpochBoundry), executeUpdateEpoch)
|
|
|
|
|
|
proc initProtocolHandler*(wl: WakuLightPush) =
|
|
|
|
proc handler(conn: Connection, proto: string) {.async, gcsafe, closure.} =
|
|
let message = await conn.readLp(MaxRpcSize.int)
|
|
let res = PushRPC.init(message)
|
|
if res.isErr():
|
|
error "failed to decode rpc"
|
|
waku_lightpush_errors.inc(labelValues = [decodeRpcFailure])
|
|
return
|
|
|
|
let rpc = res.get()
|
|
|
|
if rpc.request != PushRequest():
|
|
info "lightpush push request"
|
|
waku_lightpush_messages.inc(labelValues = ["PushRequest"])
|
|
|
|
let
|
|
pubSubTopic = rpc.request.pubSubTopic
|
|
message = rpc.request.message
|
|
debug "PushRequest", pubSubTopic=pubSubTopic, msg=message
|
|
|
|
var response: PushResponse
|
|
|
|
if wl.wakuDandelionStem.isSome and wl.wakuDandelionStem.get().isStemState: # Node is in Dandelion Stem State
|
|
let rpc = PushRequest(pubSubTopic: pubsubTopic, message: message)
|
|
discard wl.request(rpc, wl.wakuDandelionStem.get().dandelionRelay)
|
|
response = PushResponse(is_success: true, info: "Totally.") # do not tell that we are in stem state
|
|
|
|
if not wl.relayReference.isNil():
|
|
let data = message.encode().buffer
|
|
|
|
# Assumimng success, should probably be extended to check for network, peers, etc
|
|
discard wl.relayReference.publish(pubSubTopic, data)
|
|
response = PushResponse(is_success: true, info: "Totally.")
|
|
else:
|
|
debug "No relay protocol present, unsuccesssful push"
|
|
response = PushResponse(is_success: false, info: "No relay protocol")
|
|
|
|
let rpc = PushRPC(requestId: rpc.requestId, response: response)
|
|
await conn.writeLp(rpc.encode().buffer)
|
|
|
|
if rpc.response != PushResponse():
|
|
waku_lightpush_messages.inc(labelValues = ["PushResponse"])
|
|
if rpc.response.isSuccess:
|
|
info "lightpush message success"
|
|
else:
|
|
info "lightpush message failure", info=rpc.response.info
|
|
|
|
wl.handler = handler
|
|
wl.codec = WakuLightPushCodec
|
|
|
|
proc init*(T: type WakuLightPush, peerManager: PeerManager, rng: ref rand.HmacDrbgContext, handler: PushRequestHandler, relay: WakuRelay = nil, dandelion: bool = false): T =
|
|
debug "init"
|
|
let rng = crypto.newRng()
|
|
|
|
var wdOption: Option[WakuDandelionStem]
|
|
if dandelion:
|
|
let wd = WakuDandelionStem()
|
|
wdOption = some(wd)
|
|
|
|
let wl = WakuLightPush(rng: rng,
|
|
peerManager: peerManager,
|
|
requestHandler: handler,
|
|
relayReference: relay,
|
|
wakuDandelionStem: wdOption)
|
|
wl.initProtocolHandler()
|
|
|
|
if dandelion:
|
|
wl.startDandelionEpochLoop()
|
|
|
|
return wl
|
|
|
|
|
|
proc setPeer*(wlp: WakuLightPush, peer: RemotePeerInfo) =
|
|
wlp.peerManager.addPeer(peer, WakuLightPushCodec)
|
|
waku_lightpush_peers.inc()
|
|
|
|
|
|
|