2022-11-04 10:52:27 +01:00
|
|
|
when (NimMajor, NimMinor) < (1, 4):
|
|
|
|
{.push raises: [Defect].}
|
|
|
|
else:
|
|
|
|
{.push raises: [].}
|
2022-10-25 14:55:31 +02:00
|
|
|
|
2024-03-16 00:08:47 +01:00
|
|
|
import std/options, stew/results, chronicles, chronos, metrics, bearssl/rand
|
2022-10-25 14:55:31 +02:00
|
|
|
import
|
2023-04-18 15:22:10 +02:00
|
|
|
../node/peer_manager,
|
|
|
|
../utils/requests,
|
2023-04-19 13:29:23 +02:00
|
|
|
../waku_core,
|
2024-01-30 07:28:21 -05:00
|
|
|
./common,
|
2022-10-25 14:55:31 +02:00
|
|
|
./protocol_metrics,
|
|
|
|
./rpc,
|
|
|
|
./rpc_codec
|
|
|
|
|
|
|
|
logScope:
|
2022-11-03 16:36:24 +01:00
|
|
|
topics = "waku lightpush client"
|
2022-10-25 14:55:31 +02:00
|
|
|
|
|
|
|
type WakuLightPushClient* = ref object
|
2024-03-16 00:08:47 +01:00
|
|
|
peerManager*: PeerManager
|
|
|
|
rng*: ref rand.HmacDrbgContext
|
2022-10-25 14:55:31 +02:00
|
|
|
|
2024-03-16 00:08:47 +01:00
|
|
|
proc new*(
|
|
|
|
T: type WakuLightPushClient, peerManager: PeerManager, rng: ref rand.HmacDrbgContext
|
|
|
|
): T =
|
2022-10-25 14:55:31 +02:00
|
|
|
WakuLightPushClient(peerManager: peerManager, rng: rng)
|
|
|
|
|
2024-03-16 00:08:47 +01:00
|
|
|
proc sendPushRequest(
|
|
|
|
wl: WakuLightPushClient, req: PushRequest, peer: PeerId | RemotePeerInfo
|
|
|
|
): Future[WakuLightPushResult[void]] {.async, gcsafe.} =
|
2022-10-25 14:55:31 +02:00
|
|
|
let connOpt = await wl.peerManager.dialPeer(peer, WakuLightPushCodec)
|
|
|
|
if connOpt.isNone():
|
|
|
|
waku_lightpush_errors.inc(labelValues = [dialFailure])
|
|
|
|
return err(dialFailure)
|
|
|
|
let connection = connOpt.get()
|
|
|
|
|
2022-11-18 20:01:01 +01:00
|
|
|
let rpc = PushRPC(requestId: generateRequestId(wl.rng), request: some(req))
|
2022-10-25 14:55:31 +02:00
|
|
|
await connection.writeLP(rpc.encode().buffer)
|
|
|
|
|
2024-02-06 17:37:42 +01:00
|
|
|
var buffer: seq[byte]
|
|
|
|
try:
|
2024-04-20 09:10:52 +05:30
|
|
|
buffer = await connection.readLp(DefaultMaxRpcSize.int)
|
2024-02-06 17:37:42 +01:00
|
|
|
except LPStreamRemoteClosedError:
|
|
|
|
return err("Exception reading: " & getCurrentExceptionMsg())
|
|
|
|
|
2022-11-07 16:24:16 +01:00
|
|
|
let decodeRespRes = PushRPC.decode(buffer)
|
2022-10-25 14:55:31 +02:00
|
|
|
if decodeRespRes.isErr():
|
|
|
|
error "failed to decode response"
|
|
|
|
waku_lightpush_errors.inc(labelValues = [decodeRpcFailure])
|
|
|
|
return err(decodeRpcFailure)
|
|
|
|
|
|
|
|
let pushResponseRes = decodeRespRes.get()
|
2022-11-18 20:01:01 +01:00
|
|
|
if pushResponseRes.response.isNone():
|
2022-10-25 14:55:31 +02:00
|
|
|
waku_lightpush_errors.inc(labelValues = [emptyResponseBodyFailure])
|
|
|
|
return err(emptyResponseBodyFailure)
|
|
|
|
|
2022-11-18 20:01:01 +01:00
|
|
|
let response = pushResponseRes.response.get()
|
2022-10-25 14:55:31 +02:00
|
|
|
if not response.isSuccess:
|
2022-11-18 20:01:01 +01:00
|
|
|
if response.info.isSome():
|
|
|
|
return err(response.info.get())
|
2022-10-25 14:55:31 +02:00
|
|
|
else:
|
|
|
|
return err("unknown failure")
|
|
|
|
|
|
|
|
return ok()
|
|
|
|
|
2024-03-16 00:08:47 +01:00
|
|
|
proc publish*(
|
|
|
|
wl: WakuLightPushClient,
|
|
|
|
pubSubTopic: PubsubTopic,
|
|
|
|
message: WakuMessage,
|
|
|
|
peer: PeerId | RemotePeerInfo,
|
|
|
|
): Future[WakuLightPushResult[void]] {.async, gcsafe.} =
|
2023-08-17 08:11:18 -04:00
|
|
|
let pushRequest = PushRequest(pubSubTopic: pubSubTopic, message: message)
|
2024-03-16 00:08:47 +01:00
|
|
|
return await wl.sendPushRequest(pushRequest, peer)
|