mirror of
https://github.com/logos-messaging/logos-messaging-nim.git
synced 2026-07-28 02:53:27 +00:00
deploy: 7cd242f0fcc04800b815b089c491af8763a617b2
This commit is contained in:
parent
476f3b8628
commit
faec25ab00
@ -42,7 +42,7 @@ proc sendPushRequest(wl: WakuLightPushClient, req: PushRequest, peer: PeerId|Rem
|
|||||||
return err(dialFailure)
|
return err(dialFailure)
|
||||||
let connection = connOpt.get()
|
let connection = connOpt.get()
|
||||||
|
|
||||||
let rpc = PushRPC(requestId: generateRequestId(wl.rng), request: req)
|
let rpc = PushRPC(requestId: generateRequestId(wl.rng), request: some(req))
|
||||||
await connection.writeLP(rpc.encode().buffer)
|
await connection.writeLP(rpc.encode().buffer)
|
||||||
|
|
||||||
var buffer = await connection.readLp(MaxRpcSize.int)
|
var buffer = await connection.readLp(MaxRpcSize.int)
|
||||||
@ -53,14 +53,14 @@ proc sendPushRequest(wl: WakuLightPushClient, req: PushRequest, peer: PeerId|Rem
|
|||||||
return err(decodeRpcFailure)
|
return err(decodeRpcFailure)
|
||||||
|
|
||||||
let pushResponseRes = decodeRespRes.get()
|
let pushResponseRes = decodeRespRes.get()
|
||||||
if pushResponseRes.response == PushResponse():
|
if pushResponseRes.response.isNone():
|
||||||
waku_lightpush_errors.inc(labelValues = [emptyResponseBodyFailure])
|
waku_lightpush_errors.inc(labelValues = [emptyResponseBodyFailure])
|
||||||
return err(emptyResponseBodyFailure)
|
return err(emptyResponseBodyFailure)
|
||||||
|
|
||||||
let response = pushResponseRes.response
|
let response = pushResponseRes.response.get()
|
||||||
if not response.isSuccess:
|
if not response.isSuccess:
|
||||||
if response.info != "":
|
if response.info.isSome():
|
||||||
return err(response.info)
|
return err(response.info.get())
|
||||||
else:
|
else:
|
||||||
return err("unknown failure")
|
return err("unknown failure")
|
||||||
|
|
||||||
|
|||||||
@ -45,27 +45,27 @@ proc initProtocolHandler*(wl: WakuLightPush) =
|
|||||||
return
|
return
|
||||||
|
|
||||||
let req = reqDecodeRes.get()
|
let req = reqDecodeRes.get()
|
||||||
if req.request == PushRequest():
|
if req.request.isNone():
|
||||||
error "invalid lightpush rpc received", error=emptyRequestBodyFailure
|
error "invalid lightpush rpc received", error=emptyRequestBodyFailure
|
||||||
waku_lightpush_errors.inc(labelValues = [emptyRequestBodyFailure])
|
waku_lightpush_errors.inc(labelValues = [emptyRequestBodyFailure])
|
||||||
return
|
return
|
||||||
|
|
||||||
waku_lightpush_messages.inc(labelValues = ["PushRequest"])
|
waku_lightpush_messages.inc(labelValues = ["PushRequest"])
|
||||||
let
|
let
|
||||||
pubSubTopic = req.request.pubSubTopic
|
pubSubTopic = req.request.get().pubSubTopic
|
||||||
message = req.request.message
|
message = req.request.get().message
|
||||||
debug "push request", peerId=conn.peerId, requestId=req.requestId, pubsubTopic=pubsubTopic
|
debug "push request", peerId=conn.peerId, requestId=req.requestId, pubsubTopic=pubsubTopic
|
||||||
|
|
||||||
var response: PushResponse
|
var response: PushResponse
|
||||||
let handleRes = await wl.pushHandler(conn.peerId, pubsubTopic, message)
|
let handleRes = await wl.pushHandler(conn.peerId, pubsubTopic, message)
|
||||||
if handleRes.isOk():
|
if handleRes.isOk():
|
||||||
response = PushResponse(is_success: true, info: "OK")
|
response = PushResponse(is_success: true, info: some("OK"))
|
||||||
else:
|
else:
|
||||||
response = PushResponse(is_success: false, info: handleRes.error)
|
response = PushResponse(is_success: false, info: some(handleRes.error))
|
||||||
waku_lightpush_errors.inc(labelValues = [messagePushFailure])
|
waku_lightpush_errors.inc(labelValues = [messagePushFailure])
|
||||||
error "pushed message handling failed", error=handleRes.error
|
error "pushed message handling failed", error=handleRes.error
|
||||||
|
|
||||||
let rpc = PushRPC(requestId: req.requestId, response: response)
|
let rpc = PushRPC(requestId: req.requestId, response: some(response))
|
||||||
await conn.writeLp(rpc.encode().buffer)
|
await conn.writeLp(rpc.encode().buffer)
|
||||||
|
|
||||||
wl.handler = handle
|
wl.handler = handle
|
||||||
|
|||||||
@ -3,6 +3,8 @@ when (NimMajor, NimMinor) < (1, 4):
|
|||||||
else:
|
else:
|
||||||
{.push raises: [].}
|
{.push raises: [].}
|
||||||
|
|
||||||
|
import
|
||||||
|
std/options
|
||||||
import
|
import
|
||||||
../waku_message
|
../waku_message
|
||||||
|
|
||||||
@ -13,9 +15,9 @@ type
|
|||||||
|
|
||||||
PushResponse* = object
|
PushResponse* = object
|
||||||
isSuccess*: bool
|
isSuccess*: bool
|
||||||
info*: string
|
info*: Option[string]
|
||||||
|
|
||||||
PushRPC* = object
|
PushRPC* = object
|
||||||
requestId*: string
|
requestId*: string
|
||||||
request*: PushRequest
|
request*: Option[PushRequest]
|
||||||
response*: PushResponse
|
response*: Option[PushResponse]
|
||||||
|
|||||||
@ -4,6 +4,8 @@ else:
|
|||||||
{.push raises: [].}
|
{.push raises: [].}
|
||||||
|
|
||||||
|
|
||||||
|
import
|
||||||
|
std/options
|
||||||
import
|
import
|
||||||
../../../common/protobuf,
|
../../../common/protobuf,
|
||||||
../waku_message,
|
../waku_message,
|
||||||
@ -27,12 +29,16 @@ proc decode*(T: type PushRequest, buffer: seq[byte]): ProtoResult[T] =
|
|||||||
var rpc = PushRequest()
|
var rpc = PushRequest()
|
||||||
|
|
||||||
var pubSubTopic: PubsubTopic
|
var pubSubTopic: PubsubTopic
|
||||||
discard ?pb.getField(1, pubSubTopic)
|
if not ?pb.getField(1, pubSubTopic):
|
||||||
rpc.pubSubTopic = pubSubTopic
|
return err(ProtoError.RequiredFieldMissing)
|
||||||
|
else:
|
||||||
|
rpc.pubSubTopic = pubSubTopic
|
||||||
|
|
||||||
var buf: seq[byte]
|
var messageBuf: seq[byte]
|
||||||
discard ?pb.getField(2, buf)
|
if not ?pb.getField(2, messageBuf):
|
||||||
rpc.message = ?WakuMessage.decode(buf)
|
return err(ProtoError.RequiredFieldMissing)
|
||||||
|
else:
|
||||||
|
rpc.message = ?WakuMessage.decode(messageBuf)
|
||||||
|
|
||||||
ok(rpc)
|
ok(rpc)
|
||||||
|
|
||||||
@ -48,15 +54,19 @@ proc encode*(rpc: PushResponse): ProtoBuffer =
|
|||||||
|
|
||||||
proc decode*(T: type PushResponse, buffer: seq[byte]): ProtoResult[T] =
|
proc decode*(T: type PushResponse, buffer: seq[byte]): ProtoResult[T] =
|
||||||
let pb = initProtoBuffer(buffer)
|
let pb = initProtoBuffer(buffer)
|
||||||
var rpc = PushResponse(isSuccess: false, info: "")
|
var rpc = PushResponse()
|
||||||
|
|
||||||
var isSuccess: uint64
|
var isSuccess: uint64
|
||||||
if ?pb.getField(1, isSuccess):
|
if not ?pb.getField(1, isSuccess):
|
||||||
|
return err(ProtoError.RequiredFieldMissing)
|
||||||
|
else:
|
||||||
rpc.isSuccess = bool(isSuccess)
|
rpc.isSuccess = bool(isSuccess)
|
||||||
|
|
||||||
var info: string
|
var info: string
|
||||||
discard ?pb.getField(2, info)
|
if not ?pb.getField(2, info):
|
||||||
rpc.info = info
|
rpc.info = none(string)
|
||||||
|
else:
|
||||||
|
rpc.info = some(info)
|
||||||
|
|
||||||
ok(rpc)
|
ok(rpc)
|
||||||
|
|
||||||
@ -65,8 +75,8 @@ proc encode*(rpc: PushRPC): ProtoBuffer =
|
|||||||
var pb = initProtoBuffer()
|
var pb = initProtoBuffer()
|
||||||
|
|
||||||
pb.write3(1, rpc.requestId)
|
pb.write3(1, rpc.requestId)
|
||||||
pb.write3(2, rpc.request.encode())
|
pb.write3(2, rpc.request.map(encode))
|
||||||
pb.write3(3, rpc.response.encode())
|
pb.write3(3, rpc.response.map(encode))
|
||||||
pb.finish3()
|
pb.finish3()
|
||||||
|
|
||||||
pb
|
pb
|
||||||
@ -76,15 +86,23 @@ proc decode*(T: type PushRPC, buffer: seq[byte]): ProtoResult[T] =
|
|||||||
var rpc = PushRPC()
|
var rpc = PushRPC()
|
||||||
|
|
||||||
var requestId: string
|
var requestId: string
|
||||||
discard ?pb.getField(1, requestId)
|
if not ?pb.getField(1, requestId):
|
||||||
rpc.requestId = requestId
|
return err(ProtoError.RequiredFieldMissing)
|
||||||
|
else:
|
||||||
|
rpc.requestId = requestId
|
||||||
|
|
||||||
var requestBuffer: seq[byte]
|
var requestBuffer: seq[byte]
|
||||||
discard ?pb.getField(2, requestBuffer)
|
if not ?pb.getField(2, requestBuffer):
|
||||||
rpc.request = ?PushRequest.decode(requestBuffer)
|
rpc.request = none(PushRequest)
|
||||||
|
else:
|
||||||
|
let request = ?PushRequest.decode(requestBuffer)
|
||||||
|
rpc.request = some(request)
|
||||||
|
|
||||||
var pushBuffer: seq[byte]
|
var responseBuffer: seq[byte]
|
||||||
discard ?pb.getField(3, pushBuffer)
|
if not ?pb.getField(3, responseBuffer):
|
||||||
rpc.response = ?PushResponse.decode(pushBuffer)
|
rpc.response = none(PushResponse)
|
||||||
|
else:
|
||||||
|
let response = ?PushResponse.decode(responseBuffer)
|
||||||
|
rpc.response = some(response)
|
||||||
|
|
||||||
ok(rpc)
|
ok(rpc)
|
||||||
Loading…
x
Reference in New Issue
Block a user