2021-04-28 16:10:44 -04:00
|
|
|
package lightpush
|
|
|
|
|
|
|
|
import (
|
|
|
|
"context"
|
|
|
|
"encoding/hex"
|
|
|
|
"errors"
|
2021-11-25 09:46:04 -04:00
|
|
|
"math"
|
2021-04-28 16:10:44 -04:00
|
|
|
|
2022-10-19 15:39:32 -04:00
|
|
|
"github.com/libp2p/go-libp2p/core/host"
|
|
|
|
"github.com/libp2p/go-libp2p/core/network"
|
|
|
|
libp2pProtocol "github.com/libp2p/go-libp2p/core/protocol"
|
2023-02-06 18:16:20 -04:00
|
|
|
"github.com/libp2p/go-msgio/pbio"
|
2022-11-09 15:53:01 -04:00
|
|
|
"github.com/waku-org/go-waku/logging"
|
|
|
|
"github.com/waku-org/go-waku/waku/v2/metrics"
|
|
|
|
"github.com/waku-org/go-waku/waku/v2/protocol"
|
2023-02-06 18:16:20 -04:00
|
|
|
"github.com/waku-org/go-waku/waku/v2/protocol/lightpush/pb"
|
|
|
|
wpb "github.com/waku-org/go-waku/waku/v2/protocol/pb"
|
2022-11-09 15:53:01 -04:00
|
|
|
"github.com/waku-org/go-waku/waku/v2/protocol/relay"
|
2022-01-18 14:17:06 -04:00
|
|
|
"go.uber.org/zap"
|
2021-04-28 16:10:44 -04:00
|
|
|
)
|
|
|
|
|
2023-07-19 12:25:35 -04:00
|
|
|
// LightPushID_v20beta1 is the current Waku LightPush protocol identifier
|
2021-09-30 11:59:51 -04:00
|
|
|
const LightPushID_v20beta1 = libp2pProtocol.ID("/vac/waku/lightpush/2.0.0-beta1")
|
2021-04-28 16:10:44 -04:00
|
|
|
|
|
|
|
var (
|
|
|
|
ErrNoPeersAvailable = errors.New("no suitable remote peers")
|
|
|
|
ErrInvalidId = errors.New("invalid request id")
|
|
|
|
)
|
|
|
|
|
2023-07-19 12:25:35 -04:00
|
|
|
// WakuLightPush is the implementation of the Waku LightPush protocol
|
2021-04-28 16:10:44 -04:00
|
|
|
type WakuLightPush struct {
|
2023-01-06 18:37:57 -04:00
|
|
|
h host.Host
|
|
|
|
relay *relay.WakuRelay
|
|
|
|
cancel context.CancelFunc
|
2021-12-06 09:43:00 +01:00
|
|
|
|
2022-05-30 11:55:30 -04:00
|
|
|
log *zap.Logger
|
2021-04-28 16:10:44 -04:00
|
|
|
}
|
|
|
|
|
2022-05-04 17:08:24 -04:00
|
|
|
// NewWakuRelay returns a new instance of Waku Lightpush struct
|
2023-04-16 20:04:12 -04:00
|
|
|
func NewWakuLightPush(relay *relay.WakuRelay, log *zap.Logger) *WakuLightPush {
|
2021-04-28 16:10:44 -04:00
|
|
|
wakuLP := new(WakuLightPush)
|
|
|
|
wakuLP.relay = relay
|
2022-01-18 14:17:06 -04:00
|
|
|
wakuLP.log = log.Named("lightpush")
|
2021-04-28 16:23:03 -04:00
|
|
|
|
2021-11-01 08:38:03 -04:00
|
|
|
return wakuLP
|
|
|
|
}
|
|
|
|
|
2023-04-16 20:04:12 -04:00
|
|
|
// Sets the host to be able to mount or consume a protocol
|
|
|
|
func (wakuLP *WakuLightPush) SetHost(h host.Host) {
|
|
|
|
wakuLP.h = h
|
|
|
|
}
|
|
|
|
|
2022-05-04 17:08:24 -04:00
|
|
|
// Start inits the lighpush protocol
|
2023-01-06 18:37:57 -04:00
|
|
|
func (wakuLP *WakuLightPush) Start(ctx context.Context) error {
|
2022-10-20 09:42:01 -04:00
|
|
|
if wakuLP.relayIsNotAvailable() {
|
2021-11-07 13:08:29 +01:00
|
|
|
return errors.New("relay is required, without it, it is only a client and cannot be started")
|
2021-10-28 14:41:17 +02:00
|
|
|
}
|
2021-04-28 16:23:03 -04:00
|
|
|
|
2023-01-06 18:37:57 -04:00
|
|
|
ctx, cancel := context.WithCancel(ctx)
|
|
|
|
|
|
|
|
wakuLP.cancel = cancel
|
|
|
|
wakuLP.h.SetStreamHandlerMatch(LightPushID_v20beta1, protocol.PrefixTextMatch(string(LightPushID_v20beta1)), wakuLP.onRequest(ctx))
|
2022-01-18 14:17:06 -04:00
|
|
|
wakuLP.log.Info("Light Push protocol started")
|
2021-11-01 08:38:03 -04:00
|
|
|
|
|
|
|
return nil
|
2021-04-28 16:10:44 -04:00
|
|
|
}
|
|
|
|
|
2022-10-20 09:42:01 -04:00
|
|
|
// relayIsNotAvailable determines if this node supports relaying messages for other lightpush clients
|
|
|
|
func (wakuLp *WakuLightPush) relayIsNotAvailable() bool {
|
2021-11-07 13:08:29 +01:00
|
|
|
return wakuLp.relay == nil
|
|
|
|
}
|
|
|
|
|
2023-01-06 18:37:57 -04:00
|
|
|
func (wakuLP *WakuLightPush) onRequest(ctx context.Context) func(s network.Stream) {
|
|
|
|
return func(s network.Stream) {
|
|
|
|
defer s.Close()
|
|
|
|
logger := wakuLP.log.With(logging.HostID("peer", s.Conn().RemotePeer()))
|
|
|
|
requestPushRPC := &pb.PushRPC{}
|
2021-04-28 16:10:44 -04:00
|
|
|
|
2023-02-06 18:16:20 -04:00
|
|
|
writer := pbio.NewDelimitedWriter(s)
|
|
|
|
reader := pbio.NewDelimitedReader(s, math.MaxInt32)
|
2021-04-28 16:10:44 -04:00
|
|
|
|
2023-01-06 18:37:57 -04:00
|
|
|
err := reader.ReadMsg(requestPushRPC)
|
|
|
|
if err != nil {
|
|
|
|
logger.Error("reading request", zap.Error(err))
|
2023-04-19 12:09:03 -04:00
|
|
|
metrics.RecordLightpushError(ctx, "decode_rpc_failure")
|
2023-01-06 18:37:57 -04:00
|
|
|
return
|
|
|
|
}
|
2021-04-28 16:10:44 -04:00
|
|
|
|
2023-01-06 18:37:57 -04:00
|
|
|
logger.Info("request received")
|
|
|
|
if requestPushRPC.Query != nil {
|
|
|
|
logger.Info("push request")
|
|
|
|
response := new(pb.PushResponse)
|
2023-01-08 14:33:30 -04:00
|
|
|
|
|
|
|
pubSubTopic := requestPushRPC.Query.PubsubTopic
|
|
|
|
message := requestPushRPC.Query.Message
|
|
|
|
|
2023-04-19 12:09:03 -04:00
|
|
|
metrics.RecordLightpushMessage(ctx, "PushRequest")
|
|
|
|
|
2023-01-08 14:33:30 -04:00
|
|
|
// TODO: Assumes success, should probably be extended to check for network, peers, etc
|
|
|
|
// It might make sense to use WithReadiness option here?
|
|
|
|
|
|
|
|
_, err := wakuLP.relay.PublishToTopic(ctx, message, pubSubTopic)
|
|
|
|
|
|
|
|
if err != nil {
|
|
|
|
logger.Error("publishing message", zap.Error(err))
|
2023-04-19 12:09:03 -04:00
|
|
|
metrics.RecordLightpushError(ctx, "message_push_failure")
|
2023-01-08 14:33:30 -04:00
|
|
|
response.Info = "Could not publish message"
|
2021-04-28 16:10:44 -04:00
|
|
|
} else {
|
2023-01-08 14:33:30 -04:00
|
|
|
response.IsSuccess = true
|
|
|
|
response.Info = "Totally" // TODO: ask about this
|
2021-04-28 16:10:44 -04:00
|
|
|
}
|
|
|
|
|
2023-01-06 18:37:57 -04:00
|
|
|
responsePushRPC := &pb.PushRPC{}
|
|
|
|
responsePushRPC.RequestId = requestPushRPC.RequestId
|
|
|
|
responsePushRPC.Response = response
|
2021-04-28 16:10:44 -04:00
|
|
|
|
2023-01-06 18:37:57 -04:00
|
|
|
err = writer.WriteMsg(responsePushRPC)
|
|
|
|
if err != nil {
|
2023-04-19 12:09:03 -04:00
|
|
|
metrics.RecordLightpushError(ctx, "response_write_failure")
|
2023-01-06 18:37:57 -04:00
|
|
|
logger.Error("writing response", zap.Error(err))
|
|
|
|
_ = s.Reset()
|
|
|
|
} else {
|
|
|
|
logger.Info("response sent")
|
|
|
|
}
|
2023-04-19 12:09:03 -04:00
|
|
|
} else {
|
|
|
|
metrics.RecordLightpushError(ctx, "empty_request_body_failure")
|
2021-04-28 16:10:44 -04:00
|
|
|
}
|
|
|
|
|
2023-01-06 18:37:57 -04:00
|
|
|
if requestPushRPC.Response != nil {
|
|
|
|
if requestPushRPC.Response.IsSuccess {
|
|
|
|
logger.Info("request success")
|
|
|
|
} else {
|
|
|
|
logger.Info("request failure", zap.String("info=", requestPushRPC.Response.Info))
|
|
|
|
}
|
2023-04-19 12:09:03 -04:00
|
|
|
} else {
|
|
|
|
metrics.RecordLightpushError(ctx, "empty_response_body_failure")
|
2021-04-28 16:10:44 -04:00
|
|
|
}
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
2023-07-19 12:25:35 -04:00
|
|
|
func (wakuLP *WakuLightPush) request(ctx context.Context, req *pb.PushRequest, opts ...Option) (*pb.PushResponse, error) {
|
|
|
|
params := new(lightPushParameters)
|
2021-11-09 19:34:04 -04:00
|
|
|
params.host = wakuLP.h
|
2022-01-18 14:17:06 -04:00
|
|
|
params.log = wakuLP.log
|
2021-04-28 16:10:44 -04:00
|
|
|
|
2023-01-08 14:33:30 -04:00
|
|
|
optList := append(DefaultOptions(wakuLP.h), opts...)
|
2021-04-28 16:10:44 -04:00
|
|
|
for _, opt := range optList {
|
|
|
|
opt(params)
|
|
|
|
}
|
|
|
|
|
|
|
|
if params.selectedPeer == "" {
|
2023-04-19 12:09:03 -04:00
|
|
|
metrics.RecordLightpushError(ctx, "peer_not_found_failure")
|
2021-04-28 16:10:44 -04:00
|
|
|
return nil, ErrNoPeersAvailable
|
|
|
|
}
|
|
|
|
|
2023-07-19 12:25:35 -04:00
|
|
|
if len(params.requestID) == 0 {
|
2021-04-28 16:10:44 -04:00
|
|
|
return nil, ErrInvalidId
|
|
|
|
}
|
|
|
|
|
2022-05-30 11:55:30 -04:00
|
|
|
logger := wakuLP.log.With(logging.HostID("peer", params.selectedPeer))
|
2022-03-03 12:04:03 -04:00
|
|
|
|
2021-09-30 11:59:51 -04:00
|
|
|
connOpt, err := wakuLP.h.NewStream(ctx, params.selectedPeer, LightPushID_v20beta1)
|
2021-04-28 16:10:44 -04:00
|
|
|
if err != nil {
|
2022-05-30 11:55:30 -04:00
|
|
|
logger.Error("creating stream to peer", zap.Error(err))
|
2023-04-19 12:09:03 -04:00
|
|
|
metrics.RecordLightpushError(ctx, "dial_failure")
|
2021-04-28 16:10:44 -04:00
|
|
|
return nil, err
|
|
|
|
}
|
|
|
|
|
|
|
|
defer connOpt.Close()
|
2021-08-13 13:56:09 +02:00
|
|
|
defer func() {
|
|
|
|
err := connOpt.Reset()
|
|
|
|
if err != nil {
|
2023-04-19 12:09:03 -04:00
|
|
|
metrics.RecordLightpushError(ctx, "dial_failure")
|
2022-05-30 11:55:30 -04:00
|
|
|
logger.Error("resetting connection", zap.Error(err))
|
2021-08-13 13:56:09 +02:00
|
|
|
}
|
|
|
|
}()
|
2021-04-28 16:10:44 -04:00
|
|
|
|
2023-07-19 12:25:35 -04:00
|
|
|
pushRequestRPC := &pb.PushRPC{RequestId: hex.EncodeToString(params.requestID), Query: req}
|
2021-04-28 16:10:44 -04:00
|
|
|
|
2023-02-06 18:16:20 -04:00
|
|
|
writer := pbio.NewDelimitedWriter(connOpt)
|
|
|
|
reader := pbio.NewDelimitedReader(connOpt, math.MaxInt32)
|
2021-04-28 16:10:44 -04:00
|
|
|
|
|
|
|
err = writer.WriteMsg(pushRequestRPC)
|
|
|
|
if err != nil {
|
2023-04-19 12:09:03 -04:00
|
|
|
metrics.RecordLightpushError(ctx, "request_write_failure")
|
2022-05-30 11:55:30 -04:00
|
|
|
logger.Error("writing request", zap.Error(err))
|
2021-04-28 16:10:44 -04:00
|
|
|
return nil, err
|
|
|
|
}
|
|
|
|
|
|
|
|
pushResponseRPC := &pb.PushRPC{}
|
|
|
|
err = reader.ReadMsg(pushResponseRPC)
|
|
|
|
if err != nil {
|
2022-05-30 11:55:30 -04:00
|
|
|
logger.Error("reading response", zap.Error(err))
|
2023-04-19 12:09:03 -04:00
|
|
|
metrics.RecordLightpushError(ctx, "decode_rpc_failure")
|
2021-04-28 16:10:44 -04:00
|
|
|
return nil, err
|
|
|
|
}
|
|
|
|
|
|
|
|
return pushResponseRPC.Response, nil
|
|
|
|
}
|
2021-10-11 18:45:54 -04:00
|
|
|
|
2022-05-04 17:08:24 -04:00
|
|
|
// Stop unmounts the lightpush protocol
|
2021-11-05 16:09:48 -04:00
|
|
|
func (wakuLP *WakuLightPush) Stop() {
|
2023-01-06 18:37:57 -04:00
|
|
|
if wakuLP.cancel == nil {
|
|
|
|
return
|
2022-10-20 09:42:01 -04:00
|
|
|
}
|
2023-01-06 18:37:57 -04:00
|
|
|
|
|
|
|
wakuLP.cancel()
|
|
|
|
wakuLP.h.RemoveStreamHandler(LightPushID_v20beta1)
|
2021-10-11 18:45:54 -04:00
|
|
|
}
|
2021-11-01 08:38:03 -04:00
|
|
|
|
2022-05-04 17:08:24 -04:00
|
|
|
// PublishToTopic is used to broadcast a WakuMessage to a pubsub topic via lightpush protocol
|
2023-07-19 12:25:35 -04:00
|
|
|
func (wakuLP *WakuLightPush) PublishToTopic(ctx context.Context, message *wpb.WakuMessage, topic string, opts ...Option) ([]byte, error) {
|
2021-11-01 08:38:03 -04:00
|
|
|
if message == nil {
|
|
|
|
return nil, errors.New("message can't be null")
|
|
|
|
}
|
|
|
|
|
|
|
|
req := new(pb.PushRequest)
|
|
|
|
req.Message = message
|
2021-11-19 16:01:52 -04:00
|
|
|
req.PubsubTopic = topic
|
2021-11-01 08:38:03 -04:00
|
|
|
|
2021-11-05 16:09:48 -04:00
|
|
|
response, err := wakuLP.request(ctx, req, opts...)
|
2021-11-01 08:38:03 -04:00
|
|
|
if err != nil {
|
|
|
|
return nil, err
|
|
|
|
}
|
|
|
|
|
|
|
|
if response.IsSuccess {
|
2023-03-04 17:51:51 -04:00
|
|
|
hash := message.Hash(topic)
|
2022-11-03 09:53:33 -04:00
|
|
|
wakuLP.log.Info("waku.lightpush published", logging.HexString("hash", hash))
|
2021-11-01 08:38:03 -04:00
|
|
|
return hash, nil
|
|
|
|
} else {
|
|
|
|
return nil, errors.New(response.Info)
|
|
|
|
}
|
|
|
|
}
|
2021-11-19 16:01:52 -04:00
|
|
|
|
2022-05-04 17:08:24 -04:00
|
|
|
// Publish is used to broadcast a WakuMessage to the default waku pubsub topic via lightpush protocol
|
2023-07-19 12:25:35 -04:00
|
|
|
func (wakuLP *WakuLightPush) Publish(ctx context.Context, message *wpb.WakuMessage, opts ...Option) ([]byte, error) {
|
2021-11-19 20:03:05 -04:00
|
|
|
return wakuLP.PublishToTopic(ctx, message, relay.DefaultWakuTopic, opts...)
|
2021-11-19 16:01:52 -04:00
|
|
|
}
|