status-go/protocol/common/message_sender.go

1181 lines
36 KiB
Go
Raw Normal View History

package common
2019-09-02 09:29:06 +00:00
import (
"context"
"crypto/ecdsa"
"database/sql"
"sync"
2019-09-02 09:29:06 +00:00
"time"
"github.com/golang/protobuf/proto"
"github.com/pkg/errors"
2020-01-02 09:10:19 +00:00
datasyncnode "github.com/vacp2p/mvds/node"
datasyncproto "github.com/vacp2p/mvds/protobuf"
"go.uber.org/zap"
"github.com/status-im/status-go/eth-node/crypto"
"github.com/status-im/status-go/eth-node/types"
"github.com/status-im/status-go/protocol/datasync"
datasyncpeer "github.com/status-im/status-go/protocol/datasync/peer"
"github.com/status-im/status-go/protocol/encryption"
2020-07-31 12:22:05 +00:00
"github.com/status-im/status-go/protocol/encryption/sharedsecret"
"github.com/status-im/status-go/protocol/protobuf"
"github.com/status-im/status-go/protocol/transport"
v1protocol "github.com/status-im/status-go/protocol/v1"
2019-09-02 09:29:06 +00:00
)
// Whisper message properties.
const (
whisperTTL = 15
whisperDefaultPoW = 0.002
// whisperLargeSizePoW is the PoWTarget for larger payload sizes
whisperLargeSizePoW = 0.000002
// largeSizeInBytes is when should we be using a lower POW.
// Roughly this is 50KB
largeSizeInBytes = 50000
whisperPoWTime = 5
2019-09-02 09:29:06 +00:00
)
// RekeyCompatibility indicates whether we should be sending
// keys in 1-to-1 messages as well as in the newer format
var RekeyCompatibility = true
// SentMessage reprent a message that has been passed to the transport layer
type SentMessage struct {
PublicKey *ecdsa.PublicKey
Spec *encryption.ProtocolMessageSpec
MessageIDs [][]byte
}
type MessageEventType uint32
const (
MessageScheduled = iota + 1
MessageSent
)
type MessageEvent struct {
Recipient *ecdsa.PublicKey
Type MessageEventType
SentMessage *SentMessage
RawMessage *RawMessage
}
type MessageSender struct {
identity *ecdsa.PrivateKey
datasync *datasync.DataSync
2021-10-21 12:39:19 +00:00
database *sql.DB
protocol *encryption.Protocol
transport *transport.Transport
logger *zap.Logger
persistence *RawMessagesPersistence
2020-07-22 07:41:40 +00:00
2021-10-21 12:39:19 +00:00
datasyncEnabled bool
2020-07-22 07:41:40 +00:00
// ephemeralKeys is a map that contains the ephemeral keys of the client, used
// to decrypt messages
ephemeralKeys map[string]*ecdsa.PrivateKey
ephemeralKeysMutex sync.Mutex
// messageEventsSubscriptions contains all the subscriptions for message events
messageEventsSubscriptions []chan<- *MessageEvent
featureFlags FeatureFlags
2020-07-31 12:22:05 +00:00
// handleSharedSecrets is a callback that is called every time a new shared secret is negotiated
handleSharedSecrets func([]*sharedsecret.Secret) error
2019-09-02 09:29:06 +00:00
}
func NewMessageSender(
2019-09-02 09:29:06 +00:00
identity *ecdsa.PrivateKey,
database *sql.DB,
enc *encryption.Protocol,
transport *transport.Transport,
2019-09-02 09:29:06 +00:00
logger *zap.Logger,
features FeatureFlags,
) (*MessageSender, error) {
dataSyncTransport := datasync.NewNodeTransport()
2019-09-02 09:29:06 +00:00
dataSyncNode, err := datasyncnode.NewPersistentNode(
database,
dataSyncTransport,
datasyncpeer.PublicKeyToPeerID(identity.PublicKey),
datasyncnode.BATCH,
datasync.CalculateSendTime,
logger,
)
if err != nil {
return nil, err
}
ds := datasync.New(dataSyncNode, dataSyncTransport, features.Datasync, logger)
2019-09-02 09:29:06 +00:00
p := &MessageSender{
2021-10-21 12:39:19 +00:00
identity: identity,
datasyncEnabled: features.Datasync,
datasync: ds,
protocol: enc,
database: database,
persistence: NewRawMessagesPersistence(database),
transport: transport,
logger: logger,
ephemeralKeys: make(map[string]*ecdsa.PrivateKey),
featureFlags: features,
2019-09-02 09:29:06 +00:00
}
// Initializing DataSync is required to encrypt and send messages.
// With DataSync enabled, messages are added to the DataSync
// but actual encrypt and send calls are postponed.
// sendDataSync is responsible for encrypting and sending postponed messages.
if features.Datasync {
// We set the max message size to 3/4 of the allowed message size, to leave
// room for encryption.
// Messages will be tried to send in any case, even if they exceed this
// value
ds.Init(p.sendDataSync, transport.MaxMessageSize()/4*3, logger)
ds.Start(datasync.DatasyncTicker)
2019-09-02 09:29:06 +00:00
}
return p, nil
}
func (s *MessageSender) Stop() {
for _, c := range s.messageEventsSubscriptions {
close(c)
}
s.messageEventsSubscriptions = nil
s.datasync.Stop() // idempotent op
2019-09-02 09:29:06 +00:00
}
func (s *MessageSender) SetHandleSharedSecrets(handler func([]*sharedsecret.Secret) error) {
s.handleSharedSecrets = handler
2020-07-31 12:22:05 +00:00
}
// SendPrivate takes encoded data, encrypts it and sends through the wire.
func (s *MessageSender) SendPrivate(
2019-09-02 09:29:06 +00:00
ctx context.Context,
2019-10-14 14:10:48 +00:00
recipient *ecdsa.PublicKey,
2021-02-23 15:47:45 +00:00
rawMessage *RawMessage,
2019-09-02 09:29:06 +00:00
) ([]byte, error) {
s.logger.Debug(
2019-09-02 09:29:06 +00:00
"sending a private message",
zap.String("public-key", types.EncodeHex(crypto.FromECDSAPub(recipient))),
zap.String("site", "SendPrivate"),
2019-09-02 09:29:06 +00:00
)
2020-07-22 07:41:40 +00:00
// Currently we don't support sending through datasync and setting custom waku fields,
// as the datasync interface is not rich enough to propagate that information, so we
// would have to add some complexity to handle this.
if rawMessage.ResendAutomatically && (rawMessage.Sender != nil || rawMessage.SkipEncryptionLayer || rawMessage.SendOnPersonalTopic) {
return nil, errors.New("setting identity, skip-encryption or personal topic and datasync not supported")
2020-07-22 07:41:40 +00:00
}
// Set sender identity if not specified
if rawMessage.Sender == nil {
rawMessage.Sender = s.identity
2020-07-22 07:41:40 +00:00
}
return s.sendPrivate(ctx, recipient, rawMessage)
}
2019-09-02 09:29:06 +00:00
2021-01-11 10:32:51 +00:00
// SendCommunityMessage takes encoded data, encrypts it and sends through the wire
// using the community topic and their key
func (s *MessageSender) SendCommunityMessage(
2021-01-11 10:32:51 +00:00
ctx context.Context,
rawMessage RawMessage,
) ([]byte, error) {
s.logger.Debug(
2021-01-11 10:32:51 +00:00
"sending a community message",
2022-05-27 09:14:40 +00:00
zap.String("communityId", types.EncodeHex(rawMessage.CommunityID)),
zap.String("site", "SendCommunityMessage"),
2021-01-11 10:32:51 +00:00
)
rawMessage.Sender = s.identity
2021-01-11 10:32:51 +00:00
2022-05-27 09:14:40 +00:00
return s.sendCommunity(ctx, &rawMessage)
2021-01-11 10:32:51 +00:00
}
// SendPubsubTopicKey sends the protected topic key for a community to a list of recipients
func (s *MessageSender) SendPubsubTopicKey(
ctx context.Context,
rawMessage *RawMessage,
) ([]byte, error) {
s.logger.Debug(
"sending the protected topic key for a community",
zap.String("communityId", types.EncodeHex(rawMessage.CommunityID)),
zap.String("site", "SendPubsubTopicKey"),
)
rawMessage.Sender = s.identity
messageID, err := s.getMessageID(rawMessage)
if err != nil {
return nil, err
}
rawMessage.ID = types.EncodeHex(messageID)
// Notify before dispatching, otherwise the dispatch subscription might happen
// earlier than the scheduled
s.notifyOnScheduledMessage(nil, rawMessage)
// Send to each recipients
for _, recipient := range rawMessage.Recipients {
_, err = s.sendPrivate(ctx, recipient, rawMessage)
if err != nil {
return nil, errors.Wrap(err, "failed to send message")
}
}
return messageID, nil
}
2020-07-22 07:41:40 +00:00
// SendGroup takes encoded data, encrypts it and sends through the wire,
// always return the messageID
func (s *MessageSender) SendGroup(
ctx context.Context,
recipients []*ecdsa.PublicKey,
2020-07-28 13:22:22 +00:00
rawMessage RawMessage,
) ([]byte, error) {
s.logger.Debug(
"sending a private group message",
zap.String("site", "SendGroup"),
)
2020-07-22 07:41:40 +00:00
// Set sender if not specified
if rawMessage.Sender == nil {
rawMessage.Sender = s.identity
}
2020-07-22 07:41:40 +00:00
// Calculate messageID first and set on raw message
wrappedMessage, err := s.wrapMessageV1(&rawMessage)
if err != nil {
return nil, errors.Wrap(err, "failed to wrap message")
}
messageID := v1protocol.MessageID(&rawMessage.Sender.PublicKey, wrappedMessage)
2020-07-22 07:41:40 +00:00
rawMessage.ID = types.EncodeHex(messageID)
// We call it only once, and we nil the function after so it doesn't get called again
if rawMessage.BeforeDispatch != nil {
if err := rawMessage.BeforeDispatch(&rawMessage); err != nil {
return nil, err
}
}
rawMessage.BeforeDispatch = nil
// Send to each recipients
for _, recipient := range recipients {
_, err = s.sendPrivate(ctx, recipient, &rawMessage)
if err != nil {
return nil, errors.Wrap(err, "failed to send message")
}
}
return messageID, nil
}
2022-05-27 09:14:40 +00:00
func (s *MessageSender) getMessageID(rawMessage *RawMessage) (types.HexBytes, error) {
wrappedMessage, err := s.wrapMessageV1(rawMessage)
if err != nil {
return nil, errors.Wrap(err, "failed to wrap message")
}
messageID := v1protocol.MessageID(&rawMessage.Sender.PublicKey, wrappedMessage)
return messageID, nil
}
func ShouldCommunityMessageBeEncrypted(msgType protobuf.ApplicationMetadataMessage_Type) bool {
return msgType == protobuf.ApplicationMetadataMessage_CHAT_MESSAGE ||
msgType == protobuf.ApplicationMetadataMessage_EDIT_MESSAGE ||
msgType == protobuf.ApplicationMetadataMessage_DELETE_MESSAGE ||
msgType == protobuf.ApplicationMetadataMessage_PIN_MESSAGE ||
msgType == protobuf.ApplicationMetadataMessage_EMOJI_REACTION
}
// sendCommunity sends a message that's to be sent in a community
// If it's a chat message, it will go to the respective topic derived by the
// chat id, if it's not a chat message, it will go to the community topic.
func (s *MessageSender) sendCommunity(
2021-01-11 10:32:51 +00:00
ctx context.Context,
rawMessage *RawMessage,
) ([]byte, error) {
2022-05-27 09:14:40 +00:00
s.logger.Debug("sending community message", zap.String("recipient", types.EncodeHex(crypto.FromECDSAPub(&rawMessage.Sender.PublicKey))))
2021-01-11 10:32:51 +00:00
// Set sender
if rawMessage.Sender == nil {
rawMessage.Sender = s.identity
}
2022-05-27 09:14:40 +00:00
messageID, err := s.getMessageID(rawMessage)
2021-01-11 10:32:51 +00:00
if err != nil {
2022-05-27 09:14:40 +00:00
return nil, err
2021-01-11 10:32:51 +00:00
}
rawMessage.ID = types.EncodeHex(messageID)
2022-05-27 09:14:40 +00:00
messageIDs := [][]byte{messageID}
2021-01-11 10:32:51 +00:00
2023-06-20 16:12:59 +00:00
if rawMessage.BeforeDispatch != nil {
if err := rawMessage.BeforeDispatch(rawMessage); err != nil {
return nil, err
}
}
2021-01-11 10:32:51 +00:00
// Notify before dispatching, otherwise the dispatch subscription might happen
// earlier than the scheduled
s.notifyOnScheduledMessage(nil, rawMessage)
2021-01-11 10:32:51 +00:00
var hash []byte
var newMessage *types.NewMessage
// Check if it's a key exchange message. In this case we send it
// to all the recipients
2022-05-27 09:14:40 +00:00
if rawMessage.CommunityKeyExMsgType != KeyExMsgNone {
forceRekey := rawMessage.CommunityKeyExMsgType == KeyExMsgRekey
// If rekeycompatibility is on, we always
// want to execute below, otherwise we execute
// only when we want to fill up old keys to a given user
if RekeyCompatibility || !forceRekey {
keyExMessageSpecs, err := s.protocol.GetKeyExMessageSpecs(rawMessage.HashRatchetGroupID, s.identity, rawMessage.Recipients, forceRekey)
if err != nil {
return nil, err
}
for i, spec := range keyExMessageSpecs {
recipient := rawMessage.Recipients[i]
_, _, err = s.sendMessageSpec(ctx, recipient, spec, messageIDs)
if err != nil {
return nil, err
}
}
2022-05-27 09:14:40 +00:00
}
if forceRekey {
var ratchet *encryption.HashRatchetKeyCompatibility
// We have just rekeyed, pull the latest
if RekeyCompatibility {
ratchet, err = s.protocol.GetCurrentKeyForGroup(rawMessage.HashRatchetGroupID)
if err != nil {
return nil, err
}
}
// We send the message over the community topic
messages, err := s.protocol.BuildHashRatchetReKeyGroupMessage(s.identity, rawMessage.Recipients, rawMessage.HashRatchetGroupID, ratchet)
2022-05-27 09:14:40 +00:00
if err != nil {
return nil, err
}
for _, m := range messages {
payload, err := proto.Marshal(m.Message)
if err != nil {
return nil, err
}
rawMessage.Payload = payload
newMessage = &types.NewMessage{
TTL: whisperTTL,
Payload: payload,
PowTarget: calculatePoW(payload),
PowTime: whisperPoWTime,
}
newMessage.Ephemeral = rawMessage.Ephemeral
newMessage.PubsubTopic = rawMessage.PubsubTopic
messageID := v1protocol.MessageID(&rawMessage.Sender.PublicKey, payload)
rawMessage.ID = types.EncodeHex(messageID)
// notify before dispatching
s.notifyOnScheduledMessage(nil, rawMessage)
_, err = s.transport.SendPublic(ctx, newMessage, types.EncodeHex(rawMessage.CommunityID))
if err != nil {
return nil, err
}
}
2022-05-27 09:14:40 +00:00
}
return nil, nil
}
wrappedMessage, err := s.wrapMessageV1(rawMessage)
if err != nil {
return nil, err
}
// If it's a chat message, we send it on the community chat topic
if ShouldCommunityMessageBeEncrypted(rawMessage.MessageType) {
messageSpec, err := s.protocol.BuildHashRatchetMessage(rawMessage.HashRatchetGroupID, wrappedMessage)
2022-05-27 09:14:40 +00:00
if err != nil {
return nil, err
}
2021-01-11 10:32:51 +00:00
payload, err := proto.Marshal(messageSpec.Message)
if err != nil {
return nil, errors.Wrap(err, "failed to marshal")
}
2022-09-21 16:05:29 +00:00
hash, newMessage, err = s.dispatchCommunityChatMessage(ctx, rawMessage, payload)
if err != nil {
return nil, err
}
sentMessage := &SentMessage{
Spec: messageSpec,
MessageIDs: messageIDs,
}
s.notifyOnSentMessage(sentMessage)
2021-10-29 14:29:28 +00:00
} else {
2021-01-11 10:32:51 +00:00
payload := wrappedMessage
pubkey, err := crypto.DecompressPubkey(rawMessage.CommunityID)
if err != nil {
return nil, errors.Wrap(err, "failed to decompress pubkey")
}
hash, newMessage, err = s.dispatchCommunityMessage(ctx, pubkey, payload, messageIDs, rawMessage.PubsubTopic)
if err != nil {
s.logger.Error("failed to send a community message", zap.Error(err))
return nil, errors.Wrap(err, "failed to send a message spec")
}
2022-05-27 09:14:40 +00:00
s.logger.Debug("sent community message ", zap.String("messageID", messageID.String()), zap.String("hash", types.EncodeHex(hash)))
2022-05-27 09:14:40 +00:00
}
s.transport.Track(messageIDs, hash, newMessage)
2021-01-11 10:32:51 +00:00
return messageID, nil
}
// sendPrivate sends data to the recipient identifying with a given public key.
func (s *MessageSender) sendPrivate(
ctx context.Context,
2019-10-14 14:10:48 +00:00
recipient *ecdsa.PublicKey,
rawMessage *RawMessage,
) ([]byte, error) {
s.logger.Debug("sending private message", zap.String("recipient", types.EncodeHex(crypto.FromECDSAPub(recipient))))
wrappedMessage, err := s.wrapMessageV1(rawMessage)
2019-09-02 09:29:06 +00:00
if err != nil {
return nil, errors.Wrap(err, "failed to wrap message")
}
messageID := v1protocol.MessageID(&rawMessage.Sender.PublicKey, wrappedMessage)
2020-07-22 07:41:40 +00:00
rawMessage.ID = types.EncodeHex(messageID)
rawMessage.PubsubTopic = transport.DefaultShardPubsubTopic() // TODO: determine which pubsub topic should be used for 1:1 messages
2020-07-22 07:41:40 +00:00
2023-06-20 16:12:59 +00:00
if rawMessage.BeforeDispatch != nil {
if err := rawMessage.BeforeDispatch(rawMessage); err != nil {
return nil, err
}
}
2020-07-22 07:41:40 +00:00
// Notify before dispatching, otherwise the dispatch subscription might happen
// earlier than the scheduled
s.notifyOnScheduledMessage(recipient, rawMessage)
2019-09-02 09:29:06 +00:00
if s.featureFlags.Datasync && rawMessage.ResendAutomatically {
// No need to call transport tracking.
// It is done in a data sync dispatch step.
datasyncID, err := s.addToDataSync(recipient, wrappedMessage)
2021-02-23 15:47:45 +00:00
if err != nil {
2019-09-02 09:29:06 +00:00
return nil, errors.Wrap(err, "failed to send message with datasync")
}
// We don't need to receive confirmations from our own devices
if !IsPubKeyEqual(recipient, &s.identity.PublicKey) {
confirmation := &RawMessageConfirmation{
DataSyncID: datasyncID,
MessageID: messageID,
PublicKey: crypto.CompressPubkey(recipient),
}
err = s.persistence.InsertPendingConfirmation(confirmation)
if err != nil {
return nil, err
}
}
} else if rawMessage.SkipEncryptionLayer {
// When SkipProtocolLayer is set we don't pass the message to the encryption layer
messageIDs := [][]byte{messageID}
hash, newMessage, err := s.sendPrivateRawMessage(ctx, rawMessage, recipient, wrappedMessage, messageIDs)
if err != nil {
s.logger.Error("failed to send a private message", zap.Error(err))
return nil, errors.Wrap(err, "failed to send a message spec")
}
s.logger.Debug("sent private message skipProtocolLayer", zap.String("messageID", messageID.String()), zap.String("hash", types.EncodeHex(hash)))
2021-10-29 14:29:28 +00:00
s.transport.Track(messageIDs, hash, newMessage)
2019-09-02 09:29:06 +00:00
} else {
2021-09-21 15:47:04 +00:00
messageSpec, err := s.protocol.BuildEncryptedMessage(rawMessage.Sender, recipient, wrappedMessage)
2019-09-02 09:29:06 +00:00
if err != nil {
return nil, errors.Wrap(err, "failed to encrypt message")
}
2020-07-31 12:22:05 +00:00
// The shared secret needs to be handle before we send a message
// otherwise the topic might not be set up before we receive a message
if s.handleSharedSecrets != nil {
err := s.handleSharedSecrets([]*sharedsecret.Secret{messageSpec.SharedSecret})
2020-07-31 12:22:05 +00:00
if err != nil {
return nil, err
}
}
2020-06-30 07:50:59 +00:00
messageIDs := [][]byte{messageID}
hash, newMessage, err := s.sendMessageSpec(ctx, recipient, messageSpec, messageIDs)
2019-09-02 09:29:06 +00:00
if err != nil {
s.logger.Error("failed to send a private message", zap.Error(err))
2019-09-02 09:29:06 +00:00
return nil, errors.Wrap(err, "failed to send a message spec")
}
2021-10-29 14:29:28 +00:00
s.logger.Debug("sent private message without datasync", zap.String("messageID", messageID.String()), zap.String("hash", types.EncodeHex(hash)))
s.transport.Track(messageIDs, hash, newMessage)
2019-09-02 09:29:06 +00:00
}
return messageID, nil
}
// sendPairInstallation sends data to the recipients, using DH
func (s *MessageSender) SendPairInstallation(
2019-10-14 14:10:48 +00:00
ctx context.Context,
recipient *ecdsa.PublicKey,
2020-07-28 13:22:22 +00:00
rawMessage RawMessage,
) ([]byte, error) {
s.logger.Debug("sending private message", zap.String("recipient", types.EncodeHex(crypto.FromECDSAPub(recipient))))
wrappedMessage, err := s.wrapMessageV1(&rawMessage)
if err != nil {
return nil, errors.Wrap(err, "failed to wrap message")
}
messageSpec, err := s.protocol.BuildDHMessage(s.identity, recipient, wrappedMessage)
if err != nil {
return nil, errors.Wrap(err, "failed to encrypt message")
}
messageID := v1protocol.MessageID(&s.identity.PublicKey, wrappedMessage)
2020-06-30 07:50:59 +00:00
messageIDs := [][]byte{messageID}
hash, newMessage, err := s.sendMessageSpec(ctx, recipient, messageSpec, messageIDs)
if err != nil {
return nil, errors.Wrap(err, "failed to send a message spec")
}
s.transport.Track(messageIDs, hash, newMessage)
return messageID, nil
}
func (s *MessageSender) encodeMembershipUpdate(
message v1protocol.MembershipUpdateMessage,
chatEntity ChatEntity,
) ([]byte, error) {
2019-10-14 14:10:48 +00:00
if chatEntity != nil {
chatEntityProtobuf := chatEntity.GetProtobuf()
2020-07-28 08:02:51 +00:00
switch chatEntityProtobuf := chatEntityProtobuf.(type) {
case *protobuf.ChatMessage:
2020-07-28 08:02:51 +00:00
message.Message = chatEntityProtobuf
case *protobuf.EmojiReaction:
2020-07-28 08:02:51 +00:00
message.EmojiReaction = chatEntityProtobuf
}
2019-10-14 14:10:48 +00:00
}
encodedMessage, err := v1protocol.EncodeMembershipUpdateMessage(message)
2019-10-14 14:10:48 +00:00
if err != nil {
return nil, errors.Wrap(err, "failed to encode membership update message")
}
return encodedMessage, nil
2019-10-14 14:10:48 +00:00
}
// EncodeMembershipUpdate takes a group and an optional chat message and returns the protobuf representation to be sent on the wire.
// All the events in a group are encoded and added to the payload
func (s *MessageSender) EncodeMembershipUpdate(
group *v1protocol.Group,
chatEntity ChatEntity,
) ([]byte, error) {
message := v1protocol.MembershipUpdateMessage{
ChatID: group.ChatID(),
Events: group.Events(),
}
return s.encodeMembershipUpdate(message, chatEntity)
}
// EncodeAbridgedMembershipUpdate takes a group and an optional chat message and returns the protobuf representation to be sent on the wire.
// Only the events relevant to the current group are encoded
func (s *MessageSender) EncodeAbridgedMembershipUpdate(
group *v1protocol.Group,
chatEntity ChatEntity,
) ([]byte, error) {
message := v1protocol.MembershipUpdateMessage{
ChatID: group.ChatID(),
Events: group.AbridgedEvents(),
}
return s.encodeMembershipUpdate(message, chatEntity)
}
2022-09-21 16:05:29 +00:00
func (s *MessageSender) dispatchCommunityChatMessage(ctx context.Context, rawMessage *RawMessage, wrappedMessage []byte) ([]byte, *types.NewMessage, error) {
newMessage := &types.NewMessage{
TTL: whisperTTL,
Payload: wrappedMessage,
PowTarget: calculatePoW(wrappedMessage),
PowTime: whisperPoWTime,
PubsubTopic: rawMessage.PubsubTopic,
}
2023-06-20 16:12:59 +00:00
if rawMessage.BeforeDispatch != nil {
if err := rawMessage.BeforeDispatch(rawMessage); err != nil {
return nil, nil, err
}
}
// notify before dispatching
s.notifyOnScheduledMessage(nil, rawMessage)
hash, err := s.transport.SendPublic(ctx, newMessage, rawMessage.LocalChatID)
if err != nil {
return nil, nil, err
}
return hash, newMessage, nil
}
// SendPublic takes encoded data, encrypts it and sends through the wire.
func (s *MessageSender) SendPublic(
ctx context.Context,
chatName string,
2020-07-28 13:22:22 +00:00
rawMessage RawMessage,
) ([]byte, error) {
2020-07-22 07:41:40 +00:00
// Set sender
if rawMessage.Sender == nil {
rawMessage.Sender = s.identity
}
2019-09-02 09:29:06 +00:00
wrappedMessage, err := s.wrapMessageV1(&rawMessage)
2019-09-02 09:29:06 +00:00
if err != nil {
return nil, errors.Wrap(err, "failed to wrap message")
}
2020-07-31 12:22:05 +00:00
var newMessage *types.NewMessage
messageSpec, err := s.protocol.BuildPublicMessage(s.identity, wrappedMessage)
if err != nil {
s.logger.Error("failed to send a public message", zap.Error(err))
return nil, errors.Wrap(err, "failed to wrap a public message in the encryption layer")
}
if !rawMessage.SkipEncryptionLayer {
2020-07-31 12:22:05 +00:00
newMessage, err = MessageSpecToWhisper(messageSpec)
if err != nil {
return nil, err
}
} else {
newMessage = &types.NewMessage{
TTL: whisperTTL,
Payload: wrappedMessage,
PowTarget: calculatePoW(wrappedMessage),
PowTime: whisperPoWTime,
}
2019-09-02 09:29:06 +00:00
}
newMessage.Ephemeral = rawMessage.Ephemeral
newMessage.PubsubTopic = rawMessage.PubsubTopic
2020-07-22 07:41:40 +00:00
messageID := v1protocol.MessageID(&rawMessage.Sender.PublicKey, wrappedMessage)
rawMessage.ID = types.EncodeHex(messageID)
2023-06-20 16:12:59 +00:00
if rawMessage.BeforeDispatch != nil {
if err := rawMessage.BeforeDispatch(&rawMessage); err != nil {
return nil, err
}
}
2020-07-22 07:41:40 +00:00
// notify before dispatching
s.notifyOnScheduledMessage(nil, &rawMessage)
2020-07-22 07:41:40 +00:00
hash, err := s.transport.SendPublic(ctx, newMessage, chatName)
2019-09-02 09:29:06 +00:00
if err != nil {
return nil, err
}
2021-10-29 14:29:28 +00:00
s.logger.Debug("sent public message", zap.String("messageID", messageID.String()), zap.String("hash", types.EncodeHex(hash)))
sentMessage := &SentMessage{
Spec: messageSpec,
MessageIDs: [][]byte{messageID},
}
s.notifyOnSentMessage(sentMessage)
s.transport.Track([][]byte{messageID}, hash, newMessage)
2019-09-02 09:29:06 +00:00
return messageID, nil
}
2021-02-23 15:47:45 +00:00
// unwrapDatasyncMessage tries to unwrap message as datasync one and in case of success
// returns cloned messages with replaced payloads
func unwrapDatasyncMessage(m *v1protocol.StatusMessage, datasync *datasync.DataSync) ([]*v1protocol.StatusMessage, [][]byte, error) {
var statusMessages []*v1protocol.StatusMessage
payloads, acks, err := datasync.UnwrapPayloadsAndAcks(
m.SigPubKey(),
m.EncryptionLayer.Payload,
2021-02-23 15:47:45 +00:00
)
if err != nil {
return nil, nil, err
}
for _, payload := range payloads {
message, err := m.Clone()
if err != nil {
return nil, nil, err
}
message.EncryptionLayer.Payload = payload
2021-02-23 15:47:45 +00:00
statusMessages = append(statusMessages, message)
}
return statusMessages, acks, nil
}
// HandleMessages expects a whisper message as input, and it will go through
2019-09-02 09:29:06 +00:00
// a series of transformations until the message is parsed into an application
// layer message, or in case of Raw methods, the processing stops at the layer
// before.
// It returns an error only if the processing of required steps failed.
func (s *MessageSender) HandleMessages(wakuMessage *types.Message) ([]*v1protocol.StatusMessage, [][]byte, error) {
logger := s.logger.With(zap.String("site", "HandleMessages"))
hlogger := logger.With(zap.ByteString("hash", wakuMessage.Hash))
2021-02-23 15:47:45 +00:00
var statusMessages []*v1protocol.StatusMessage
2022-09-21 16:05:29 +00:00
var acks [][]byte
2019-09-02 09:29:06 +00:00
response, err := s.handleMessage(wakuMessage)
2019-09-02 09:29:06 +00:00
if err != nil {
// Hash ratchet with a group id not found yet, save the message for future processing
if err == encryption.ErrHashRatchetGroupIDNotFound && len(response.Message.EncryptionLayer.HashRatchetInfo) == 1 {
info := response.Message.EncryptionLayer.HashRatchetInfo[0]
return nil, nil, s.persistence.SaveHashRatchetMessage(info.GroupID, info.KeyID, wakuMessage)
}
2019-09-02 09:29:06 +00:00
2022-09-21 16:05:29 +00:00
return nil, nil, err
}
statusMessages = append(statusMessages, response.Messages()...)
acks = append(acks, response.DatasyncAcks...)
2022-09-21 16:05:29 +00:00
// Process queued hash ratchet messages
for _, hashRatchetInfo := range response.Message.EncryptionLayer.HashRatchetInfo {
messages, err := s.persistence.GetHashRatchetMessages(hashRatchetInfo.KeyID)
2022-09-21 16:05:29 +00:00
if err != nil {
return nil, nil, err
}
var processedIds [][]byte
2022-09-21 16:05:29 +00:00
for _, message := range messages {
response, err := s.handleMessage(message)
if err != nil {
hlogger.Debug("failed to handle hash ratchet message", zap.Error(err))
continue
}
statusMessages = append(statusMessages, response.Messages()...)
acks = append(acks, response.DatasyncAcks...)
processedIds = append(processedIds, message.Hash)
2022-09-21 16:05:29 +00:00
}
err = s.persistence.DeleteHashRatchetMessages(processedIds)
if err != nil {
s.logger.Warn("failed to delete hash ratchet messages", zap.Error(err))
return nil, nil, err
}
}
return statusMessages, acks, nil
}
type handleMessageResponse struct {
Message *v1protocol.StatusMessage
DatasyncMessages []*v1protocol.StatusMessage
DatasyncAcks [][]byte
}
func (h *handleMessageResponse) Messages() []*v1protocol.StatusMessage {
if len(h.DatasyncMessages) > 0 {
return h.DatasyncMessages
}
return []*v1protocol.StatusMessage{h.Message}
}
func (s *MessageSender) handleMessage(wakuMessage *types.Message) (*handleMessageResponse, error) {
logger := s.logger.With(zap.String("site", "handleMessage"))
hlogger := logger.With(zap.ByteString("hash", wakuMessage.Hash))
response := &handleMessageResponse{
Message: &v1protocol.StatusMessage{},
DatasyncMessages: []*v1protocol.StatusMessage{},
DatasyncAcks: [][]byte{},
}
err := response.Message.HandleTransportLayer(wakuMessage)
if err != nil {
hlogger.Error("failed to handle transport layer message", zap.Error(err))
return nil, err
}
err = s.handleEncryptionLayer(context.Background(), response.Message)
if err != nil {
hlogger.Debug("failed to handle an encryption message", zap.Error(err))
// Hash ratchet with a group id not found yet, stop processing
if err == encryption.ErrHashRatchetGroupIDNotFound {
return response, err
}
2022-09-21 16:05:29 +00:00
}
datasyncMessages, as, err := unwrapDatasyncMessage(response.Message, s.datasync)
2019-09-02 09:29:06 +00:00
if err != nil {
hlogger.Debug("failed to handle datasync message", zap.Error(err))
2022-09-21 16:05:29 +00:00
} else {
response.DatasyncMessages = append(response.DatasyncMessages, datasyncMessages...)
response.DatasyncAcks = append(response.DatasyncAcks, as...)
2019-09-02 09:29:06 +00:00
}
for _, msg := range response.Messages() {
err := msg.HandleApplicationLayer()
2019-09-02 09:29:06 +00:00
if err != nil {
hlogger.Error("failed to handle application metadata layer message", zap.Error(err))
}
}
return response, nil
2019-09-02 09:29:06 +00:00
}
2020-07-22 07:41:40 +00:00
// fetchDecryptionKey returns the private key associated with this public key, and returns true if it's an ephemeral key
func (s *MessageSender) fetchDecryptionKey(destination *ecdsa.PublicKey) (*ecdsa.PrivateKey, bool) {
destinationID := types.EncodeHex(crypto.FromECDSAPub(destination))
2020-07-22 07:41:40 +00:00
s.ephemeralKeysMutex.Lock()
decryptionKey, ok := s.ephemeralKeys[destinationID]
s.ephemeralKeysMutex.Unlock()
2020-07-22 07:41:40 +00:00
// the key is not there, fallback on identity
if !ok {
return s.identity, false
}
2020-07-22 07:41:40 +00:00
return decryptionKey, true
}
func (s *MessageSender) handleEncryptionLayer(ctx context.Context, message *v1protocol.StatusMessage) error {
logger := s.logger.With(zap.String("site", "handleEncryptionLayer"))
2020-07-22 07:41:40 +00:00
publicKey := message.SigPubKey()
// if it's an ephemeral key, we don't negotiate a topic
decryptionKey, skipNegotiation := s.fetchDecryptionKey(message.TransportLayer.Dst)
2019-09-02 09:29:06 +00:00
err := message.HandleEncryptionLayer(decryptionKey, publicKey, s.protocol, skipNegotiation)
2020-07-22 07:41:40 +00:00
// if it's an ephemeral key, we don't have to handle a device not found error
if err == encryption.ErrDeviceNotFound && !skipNegotiation {
if err := s.handleErrDeviceNotFound(ctx, publicKey); err != nil {
logger.Error("failed to handle ErrDeviceNotFound", zap.Error(err))
2019-09-02 09:29:06 +00:00
}
}
if err != nil {
logger.Error("failed to handle an encrypted message", zap.Error(err))
return err
2019-09-02 09:29:06 +00:00
}
return nil
}
func (s *MessageSender) handleErrDeviceNotFound(ctx context.Context, publicKey *ecdsa.PublicKey) error {
2019-09-02 09:29:06 +00:00
now := time.Now().Unix()
advertise, err := s.protocol.ShouldAdvertiseBundle(publicKey, now)
2019-09-02 09:29:06 +00:00
if err != nil {
return err
}
if !advertise {
return nil
}
messageSpec, err := s.protocol.BuildBundleAdvertiseMessage(s.identity, publicKey)
2019-09-02 09:29:06 +00:00
if err != nil {
return err
}
ctx, cancel := context.WithTimeout(ctx, time.Second)
defer cancel()
2020-06-30 07:50:59 +00:00
// We don't pass an array of messageIDs as no action needs to be taken
// when sending a bundle
_, _, err = s.sendMessageSpec(ctx, publicKey, messageSpec, nil)
2019-09-02 09:29:06 +00:00
if err != nil {
return err
}
s.protocol.ConfirmBundleAdvertisement(publicKey, now)
2019-09-02 09:29:06 +00:00
return nil
}
func (s *MessageSender) wrapMessageV1(rawMessage *RawMessage) ([]byte, error) {
wrappedMessage, err := v1protocol.WrapMessageV1(rawMessage.Payload, rawMessage.MessageType, rawMessage.Sender)
2019-09-02 09:29:06 +00:00
if err != nil {
return nil, errors.Wrap(err, "failed to wrap message")
2019-09-02 09:29:06 +00:00
}
return wrappedMessage, nil
2019-09-02 09:29:06 +00:00
}
func (s *MessageSender) addToDataSync(publicKey *ecdsa.PublicKey, message []byte) ([]byte, error) {
groupID := datasync.ToOneToOneGroupID(&s.identity.PublicKey, publicKey)
2019-09-02 09:29:06 +00:00
peerID := datasyncpeer.PublicKeyToPeerID(*publicKey)
exist, err := s.datasync.IsPeerInGroup(groupID, peerID)
2019-09-02 09:29:06 +00:00
if err != nil {
2021-02-23 15:47:45 +00:00
return nil, errors.Wrap(err, "failed to check if peer is in group")
2019-09-02 09:29:06 +00:00
}
if !exist {
if err := s.datasync.AddPeer(groupID, peerID); err != nil {
2021-02-23 15:47:45 +00:00
return nil, errors.Wrap(err, "failed to add peer")
2019-09-02 09:29:06 +00:00
}
}
id, err := s.datasync.AppendMessage(groupID, message)
2019-09-02 09:29:06 +00:00
if err != nil {
2021-02-23 15:47:45 +00:00
return nil, errors.Wrap(err, "failed to append message to datasync")
2019-09-02 09:29:06 +00:00
}
2021-02-23 15:47:45 +00:00
return id[:], nil
2019-09-02 09:29:06 +00:00
}
// sendDataSync sends a message scheduled by the data sync layer.
// Data Sync layer calls this method "dispatch" function.
func (s *MessageSender) sendDataSync(ctx context.Context, publicKey *ecdsa.PublicKey, marshalledDatasyncPayload []byte, payload *datasyncproto.Payload) error {
// Calculate the messageIDs
2019-09-02 09:29:06 +00:00
messageIDs := make([][]byte, 0, len(payload.Messages))
2021-10-29 14:29:28 +00:00
hexMessageIDs := make([]string, 0, len(payload.Messages))
2019-09-02 09:29:06 +00:00
for _, payload := range payload.Messages {
2021-10-29 14:29:28 +00:00
mid := v1protocol.MessageID(&s.identity.PublicKey, payload.Body)
messageIDs = append(messageIDs, mid)
hexMessageIDs = append(hexMessageIDs, mid.String())
2019-09-02 09:29:06 +00:00
}
2021-09-21 15:47:04 +00:00
messageSpec, err := s.protocol.BuildEncryptedMessage(s.identity, publicKey, marshalledDatasyncPayload)
2019-09-02 09:29:06 +00:00
if err != nil {
return errors.Wrap(err, "failed to encrypt message")
}
2020-07-31 12:22:05 +00:00
// The shared secret needs to be handle before we send a message
// otherwise the topic might not be set up before we receive a message
if s.handleSharedSecrets != nil {
err := s.handleSharedSecrets([]*sharedsecret.Secret{messageSpec.SharedSecret})
2020-07-31 12:22:05 +00:00
if err != nil {
return err
}
}
hash, newMessage, err := s.sendMessageSpec(ctx, publicKey, messageSpec, messageIDs)
2019-09-02 09:29:06 +00:00
if err != nil {
s.logger.Error("failed to send a datasync message", zap.Error(err))
2019-09-02 09:29:06 +00:00
return err
}
2021-10-29 14:29:28 +00:00
s.logger.Debug("sent private messages", zap.Any("messageIDs", hexMessageIDs), zap.String("hash", types.EncodeHex(hash)))
s.transport.Track(messageIDs, hash, newMessage)
2019-09-02 09:29:06 +00:00
return nil
}
2020-07-31 12:22:05 +00:00
// sendPrivateRawMessage sends a message not wrapped in an encryption layer
func (s *MessageSender) sendPrivateRawMessage(ctx context.Context, rawMessage *RawMessage, publicKey *ecdsa.PublicKey, payload []byte, messageIDs [][]byte) ([]byte, *types.NewMessage, error) {
newMessage := &types.NewMessage{
TTL: whisperTTL,
Payload: payload,
PowTarget: calculatePoW(payload),
PowTime: whisperPoWTime,
PubsubTopic: rawMessage.PubsubTopic,
}
var hash []byte
var err error
if rawMessage.SendOnPersonalTopic {
hash, err = s.transport.SendPrivateOnPersonalTopic(ctx, newMessage, publicKey)
} else {
hash, err = s.transport.SendPrivateWithPartitioned(ctx, newMessage, publicKey)
}
if err != nil {
return nil, nil, err
}
return hash, newMessage, nil
}
// sendCommunityMessage sends a message not wrapped in an encryption layer
2021-01-11 10:32:51 +00:00
// to a community
func (s *MessageSender) dispatchCommunityMessage(ctx context.Context, publicKey *ecdsa.PublicKey, payload []byte, messageIDs [][]byte, pubsubTopic string) ([]byte, *types.NewMessage, error) {
2021-01-11 10:32:51 +00:00
newMessage := &types.NewMessage{
TTL: whisperTTL,
Payload: payload,
PowTarget: calculatePoW(payload),
PowTime: whisperPoWTime,
PubsubTopic: pubsubTopic,
2021-01-11 10:32:51 +00:00
}
hash, err := s.transport.SendCommunityMessage(ctx, newMessage, publicKey)
2021-01-11 10:32:51 +00:00
if err != nil {
return nil, nil, err
}
return hash, newMessage, nil
}
// sendMessageSpec analyses the spec properties and selects a proper transport method.
func (s *MessageSender) sendMessageSpec(ctx context.Context, publicKey *ecdsa.PublicKey, messageSpec *encryption.ProtocolMessageSpec, messageIDs [][]byte) ([]byte, *types.NewMessage, error) {
newMessage, err := MessageSpecToWhisper(messageSpec)
2019-09-02 09:29:06 +00:00
if err != nil {
return nil, nil, err
}
logger := s.logger.With(zap.String("site", "sendMessageSpec"))
2019-09-02 09:29:06 +00:00
var hash []byte
2020-07-31 12:22:05 +00:00
// process shared secret
if messageSpec.AgreedSecret {
2019-09-02 09:29:06 +00:00
logger.Debug("sending using shared secret")
hash, err = s.transport.SendPrivateWithSharedSecret(ctx, newMessage, publicKey, messageSpec.SharedSecret.Key)
2020-07-31 12:22:05 +00:00
} else {
2019-09-02 09:29:06 +00:00
logger.Debug("sending partitioned topic")
hash, err = s.transport.SendPrivateWithPartitioned(ctx, newMessage, publicKey)
2019-09-02 09:29:06 +00:00
}
if err != nil {
return nil, nil, err
}
sentMessage := &SentMessage{
PublicKey: publicKey,
Spec: messageSpec,
MessageIDs: messageIDs,
}
2020-06-30 07:50:59 +00:00
s.notifyOnSentMessage(sentMessage)
2020-07-22 07:41:40 +00:00
return hash, newMessage, nil
}
func (s *MessageSender) SubscribeToMessageEvents() <-chan *MessageEvent {
c := make(chan *MessageEvent, 100)
s.messageEventsSubscriptions = append(s.messageEventsSubscriptions, c)
2020-07-22 07:41:40 +00:00
return c
}
func (s *MessageSender) notifyOnSentMessage(sentMessage *SentMessage) {
event := &MessageEvent{
Type: MessageSent,
SentMessage: sentMessage,
}
// Publish on channels, drop if buffer is full
for _, c := range s.messageEventsSubscriptions {
select {
case c <- event:
default:
s.logger.Warn("message events subscription channel full when publishing sent event, dropping message")
2020-06-30 07:50:59 +00:00
}
}
2019-09-02 09:29:06 +00:00
}
func (s *MessageSender) notifyOnScheduledMessage(recipient *ecdsa.PublicKey, message *RawMessage) {
event := &MessageEvent{
Recipient: recipient,
Type: MessageScheduled,
RawMessage: message,
}
2020-07-22 07:41:40 +00:00
// Publish on channels, drop if buffer is full
for _, c := range s.messageEventsSubscriptions {
2020-07-22 07:41:40 +00:00
select {
case c <- event:
2020-07-22 07:41:40 +00:00
default:
s.logger.Warn("message events subscription channel full when publishing scheduled event, dropping message")
2020-07-22 07:41:40 +00:00
}
}
}
func (s *MessageSender) JoinPublic(id string) (*transport.Filter, error) {
return s.transport.JoinPublic(id)
2020-07-09 16:52:26 +00:00
}
2020-07-22 07:41:40 +00:00
// AddEphemeralKey adds an ephemeral key that we will be listening to
// note that we never removed them from now, as waku/whisper does not
// recalculate topics on removal, so effectively there's no benefit.
// On restart they will be gone.
func (s *MessageSender) AddEphemeralKey(privateKey *ecdsa.PrivateKey) (*transport.Filter, error) {
s.ephemeralKeysMutex.Lock()
s.ephemeralKeys[types.EncodeHex(crypto.FromECDSAPub(&privateKey.PublicKey))] = privateKey
s.ephemeralKeysMutex.Unlock()
return s.transport.LoadKeyFilters(privateKey)
}
func MessageSpecToWhisper(spec *encryption.ProtocolMessageSpec) (*types.NewMessage, error) {
var newMessage *types.NewMessage
2019-09-02 09:29:06 +00:00
payload, err := proto.Marshal(spec.Message)
if err != nil {
return newMessage, err
}
newMessage = &types.NewMessage{
2019-09-02 09:29:06 +00:00
TTL: whisperTTL,
Payload: payload,
PowTarget: calculatePoW(payload),
2019-09-02 09:29:06 +00:00
PowTime: whisperPoWTime,
}
return newMessage, nil
}
// calculatePoW returns the PoWTarget to be used.
// We check the size and arbitrarily set it to a lower PoW if the packet is
// greater than 50KB. We do this as the defaultPoW is too high for clients to send
// large messages.
func calculatePoW(payload []byte) float64 {
if len(payload) > largeSizeInBytes {
return whisperLargeSizePoW
}
return whisperDefaultPoW
}
2021-10-21 12:39:19 +00:00
func (s *MessageSender) StopDatasync() {
s.datasync.Stop()
}
func (s *MessageSender) StartDatasync() {
dataSyncTransport := datasync.NewNodeTransport()
dataSyncNode, err := datasyncnode.NewPersistentNode(
s.database,
dataSyncTransport,
datasyncpeer.PublicKeyToPeerID(s.identity.PublicKey),
datasyncnode.BATCH,
datasync.CalculateSendTime,
s.logger,
)
if err != nil {
return
}
ds := datasync.New(dataSyncNode, dataSyncTransport, true, s.logger)
if s.datasyncEnabled {
ds.Init(s.sendDataSync, s.transport.MaxMessageSize()/4*3, s.logger)
ds.Start(datasync.DatasyncTicker)
}
s.datasync = ds
}
// GetCurrentKeyForGroup returns the latest key timestampID belonging to a key group
func (s *MessageSender) GetCurrentKeyForGroup(groupID []byte) (*encryption.HashRatchetKeyCompatibility, error) {
return s.protocol.GetCurrentKeyForGroup(groupID)
}
// GetKeyIDsForGroup returns a slice of key IDs belonging to a given group ID
func (s *MessageSender) GetKeysForGroup(groupID []byte) ([]*encryption.HashRatchetKeyCompatibility, error) {
return s.protocol.GetKeysForGroup(groupID)
}